From 0c62a61c6422e6bea2c5eb12bae3fa482db29d83 Mon Sep 17 00:00:00 2001 From: Evgeny Malygin Date: Fri, 24 Jul 2026 13:20:29 -0400 Subject: [PATCH] Fix: CloseQueueStrategy leaks the queue when the close-request write fails Signed-off-by: Evgeny Malygin --- .../bloomberg/bmq/impl/CloseQueueStrategy.java | 12 +++++++++++- .../bloomberg/bmq/impl/BrokerSessionTest.java | 17 +++++++---------- 2 files changed, 18 insertions(+), 11 deletions(-) diff --git a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/CloseQueueStrategy.java b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/CloseQueueStrategy.java index f767bc5b..38613b29 100644 --- a/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/CloseQueueStrategy.java +++ b/bmq-sdk/src/main/java/com/bloomberg/bmq/impl/CloseQueueStrategy.java @@ -261,7 +261,17 @@ private void handleConfigureStatus(StatusCategory confStatusCategory) { if (genericResult.isFailure()) { logger.error("Failed to send closeQueue request. RC={}", genericResult); - setQueueState(QueueState.e_CLOSED); + if (scenario.isDefault()) { + // If no connection the queue now is closed. + onFullyClosed(); + if (genericResult.isNotConnected()) { + genericResult = GenericResult.SUCCESS; + } + } else { + // The late scenarios have already removed the queue from the + // active maps, so just mark it closed (mirrors onCloseResponse()). + setQueueState(QueueState.e_CLOSED); + } resultHook(genericResult); } diff --git a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/BrokerSessionTest.java b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/BrokerSessionTest.java index 8836883e..8094bf07 100644 --- a/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/BrokerSessionTest.java +++ b/bmq-sdk/src/test/java/com/bloomberg/bmq/impl/BrokerSessionTest.java @@ -1254,7 +1254,9 @@ void closeQueueRequestFailsTest() { obj.sendConfigureStreamResponse(configureRequest); - assertEquals(CloseQueueResult.NOT_CONNECTED, closeFuture.get(FUTURE_TIMEOUT)); + // With no connection the queue is considered closed, so close + // reports success. + assertEquals(CloseQueueResult.SUCCESS, closeFuture.get(FUTURE_TIMEOUT)); // Check close request obj.verifyCloseQueueRequest(true); @@ -1266,21 +1268,16 @@ void closeQueueRequestFailsTest() { // Check event obj.verifyQueueControlEvent( - QueueControlEvent.Type.e_QUEUE_CLOSE_RESULT, CloseQueueResult.NOT_CONNECTED); + QueueControlEvent.Type.e_QUEUE_CLOSE_RESULT, CloseQueueResult.SUCCESS); assertEquals(QueueState.e_CLOSED, queue.getState()); // Open the queue obj.connection().resetWriteStatus(); // SUCCESS - // Due to the bug, queue is not closed fully, so we cannot reopen it - logger.info("[BUG] Current queue cannot be reopened"); - assertEquals( - OpenQueueResult.QUEUE_ID_NOT_UNIQUE, - queue.open(queueOptions, SEQUENCE_TIMEOUT)); - - // Create and open another queue - queue = obj.createQueue(createUri(), flags); + // The queue is fully closed and no longer registered in the active + // queue map, so the same queue can be reopened. + assertNull(obj.queueManager().findByQueueId(queue.getFullQueueId())); obj.openQueue(queue, queueOptions);