This is an automated email from the ASF dual-hosted git repository.

Jackie-Jiang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pinot.git


The following commit(s) were added to refs/heads/master by this push:
     new 4cfedb07f85 Make BrokerRoutingManagerConcurrencyTest deterministic 
(#19356)
4cfedb07f85 is described below

commit 4cfedb07f851ef7bec019872c5c472766f4a46f0
Author: Xiaotian (Jackie) Jiang <[email protected]>
AuthorDate: Tue Aug 25 11:14:03 2026 -0700

    Make BrokerRoutingManagerConcurrencyTest deterministic (#19356)
---
 .../BrokerRoutingManagerConcurrencyTest.java       | 681 ++++++++++-----------
 1 file changed, 322 insertions(+), 359 deletions(-)

diff --git 
a/pinot-broker/src/test/java/org/apache/pinot/broker/routing/manager/BrokerRoutingManagerConcurrencyTest.java
 
b/pinot-broker/src/test/java/org/apache/pinot/broker/routing/manager/BrokerRoutingManagerConcurrencyTest.java
index 174c0f2c79f..dfa1c51531a 100644
--- 
a/pinot-broker/src/test/java/org/apache/pinot/broker/routing/manager/BrokerRoutingManagerConcurrencyTest.java
+++ 
b/pinot-broker/src/test/java/org/apache/pinot/broker/routing/manager/BrokerRoutingManagerConcurrencyTest.java
@@ -21,16 +21,19 @@ package org.apache.pinot.broker.routing.manager;
 import java.lang.reflect.Field;
 import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.Collections;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
-import java.util.concurrent.Future;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicReference;
+import java.util.concurrent.locks.ReadWriteLock;
+import java.util.function.BiConsumer;
 import org.apache.helix.HelixConstants.ChangeType;
 import org.apache.helix.model.ExternalView;
 import org.apache.helix.model.IdealState;
@@ -58,7 +61,6 @@ import org.apache.pinot.spi.utils.builder.TableConfigBuilder;
 import org.apache.pinot.spi.utils.builder.TableNameBuilder;
 import org.mockito.Mock;
 import org.mockito.Mockito;
-import org.testng.Assert;
 import org.testng.annotations.AfterClass;
 import org.testng.annotations.BeforeClass;
 import org.testng.annotations.Test;
@@ -66,6 +68,7 @@ import org.testng.annotations.Test;
 import static org.mockito.ArgumentMatchers.anyInt;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.when;
+import static org.testng.Assert.*;
 
 /// Test class to validate concurrency and race condition handling in 
BrokerRoutingManager,
 /// specifically focusing on TimeBoundaryManager coordination for hybrid 
tables.
@@ -88,7 +91,8 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
   private BrokerRoutingManager _routingManager;
 
   @BeforeClass
-  public void setUp() throws Exception {
+  public void setUp()
+      throws Exception {
     // Start ZooKeeper and initialize the test infrastructure
     startZk();
     startController();
@@ -102,7 +106,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
     
Mockito.when(_pinotConfig.getProperty(Mockito.eq("pinot.broker.adaptive.server.selector.type")))
         .thenReturn("UNIFORM_RANDOM");
     Mockito.when(_pinotConfig.getProperty(
-        
Mockito.eq(CommonConstants.Broker.CONFIG_OF_ROUTING_ASSIGNMENT_CHANGE_PROCESS_PARALLELISM),
 anyInt()))
+            
Mockito.eq(CommonConstants.Broker.CONFIG_OF_ROUTING_ASSIGNMENT_CHANGE_PROCESS_PARALLELISM),
 anyInt()))
         .thenReturn(10);
     Mockito.when(_pinotConfig.getProperty(Mockito.anyString(), 
Mockito.anyString()))
         .thenAnswer(invocation -> invocation.getArgument(1)); // Return 
default value
@@ -168,7 +172,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
     try {
       _routingManager.processClusterChange(ChangeType.INSTANCE_CONFIG);
     } catch (Exception e) {
-      Assert.fail("Direct call to processClusterChange failed", e);
+      fail("Direct call to processClusterChange failed", e);
     }
   }
 
@@ -182,7 +186,18 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       Map<?, ?> routingEntryMap = (Map<?, ?>) 
routingEntryMapField.get(_routingManager);
       routingEntryMap.clear();
     } catch (Exception e) {
-      Assert.fail("Failed to clear routing entries", e);
+      fail("Failed to clear routing entries", e);
+    }
+  }
+
+  private ReadWriteLock getGlobalLock() {
+    try {
+      // Access the private lock from the parent class 
(BaseBrokerRoutingManager)
+      Field globalLockField = 
BaseBrokerRoutingManager.class.getDeclaredField("_globalLock");
+      globalLockField.setAccessible(true);
+      return (ReadWriteLock) globalLockField.get(_routingManager);
+    } catch (Exception e) {
+      throw new IllegalStateException("Failed to access the global lock", e);
     }
   }
 
@@ -247,13 +262,13 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
   private void validateDisabledInstanceNotInRouting(String tableNameWithType, 
String disabledInstance) {
     try {
       Object routingEntry = getRoutingEntry(tableNameWithType);
-      Assert.assertNotNull(routingEntry, "Routing entry should exist for 
table: " + tableNameWithType);
+      assertNotNull(routingEntry, "Routing entry should exist for table: " + 
tableNameWithType);
 
       // Get the InstanceSelector from the routing entry
       java.lang.reflect.Field instanceSelectorField = 
routingEntry.getClass().getDeclaredField("_instanceSelector");
       instanceSelectorField.setAccessible(true);
       Object instanceSelector = instanceSelectorField.get(routingEntry);
-      Assert.assertNotNull(instanceSelector, "InstanceSelector should exist");
+      assertNotNull(instanceSelector, "InstanceSelector should exist");
 
       // Get the _enabledInstances field from BaseInstanceSelector
       java.lang.reflect.Field enabledInstancesField =
@@ -262,28 +277,26 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       Object enabledInstancesObj = enabledInstancesField.get(instanceSelector);
 
       if (enabledInstancesObj != null) {
-        @SuppressWarnings("unchecked")
-        java.util.Set<String> enabledInstances = (java.util.Set<String>) 
enabledInstancesObj;
-        Assert.assertFalse(enabledInstances.contains(disabledInstance),
+        Set<?> enabledInstances = (Set<?>) enabledInstancesObj;
+        assertFalse(enabledInstances.contains(disabledInstance),
             "Disabled instance " + disabledInstance + " should NOT be in 
enabled instances for table "
                 + tableNameWithType + ". Enabled instances: " + 
enabledInstances);
       }
     } catch (Exception e) {
-      Assert.fail("Failed to validate disabled instance exclusion for table " 
+ tableNameWithType + ": "
-          + e.getMessage());
+      fail("Failed to validate disabled instance exclusion for table " + 
tableNameWithType + ": " + e.getMessage());
     }
   }
 
   private void validateEnabledInstanceInRouting(String tableNameWithType, 
String enabledInstance) {
     try {
       Object routingEntry = getRoutingEntry(tableNameWithType);
-      Assert.assertNotNull(routingEntry, "Routing entry should exist for 
table: " + tableNameWithType);
+      assertNotNull(routingEntry, "Routing entry should exist for table: " + 
tableNameWithType);
 
       // Get the InstanceSelector from the routing entry
       java.lang.reflect.Field instanceSelectorField = 
routingEntry.getClass().getDeclaredField("_instanceSelector");
       instanceSelectorField.setAccessible(true);
       Object instanceSelector = instanceSelectorField.get(routingEntry);
-      Assert.assertNotNull(instanceSelector, "InstanceSelector should exist");
+      assertNotNull(instanceSelector, "InstanceSelector should exist");
 
       // Get the _enabledInstances field from BaseInstanceSelector
       java.lang.reflect.Field enabledInstancesField =
@@ -291,15 +304,13 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       enabledInstancesField.setAccessible(true);
       Object enabledInstancesObj = enabledInstancesField.get(instanceSelector);
 
-      Assert.assertNotNull(enabledInstancesObj, "Enabled instances should not 
be null for table " + tableNameWithType);
-      @SuppressWarnings("unchecked")
-      java.util.Set<String> enabledInstances = (java.util.Set<String>) 
enabledInstancesObj;
-      Assert.assertTrue(enabledInstances.contains(enabledInstance),
-          "Enabled instance " + enabledInstance + " should be in enabled 
instances for table "
-              + tableNameWithType + ". Enabled instances: " + 
enabledInstances);
+      assertNotNull(enabledInstancesObj, "Enabled instances should not be null 
for table " + tableNameWithType);
+      Set<?> enabledInstances = (Set<?>) enabledInstancesObj;
+      assertTrue(enabledInstances.contains(enabledInstance),
+          "Enabled instance " + enabledInstance + " should be in enabled 
instances for table " + tableNameWithType
+              + ". Enabled instances: " + enabledInstances);
     } catch (Exception e) {
-      Assert.fail("Failed to validate enabled instance inclusion for table " + 
tableNameWithType + ": "
-          + e.getMessage());
+      fail("Failed to validate enabled instance inclusion for table " + 
tableNameWithType + ": " + e.getMessage());
     }
   }
 
@@ -321,7 +332,8 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
   /// 2. TimeBoundaryManager is properly set on the offline table when 
realtime is built
   /// 3. No race conditions occur that would leave the offline table without a 
TimeBoundaryManager
   @Test
-  public void testConcurrentHybridTableBuildNoTimeBoundaryManagerRace() throws 
Exception {
+  public void testConcurrentHybridTableBuildNoTimeBoundaryManagerRace()
+      throws Exception {
     // Clean any existing routing entries to ensure test isolation
     clearRoutingEntries();
 
@@ -335,8 +347,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
 
     try {
       // Build OFFLINE table in thread 1
-      @SuppressWarnings("unused")
-      Future<?> offlineTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           _routingManager.buildRouting(OFFLINE_TABLE_NAME);
@@ -348,8 +359,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Build REALTIME table in thread 2
-      @SuppressWarnings("unused")
-      Future<?> realtimeTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           _routingManager.buildRouting(REALTIME_TABLE_NAME);
@@ -364,32 +374,32 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       startLatch.countDown();
 
       // Wait for both threads to complete
-      Assert.assertTrue(finishLatch.await(30, TimeUnit.SECONDS), "Threads 
didn't complete in time");
+      assertTrue(finishLatch.await(30, TimeUnit.SECONDS), "Threads didn't 
complete in time");
 
       // Check if any thread failed
       if (offlineException.get() != null) {
-        Assert.fail("Offline table build failed", offlineException.get());
+        fail("Offline table build failed", offlineException.get());
       }
       if (realtimeException.get() != null) {
-        Assert.fail("Realtime table build failed", realtimeException.get());
+        fail("Realtime table build failed", realtimeException.get());
       }
 
       // Verify both tables exist
-      Assert.assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME), 
"Offline table routing should exist");
-      Assert.assertTrue(_routingManager.routingExists(REALTIME_TABLE_NAME), 
"Realtime table routing should exist");
+      assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME), "Offline 
table routing should exist");
+      assertTrue(_routingManager.routingExists(REALTIME_TABLE_NAME), "Realtime 
table routing should exist");
 
       // CRITICAL VERIFICATION: Check TimeBoundaryManager coordination
       // The offline table should have a TimeBoundaryManager when realtime 
table exists
       Object offlineEntry = getRoutingEntry(OFFLINE_TABLE_NAME);
-      Assert.assertNotNull(offlineEntry, "Offline routing entry should exist");
+      assertNotNull(offlineEntry, "Offline routing entry should exist");
 
       // If realtime table was built, offline should have TimeBoundaryManager
       TimeBoundaryManager timeBoundaryManager = 
getTimeBoundaryManager(offlineEntry);
-      Assert.assertNotNull(timeBoundaryManager, "Offline table should have 
TimeBoundaryManager when realtime table "
-          + "exists - this indicates a race condition in cross-table 
coordination");
+      assertNotNull(timeBoundaryManager, "Offline table should have 
TimeBoundaryManager when realtime table exists - "
+          + "this indicates a race condition in cross-table coordination");
     } finally {
       executor.shutdown();
-      Assert.assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), 
"Executor didn't shutdown in time");
+      assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), "Executor 
didn't shutdown in time");
     }
   }
 
@@ -401,21 +411,20 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
 
     // Step 1: Build OFFLINE table first - should not have TimeBoundaryManager
     _routingManager.buildRouting(OFFLINE_TABLE_NAME);
-    Assert.assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME), 
"Offline table routing should exist");
+    assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME), "Offline 
table routing should exist");
 
     Object offlineEntry = getRoutingEntry(OFFLINE_TABLE_NAME);
     TimeBoundaryManager timeBoundaryManager = 
getTimeBoundaryManager(offlineEntry);
-    Assert.assertNull(timeBoundaryManager,
-        "Offline table should not have TimeBoundaryManager when realtime 
doesn't exist");
+    assertNull(timeBoundaryManager, "Offline table should not have 
TimeBoundaryManager when realtime doesn't exist");
 
     // Step 2: Build REALTIME table - should add TimeBoundaryManager to 
existing offline table
     _routingManager.buildRouting(REALTIME_TABLE_NAME);
-    Assert.assertTrue(_routingManager.routingExists(REALTIME_TABLE_NAME));
+    assertTrue(_routingManager.routingExists(REALTIME_TABLE_NAME));
 
     // Verify TimeBoundaryManager was added to offline table
     offlineEntry = getRoutingEntry(OFFLINE_TABLE_NAME);
     timeBoundaryManager = getTimeBoundaryManager(offlineEntry);
-    Assert.assertNotNull(timeBoundaryManager, "Offline table should have 
TimeBoundaryManager after realtime is built");
+    assertNotNull(timeBoundaryManager, "Offline table should have 
TimeBoundaryManager after realtime is built");
   }
 
   @Test
@@ -436,6 +445,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
     Class<?> baseClass = manager.getClass().getSuperclass();
     Field startTimesField = 
baseClass.getDeclaredField("_routingTableBuildStartTimeMs");
     startTimesField.setAccessible(true);
+    //noinspection unchecked
     Map<String, Long> startTimes = (Map<String, Long>) 
startTimesField.get(manager);
     if (startTimes == null) {
       startTimes = new ConcurrentHashMap<>();
@@ -447,23 +457,24 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
     manager.buildRouting(tableNameWithType);
 
     // Ensure routing was not created and the last start time was not 
overwritten
-    Assert.assertFalse(manager.routingExists(tableNameWithType));
-    Assert.assertEquals(startTimes.get(tableNameWithType).longValue(), 
futureStart);
+    assertFalse(manager.routingExists(tableNameWithType));
+    assertEquals(startTimes.get(tableNameWithType).longValue(), futureStart);
   }
 
   /// Test concurrent interactions between processSegmentAssignmentChange and 
buildRouting.
   /// This validates that the global read lock (for 
processSegmentAssignmentChange) and
   /// per-table locks (for buildRouting) work correctly together without 
deadlocks.
   @Test
-  public void testConcurrentProcessSegmentAssignmentChangeAndBuildRouting() 
throws Exception {
+  public void testConcurrentProcessSegmentAssignmentChangeAndBuildRouting()
+      throws Exception {
     clearRoutingEntries();
 
     // First, build initial routing entries for both tables
     _routingManager.buildRouting(OFFLINE_TABLE_NAME);
     _routingManager.buildRouting(REALTIME_TABLE_NAME);
 
-    Assert.assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME), 
"Initial offline routing should exist");
-    Assert.assertTrue(_routingManager.routingExists(REALTIME_TABLE_NAME), 
"Initial realtime routing should exist");
+    assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME), "Initial 
offline routing should exist");
+    assertTrue(_routingManager.routingExists(REALTIME_TABLE_NAME), "Initial 
realtime routing should exist");
     ExecutorService executor = Executors.newFixedThreadPool(3);
     CountDownLatch startLatch = new CountDownLatch(1);
     CountDownLatch finishLatch = new CountDownLatch(3);
@@ -475,7 +486,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
     try {
       // Thread 1: Process segment assignment change (takes global read lock, 
per-raw-table-name lock for each table one
       // at a time)
-      Future<?> segmentAssignmentTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           // This should take global read lock
@@ -488,7 +499,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 2: Build routing for offline table (takes per-table lock)
-      Future<?> buildOfflineTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           // This should take per-raw-table-name lock for OFFLINE table
@@ -501,7 +512,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 3: Build routing for realtime table (takes same per-table lock)
-      Future<?> buildRealtimeTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           // This should take per-raw-table-name lock for REALTIME table
@@ -517,37 +528,37 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       startLatch.countDown();
 
       // Wait for completion with timeout
-      Assert.assertTrue(finishLatch.await(10, TimeUnit.SECONDS), "All tasks 
should complete within timeout");
+      assertTrue(finishLatch.await(10, TimeUnit.SECONDS), "All tasks should 
complete within timeout");
 
       // Verify no exceptions occurred
       if (segmentAssignmentException.get() != null) {
-        Assert.fail("Segment assignment change failed: " + 
segmentAssignmentException.get().getMessage());
+        fail("Segment assignment change failed: " + 
segmentAssignmentException.get().getMessage());
       }
       if (buildRoutingOfflineException.get() != null) {
-        Assert.fail("Offline table build failed: " + 
buildRoutingOfflineException.get().getMessage());
+        fail("Offline table build failed: " + 
buildRoutingOfflineException.get().getMessage());
       }
       if (buildRoutingRealtimeException.get() != null) {
-        Assert.fail("Realtime table build failed: " + 
buildRoutingRealtimeException.get().getMessage());
+        fail("Realtime table build failed: " + 
buildRoutingRealtimeException.get().getMessage());
       }
 
       // Verify routing entries still exist after concurrent operations
-      Assert.assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME),
+      assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME),
           "Offline routing should exist after concurrent operations");
-      Assert.assertTrue(_routingManager.routingExists(REALTIME_TABLE_NAME),
+      assertTrue(_routingManager.routingExists(REALTIME_TABLE_NAME),
           "Realtime routing should exist after concurrent operations");
 
       // CRITICAL VERIFICATION: Check TimeBoundaryManager coordination
       // The offline table should have a TimeBoundaryManager when realtime 
table exists
       Object offlineEntry = getRoutingEntry(OFFLINE_TABLE_NAME);
-      Assert.assertNotNull(offlineEntry, "Offline routing entry should exist");
+      assertNotNull(offlineEntry, "Offline routing entry should exist");
 
       // If realtime table was built, offline should have TimeBoundaryManager
       TimeBoundaryManager timeBoundaryManager = 
getTimeBoundaryManager(offlineEntry);
-      Assert.assertNotNull(timeBoundaryManager, "Offline table should have 
TimeBoundaryManager when realtime table "
-          + "exists - this indicates a race condition in cross-table 
coordination");
+      assertNotNull(timeBoundaryManager, "Offline table should have 
TimeBoundaryManager when realtime table exists - "
+          + "this indicates a race condition in cross-table coordination");
     } finally {
       executor.shutdown();
-      Assert.assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), 
"Executor didn't shutdown in time");
+      assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), "Executor 
didn't shutdown in time");
     }
   }
 
@@ -557,7 +568,8 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
   /// per-table locks (for buildRouting) work correctly together, especially 
when buildRouting creates new routing
   /// entries.
   @Test
-  public void 
testConcurrentProcessInstanceConfigChangeAndBuildRoutingNewTable() throws 
Exception {
+  public void 
testConcurrentProcessInstanceConfigChangeAndBuildRoutingNewTable()
+      throws Exception {
     clearRoutingEntries();
 
     // Add additional server instances for this test
@@ -593,10 +605,8 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
     ZKMetadataProvider.setSchema(_propertyStore, newSchema);
 
     // Create IdealState and ExternalView for the new tables with both servers
-    createIdealStateAndExternalViewWithMultipleServers(newOfflineTable, 
enabledServerInstance,
-        disabledServerInstance);
-    createIdealStateAndExternalViewWithMultipleServers(newRealtimeTable, 
enabledServerInstance,
-        disabledServerInstance);
+    createIdealStateAndExternalViewWithMultipleServers(newOfflineTable, 
enabledServerInstance, disabledServerInstance);
+    createIdealStateAndExternalViewWithMultipleServers(newRealtimeTable, 
enabledServerInstance, disabledServerInstance);
 
     // Create segment metadata
     createSegmentMetadata(newOfflineTable, "newSegment_0", 
System.currentTimeMillis());
@@ -616,7 +626,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
     try {
       // Thread 1: Process instance config change (takes global write lock and 
per-raw-table-name locks for each table
       // one at a time)
-      Future<?> instanceConfigTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           // This should take global write lock and process the disabled 
instance
@@ -629,7 +639,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 2: Build routing for new offline table (takes per-table lock 
and adds new entry)
-      Future<?> buildNewOfflineTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           // This should take per-table lock and create new routing entry
@@ -642,7 +652,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 3: Build routing for new realtime table (takes per-table lock 
and adds new entry)
-      Future<?> buildNewRealtimeTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           // This should take per-table lock and create new routing entry
@@ -658,44 +668,42 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       startLatch.countDown();
 
       // Wait for completion with timeout
-      Assert.assertTrue(finishLatch.await(15, TimeUnit.SECONDS), "All tasks 
should complete within timeout");
+      assertTrue(finishLatch.await(15, TimeUnit.SECONDS), "All tasks should 
complete within timeout");
 
       // Verify no exceptions occurred
       if (instanceConfigException.get() != null) {
-        Assert.fail("Instance config change failed: " + 
instanceConfigException.get().getMessage());
+        fail("Instance config change failed: " + 
instanceConfigException.get().getMessage());
       }
       if (buildNewOfflineException.get() != null) {
-        Assert.fail("New offline table build failed: " + 
buildNewOfflineException.get().getMessage());
+        fail("New offline table build failed: " + 
buildNewOfflineException.get().getMessage());
       }
       if (buildNewRealtimeException.get() != null) {
-        Assert.fail("New realtime table build failed: " + 
buildNewRealtimeException.get().getMessage());
+        fail("New realtime table build failed: " + 
buildNewRealtimeException.get().getMessage());
       }
 
       // Verify new routing entries were created successfully
-      Assert.assertTrue(_routingManager.routingExists(newOfflineTable),
+      assertTrue(_routingManager.routingExists(newOfflineTable),
           "New offline routing should exist after concurrent operations");
-      Assert.assertTrue(_routingManager.routingExists(newRealtimeTable),
+      assertTrue(_routingManager.routingExists(newRealtimeTable),
           "New realtime routing should exist after concurrent operations");
 
       // Verify TimeBoundaryManager coordination for the new hybrid table
       Object newOfflineEntry = getRoutingEntry(newOfflineTable);
-      Assert.assertNotNull(newOfflineEntry, "New offline routing entry should 
exist");
+      assertNotNull(newOfflineEntry, "New offline routing entry should exist");
       TimeBoundaryManager timeBoundaryManager = 
getTimeBoundaryManager(newOfflineEntry);
-      Assert.assertNotNull(timeBoundaryManager,
-          "New offline table should have TimeBoundaryManager when realtime 
exists");
+      assertNotNull(timeBoundaryManager, "New offline table should have 
TimeBoundaryManager when realtime exists");
 
       // CRITICAL: Verify that the disabled instance is NOT included in the 
routing entries
       validateDisabledInstanceNotInRouting(newOfflineTable, 
disabledServerInstance);
       validateDisabledInstanceNotInRouting(newRealtimeTable, 
disabledServerInstance);
     } finally {
       executor.shutdown();
-      Assert.assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), 
"Executor didn't shutdown in time");
+      assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), "Executor 
didn't shutdown in time");
     }
   }
 
   private TableConfig createTableConfig(String tableNameWithType, TableType 
tableType) {
-    return new TableConfigBuilder(tableType)
-        .setTableName(TableNameBuilder.extractRawTableName(tableNameWithType))
+    return new 
TableConfigBuilder(tableType).setTableName(TableNameBuilder.extractRawTableName(tableNameWithType))
         .setTimeColumnName("timestamp")
         .build();
   }
@@ -713,8 +721,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       Class<?> baseClass = _routingManager.getClass().getSuperclass();
       java.lang.reflect.Field field = 
baseClass.getDeclaredField("_routingEntryMap");
       field.setAccessible(true);
-      @SuppressWarnings("unchecked")
-      java.util.Map<String, Object> routingEntryMap = (java.util.Map<String, 
Object>) field.get(_routingManager);
+      Map<?, ?> routingEntryMap = (Map<?, ?>) field.get(_routingManager);
       return routingEntryMap.get(tableNameWithType);
     } catch (Exception e) {
       throw new RuntimeException("Failed to access routing entry", e);
@@ -738,7 +745,8 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
   /// (global read lock + per-table lock).
   /// This validates that global write lock properly blocks global read lock 
operations.
   @Test
-  public void testConcurrentExcludeServerAndBuildRouting() throws Exception {
+  public void testConcurrentExcludeServerAndBuildRouting()
+      throws Exception {
     clearRoutingEntries();
 
     String disabledServerInstance = "Server_localhost_8001"; // We'll disable 
this one
@@ -747,8 +755,8 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
     _routingManager.buildRouting(OFFLINE_TABLE_NAME);
     _routingManager.buildRouting(REALTIME_TABLE_NAME);
 
-    Assert.assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME), 
"Initial offline routing should exist");
-    Assert.assertTrue(_routingManager.routingExists(REALTIME_TABLE_NAME), 
"Initial realtime routing should exist");
+    assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME), "Initial 
offline routing should exist");
+    assertTrue(_routingManager.routingExists(REALTIME_TABLE_NAME), "Initial 
realtime routing should exist");
 
     ExecutorService executor = Executors.newFixedThreadPool(3);
     CountDownLatch startLatch = new CountDownLatch(1);
@@ -764,7 +772,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
 
     try {
       // Thread 1: Exclude server from routing (takes global write lock)
-      Future<?> excludeServerTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           // This should take global write lock
@@ -777,7 +785,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 2: Build routing for offline table (takes global read lock + 
per-table lock)
-      Future<?> buildOfflineTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           // This should take global read lock + per-table lock
@@ -790,7 +798,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 3: Build routing for realtime table (takes global read lock + 
per-table lock)
-      Future<?> buildRealtimeTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           // This should take global read lock + different per-table lock
@@ -806,41 +814,41 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       startLatch.countDown();
 
       // Wait for completion with timeout
-      Assert.assertTrue(finishLatch.await(10, TimeUnit.SECONDS), "All tasks 
should complete within timeout");
+      assertTrue(finishLatch.await(10, TimeUnit.SECONDS), "All tasks should 
complete within timeout");
 
       // Verify no exceptions occurred
       if (excludeServerException.get() != null) {
-        Assert.fail("Exclude server failed: " + 
excludeServerException.get().getMessage());
+        fail("Exclude server failed: " + 
excludeServerException.get().getMessage());
       }
       if (buildOfflineException.get() != null) {
-        Assert.fail("Build offline routing failed: " + 
buildOfflineException.get().getMessage());
+        fail("Build offline routing failed: " + 
buildOfflineException.get().getMessage());
       }
       if (buildRealtimeException.get() != null) {
-        Assert.fail("Build realtime routing failed: " + 
buildRealtimeException.get().getMessage());
+        fail("Build realtime routing failed: " + 
buildRealtimeException.get().getMessage());
       }
 
       // Verify routing entries still exist after operations
-      Assert.assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME),
+      assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME),
           "Offline routing should exist after exclude server operation");
-      Assert.assertTrue(_routingManager.routingExists(REALTIME_TABLE_NAME),
+      assertTrue(_routingManager.routingExists(REALTIME_TABLE_NAME),
           "Realtime routing should exist after exclude server operation");
 
       // CRITICAL VERIFICATION: Check TimeBoundaryManager coordination
       // The offline table should have a TimeBoundaryManager when realtime 
table exists
       Object offlineEntry = getRoutingEntry(OFFLINE_TABLE_NAME);
-      Assert.assertNotNull(offlineEntry, "Offline routing entry should exist");
+      assertNotNull(offlineEntry, "Offline routing entry should exist");
 
       // If realtime table was built, offline should have TimeBoundaryManager
       TimeBoundaryManager timeBoundaryManager = 
getTimeBoundaryManager(offlineEntry);
-      Assert.assertNotNull(timeBoundaryManager, "Offline table should have 
TimeBoundaryManager when realtime table "
-          + "exists - this indicates a race condition in cross-table 
coordination");
+      assertNotNull(timeBoundaryManager, "Offline table should have 
TimeBoundaryManager when realtime table exists - "
+          + "this indicates a race condition in cross-table coordination");
 
       // CRITICAL: Verify that the disabled instance is NOT included in the 
routing entries
       validateDisabledInstanceNotInRouting(OFFLINE_TABLE_NAME, 
disabledServerInstance);
       validateDisabledInstanceNotInRouting(REALTIME_TABLE_NAME, 
disabledServerInstance);
     } finally {
       executor.shutdown();
-      Assert.assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), 
"Executor didn't shutdown in time");
+      assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), "Executor 
didn't shutdown in time");
     }
   }
 
@@ -848,7 +856,8 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
   /// (global read lock + per-table lock).
   /// This validates proper coordination between global write operations and 
segment refresh operations.
   @Test
-  public void testConcurrentIncludeServerAndRefreshSegment() throws Exception {
+  public void testConcurrentIncludeServerAndRefreshSegment()
+      throws Exception {
     clearRoutingEntries();
 
     String includedServerInstance = "Server_localhost_8001"; // We'll include 
this one
@@ -874,7 +883,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
 
     try {
       // Thread 1: Include server to routing (takes global write lock)
-      Future<?> includeServerTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           // This should take global write lock
@@ -887,7 +896,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 2: Refresh segment for offline table (takes global read lock + 
per-table lock)
-      Future<?> refreshOfflineTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           // This should take global read lock + per-table lock
@@ -900,7 +909,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 3: Refresh segment for realtime table (takes global read lock 
+ per-table lock)
-      Future<?> refreshRealtimeTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           // This should take global read lock + different per-table lock
@@ -916,17 +925,17 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       startLatch.countDown();
 
       // Wait for completion with timeout
-      Assert.assertTrue(finishLatch.await(10, TimeUnit.SECONDS), "All tasks 
should complete within timeout");
+      assertTrue(finishLatch.await(10, TimeUnit.SECONDS), "All tasks should 
complete within timeout");
 
       // Verify no exceptions occurred
       if (includeServerException.get() != null) {
-        Assert.fail("Include server failed: " + 
includeServerException.get().getMessage());
+        fail("Include server failed: " + 
includeServerException.get().getMessage());
       }
       if (refreshOfflineException.get() != null) {
-        Assert.fail("Refresh offline segment failed: " + 
refreshOfflineException.get().getMessage());
+        fail("Refresh offline segment failed: " + 
refreshOfflineException.get().getMessage());
       }
       if (refreshRealtimeException.get() != null) {
-        Assert.fail("Refresh realtime segment failed: " + 
refreshRealtimeException.get().getMessage());
+        fail("Refresh realtime segment failed: " + 
refreshRealtimeException.get().getMessage());
       }
 
       // CRITICAL: Verify that the included instance IS now included in the 
routing entries
@@ -934,7 +943,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       validateEnabledInstanceInRouting(REALTIME_TABLE_NAME, 
includedServerInstance);
     } finally {
       executor.shutdown();
-      Assert.assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), 
"Executor didn't shutdown in time");
+      assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), "Executor 
didn't shutdown in time");
     }
   }
 
@@ -942,7 +951,8 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
   /// modifications. This validates that query path operations can execute 
concurrently and are not blocked by
   /// routing modifications.
   @Test
-  public void testConcurrentQueryOperationsDuringRoutingModifications() throws 
Exception {
+  public void testConcurrentQueryOperationsDuringRoutingModifications()
+      throws Exception {
     clearRoutingEntries();
 
     // Build initial routing entries
@@ -968,7 +978,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
 
     try {
       // Thread 1: Build routing (takes global read lock + per-table lock)
-      Future<?> buildRoutingTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           for (int i = 0; i < 5; i++) {
@@ -983,7 +993,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 2: Refresh segment (takes global read lock + per-table lock)
-      Future<?> refreshSegmentTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           for (int i = 0; i < 5; i++) {
@@ -998,7 +1008,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 3: Get routing table (read-only, no locks in query path)
-      Future<?> getRoutingTableTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           for (int i = 0; i < 10; i++) {
@@ -1013,7 +1023,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 4: Get time boundary info (read-only, no locks in query path)
-      Future<?> getTimeBoundaryTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           for (int i = 0; i < 10; i++) {
@@ -1028,7 +1038,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 5: Get query timeout (read-only, no locks in query path)
-      Future<?> getQueryTimeoutTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           for (int i = 0; i < 10; i++) {
@@ -1044,7 +1054,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 6: Remove routing (takes global read lock + per-table lock)
-      Future<?> removeRoutingTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(50); // Let other operations run first
@@ -1060,54 +1070,59 @@ public class BrokerRoutingManagerConcurrencyTest 
extends ControllerTest {
       startLatch.countDown();
 
       // Wait for completion with timeout
-      Assert.assertTrue(finishLatch.await(15, TimeUnit.SECONDS), "All tasks 
should complete within timeout");
+      assertTrue(finishLatch.await(15, TimeUnit.SECONDS), "All tasks should 
complete within timeout");
 
       // Verify no exceptions occurred
       if (buildRoutingException.get() != null) {
-        Assert.fail("Build routing failed: " + 
buildRoutingException.get().getMessage());
+        fail("Build routing failed: " + 
buildRoutingException.get().getMessage());
       }
       if (refreshSegmentException.get() != null) {
-        Assert.fail("Refresh segment failed: " + 
refreshSegmentException.get().getMessage());
+        fail("Refresh segment failed: " + 
refreshSegmentException.get().getMessage());
       }
       if (getRoutingTableException.get() != null) {
-        Assert.fail("Get routing table failed: " + 
getRoutingTableException.get().getMessage());
+        fail("Get routing table failed: " + 
getRoutingTableException.get().getMessage());
       }
       if (getTimeBoundaryException.get() != null) {
-        Assert.fail("Get time boundary failed: " + 
getTimeBoundaryException.get().getMessage());
+        fail("Get time boundary failed: " + 
getTimeBoundaryException.get().getMessage());
       }
       if (getQueryTimeoutException.get() != null) {
-        Assert.fail("Get query timeout failed: " + 
getQueryTimeoutException.get().getMessage());
+        fail("Get query timeout failed: " + 
getQueryTimeoutException.get().getMessage());
       }
       if (removeRoutingException.get() != null) {
-        Assert.fail("Remove routing failed: " + 
removeRoutingException.get().getMessage());
+        fail("Remove routing failed: " + 
removeRoutingException.get().getMessage());
       }
 
       // Verify offline routing still exists but realtime was removed
-      Assert.assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME),
+      assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME),
           "Offline routing should still exist after concurrent operations");
-      Assert.assertFalse(_routingManager.routingExists(REALTIME_TABLE_NAME),
+      assertFalse(_routingManager.routingExists(REALTIME_TABLE_NAME),
           "Realtime routing should not exist after concurrent operations");
 
       // CRITICAL VERIFICATION: Check TimeBoundaryManager coordination
       // The offline table should have a TimeBoundaryManager when realtime 
table exists
       Object offlineEntry = getRoutingEntry(OFFLINE_TABLE_NAME);
-      Assert.assertNotNull(offlineEntry, "Offline routing entry should exist");
+      assertNotNull(offlineEntry, "Offline routing entry should exist");
 
       // If realtime table wasn't built, offline shouldn't have 
TimeBoundaryManager
       TimeBoundaryManager timeBoundaryManager = 
getTimeBoundaryManager(offlineEntry);
-      Assert.assertNull(timeBoundaryManager, "Offline table shouldn't have 
TimeBoundaryManager when realtime table "
-          + "doesn't exist - this indicates a race condition in cross-table 
coordination");
+      assertNull(timeBoundaryManager, "Offline table shouldn't have 
TimeBoundaryManager when realtime table doesn't "
+          + "exist - this indicates a race condition in cross-table 
coordination");
     } finally {
       executor.shutdown();
-      Assert.assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), 
"Executor didn't shutdown in time");
+      assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), "Executor 
didn't shutdown in time");
     }
   }
 
-  /// Test concurrent interactions between two global write lock methods: 
processInstanceConfigChange and
-  /// includeServerToRouting. This validates that global write lock methods 
are properly serialized and don't
-  /// cause deadlocks or race conditions.
+  /// Test that the global write lock methods (processInstanceConfigChange, 
includeServerToRouting and
+  /// excludeServerFromRouting) are serialized by the global write lock and 
don't deadlock.
+  ///
+  /// The test holds the global write lock while invoking all three methods 
concurrently: none of them may complete
+  /// while the lock is held, and all of them must complete once it is 
released. This is deterministic, unlike
+  /// observing invocation order from the caller side, which cannot tell 
whether the callee has actually entered its
+  /// critical section.
   @Test
-  public void testConcurrentGlobalWriteLockMethods() throws Exception {
+  public void testConcurrentGlobalWriteLockMethods()
+      throws Exception {
     clearRoutingEntries();
 
     // Build initial routing entries
@@ -1118,120 +1133,59 @@ public class BrokerRoutingManagerConcurrencyTest 
extends ControllerTest {
     _routingManager.excludeServerFromRouting("Server_localhost_8000");
 
     ExecutorService executor = Executors.newFixedThreadPool(3);
-    CountDownLatch startLatch = new CountDownLatch(1);
+    CountDownLatch startedLatch = new CountDownLatch(3);
     CountDownLatch finishLatch = new CountDownLatch(3);
+    List<String> completedOperations = Collections.synchronizedList(new 
ArrayList<>());
+    Map<String, Exception> operationExceptions = new ConcurrentHashMap<>();
+    BiConsumer<String, Runnable> submitOperation = (operationName, operation) 
-> executor.submit(() -> {
+      try {
+        startedLatch.countDown();
+        operation.run();
+        completedOperations.add(operationName);
+      } catch (Exception e) {
+        operationExceptions.put(operationName, e);
+      } finally {
+        finishLatch.countDown();
+      }
+    });
 
-    AtomicReference<Exception> processInstanceConfigException = new 
AtomicReference<>();
-    AtomicReference<Exception> includeServerException = new 
AtomicReference<>();
-    AtomicReference<Exception> excludeServerException = new 
AtomicReference<>();
-
-    // Track execution order to verify serialization
-    List<String> executionOrder = new ArrayList<>();
-
+    ReadWriteLock globalLock = getGlobalLock();
     try {
-      // Thread 1: Process instance config change (takes global write lock)
-      Future<?> processInstanceConfigTask = executor.submit(() -> {
-        try {
-          startLatch.await();
-          synchronized (executionOrder) {
-            executionOrder.add("processInstanceConfigChange_start");
-          }
-          // This should take global write lock
-          _routingManager.processClusterChange(ChangeType.INSTANCE_CONFIG);
-          synchronized (executionOrder) {
-            executionOrder.add("processInstanceConfigChange_end");
-          }
-        } catch (Exception e) {
-          processInstanceConfigException.set(e);
-        } finally {
-          finishLatch.countDown();
-        }
-      });
-
-      // Thread 2: Include server to routing (takes global write lock)
-      Future<?> includeServerTask = executor.submit(() -> {
-        try {
-          startLatch.await();
-          Thread.sleep(10); // Small delay to encourage different ordering
-          synchronized (executionOrder) {
-            executionOrder.add("includeServerToRouting_start");
-          }
-          // This should take global write lock and be serialized with 
processInstanceConfigChange
-          _routingManager.includeServerToRouting("Server_localhost_8000");
-          synchronized (executionOrder) {
-            executionOrder.add("includeServerToRouting_end");
-          }
-        } catch (Exception e) {
-          includeServerException.set(e);
-        } finally {
-          finishLatch.countDown();
-        }
-      });
-
-      // Thread 3: Exclude another server (takes global write lock)
-      Future<?> excludeServerTask = executor.submit(() -> {
-        try {
-          startLatch.await();
-          Thread.sleep(20); // Small delay to encourage different ordering
-          synchronized (executionOrder) {
-            executionOrder.add("excludeServerFromRouting_start");
-          }
-          // This should take global write lock and be serialized with other 
global write operations
-          _routingManager.excludeServerFromRouting("Server_localhost_8001");
-          synchronized (executionOrder) {
-            executionOrder.add("excludeServerFromRouting_end");
-          }
-        } catch (Exception e) {
-          excludeServerException.set(e);
-        } finally {
-          finishLatch.countDown();
-        }
-      });
-
-      // Start all threads simultaneously
-      startLatch.countDown();
-
-      // Wait for completion with timeout
-      Assert.assertTrue(finishLatch.await(15, TimeUnit.SECONDS), "All tasks 
should complete within timeout");
-
-      // Verify no exceptions occurred
-      if (processInstanceConfigException.get() != null) {
-        Assert.fail("Process instance config failed: " + 
processInstanceConfigException.get().getMessage());
-      }
-      if (includeServerException.get() != null) {
-        Assert.fail("Include server failed: " + 
includeServerException.get().getMessage());
-      }
-      if (excludeServerException.get() != null) {
-        Assert.fail("Exclude server failed: " + 
excludeServerException.get().getMessage());
+      // Hold the global write lock so that each method blocks right after 
being invoked
+      globalLock.writeLock().lock();
+      try {
+        submitOperation.accept("processInstanceConfigChange",
+            () -> 
_routingManager.processClusterChange(ChangeType.INSTANCE_CONFIG));
+        submitOperation.accept("includeServerToRouting",
+            () -> 
_routingManager.includeServerToRouting("Server_localhost_8000"));
+        submitOperation.accept("excludeServerFromRouting",
+            () -> 
_routingManager.excludeServerFromRouting("Server_localhost_8001"));
+
+        // Wait until all three operations have been invoked, then verify that 
none of them completes while the
+        // global write lock is held
+        assertTrue(startedLatch.await(15, TimeUnit.SECONDS), "All operations 
should have been invoked");
+        assertFalse(finishLatch.await(500, TimeUnit.MILLISECONDS),
+            "No operation should complete while the global write lock is 
held");
+        assertTrue(completedOperations.isEmpty(),
+            "No operation should complete while the global write lock is held, 
completed: " + completedOperations);
+      } finally {
+        globalLock.writeLock().unlock();
       }
 
-      // Verify that operations were properly serialized (no interleaving of 
start/end events)
-      Assert.assertEquals(executionOrder.size(), 6, "Should have 6 execution 
events (3 starts, 3 ends)");
-
-      // Validate that each operation completed before the next one started
-      // (no _start event should occur between another operation's _start and 
_end)
-      for (int i = 0; i < executionOrder.size(); i += 2) {
-        String startEvent = executionOrder.get(i);
-        String endEvent = executionOrder.get(i + 1);
-        Assert.assertTrue(startEvent.endsWith("_start"), "Event at position " 
+ i + " should be a start event");
-        Assert.assertTrue(endEvent.endsWith("_end"), "Event at position " + (i 
+ 1) + " should be an end event");
-
-        // Extract operation name (everything before the last underscore)
-        String startOperation = startEvent.substring(0, 
startEvent.lastIndexOf("_"));
-        String endOperation = endEvent.substring(0, endEvent.lastIndexOf("_"));
-        Assert.assertEquals(startOperation, endOperation, "Start and end 
events should be for the same operation");
-      }
+      // Once the lock is released, all operations must complete without 
deadlock
+      assertTrue(finishLatch.await(15, TimeUnit.SECONDS),
+          "All operations should complete after the global write lock is 
released");
+      assertTrue(operationExceptions.isEmpty(), "Operations failed: " + 
operationExceptions);
+      assertEquals(completedOperations.size(), 3, "All operations should have 
completed");
 
       // Verify routing entries still exist and are properly updated
-      Assert.assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME),
+      assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME),
           "Offline routing should exist after concurrent global write 
operations");
-      Assert.assertTrue(_routingManager.routingExists(REALTIME_TABLE_NAME),
+      assertTrue(_routingManager.routingExists(REALTIME_TABLE_NAME),
           "Realtime routing should exist after concurrent global write 
operations");
-
-      System.out.println("Global write lock methods executed in order: " + 
executionOrder);
     } finally {
       executor.shutdown();
-      Assert.assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), 
"Executor didn't shutdown in time");
+      assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), "Executor 
didn't shutdown in time");
     }
   }
 
@@ -1239,7 +1193,8 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
   /// This validates proper coordination between logical table operations, 
regular table operations,
   /// and global write operations. Uses a hybrid logical table configuration 
with both offline and realtime tables.
   @Test
-  public void testConcurrentLogicalTableBuildAndRegularBuild() throws 
Exception {
+  public void testConcurrentLogicalTableBuildAndRegularBuild()
+      throws Exception {
     clearRoutingEntries();
 
     String logicalTableName = "testLogicalTable";
@@ -1247,10 +1202,9 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
     String physicalRealtimeTable = "testLogicalTable_REALTIME";
 
     // Create hybrid logical table config with both offline and realtime tables
-    LogicalTableConfig logicalTableConfig = 
createLogicalTableConfig(logicalTableName,
-        Map.of(
-            "testLogicalTable", new PhysicalTableConfig()
-        ), physicalOfflineTable, physicalRealtimeTable);
+    LogicalTableConfig logicalTableConfig =
+        createLogicalTableConfig(logicalTableName, Map.of("testLogicalTable", 
new PhysicalTableConfig()),
+            physicalOfflineTable, physicalRealtimeTable);
     ZKMetadataProvider.setLogicalTableConfig(_propertyStore, 
logicalTableConfig);
 
     // Create physical table configs and schemas
@@ -1284,7 +1238,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
 
     try {
       // Thread 1: Build routing for logical table (global read lock + 
per-table locks for both physical tables)
-      Future<?> logicalBuildTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           // This should take global read lock + per-table locks for both 
physical tables
@@ -1297,7 +1251,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 2: Build routing for regular table (global read lock + 
different per-table lock)
-      Future<?> regularBuildTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(5); // Small delay to encourage interleaving
@@ -1311,7 +1265,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 3: Refresh segment on offline physical table (global read lock 
+ same per-table lock as logical)
-      Future<?> refreshOfflineTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(10); // Small delay to encourage interleaving
@@ -1325,7 +1279,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 4: Refresh segment on realtime physical table (global read 
lock + different per-table lock)
-      Future<?> refreshRealtimeTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(12); // Small delay to encourage interleaving
@@ -1339,7 +1293,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 5: Exclude server from routing (global write lock - should 
serialize with all read operations)
-      Future<?> excludeServerTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(15); // Small delay to encourage interleaving
@@ -1356,34 +1310,34 @@ public class BrokerRoutingManagerConcurrencyTest 
extends ControllerTest {
       startLatch.countDown();
 
       // Wait for completion with timeout
-      Assert.assertTrue(finishLatch.await(15, TimeUnit.SECONDS), "All tasks 
should complete within timeout");
+      assertTrue(finishLatch.await(15, TimeUnit.SECONDS), "All tasks should 
complete within timeout");
 
       // Verify no exceptions occurred
       if (logicalBuildException.get() != null) {
-        Assert.fail("Logical table build failed: " + 
logicalBuildException.get().getMessage());
+        fail("Logical table build failed: " + 
logicalBuildException.get().getMessage());
       }
       if (regularBuildException.get() != null) {
-        Assert.fail("Regular table build failed: " + 
regularBuildException.get().getMessage());
+        fail("Regular table build failed: " + 
regularBuildException.get().getMessage());
       }
       if (refreshOfflineException.get() != null) {
-        Assert.fail("Refresh offline segment failed: " + 
refreshOfflineException.get().getMessage());
+        fail("Refresh offline segment failed: " + 
refreshOfflineException.get().getMessage());
       }
       if (refreshRealtimeException.get() != null) {
-        Assert.fail("Refresh realtime segment failed: " + 
refreshRealtimeException.get().getMessage());
+        fail("Refresh realtime segment failed: " + 
refreshRealtimeException.get().getMessage());
       }
       if (excludeServerException.get() != null) {
-        Assert.fail("Exclude server failed: " + 
excludeServerException.get().getMessage());
+        fail("Exclude server failed: " + 
excludeServerException.get().getMessage());
       }
 
       // Verify routing entries exist for regular table
-      Assert.assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME),
+      assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME),
           "Regular table routing should exist after concurrent operations");
 
       // Verify that the excluded server is not in routing for any existing 
tables
       validateDisabledInstanceNotInRouting(OFFLINE_TABLE_NAME, 
"Server_localhost_8001");
 
       // Verify routing entries exist for regular table
-      Assert.assertTrue(_routingManager.routingExists(physicalOfflineTable),
+      assertTrue(_routingManager.routingExists(physicalOfflineTable),
           "Regular table routing should exist after concurrent operations");
 
       // Verify that the excluded server is not in routing for any existing 
tables
@@ -1396,18 +1350,18 @@ public class BrokerRoutingManagerConcurrencyTest 
extends ControllerTest {
       // CRITICAL VERIFICATION: Check TimeBoundaryManager coordination
       // If the logical table build actually created routing for physical 
tables, verify TimeBoundaryManager
       Object offlineEntry = getRoutingEntry(physicalOfflineTable);
-      Assert.assertNotNull(offlineEntry, "Physical offline routing entry 
should exist");
+      assertNotNull(offlineEntry, "Physical offline routing entry should 
exist");
 
       // Physical offline table should have TimeBoundaryManager due to logical 
table setup
       TimeBoundaryManager timeBoundaryManager = 
getTimeBoundaryManager(offlineEntry);
-      Assert.assertNotNull(timeBoundaryManager, "Physical offline table should 
have TimeBoundaryManager if part of "
-          + "a logical table - this indicates a race condition in cross-table 
coordination");
+      assertNotNull(timeBoundaryManager, "Physical offline table should have 
TimeBoundaryManager if part of a logical "
+          + "table - this indicates a race condition in cross-table 
coordination");
 
-      Assert.assertFalse(_routingManager.routingExists(physicalRealtimeTable),
+      assertFalse(_routingManager.routingExists(physicalRealtimeTable),
           "Physical realtime routing entry should not exist since we never 
built routing entry for it");
     } finally {
       executor.shutdown();
-      Assert.assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), 
"Executor didn't shutdown in time");
+      assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), "Executor 
didn't shutdown in time");
     }
   }
 
@@ -1415,7 +1369,8 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
   /// This validates proper coordination between logical table operations, 
regular table operations,
   /// and global write operations. Uses a hybrid logical table configuration 
with both offline and realtime tables.
   @Test
-  public void testConcurrentLogicalTableBuildAndRegularBuildAndRealtimeBuild() 
throws Exception {
+  public void testConcurrentLogicalTableBuildAndRegularBuildAndRealtimeBuild()
+      throws Exception {
     clearRoutingEntries();
 
     String logicalTableName = "testLogicalTable";
@@ -1423,10 +1378,9 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
     String physicalRealtimeTable = "testLogicalTable_REALTIME";
 
     // Create hybrid logical table config with both offline and realtime tables
-    LogicalTableConfig logicalTableConfig = 
createLogicalTableConfig(logicalTableName,
-        Map.of(
-            "testLogicalTable", new PhysicalTableConfig()
-        ), physicalOfflineTable, physicalRealtimeTable);
+    LogicalTableConfig logicalTableConfig =
+        createLogicalTableConfig(logicalTableName, Map.of("testLogicalTable", 
new PhysicalTableConfig()),
+            physicalOfflineTable, physicalRealtimeTable);
     ZKMetadataProvider.setLogicalTableConfig(_propertyStore, 
logicalTableConfig);
 
     // Create physical table configs and schemas
@@ -1461,7 +1415,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
 
     try {
       // Thread 1: Build routing for logical table (global read lock + 
per-table locks for both physical tables)
-      Future<?> logicalBuildTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           // This should take global read lock + per-table locks for both 
physical tables
@@ -1474,7 +1428,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 2: Build routing for regular table (global read lock + 
different per-table lock)
-      Future<?> regularBuildTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(5); // Small delay to encourage interleaving
@@ -1488,7 +1442,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 3: Refresh segment on offline physical table (global read lock 
+ same per-table lock as logical)
-      Future<?> refreshOfflineTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(10); // Small delay to encourage interleaving
@@ -1502,7 +1456,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 4: Refresh segment on realtime physical table (global read 
lock + different per-table lock)
-      Future<?> refreshRealtimeTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(12); // Small delay to encourage interleaving
@@ -1516,7 +1470,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 5: Exclude server from routing (global write lock - should 
serialize with all read operations)
-      Future<?> excludeServerTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(15); // Small delay to encourage interleaving
@@ -1530,7 +1484,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 6: Build routing for physical realtime table (global read lock 
+ different per-table lock)
-      Future<?> physicalRealtimeBuildTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(5); // Small delay to encourage interleaving
@@ -1547,37 +1501,37 @@ public class BrokerRoutingManagerConcurrencyTest 
extends ControllerTest {
       startLatch.countDown();
 
       // Wait for completion with timeout
-      Assert.assertTrue(finishLatch.await(15, TimeUnit.SECONDS), "All tasks 
should complete within timeout");
+      assertTrue(finishLatch.await(15, TimeUnit.SECONDS), "All tasks should 
complete within timeout");
 
       // Verify no exceptions occurred
       if (logicalBuildException.get() != null) {
-        Assert.fail("Logical table build failed: " + 
logicalBuildException.get().getMessage());
+        fail("Logical table build failed: " + 
logicalBuildException.get().getMessage());
       }
       if (regularBuildException.get() != null) {
-        Assert.fail("Regular table build failed: " + 
regularBuildException.get().getMessage());
+        fail("Regular table build failed: " + 
regularBuildException.get().getMessage());
       }
       if (refreshOfflineException.get() != null) {
-        Assert.fail("Refresh offline segment failed: " + 
refreshOfflineException.get().getMessage());
+        fail("Refresh offline segment failed: " + 
refreshOfflineException.get().getMessage());
       }
       if (refreshRealtimeException.get() != null) {
-        Assert.fail("Refresh realtime segment failed: " + 
refreshRealtimeException.get().getMessage());
+        fail("Refresh realtime segment failed: " + 
refreshRealtimeException.get().getMessage());
       }
       if (excludeServerException.get() != null) {
-        Assert.fail("Exclude server failed: " + 
excludeServerException.get().getMessage());
+        fail("Exclude server failed: " + 
excludeServerException.get().getMessage());
       }
       if (realtimeBuildException.get() != null) {
-        Assert.fail("Realtime table build failed: " + 
realtimeBuildException.get().getMessage());
+        fail("Realtime table build failed: " + 
realtimeBuildException.get().getMessage());
       }
 
       // Verify routing entries exist for regular table
-      Assert.assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME),
+      assertTrue(_routingManager.routingExists(OFFLINE_TABLE_NAME),
           "Regular table routing should exist after concurrent operations");
 
       // Verify that the excluded server is not in routing for any existing 
tables
       validateDisabledInstanceNotInRouting(OFFLINE_TABLE_NAME, 
"Server_localhost_8001");
 
       // Verify routing entries exist for regular table
-      Assert.assertTrue(_routingManager.routingExists(physicalOfflineTable),
+      assertTrue(_routingManager.routingExists(physicalOfflineTable),
           "Regular table routing should exist after concurrent operations");
 
       // Verify that the excluded server is not in routing for any existing 
tables
@@ -1590,25 +1544,26 @@ public class BrokerRoutingManagerConcurrencyTest 
extends ControllerTest {
       // CRITICAL VERIFICATION: Check TimeBoundaryManager coordination
       // If the logical table build actually created routing for physical 
tables, verify TimeBoundaryManager
       Object offlineEntry = getRoutingEntry(physicalOfflineTable);
-      Assert.assertNotNull(offlineEntry, "Physical offline routing entry 
should exist");
+      assertNotNull(offlineEntry, "Physical offline routing entry should 
exist");
 
       // Physical offline table should have TimeBoundaryManager due to logical 
table setup
       TimeBoundaryManager timeBoundaryManager = 
getTimeBoundaryManager(offlineEntry);
-      Assert.assertNotNull(timeBoundaryManager, "Physical offline table should 
have TimeBoundaryManager if part of "
-          + "a logical table - this indicates a race condition in cross-table 
coordination");
+      assertNotNull(timeBoundaryManager, "Physical offline table should have 
TimeBoundaryManager if part of a logical "
+          + "table - this indicates a race condition in cross-table 
coordination");
 
-      Assert.assertTrue(_routingManager.routingExists(physicalRealtimeTable),
+      assertTrue(_routingManager.routingExists(physicalRealtimeTable),
           "Physical realtime routing entry should exist since we built routing 
entry for it");
     } finally {
       executor.shutdown();
-      Assert.assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), 
"Executor didn't shutdown in time");
+      assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS), "Executor 
didn't shutdown in time");
     }
   }
 
   /// Test concurrent interactions with logical table containing multiple 
physical tables.
   /// This validates coordination when logical table operations affect 
multiple per-table locks.
   @Test
-  public void testConcurrentMultiPhysicalTableLogicalOperations() throws 
Exception {
+  public void testConcurrentMultiPhysicalTableLogicalOperations()
+      throws Exception {
     clearRoutingEntries();
 
     String logicalTableName = "testMultiLogicalTable";
@@ -1618,18 +1573,14 @@ public class BrokerRoutingManagerConcurrencyTest 
extends ControllerTest {
 
     // Create logical table config with multiple physical tables
     LogicalTableConfig logicalTableConfig = 
createLogicalTableConfig(logicalTableName,
-        Map.of(
-            "multiPhysical1", new PhysicalTableConfig(),
-            "multiPhysical2", new PhysicalTableConfig(),
-            "multiPhysical3", new PhysicalTableConfig()
-        ),
-        physicalTable1, physicalTable3);
+        Map.of("multiPhysical1", new PhysicalTableConfig(), "multiPhysical2", 
new PhysicalTableConfig(),
+            "multiPhysical3", new PhysicalTableConfig()), physicalTable1, 
physicalTable3);
     ZKMetadataProvider.setLogicalTableConfig(_propertyStore, 
logicalTableConfig);
 
     // Create physical table configs and schemas
     for (String tableNameWithType : Arrays.asList(physicalTable1, 
physicalTable2, physicalTable3)) {
-      TableType tableType = 
TableNameBuilder.isOfflineTableResource(tableNameWithType)
-          ? TableType.OFFLINE : TableType.REALTIME;
+      TableType tableType =
+          TableNameBuilder.isOfflineTableResource(tableNameWithType) ? 
TableType.OFFLINE : TableType.REALTIME;
       TableConfig physicalTableConfig = createTableConfig(tableNameWithType, 
tableType);
       ZKMetadataProvider.setTableConfig(_propertyStore, physicalTableConfig);
       ZKMetadataProvider.setSchema(_propertyStore, createMockSchema());
@@ -1639,6 +1590,8 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
     ExecutorService executor = Executors.newFixedThreadPool(5);
     CountDownLatch startLatch = new CountDownLatch(1);
     CountDownLatch finishLatch = new CountDownLatch(5);
+    // Released once the logical table build finishes, so that the logical 
table remove runs after it
+    CountDownLatch logicalBuildDoneLatch = new CountDownLatch(1);
 
     AtomicReference<Exception> logicalBuildException = new AtomicReference<>();
     AtomicReference<Exception> logicalRemoveException = new 
AtomicReference<>();
@@ -1648,7 +1601,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
 
     try {
       // Thread 1: Build routing for logical table (global read lock + 
multiple per-table locks)
-      Future<?> logicalBuildTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           // This should take global read lock + per-table locks for all 3 
physical tables
@@ -1656,15 +1609,20 @@ public class BrokerRoutingManagerConcurrencyTest 
extends ControllerTest {
         } catch (Exception e) {
           logicalBuildException.set(e);
         } finally {
+          logicalBuildDoneLatch.countDown();
           finishLatch.countDown();
         }
       });
 
       // Thread 2: Remove routing for logical table (global read lock + 
multiple per-table locks)
-      Future<?> logicalRemoveTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
-          Thread.sleep(20); // Delay to let build start first
+          // Wait for the logical table build to finish so that the remove 
detaches the time boundary manager the
+          // build attached; without this ordering the final time boundary 
manager state would depend on scheduling
+          if (!logicalBuildDoneLatch.await(15, TimeUnit.SECONDS)) {
+            throw new IllegalStateException("Timed out waiting for the logical 
table build to finish");
+          }
           // This should take global read lock + per-table locks for all 3 
physical tables
           _routingManager.removeRoutingForLogicalTable(logicalTableName);
         } catch (Exception e) {
@@ -1675,7 +1633,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 3: Build routing for one of the physical tables directly 
(competing per-table lock)
-      Future<?> regularBuild1Task = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(5);
@@ -1689,7 +1647,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 4: Build routing for another physical table (different 
competing per-table lock)
-      Future<?> regularBuild2Task = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(10);
@@ -1703,7 +1661,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 5: Include server to routing (global write lock - should block 
all read operations)
-      Future<?> includeServerTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(15);
@@ -1720,57 +1678,57 @@ public class BrokerRoutingManagerConcurrencyTest 
extends ControllerTest {
       startLatch.countDown();
 
       // Wait for completion with timeout
-      Assert.assertTrue(finishLatch.await(20, TimeUnit.SECONDS), "All tasks 
should complete within timeout");
+      assertTrue(finishLatch.await(20, TimeUnit.SECONDS), "All tasks should 
complete within timeout");
 
       // Verify no exceptions occurred
       if (logicalBuildException.get() != null) {
-        Assert.fail("Logical table build failed: " + 
logicalBuildException.get().getMessage());
+        fail("Logical table build failed: " + 
logicalBuildException.get().getMessage());
       }
       if (logicalRemoveException.get() != null) {
-        Assert.fail("Logical table remove failed: " + 
logicalRemoveException.get().getMessage());
+        fail("Logical table remove failed: " + 
logicalRemoveException.get().getMessage());
       }
       if (regularBuild1Exception.get() != null) {
-        Assert.fail("Regular table 1 build failed: " + 
regularBuild1Exception.get().getMessage());
+        fail("Regular table 1 build failed: " + 
regularBuild1Exception.get().getMessage());
       }
       if (regularBuild2Exception.get() != null) {
-        Assert.fail("Regular table 2 build failed: " + 
regularBuild2Exception.get().getMessage());
+        fail("Regular table 2 build failed: " + 
regularBuild2Exception.get().getMessage());
       }
       if (includeServerException.get() != null) {
-        Assert.fail("Include server failed: " + 
includeServerException.get().getMessage());
+        fail("Include server failed: " + 
includeServerException.get().getMessage());
       }
 
+      // CRITICAL VERIFICATION: Check TimeBoundaryManager coordination before 
rebuilding any routing, because a
+      // rebuild recomputes the time boundary manager and would mask what the 
concurrent phase produced. The logical
+      // table build attached the manager and the logical table remove (which 
runs after the build) detached it;
+      // rebuilding the physical table directly cannot re-attach it because no 
realtime counterpart exists.
+      Object offlineEntry = getRoutingEntry(physicalTable1);
+      assertNotNull(offlineEntry, "Physical offline routing entry should 
exist");
+      TimeBoundaryManager timeBoundaryManager = 
getTimeBoundaryManager(offlineEntry);
+      assertNull(timeBoundaryManager, "Physical offline table shouldn't have 
TimeBoundaryManager if part of a logical "
+          + "table since logical table was removed - this indicates a race 
condition in cross-table coordination");
+
       // Verify that routing can be built for all tables after operations
       _routingManager.buildRouting(physicalTable1);
       _routingManager.buildRouting(physicalTable2);
       _routingManager.buildRouting(physicalTable3);
 
-      Assert.assertTrue(_routingManager.routingExists(physicalTable1),
+      assertTrue(_routingManager.routingExists(physicalTable1),
           "Physical table 1 routing should be buildable after concurrent 
operations");
-      Assert.assertTrue(_routingManager.routingExists(physicalTable2),
+      assertTrue(_routingManager.routingExists(physicalTable2),
           "Physical table 2 routing should be buildable after concurrent 
operations");
-      Assert.assertTrue(_routingManager.routingExists(physicalTable3),
+      assertTrue(_routingManager.routingExists(physicalTable3),
           "Physical table 3 routing should be buildable after concurrent 
operations");
-
-      // CRITICAL VERIFICATION: Check TimeBoundaryManager coordination
-      // If the logical table build actually created routing for physical 
tables, verify TimeBoundaryManager
-      Object offlineEntry = getRoutingEntry(physicalTable1);
-      Assert.assertNotNull(offlineEntry, "Physical offline routing entry 
should exist");
-
-      // Physical offline table shouldn't have TimeBoundaryManager due to 
logical table setup and then removal
-      TimeBoundaryManager timeBoundaryManager = 
getTimeBoundaryManager(offlineEntry);
-      Assert.assertNull(timeBoundaryManager, "Physical offline table shouldn't 
have TimeBoundaryManager if part of "
-          + "a logical table since logical table was removed - this indicates 
a race condition in cross-table "
-          + "coordination");
     } finally {
       executor.shutdown();
-      Assert.assertTrue(executor.awaitTermination(15, TimeUnit.SECONDS), 
"Executor didn't shutdown in time");
+      assertTrue(executor.awaitTermination(15, TimeUnit.SECONDS), "Executor 
didn't shutdown in time");
     }
   }
 
   /// Test concurrent interactions with logical table containing multiple 
physical tables.
   /// This validates coordination when logical table operations affect 
multiple per-table locks.
   @Test
-  public void 
testConcurrentMultiPhysicalTableLogicalOperationsWithRealtimeBuild() throws 
Exception {
+  public void 
testConcurrentMultiPhysicalTableLogicalOperationsWithRealtimeBuild()
+      throws Exception {
     clearRoutingEntries();
 
     String logicalTableName = "testMultiLogicalTable";
@@ -1781,18 +1739,14 @@ public class BrokerRoutingManagerConcurrencyTest 
extends ControllerTest {
 
     // Create logical table config with multiple physical tables
     LogicalTableConfig logicalTableConfig = 
createLogicalTableConfig(logicalTableName,
-        Map.of(
-            "multiPhysical1", new PhysicalTableConfig(),
-            "multiPhysical2", new PhysicalTableConfig(),
-            "multiPhysical3", new PhysicalTableConfig()
-        ),
-        physicalTable1, physicalTable4);
+        Map.of("multiPhysical1", new PhysicalTableConfig(), "multiPhysical2", 
new PhysicalTableConfig(),
+            "multiPhysical3", new PhysicalTableConfig()), physicalTable1, 
physicalTable4);
     ZKMetadataProvider.setLogicalTableConfig(_propertyStore, 
logicalTableConfig);
 
     // Create physical table configs and schemas
     for (String tableNameWithType : Arrays.asList(physicalTable1, 
physicalTable2, physicalTable3, physicalTable4)) {
-      TableType tableType = 
TableNameBuilder.isOfflineTableResource(tableNameWithType)
-          ? TableType.OFFLINE : TableType.REALTIME;
+      TableType tableType =
+          TableNameBuilder.isOfflineTableResource(tableNameWithType) ? 
TableType.OFFLINE : TableType.REALTIME;
       TableConfig physicalTableConfig = createTableConfig(tableNameWithType, 
tableType);
       ZKMetadataProvider.setTableConfig(_propertyStore, physicalTableConfig);
 
@@ -1807,6 +1761,9 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
     CountDownLatch startLatch = new CountDownLatch(1);
     CountDownLatch finishLatch = new CountDownLatch(6);
 
+    // Released once the logical table build finishes, so that the logical 
table remove runs after it
+    CountDownLatch logicalBuildDoneLatch = new CountDownLatch(1);
+
     AtomicReference<Exception> logicalBuildException = new AtomicReference<>();
     AtomicReference<Exception> logicalRemoveException = new 
AtomicReference<>();
     AtomicReference<Exception> regularBuild1Exception = new 
AtomicReference<>();
@@ -1816,7 +1773,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
 
     try {
       // Thread 1: Build routing for logical table (global read lock + 
multiple per-table locks)
-      Future<?> logicalBuildTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           // This should take global read lock + per-table locks for all 3 
physical tables
@@ -1824,15 +1781,20 @@ public class BrokerRoutingManagerConcurrencyTest 
extends ControllerTest {
         } catch (Exception e) {
           logicalBuildException.set(e);
         } finally {
+          logicalBuildDoneLatch.countDown();
           finishLatch.countDown();
         }
       });
 
       // Thread 2: Remove routing for logical table (global read lock + 
multiple per-table locks)
-      Future<?> logicalRemoveTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
-          Thread.sleep(25); // Delay to let build start first
+          // Wait for the logical table build to finish so that the remove 
runs after the build, matching the
+          // built-then-removed scenario under test
+          if (!logicalBuildDoneLatch.await(15, TimeUnit.SECONDS)) {
+            throw new IllegalStateException("Timed out waiting for the logical 
table build to finish");
+          }
           // This should take global read lock + per-table locks for all 3 
physical tables
           _routingManager.removeRoutingForLogicalTable(logicalTableName);
         } catch (Exception e) {
@@ -1843,7 +1805,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 3: Build routing for one of the physical tables directly 
(competing per-table lock)
-      Future<?> regularBuild1Task = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(5);
@@ -1857,7 +1819,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 4: Build routing for another physical table (different 
competing per-table lock)
-      Future<?> regularBuild2Task = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(10);
@@ -1871,7 +1833,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 5: Include server to routing (global write lock - should block 
all read operations)
-      Future<?> includeServerTask = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(15);
@@ -1885,7 +1847,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
       });
 
       // Thread 6: Build routing for realtime physical table (competing 
per-table lock)
-      Future<?> regularBuild4Task = executor.submit(() -> {
+      executor.submit(() -> {
         try {
           startLatch.await();
           Thread.sleep(20);
@@ -1902,55 +1864,57 @@ public class BrokerRoutingManagerConcurrencyTest 
extends ControllerTest {
       startLatch.countDown();
 
       // Wait for completion with timeout
-      Assert.assertTrue(finishLatch.await(20, TimeUnit.SECONDS), "All tasks 
should complete within timeout");
+      assertTrue(finishLatch.await(20, TimeUnit.SECONDS), "All tasks should 
complete within timeout");
 
       // Verify no exceptions occurred
       if (logicalBuildException.get() != null) {
-        Assert.fail("Logical table build failed: " + 
logicalBuildException.get().getMessage());
+        fail("Logical table build failed: " + 
logicalBuildException.get().getMessage());
       }
       if (logicalRemoveException.get() != null) {
-        Assert.fail("Logical table remove failed: " + 
logicalRemoveException.get().getMessage());
+        fail("Logical table remove failed: " + 
logicalRemoveException.get().getMessage());
       }
       if (regularBuild1Exception.get() != null) {
-        Assert.fail("Regular table 1 build failed: " + 
regularBuild1Exception.get().getMessage());
+        fail("Regular table 1 build failed: " + 
regularBuild1Exception.get().getMessage());
       }
       if (regularBuild2Exception.get() != null) {
-        Assert.fail("Regular table 2 build failed: " + 
regularBuild2Exception.get().getMessage());
+        fail("Regular table 2 build failed: " + 
regularBuild2Exception.get().getMessage());
       }
       if (includeServerException.get() != null) {
-        Assert.fail("Include server failed: " + 
includeServerException.get().getMessage());
+        fail("Include server failed: " + 
includeServerException.get().getMessage());
       }
       if (realtimeBuildException.get() != null) {
-        Assert.fail("Regular table 4 build failed: " + 
realtimeBuildException.get().getMessage());
+        fail("Regular table 4 build failed: " + 
realtimeBuildException.get().getMessage());
       }
 
+      // CRITICAL VERIFICATION: Check TimeBoundaryManager coordination before 
rebuilding any routing, because a
+      // rebuild recomputes the time boundary manager and would mask what the 
concurrent phase produced. Since
+      // physicalTable4 (realtime) routing exists and shares the raw table 
name (multiPhysical1) with physicalTable1,
+      // the offline table must end up with a time boundary manager regardless 
of how the operations interleave:
+      // building the realtime routing attaches it to the offline counterpart, 
rebuilding the offline routing attaches
+      // it while the realtime routing exists, and the logical table remove 
skips detaching it for hybrid physical
+      // tables.
+      Object offlineEntry = getRoutingEntry(physicalTable1);
+      assertNotNull(offlineEntry, "Physical offline routing entry should 
exist");
+      TimeBoundaryManager timeBoundaryManager = 
getTimeBoundaryManager(offlineEntry);
+      assertNotNull(timeBoundaryManager, "Physical offline table should have 
TimeBoundaryManager when realtime table "
+          + "exists - this indicates proper hybrid table coordination");
+
       // Verify that routing can be built for all tables after operations
       _routingManager.buildRouting(physicalTable1);
       _routingManager.buildRouting(physicalTable2);
       _routingManager.buildRouting(physicalTable3);
 
-      Assert.assertTrue(_routingManager.routingExists(physicalTable1),
+      assertTrue(_routingManager.routingExists(physicalTable1),
           "Physical table 1 routing should be buildable after concurrent 
operations");
-      Assert.assertTrue(_routingManager.routingExists(physicalTable2),
+      assertTrue(_routingManager.routingExists(physicalTable2),
           "Physical table 2 routing should be buildable after concurrent 
operations");
-      Assert.assertTrue(_routingManager.routingExists(physicalTable3),
+      assertTrue(_routingManager.routingExists(physicalTable3),
           "Physical table 3 routing should be buildable after concurrent 
operations");
-      Assert.assertTrue(_routingManager.routingExists(physicalTable4),
+      assertTrue(_routingManager.routingExists(physicalTable4),
           "Physical table 4 routing should be buildable after concurrent 
operations");
-
-      // CRITICAL VERIFICATION: Check TimeBoundaryManager coordination
-      // The logical table was built then removed, so check final state
-      Object offlineEntry = getRoutingEntry(physicalTable1);
-      Assert.assertNotNull(offlineEntry, "Physical offline routing entry 
should exist");
-
-      // Since physicalTable4 (realtime) exists and they share the same raw 
table name (multiPhysical1),
-      // the offline table should have a TimeBoundaryManager for hybrid table 
coordination
-      TimeBoundaryManager timeBoundaryManager = 
getTimeBoundaryManager(offlineEntry);
-      Assert.assertNotNull(timeBoundaryManager, "Physical offline table should 
have TimeBoundaryManager when "
-          + "realtime table exists - this indicates proper hybrid table 
coordination");
     } finally {
       executor.shutdown();
-      Assert.assertTrue(executor.awaitTermination(15, TimeUnit.SECONDS), 
"Executor didn't shutdown in time");
+      assertTrue(executor.awaitTermination(15, TimeUnit.SECONDS), "Executor 
didn't shutdown in time");
     }
   }
 
@@ -1961,8 +1925,7 @@ public class BrokerRoutingManagerConcurrencyTest extends 
ControllerTest {
     parameters.put("includedTables", List.of(refOfflineTableName));
     TimeBoundaryConfig timeBoundaryConfig = new TimeBoundaryConfig("min", 
parameters);
 
-    return new LogicalTableConfigBuilder()
-        .setTableName(logicalTableName)
+    return new LogicalTableConfigBuilder().setTableName(logicalTableName)
         .setPhysicalTableConfigMap(physicalTableConfigMap)
         .setBrokerTenant("DefaultTenant")
         .setRefOfflineTableName(refOfflineTableName)


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to