From 80b65880b68e51a261cea93b865d6be23bad7511 Mon Sep 17 00:00:00 2001 From: Matt Pavlovich Date: Tue, 4 Aug 2026 10:04:38 -0500 Subject: [PATCH 1/2] Add test proving FutureBrokerInfo.get(timeout) ignores its timeout The timed get loop's condition uses || between the not-disposed check and the deadline check, so the loop runs until disposal regardless of the caller's timeout. Widens FutureBrokerInfo to package-private so the test drives the real class. Two of four scenarios fail until the condition is corrected. --- .../DemandForwardingBridgeSupport.java | 2 +- .../network/FutureBrokerInfoTimeoutTest.java | 160 ++++++++++++++++++ 2 files changed, 161 insertions(+), 1 deletion(-) create mode 100644 activemq-broker/src/test/java/org/apache/activemq/network/FutureBrokerInfoTimeoutTest.java diff --git a/activemq-broker/src/main/java/org/apache/activemq/network/DemandForwardingBridgeSupport.java b/activemq-broker/src/main/java/org/apache/activemq/network/DemandForwardingBridgeSupport.java index c8bb1386c5f..ca4b2d1fe1a 100644 --- a/activemq-broker/src/main/java/org/apache/activemq/network/DemandForwardingBridgeSupport.java +++ b/activemq-broker/src/main/java/org/apache/activemq/network/DemandForwardingBridgeSupport.java @@ -2007,7 +2007,7 @@ public void resetStats() { * Used to allow for async tasks to await receipt of the BrokerInfo from the local and * remote sides of the network bridge. */ - private static class FutureBrokerInfo implements Future { + static class FutureBrokerInfo implements Future { private final CountDownLatch slot = new CountDownLatch(1); private final AtomicBoolean disposed; diff --git a/activemq-broker/src/test/java/org/apache/activemq/network/FutureBrokerInfoTimeoutTest.java b/activemq-broker/src/test/java/org/apache/activemq/network/FutureBrokerInfoTimeoutTest.java new file mode 100644 index 00000000000..e3476a95a54 --- /dev/null +++ b/activemq-broker/src/test/java/org/apache/activemq/network/FutureBrokerInfoTimeoutTest.java @@ -0,0 +1,160 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.activemq.network; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertTrue; + +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.activemq.command.BrokerInfo; +import org.apache.activemq.network.DemandForwardingBridgeSupport.FutureBrokerInfo; +import org.junit.Test; + +/** + * Verifies that {@link DemandForwardingBridgeSupport.FutureBrokerInfo} + * honors the timeout passed to {@link FutureBrokerInfo#get(long, TimeUnit)}. + * + * The timed get loop must exit when EITHER the bridge is disposed OR the + * deadline expires. A faulty OR condition on the two checks keeps the loop + * alive as long as the bridge is not disposed, ignoring the caller's timeout + * entirely and parking bridge start threads indefinitely when the peer never + * delivers its BrokerInfo. + */ +public class FutureBrokerInfoTimeoutTest { + + /** + * No info, not disposed: get(200ms) must throw TimeoutException promptly + * rather than blocking until disposal. + */ + @Test(timeout = 10000) + public void testGetTimedThrowsTimeoutExceptionWithinTimeout() throws Exception { + AtomicBoolean disposed = new AtomicBoolean(false); + FutureBrokerInfo future = new FutureBrokerInfo(null, disposed); + + AtomicBoolean timedOut = new AtomicBoolean(false); + AtomicLong elapsed = new AtomicLong(-1); + + Thread t = new Thread(() -> { + long start = System.currentTimeMillis(); + try { + future.get(200, TimeUnit.MILLISECONDS); + } catch (TimeoutException e) { + timedOut.set(true); + } catch (Exception ignored) { + } finally { + elapsed.set(System.currentTimeMillis() - start); + } + }); + t.setDaemon(true); + t.start(); + t.join(3000); + + try { + assertFalse("get(200ms) should have returned within 3s but is still blocked " + + "- the timeout is being ignored", t.isAlive()); + } finally { + // unblock the daemon thread if the timeout was ignored + disposed.set(true); + } + + assertTrue("Expected TimeoutException from get(200ms)", timedOut.get()); + assertTrue("Timed out too early: " + elapsed.get() + "ms", elapsed.get() >= 180); + } + + /** + * Already disposed, no info: get with a long timeout must not wait out the + * full timeout - disposal exits the wait immediately and surfaces as + * TimeoutException (info is absent). + */ + @Test(timeout = 10000) + public void testGetTimedExitsPromptlyWhenAlreadyDisposed() throws Exception { + AtomicBoolean disposed = new AtomicBoolean(true); + FutureBrokerInfo future = new FutureBrokerInfo(null, disposed); + + AtomicBoolean timedOut = new AtomicBoolean(false); + + Thread t = new Thread(() -> { + try { + future.get(60, TimeUnit.SECONDS); + } catch (TimeoutException e) { + timedOut.set(true); + } catch (Exception ignored) { + } + }); + t.setDaemon(true); + t.start(); + t.join(3000); + + assertFalse("get(60s) on a disposed future should return immediately " + + "- it must not wait out the deadline", t.isAlive()); + assertTrue("Expected TimeoutException on disposed future with no info", timedOut.get()); + } + + /** + * Happy path: info already present - get returns it regardless of timeout. + */ + @Test(timeout = 10000) + public void testGetTimedReturnsInfoWhenAlreadySet() throws Exception { + AtomicBoolean disposed = new AtomicBoolean(false); + FutureBrokerInfo future = new FutureBrokerInfo(null, disposed); + + BrokerInfo brokerInfo = new BrokerInfo(); + brokerInfo.setBrokerName("test-broker"); + future.set(brokerInfo); + + BrokerInfo result = future.get(200, TimeUnit.MILLISECONDS); + assertNotNull(result); + assertEquals("test-broker", result.getBrokerName()); + } + + /** + * Info arrives mid-wait: get returns it promptly, well before the timeout. + */ + @Test(timeout = 10000) + public void testGetTimedReturnsPromptlyWhenInfoSetDuringWait() throws Exception { + AtomicBoolean disposed = new AtomicBoolean(false); + FutureBrokerInfo future = new FutureBrokerInfo(null, disposed); + + BrokerInfo brokerInfo = new BrokerInfo(); + brokerInfo.setBrokerName("late-broker"); + + AtomicReference result = new AtomicReference<>(); + Thread t = new Thread(() -> { + try { + result.set(future.get(30, TimeUnit.SECONDS)); + } catch (Exception ignored) { + } + }); + t.setDaemon(true); + t.start(); + + Thread.sleep(100); + future.set(brokerInfo); + + t.join(3000); + assertFalse("get should have returned promptly after set()", t.isAlive()); + assertNotNull(result.get()); + assertEquals("late-broker", result.get().getBrokerName()); + } +} From 7311a15f423d34ea0993f2885a1d02d5af111684 Mon Sep 17 00:00:00 2001 From: Matt Pavlovich Date: Tue, 4 Aug 2026 10:05:05 -0500 Subject: [PATCH 2/2] Fix FutureBrokerInfo.get(timeout) to honor its timeout Correct the loop condition from || to && so the timed get exits when EITHER the bridge is disposed OR the deadline expires. Previously a peer that never delivered its BrokerInfo parked the bridge start thread until disposal, ignoring the caller's timeout. --- .../apache/activemq/network/DemandForwardingBridgeSupport.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/activemq-broker/src/main/java/org/apache/activemq/network/DemandForwardingBridgeSupport.java b/activemq-broker/src/main/java/org/apache/activemq/network/DemandForwardingBridgeSupport.java index ca4b2d1fe1a..0ca4cc2fcea 100644 --- a/activemq-broker/src/main/java/org/apache/activemq/network/DemandForwardingBridgeSupport.java +++ b/activemq-broker/src/main/java/org/apache/activemq/network/DemandForwardingBridgeSupport.java @@ -2058,7 +2058,7 @@ public BrokerInfo get(long timeout, TimeUnit unit) throws InterruptedException, if (info == null) { long deadline = System.currentTimeMillis() + unit.toMillis(timeout); - while (!disposed.get() || System.currentTimeMillis() - deadline < 0) { + while (!disposed.get() && System.currentTimeMillis() - deadline < 0) { if (slot.await(1, TimeUnit.MILLISECONDS)) { break; }