From ce1baa4e255009084a4454b8e1ec7e193f984fde Mon Sep 17 00:00:00 2001 From: Evgeny Malygin Date: Fri, 18 Sep 2026 16:48:19 -0400 Subject: [PATCH 1/2] UT: add test for queue lookup by appId/subQId Signed-off-by: Evgeny Malygin --- .../bloomberg/bmq/impl/QueueManagerTest.java | 74 +++++++++++++++++++ 1 file changed, 74 insertions(+) diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/QueueManagerTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/QueueManagerTest.java index 4506efd6..51d5cfe3 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/QueueManagerTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/QueueManagerTest.java @@ -17,6 +17,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.fail; @@ -279,6 +280,79 @@ void removeExpiredQueueTest() { logger.info("=================================================================="); } + /** + * Test that an expired queue is looked up by both its canonical URI and its appId. + * + *

The same appId may be used by queues with different canonical URIs, and each of them has + * its own subQueueId. {@link LateResponseHandler} rebuilds a {@link QueueId} from the qId of a + * late configure response and the subQueueId found by appId, so the lookup must give the + * subQueueId of the queue the response belongs to. + * + *

Test steps: + * + *

    + *
  1. create queue manager instance + *
  2. insert two active queues of one canonical URI with appIds "foo" and "bar", so that + * "bar" gets subQueueId 2 + *
  3. insert an active queue of another canonical URI with appId "bar", which gets subQueueId + * 1 + *
  4. move all three queues to expired queues, as a local timeout does + *
  5. for each queue, rebuild a QueueId from its qId and the subQueueId found by its appId, + * the way LateResponseHandler does + *
  6. check that every rebuilt QueueId resolves to the queue it was built from + *
+ */ + @Test + void lookupByUriAndAppIdTest() { + logger.info("===================================================="); + logger.info("BEGIN Testing QueueManager lookupByUriAndAppIdTest."); + logger.info("===================================================="); + + QueueManager obj = QueueManager.createInstance(); + + Uri fooOfQueueA = new Uri("bmq://ts.trades.myapp/queueA?id=foo"); + Uri barOfQueueA = new Uri("bmq://ts.trades.myapp/queueA?id=bar"); + Uri barOfQueueB = new Uri("bmq://ts.trades.myapp/queueB?id=bar"); + + QueueImpl queueAFoo = insertActive(obj, fooOfQueueA); + QueueImpl queueABar = insertActive(obj, barOfQueueA); + QueueImpl queueBBar = insertActive(obj, barOfQueueB); + + // Two appIds of one canonical URI, so the second one gets subQueueId 2. + assertEquals(1, queueAFoo.getSubQueueId()); + assertEquals(2, queueABar.getSubQueueId()); + + // Another canonical URI starts its subQueueIds anew. + assertEquals(1, queueBBar.getSubQueueId()); + + assertTrue(obj.insertExpired(queueAFoo)); + assertTrue(obj.insertExpired(queueABar)); + assertTrue(obj.insertExpired(queueBBar)); + + assertEquals(queueAFoo, findExpiredByAppId(obj, queueAFoo)); + assertEquals(queueABar, findExpiredByAppId(obj, queueABar)); + assertEquals(queueBBar, findExpiredByAppId(obj, queueBBar)); + + logger.info("=================================================="); + logger.info("END Testing QueueManager lookupByUriAndAppIdTest."); + logger.info("=================================================="); + } + + private QueueImpl insertActive(QueueManager manager, Uri uri) { + QueueImpl queue = createQueue(session, uri, 0L); + QueueId queueId = manager.generateNextQueueId(uri); + queue.setQueueId(queueId.getQId()).setSubQueueId(queueId.getSubQId()); + assertTrue(manager.insert(queue)); + return queue; + } + + /** Resolve an expired queue the way LateResponseHandler does. */ + private QueueImpl findExpiredByAppId(QueueManager manager, QueueImpl queue) { + Integer subQId = manager.findSubQId(queue.getUri().id()); + assertNotNull(subQId); + return manager.findExpiredByQueueId(QueueId.createInstance(queue.getQueueId(), subQId)); + } + /** * Test for queue manager to test increment queue substream count * From a866edeb76dd79d084ac4c956418005f36a6b8c5 Mon Sep 17 00:00:00 2001 From: Evgeny Malygin Date: Fri, 18 Sep 2026 17:00:04 -0400 Subject: [PATCH 2/2] Fix: scan expired queues by both qId/appId Signed-off-by: Evgeny Malygin --- .../bmq/impl/LateResponseHandler.java | 2 +- .../com/bloomberg/bmq/impl/QueueManager.java | 52 +++++++++++-------- .../bloomberg/bmq/impl/QueueStateManager.java | 4 +- .../bloomberg/bmq/impl/QueueManagerTest.java | 2 +- 4 files changed, 35 insertions(+), 25 deletions(-) diff --git a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/LateResponseHandler.java b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/LateResponseHandler.java index 7216990e..a86ce29f 100644 --- a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/LateResponseHandler.java +++ b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/LateResponseHandler.java @@ -106,7 +106,7 @@ private void handleConfigureStreamResponse(ConfigureStreamResponse configureStre } int qId = originalRequest.id(); - Integer subQId = queueStateManager.findSubQId(parameters.appId()); + Integer subQId = queueStateManager.findSubQId(qId, parameters.appId()); if (subQId == null) { subQId = QueueId.k_DEFAULT_SUBQUEUE_ID; } diff --git a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/QueueManager.java b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/QueueManager.java index e8034c60..af6c80ba 100644 --- a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/QueueManager.java +++ b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/QueueManager.java @@ -22,6 +22,7 @@ import java.util.Collection; import java.util.HashMap; import java.util.Map; +import java.util.Objects; import java.util.TreeMap; import java.util.concurrent.atomic.AtomicInteger; import java.util.stream.Collectors; @@ -49,7 +50,6 @@ public class QueueManager { private Map expiredQueueMap; private Map subscriptionIdMap; - private Map appId_subQId_Map; private AtomicInteger nextQueueId; private QueueManager() { @@ -58,7 +58,6 @@ private QueueManager() { keyQueueIdMap = new HashMap<>(); expiredQueueMap = new HashMap<>(); subscriptionIdMap = new HashMap<>(); - appId_subQId_Map = new HashMap<>(); lock = new Object(); nextQueueId = new AtomicInteger(0); } @@ -71,7 +70,6 @@ public void reset() { keyQueueIdMap = new HashMap<>(); expiredQueueMap = new HashMap<>(); subscriptionIdMap = new HashMap<>(); - appId_subQId_Map = new HashMap<>(); nextQueueId = new AtomicInteger(0); } } @@ -147,14 +145,6 @@ public boolean insertExpired(QueueImpl queue) { return false; } expiredQueueMap.put(queueId, queue); - - SubQueueIdInfo subQueueIdInfo = queue.getParameters().getSubIdInfo(); - if (subQueueIdInfo != null) { - String appId = subQueueIdInfo.appId(); - if (appId != null) { - appId_subQId_Map.put(appId, queueId.getSubQId()); - } - } } return true; } @@ -236,14 +226,6 @@ public boolean removeExpired(QueueImpl queue) { "Wrong QueueIds: %d != %d", queue.getQueueId(), qHandle.getQueueId())); } - - SubQueueIdInfo subQueueIdInfo = queue.getParameters().getSubIdInfo(); - if (subQueueIdInfo != null) { - String appId = subQueueIdInfo.appId(); - if (appId != null) { - appId_subQId_Map.remove(appId); - } - } } return true; } @@ -294,9 +276,37 @@ public QueueImpl findByUri(Uri uri) { } } - public Integer findSubQId(String appId) { + /** + * Find subQueueId of an expired queue by queue id and appId. + * + *

All subQueues of one canonical URI share the queue id, and each of them has its own appId, + * so the two together identify the queue. The same appId may be reopened with a greater + * subQueueId while the previous one is still expired; in that case the oldest subQueueId is + * returned, as responses come in the order of the requests. + * + *

Thread safe. + * + * @param qId queue id used for search + * @param appId appId used for search + * @return subQueueId if the queue is found, otherwise null. + */ + public Integer findSubQId(int qId, String appId) { synchronized (lock) { - return appId_subQId_Map.get(appId); + Integer subQId = null; + for (Map.Entry entry : expiredQueueMap.entrySet()) { + QueueId queueId = entry.getKey(); + if (queueId.getQId() != qId) { + continue; + } + SubQueueIdInfo subQueueIdInfo = entry.getValue().getParameters().getSubIdInfo(); + if (subQueueIdInfo == null || !Objects.equals(subQueueIdInfo.appId(), appId)) { + continue; + } + if (subQId == null || queueId.getSubQId() < subQId) { + subQId = queueId.getSubQId(); + } + } + return subQId; } } diff --git a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/QueueStateManager.java b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/QueueStateManager.java index ec5c39ce..8780381e 100644 --- a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/QueueStateManager.java +++ b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/QueueStateManager.java @@ -225,8 +225,8 @@ private QueueStateManager() { queueManager = QueueManager.createInstance(); } - public Integer findSubQId(String appId) { - return queueManager.findSubQId(appId); + public Integer findSubQId(int qId, String appId) { + return queueManager.findSubQId(qId, appId); } public QueueImpl findByQueueId(QueueId queueId) { diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/QueueManagerTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/QueueManagerTest.java index 51d5cfe3..d5fefecd 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/QueueManagerTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/QueueManagerTest.java @@ -348,7 +348,7 @@ private QueueImpl insertActive(QueueManager manager, Uri uri) { /** Resolve an expired queue the way LateResponseHandler does. */ private QueueImpl findExpiredByAppId(QueueManager manager, QueueImpl queue) { - Integer subQId = manager.findSubQId(queue.getUri().id()); + Integer subQId = manager.findSubQId(queue.getQueueId(), queue.getUri().id()); assertNotNull(subQId); return manager.findExpiredByQueueId(QueueId.createInstance(queue.getQueueId(), subQId)); }