|
18 | 18 | */ |
19 | 19 | package org.apache.pulsar.broker.admin; |
20 | 20 |
|
| 21 | +import static org.apache.pulsar.common.policies.data.PoliciesUtil.getBundles; |
21 | 22 | import static org.mockito.Mockito.doReturn; |
22 | 23 | import static org.mockito.Mockito.doThrow; |
23 | 24 | import static org.mockito.Mockito.mock; |
|
38 | 39 | import org.apache.pulsar.broker.web.PulsarWebResource; |
39 | 40 | import org.apache.pulsar.broker.web.RestException; |
40 | 41 | import org.apache.pulsar.common.naming.NamespaceName; |
| 42 | +import org.apache.pulsar.common.policies.data.BookieAffinityGroupData; |
| 43 | +import org.apache.pulsar.common.policies.data.BundlesData; |
41 | 44 | import org.apache.pulsar.common.policies.data.ClusterData; |
42 | 45 | import org.apache.pulsar.common.policies.data.DelayedDeliveryPolicies; |
43 | 46 | import org.apache.pulsar.common.policies.data.DispatchRate; |
| 47 | +import org.apache.pulsar.common.policies.data.Policies; |
44 | 48 | import org.apache.pulsar.common.policies.data.PolicyName; |
45 | 49 | import org.apache.pulsar.common.policies.data.PolicyOperation; |
46 | 50 | import org.apache.pulsar.common.policies.data.PublishRate; |
@@ -236,4 +240,149 @@ public void testOperationDelayedDelivery() throws Exception { |
236 | 240 | assertEquals(e.getResponse().getStatus(), Response.Status.PRECONDITION_FAILED.getStatusCode()); |
237 | 241 | } |
238 | 242 | } |
| 243 | + |
| 244 | + @Test |
| 245 | + public void testSetBookieAffinityGroupWithEmptyPolicies() throws Exception { |
| 246 | + // 1. create namespace with empty policies |
| 247 | + String setBookieAffinityGroupNs = "test-set-bookie-affinity-group-ns"; |
| 248 | + asyncRequests(response -> namespaces.createNamespace(response, testTenant, setBookieAffinityGroupNs, null)); |
| 249 | + |
| 250 | + // 2.set bookie affinity group |
| 251 | + String primaryAffinityGroup = "primary-affinity-group"; |
| 252 | + String secondaryAffinityGroup = "secondary-affinity-group"; |
| 253 | + BookieAffinityGroupData bookieAffinityGroupDataReq = |
| 254 | + BookieAffinityGroupData.builder().bookkeeperAffinityGroupPrimary(primaryAffinityGroup) |
| 255 | + .bookkeeperAffinityGroupSecondary(secondaryAffinityGroup).build(); |
| 256 | + namespaces.setBookieAffinityGroup(testTenant, setBookieAffinityGroupNs, bookieAffinityGroupDataReq); |
| 257 | + |
| 258 | + // 3.query namespace num bundles, should be conf.getDefaultNumberOfNamespaceBundles() |
| 259 | + BundlesData bundlesData = (BundlesData) asyncRequests( |
| 260 | + response -> namespaces.getBundlesData(response, testTenant, setBookieAffinityGroupNs)); |
| 261 | + assertEquals(bundlesData.getNumBundles(), conf.getDefaultNumberOfNamespaceBundles()); |
| 262 | + |
| 263 | + // 4.assert namespace bookie affinity group |
| 264 | + BookieAffinityGroupData bookieAffinityGroupDataResp = |
| 265 | + namespaces.getBookieAffinityGroup(testTenant, setBookieAffinityGroupNs); |
| 266 | + assertEquals(bookieAffinityGroupDataResp, bookieAffinityGroupDataReq); |
| 267 | + } |
| 268 | + |
| 269 | + @Test |
| 270 | + public void testSetBookieAffinityGroupWithExistBundlePolicies() throws Exception { |
| 271 | + // 1. create namespace with specified num bundles |
| 272 | + String setBookieAffinityGroupNs = "test-set-bookie-affinity-group-ns"; |
| 273 | + Policies policies = new Policies(); |
| 274 | + policies.bundles = getBundles(10); |
| 275 | + asyncRequests(response -> namespaces.createNamespace(response, testTenant, setBookieAffinityGroupNs, policies)); |
| 276 | + |
| 277 | + // 2.set bookie affinity group |
| 278 | + String primaryAffinityGroup = "primary-affinity-group"; |
| 279 | + String secondaryAffinityGroup = "secondary-affinity-group"; |
| 280 | + BookieAffinityGroupData bookieAffinityGroupDataReq = |
| 281 | + BookieAffinityGroupData.builder().bookkeeperAffinityGroupPrimary(primaryAffinityGroup) |
| 282 | + .bookkeeperAffinityGroupSecondary(secondaryAffinityGroup).build(); |
| 283 | + namespaces.setBookieAffinityGroup(testTenant, setBookieAffinityGroupNs, bookieAffinityGroupDataReq); |
| 284 | + |
| 285 | + // 3.query namespace num bundles, should be policies.bundles, which we set before |
| 286 | + BundlesData bundlesData = (BundlesData) asyncRequests( |
| 287 | + response -> namespaces.getBundlesData(response, testTenant, setBookieAffinityGroupNs)); |
| 288 | + assertEquals(bundlesData, policies.bundles); |
| 289 | + |
| 290 | + // 4.assert namespace bookie affinity group |
| 291 | + BookieAffinityGroupData bookieAffinityGroupDataResp = |
| 292 | + namespaces.getBookieAffinityGroup(testTenant, setBookieAffinityGroupNs); |
| 293 | + assertEquals(bookieAffinityGroupDataResp, bookieAffinityGroupDataReq); |
| 294 | + } |
| 295 | + |
| 296 | + @Test |
| 297 | + public void testSetNamespaceAntiAffinityGroupWithEmptyPolicies() throws Exception { |
| 298 | + // 1. create namespace with empty policies |
| 299 | + String setNamespaceAntiAffinityGroupNs = "test-set-namespace-anti-affinity-group-ns"; |
| 300 | + asyncRequests( |
| 301 | + response -> namespaces.createNamespace(response, testTenant, setNamespaceAntiAffinityGroupNs, null)); |
| 302 | + |
| 303 | + // 2.set namespace anti affinity group |
| 304 | + String namespaceAntiAffinityGroupReq = "namespace-anti-affinity-group"; |
| 305 | + namespaces.setNamespaceAntiAffinityGroup(testTenant, setNamespaceAntiAffinityGroupNs, |
| 306 | + namespaceAntiAffinityGroupReq); |
| 307 | + |
| 308 | + // 3.query namespace num bundles, should be conf.getDefaultNumberOfNamespaceBundles() |
| 309 | + BundlesData bundlesData = (BundlesData) asyncRequests( |
| 310 | + response -> namespaces.getBundlesData(response, testTenant, setNamespaceAntiAffinityGroupNs)); |
| 311 | + assertEquals(bundlesData.getNumBundles(), conf.getDefaultNumberOfNamespaceBundles()); |
| 312 | + |
| 313 | + // 4.assert namespace anti affinity group |
| 314 | + String namespaceAntiAffinityGroupResp = |
| 315 | + namespaces.getNamespaceAntiAffinityGroup(testTenant, setNamespaceAntiAffinityGroupNs); |
| 316 | + assertEquals(namespaceAntiAffinityGroupResp, namespaceAntiAffinityGroupReq); |
| 317 | + } |
| 318 | + |
| 319 | + @Test |
| 320 | + public void testSetNamespaceAntiAffinityGroupWithExistBundlePolicies() throws Exception { |
| 321 | + // 1. create namespace with specified num bundles |
| 322 | + String setNamespaceAntiAffinityGroupNs = "test-set-namespace-anti-affinity-group-ns"; |
| 323 | + Policies policies = new Policies(); |
| 324 | + policies.bundles = getBundles(10); |
| 325 | + asyncRequests(response -> namespaces.createNamespace(response, testTenant, setNamespaceAntiAffinityGroupNs, |
| 326 | + policies)); |
| 327 | + |
| 328 | + // 2.set namespace anti affinity group |
| 329 | + String namespaceAntiAffinityGroupReq = "namespace-anti-affinity-group"; |
| 330 | + namespaces.setNamespaceAntiAffinityGroup(testTenant, setNamespaceAntiAffinityGroupNs, |
| 331 | + namespaceAntiAffinityGroupReq); |
| 332 | + |
| 333 | + // 3.query namespace num bundles, should be policies.bundles, which we set before |
| 334 | + BundlesData bundlesData = (BundlesData) asyncRequests( |
| 335 | + response -> namespaces.getBundlesData(response, testTenant, setNamespaceAntiAffinityGroupNs)); |
| 336 | + assertEquals(bundlesData, policies.bundles); |
| 337 | + |
| 338 | + // 4.assert namespace anti affinity group |
| 339 | + String namespaceAntiAffinityGroupResp = |
| 340 | + namespaces.getNamespaceAntiAffinityGroup(testTenant, setNamespaceAntiAffinityGroupNs); |
| 341 | + assertEquals(namespaceAntiAffinityGroupResp, namespaceAntiAffinityGroupReq); |
| 342 | + } |
| 343 | + |
| 344 | + @Test |
| 345 | + public void testEnableMigrationWithEmptyPolicies() throws Exception { |
| 346 | + // 1. create namespace with empty policies |
| 347 | + String enableMigrationGroupNs = "test-set-namespace-enable-migration-ns"; |
| 348 | + asyncRequests(response -> namespaces.createNamespace(response, testTenant, enableMigrationGroupNs, null)); |
| 349 | + |
| 350 | + // 2.set enable migration |
| 351 | + boolean enableMigrationReq = true; |
| 352 | + namespaces.enableMigration(testTenant, enableMigrationGroupNs, enableMigrationReq); |
| 353 | + |
| 354 | + // 3.query namespace num bundles, should be conf.getDefaultNumberOfNamespaceBundles() |
| 355 | + BundlesData bundlesData = (BundlesData) asyncRequests( |
| 356 | + response -> namespaces.getBundlesData(response, testTenant, enableMigrationGroupNs)); |
| 357 | + assertEquals(bundlesData.getNumBundles(), conf.getDefaultNumberOfNamespaceBundles()); |
| 358 | + |
| 359 | + // 4.assert namespace enable migration |
| 360 | + Policies policiesResp = (Policies) asyncRequests( |
| 361 | + response -> namespaces.getPolicies(response, testTenant, enableMigrationGroupNs)); |
| 362 | + assertEquals(policiesResp.migrated, enableMigrationReq); |
| 363 | + } |
| 364 | + |
| 365 | + @Test |
| 366 | + public void testEnableMigrationWithExistBundlePolicies() throws Exception { |
| 367 | + // 1. create namespace with specified num bundles |
| 368 | + String enableMigrationGroupNs = "test-set-namespace-enable-migration-ns"; |
| 369 | + Policies policiesReq = new Policies(); |
| 370 | + policiesReq.bundles = getBundles(10); |
| 371 | + asyncRequests( |
| 372 | + response -> namespaces.createNamespace(response, testTenant, enableMigrationGroupNs, policiesReq)); |
| 373 | + |
| 374 | + // 2.set enable migration |
| 375 | + boolean enableMigrationReq = true; |
| 376 | + namespaces.enableMigration(testTenant, enableMigrationGroupNs, enableMigrationReq); |
| 377 | + |
| 378 | + // 3.query namespace num bundles, should be policies.bundles, which we set before |
| 379 | + BundlesData bundlesData = (BundlesData) asyncRequests( |
| 380 | + response -> namespaces.getBundlesData(response, testTenant, enableMigrationGroupNs)); |
| 381 | + assertEquals(bundlesData, policiesReq.bundles); |
| 382 | + |
| 383 | + // 4.assert namespace enable migration |
| 384 | + Policies policiesResp = (Policies) asyncRequests( |
| 385 | + response -> namespaces.getPolicies(response, testTenant, enableMigrationGroupNs)); |
| 386 | + assertEquals(policiesResp.migrated, enableMigrationReq); |
| 387 | + } |
239 | 388 | } |
0 commit comments