diff --git a/java/org/apache/catalina/tribes/ChannelReceiver.java b/java/org/apache/catalina/tribes/ChannelReceiver.java
index 24105d6e6782..706694dff58e 100644
--- a/java/org/apache/catalina/tribes/ChannelReceiver.java
+++ b/java/org/apache/catalina/tribes/ChannelReceiver.java
@@ -17,6 +17,7 @@
package org.apache.catalina.tribes;
import java.io.IOException;
+import java.util.concurrent.TimeUnit;
/**
* The ChannelReceiver interface is the data receiver component at the bottom layer, the IO layer (for
@@ -29,6 +30,11 @@ public interface ChannelReceiver extends Heartbeat {
*/
int MAX_UDP_SIZE = 65535;
+ /**
+ * Default timeout in milliseconds for {@link #waitForReady(long, TimeUnit)}.
+ */
+ long DEFAULT_READY_TIMEOUT_MS = 5000;
+
/**
* Start listening for incoming messages on the host/port
*
@@ -41,6 +47,23 @@ public interface ChannelReceiver extends Heartbeat {
*/
void stop();
+ /**
+ * Wait until the receiver is ready to accept connections, or the timeout expires.
+ *
+ * The default implementation returns immediately, preserving backward compatibility
+ * for receivers that do not implement readiness signaling. Implementations that
+ * start background listener threads should override this method to block until
+ * the listener thread has entered its accept/select loop.
+ *
+ * @param timeout the maximum time to wait
+ * @param unit the time unit of the timeout argument
+ * @return {@code true} if the receiver is ready; {@code false} if the timeout elapsed
+ * @throws InterruptedException if the current thread is interrupted while waiting
+ */
+ default boolean waitForReady(long timeout, TimeUnit unit) throws InterruptedException {
+ return true;
+ }
+
/**
* String representation of the IPv4 or IPv6 address that this host is listening to.
*
diff --git a/java/org/apache/catalina/tribes/group/ChannelCoordinator.java b/java/org/apache/catalina/tribes/group/ChannelCoordinator.java
index 3bf64ebfffa4..2fee33255af4 100644
--- a/java/org/apache/catalina/tribes/group/ChannelCoordinator.java
+++ b/java/org/apache/catalina/tribes/group/ChannelCoordinator.java
@@ -16,6 +16,8 @@
*/
package org.apache.catalina.tribes.group;
+import java.util.concurrent.TimeUnit;
+
import org.apache.catalina.tribes.Channel;
import org.apache.catalina.tribes.ChannelException;
import org.apache.catalina.tribes.ChannelMessage;
@@ -160,7 +162,22 @@ protected synchronized void internalStart(int svc) throws ChannelException {
clusterReceiver.setMessageListener(this);
clusterReceiver.setChannel(getChannel());
clusterReceiver.start();
- // synchronize, big time FIXME
+ // Wait for the receiver's background thread to enter the listen loop
+ // before reading the local member. Without this synchronization, there
+ // is a race window where start() has returned but the listener thread
+ // has not yet initialized, potentially causing getLocalMember() to
+ // observe an incomplete or null member state.
+ try {
+ boolean ready = clusterReceiver.waitForReady(
+ ChannelReceiver.DEFAULT_READY_TIMEOUT_MS, TimeUnit.MILLISECONDS);
+ if (!ready) {
+ throw new ChannelException(sm.getString("channelCoordinator.receiverNotReady",
+ Long.toString(ChannelReceiver.DEFAULT_READY_TIMEOUT_MS)));
+ }
+ } catch (InterruptedException ie) {
+ Thread.currentThread().interrupt();
+ throw new ChannelException(sm.getString("channelCoordinator.receiverWaitInterrupted"), ie);
+ }
Member localMember = getChannel().getLocalMember(false);
if (localMember instanceof StaticMember staticMember) {
// static member
diff --git a/java/org/apache/catalina/tribes/group/LocalStrings.properties b/java/org/apache/catalina/tribes/group/LocalStrings.properties
index 0f21a42e7ffa..7405bc73bab0 100644
--- a/java/org/apache/catalina/tribes/group/LocalStrings.properties
+++ b/java/org/apache/catalina/tribes/group/LocalStrings.properties
@@ -16,6 +16,8 @@
channelCoordinator.alreadyStarted=Channel already started for level:[{0}]
channelCoordinator.invalid.startLevel=Invalid start level, valid levels are:SND_RX_SEQ,SND_TX_SEQ,MBR_TX_SEQ,MBR_RX_SEQ
channelCoordinator.invalidState.notStopped=Configuration may not be changed until the channel has been fully stopped
+channelCoordinator.receiverNotReady=Channel receiver did not become ready within [{0}] ms during startup.
+channelCoordinator.receiverWaitInterrupted=Interrupted while waiting for channel receiver to become ready during startup.
groupChannel.listener.alreadyExist=Listener already exists:[{0}][{1}]
groupChannel.noDestination=No destination given
diff --git a/java/org/apache/catalina/tribes/transport/nio/NioReceiver.java b/java/org/apache/catalina/tribes/transport/nio/NioReceiver.java
index 8602710505fc..17b2adf93583 100644
--- a/java/org/apache/catalina/tribes/transport/nio/NioReceiver.java
+++ b/java/org/apache/catalina/tribes/transport/nio/NioReceiver.java
@@ -30,6 +30,8 @@
import java.util.Iterator;
import java.util.Set;
import java.util.concurrent.ConcurrentLinkedDeque;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import org.apache.catalina.tribes.io.ObjectReader;
@@ -59,6 +61,13 @@ public class NioReceiver extends ReceiverBase implements Runnable, NioReceiverMB
private ServerSocketChannel serverChannel = null;
private DatagramChannel datagramChannel = null;
+ /**
+ * Latch that counts down when the listener thread has entered the select loop.
+ * A count of 0 means the receiver is ready (or not started). A count of 1 means
+ * the listener thread is still initializing.
+ */
+ private volatile CountDownLatch readyLatch = new CountDownLatch(0);
+
/**
* Queue of events to be processed by the selector thread.
*/
@@ -92,6 +101,9 @@ public void start() throws IOException {
try {
getBind();
bind();
+ // Create a fresh latch with count 1 before launching the listener thread.
+ // The latch will be counted down in listen() after setListen(true).
+ readyLatch = new CountDownLatch(1);
String channelName = "";
if (getChannel().getName() != null) {
channelName = "[" + getChannel().getName() + "]";
@@ -100,6 +112,8 @@ public void start() throws IOException {
t.setDaemon(true);
t.start();
} catch (Exception e) {
+ // Reset latch to avoid blocking callers if start fails
+ readyLatch = new CountDownLatch(0);
log.fatal(sm.getString("nioReceiver.start.fail"), e);
if (e instanceof IOException) {
throw (IOException) e;
@@ -109,6 +123,19 @@ public void start() throws IOException {
}
}
+ /**
+ * Wait until the receiver's listener thread has entered the select loop.
+ *
+ * @param timeout the maximum time to wait
+ * @param unit the time unit of the timeout argument
+ * @return {@code true} if the receiver is ready; {@code false} if the timeout elapsed
+ * @throws InterruptedException if the current thread is interrupted while waiting
+ */
+ @Override
+ public boolean waitForReady(long timeout, TimeUnit unit) throws InterruptedException {
+ return readyLatch.await(timeout, unit);
+ }
+
@Override
public AbstractRxTask createRxTask() {
NioReplicationTask thread = new NioReplicationTask(this, this);
@@ -309,6 +336,10 @@ protected void listen() throws Exception {
setListen(true);
+ // Signal that the listener thread has entered the listen loop and is
+ // ready to accept connections. This must happen after setListen(true).
+ readyLatch.countDown();
+
// Avoid NPEs if selector is set to null on stop.
Selector selector = this.selector.get();
@@ -399,6 +430,9 @@ protected void listen() throws Exception {
*/
protected void stopListening() {
setListen(false);
+ // Reset the latch so that a subsequent start() can create a fresh one.
+ // A count of 0 means "not waiting" / "already ready".
+ readyLatch = new CountDownLatch(0);
Selector selector = this.selector.get();
if (selector != null) {
try {
diff --git a/test/org/apache/catalina/tribes/group/TestChannelCoordinatorStartupRace.java b/test/org/apache/catalina/tribes/group/TestChannelCoordinatorStartupRace.java
new file mode 100644
index 000000000000..e9316ac139fd
--- /dev/null
+++ b/test/org/apache/catalina/tribes/group/TestChannelCoordinatorStartupRace.java
@@ -0,0 +1,1059 @@
+/*
+ * 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.catalina.tribes.group;
+
+import java.lang.reflect.Field;
+import java.lang.reflect.Method;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.junit.Assert;
+import org.junit.Test;
+
+import org.apache.catalina.tribes.Channel;
+import org.apache.catalina.tribes.ChannelReceiver;
+import org.apache.catalina.tribes.Member;
+import org.apache.catalina.tribes.MembershipService;
+import org.apache.catalina.tribes.membership.StaticMember;
+
+/**
+ * Unit tests for the ChannelCoordinator startup race condition fix.
+ */
+public class TestChannelCoordinatorStartupRace {
+
+ @Test
+ public void testWaitForReadyDefaultMethodReturnsImmediately() throws Exception {
+ ChannelReceiver receiver = new ChannelReceiver() {
+ @Override
+ public void start() {
+ }
+
+ @Override
+ public void stop() {
+ }
+
+ @Override
+ public String getHost() {
+ return "127.0.0.1";
+ }
+
+ @Override
+ public int getPort() {
+ return 4000;
+ }
+
+ @Override
+ public int getSecurePort() {
+ return -1;
+ }
+
+ @Override
+ public int getUdpPort() {
+ return -1;
+ }
+
+ @Override
+ public void setMessageListener(org.apache.catalina.tribes.MessageListener listener) {
+ }
+
+ @Override
+ public org.apache.catalina.tribes.MessageListener getMessageListener() {
+ return null;
+ }
+
+ @Override
+ public org.apache.catalina.tribes.Channel getChannel() {
+ return null;
+ }
+
+ @Override
+ public void setChannel(org.apache.catalina.tribes.Channel channel) {
+ }
+
+ @Override
+ public void heartbeat() {
+ }
+ };
+
+ long start = System.nanoTime();
+ boolean ready = receiver.waitForReady(5000, TimeUnit.MILLISECONDS);
+ long elapsed = System.nanoTime() - start;
+
+ Assert.assertTrue("Default waitForReady should return true", ready);
+ Assert.assertTrue("Default waitForReady should return immediately",
+ elapsed < TimeUnit.SECONDS.toNanos(1));
+ }
+
+ @Test
+ public void testNioReceiverReadyLatchContract() throws Exception {
+ Class> nioReceiverClass = Class.forName(
+ "org.apache.catalina.tribes.transport.nio.NioReceiver");
+ Field latchField = nioReceiverClass.getDeclaredField("readyLatch");
+ latchField.setAccessible(true);
+
+ Object receiver = nioReceiverClass.getDeclaredConstructor().newInstance();
+ CountDownLatch initialLatch = (CountDownLatch) latchField.get(receiver);
+ Assert.assertEquals("Initial readyLatch should have count 0", 0, initialLatch.getCount());
+
+ Method waitForReady = nioReceiverClass.getMethod("waitForReady", long.class, TimeUnit.class);
+ boolean ready = (boolean) waitForReady.invoke(receiver, 100, TimeUnit.MILLISECONDS);
+ Assert.assertTrue("waitForReady on unstarted receiver should return true", ready);
+
+ CountDownLatch freshLatch = new CountDownLatch(1);
+ latchField.set(receiver, freshLatch);
+ Assert.assertEquals("After simulated start(), latch should have count 1", 1, freshLatch.getCount());
+
+ AtomicBoolean waitResult = new AtomicBoolean(false);
+ Thread waiter = new Thread(() -> {
+ try {
+ waitResult.set((boolean) waitForReady.invoke(receiver, 500, TimeUnit.MILLISECONDS));
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ });
+ waiter.start();
+
+ Thread.sleep(100);
+ Assert.assertFalse("waitForReady should still be blocking before countdown", waitResult.get());
+
+ freshLatch.countDown();
+
+ waiter.join(2000);
+ Assert.assertTrue("waitForReady should return true after countdown", waitResult.get());
+ Assert.assertEquals("Latch should have count 0 after countdown", 0, freshLatch.getCount());
+ }
+
+ @Test
+ public void testChannelCoordinatorWaitsForReceiverBeforeLocalMember() throws Exception {
+ AtomicBoolean waitForReadyCalled = new AtomicBoolean(false);
+ AtomicBoolean localMemberAccessedBeforeReady = new AtomicBoolean(false);
+ AtomicReference operationOrder = new AtomicReference<>("");
+
+ ChannelReceiver mockReceiver = new ChannelReceiver() {
+ @Override
+ public void start() {
+ operationOrder.set(operationOrder.get() + "start();");
+ }
+
+ @Override
+ public void stop() {
+ }
+
+ @Override
+ public String getHost() {
+ return "127.0.0.1";
+ }
+
+ @Override
+ public int getPort() {
+ return 4000;
+ }
+
+ @Override
+ public int getSecurePort() {
+ return -1;
+ }
+
+ @Override
+ public int getUdpPort() {
+ return -1;
+ }
+
+ @Override
+ public void setMessageListener(org.apache.catalina.tribes.MessageListener listener) {
+ }
+
+ @Override
+ public org.apache.catalina.tribes.MessageListener getMessageListener() {
+ return null;
+ }
+
+ @Override
+ public org.apache.catalina.tribes.Channel getChannel() {
+ return null;
+ }
+
+ @Override
+ public void setChannel(org.apache.catalina.tribes.Channel channel) {
+ }
+
+ @Override
+ public void heartbeat() {
+ }
+
+ @Override
+ public boolean waitForReady(long timeout, TimeUnit unit) {
+ waitForReadyCalled.set(true);
+ operationOrder.set(operationOrder.get() + "waitForReady();");
+ return true;
+ }
+ };
+
+ StaticMember localMember = new StaticMember();
+ localMember.setHost("127.0.0.1");
+ localMember.setPort(4000);
+
+ MembershipService mockMembership = new MembershipService() {
+ @Override
+ public void setProperties(java.util.Properties properties) {
+ }
+
+ @Override
+ public java.util.Properties getProperties() {
+ return new java.util.Properties();
+ }
+
+ @Override
+ public void start() {
+ }
+
+ @Override
+ public void start(int level) {
+ }
+
+ @Override
+ public void stop(int level) {
+ }
+
+ @Override
+ public boolean hasMembers() {
+ return false;
+ }
+
+ @Override
+ public Member getMember(Member mbr) {
+ return null;
+ }
+
+ @Override
+ public Member[] getMembers() {
+ return new Member[0];
+ }
+
+ @Override
+ public Member getLocalMember(boolean incAliveTime) {
+ operationOrder.set(operationOrder.get() + "getLocalMember();");
+ if (!waitForReadyCalled.get()) {
+ localMemberAccessedBeforeReady.set(true);
+ }
+ return localMember;
+ }
+
+ @Override
+ public String[] getMembersByName() {
+ return new String[0];
+ }
+
+ @Override
+ public Member findMemberByName(String name) {
+ return null;
+ }
+
+ @Override
+ public void setLocalMemberProperties(String listenHost, int listenPort,
+ int securePort, int udpPort) {
+ operationOrder.set(operationOrder.get() + "setLocalMemberProperties();");
+ }
+
+ @Override
+ public void setMembershipListener(org.apache.catalina.tribes.MembershipListener listener) {
+ }
+
+ @Override
+ public void removeMembershipListener() {
+ }
+
+ @Override
+ public void setPayload(byte[] payload) {
+ }
+
+ @Override
+ public void setDomain(byte[] domain) {
+ }
+
+ @Override
+ public void broadcast(org.apache.catalina.tribes.ChannelMessage message) {
+ }
+
+ @Override
+ public org.apache.catalina.tribes.Channel getChannel() {
+ return null;
+ }
+
+ @Override
+ public void setChannel(org.apache.catalina.tribes.Channel channel) {
+ }
+
+ @Override
+ public org.apache.catalina.tribes.MembershipProvider getMembershipProvider() {
+ return null;
+ }
+ };
+
+ org.apache.catalina.tribes.ChannelSender mockSender =
+ new org.apache.catalina.tribes.ChannelSender() {
+ @Override
+ public void start() {
+ }
+
+ @Override
+ public void stop() {
+ }
+
+ @Override
+ public void heartbeat() {
+ }
+
+ @Override
+ public void add(org.apache.catalina.tribes.Member member) {
+ }
+
+ @Override
+ public void remove(org.apache.catalina.tribes.Member member) {
+ }
+
+ @Override
+ public void sendMessage(org.apache.catalina.tribes.ChannelMessage message,
+ org.apache.catalina.tribes.Member[] destination) {
+ }
+
+ @Override
+ public org.apache.catalina.tribes.Channel getChannel() {
+ return null;
+ }
+
+ @Override
+ public void setChannel(org.apache.catalina.tribes.Channel channel) {
+ }
+ };
+
+ ChannelCoordinator coordinator = new ChannelCoordinator(
+ mockReceiver, mockSender, mockMembership);
+
+ Class> interceptorBase = Class.forName(
+ "org.apache.catalina.tribes.group.ChannelInterceptorBase");
+ Field channelField = interceptorBase.getDeclaredField("channel");
+ channelField.setAccessible(true);
+
+ org.apache.catalina.tribes.Channel mockChannel =
+ new org.apache.catalina.tribes.Channel() {
+ @Override
+ public void addInterceptor(org.apache.catalina.tribes.ChannelInterceptor interceptor) {
+ }
+
+ @Override
+ public void start(int svc) {
+ }
+
+ @Override
+ public void stop(int svc) {
+ }
+
+ @Override
+ public org.apache.catalina.tribes.UniqueId send(Member[] destination,
+ java.io.Serializable msg, int options) {
+ return null;
+ }
+
+ @Override
+ public org.apache.catalina.tribes.UniqueId send(Member[] destination,
+ java.io.Serializable msg, int options,
+ org.apache.catalina.tribes.ErrorHandler handler) {
+ return null;
+ }
+
+ @Override
+ public void heartbeat() {
+ }
+
+ @Override
+ public void setHeartbeat(boolean enable) {
+ }
+
+ @Override
+ public void addMembershipListener(
+ org.apache.catalina.tribes.MembershipListener listener) {
+ }
+
+ @Override
+ public void addChannelListener(org.apache.catalina.tribes.ChannelListener listener) {
+ }
+
+ @Override
+ public void removeMembershipListener(
+ org.apache.catalina.tribes.MembershipListener listener) {
+ }
+
+ @Override
+ public void removeChannelListener(org.apache.catalina.tribes.ChannelListener listener) {
+ }
+
+ @Override
+ public boolean hasMembers() {
+ return false;
+ }
+
+ @Override
+ public Member[] getMembers() {
+ return new Member[0];
+ }
+
+ @Override
+ public Member getLocalMember(boolean incAlive) {
+ return mockMembership.getLocalMember(incAlive);
+ }
+
+ @Override
+ public Member getMember(Member mbr) {
+ return null;
+ }
+
+ @Override
+ public String getName() {
+ return "test-channel";
+ }
+
+ @Override
+ public void setName(String name) {
+ }
+
+ @Override
+ public java.util.concurrent.ScheduledExecutorService getUtilityExecutor() {
+ return null;
+ }
+
+ @Override
+ public void setUtilityExecutor(
+ java.util.concurrent.ScheduledExecutorService utilityExecutor) {
+ }
+ };
+ channelField.set(coordinator, mockChannel);
+
+ Method internalStart = ChannelCoordinator.class.getDeclaredMethod("internalStart", int.class);
+ internalStart.setAccessible(true);
+ internalStart.invoke(coordinator, Channel.SND_RX_SEQ);
+
+ Assert.assertTrue("waitForReady must be called during startup", waitForReadyCalled.get());
+ Assert.assertFalse("getLocalMember must NOT be accessed before waitForReady",
+ localMemberAccessedBeforeReady.get());
+
+ String order = operationOrder.get();
+ Assert.assertEquals(
+ "start();waitForReady();getLocalMember();setLocalMemberProperties();", order);
+ }
+
+ @Test
+ public void testChannelCoordinatorThrowsWhenReceiverNotReady() throws Exception {
+ ChannelReceiver slowReceiver = new ChannelReceiver() {
+ @Override
+ public void start() {
+ }
+
+ @Override
+ public void stop() {
+ }
+
+ @Override
+ public String getHost() {
+ return "127.0.0.1";
+ }
+
+ @Override
+ public int getPort() {
+ return 4000;
+ }
+
+ @Override
+ public int getSecurePort() {
+ return -1;
+ }
+
+ @Override
+ public int getUdpPort() {
+ return -1;
+ }
+
+ @Override
+ public void setMessageListener(org.apache.catalina.tribes.MessageListener listener) {
+ }
+
+ @Override
+ public org.apache.catalina.tribes.MessageListener getMessageListener() {
+ return null;
+ }
+
+ @Override
+ public org.apache.catalina.tribes.Channel getChannel() {
+ return null;
+ }
+
+ @Override
+ public void setChannel(org.apache.catalina.tribes.Channel channel) {
+ }
+
+ @Override
+ public void heartbeat() {
+ }
+
+ @Override
+ public boolean waitForReady(long timeout, TimeUnit unit) {
+ return false;
+ }
+ };
+
+ MembershipService mockMembership = new MembershipService() {
+ @Override
+ public void setProperties(java.util.Properties properties) {
+ }
+
+ @Override
+ public java.util.Properties getProperties() {
+ return new java.util.Properties();
+ }
+
+ @Override
+ public void start() {
+ }
+
+ @Override
+ public void start(int level) {
+ }
+
+ @Override
+ public void stop(int level) {
+ }
+
+ @Override
+ public boolean hasMembers() {
+ return false;
+ }
+
+ @Override
+ public Member getMember(Member mbr) {
+ return null;
+ }
+
+ @Override
+ public Member[] getMembers() {
+ return new Member[0];
+ }
+
+ @Override
+ public Member getLocalMember(boolean incAliveTime) {
+ Assert.fail("getLocalMember should NOT be called when receiver is not ready");
+ return null;
+ }
+
+ @Override
+ public String[] getMembersByName() {
+ return new String[0];
+ }
+
+ @Override
+ public Member findMemberByName(String name) {
+ return null;
+ }
+
+ @Override
+ public void setLocalMemberProperties(String listenHost, int listenPort,
+ int securePort, int udpPort) {
+ }
+
+ @Override
+ public void setMembershipListener(org.apache.catalina.tribes.MembershipListener listener) {
+ }
+
+ @Override
+ public void removeMembershipListener() {
+ }
+
+ @Override
+ public void setPayload(byte[] payload) {
+ }
+
+ @Override
+ public void setDomain(byte[] domain) {
+ }
+
+ @Override
+ public void broadcast(org.apache.catalina.tribes.ChannelMessage message) {
+ }
+
+ @Override
+ public org.apache.catalina.tribes.Channel getChannel() {
+ return null;
+ }
+
+ @Override
+ public void setChannel(org.apache.catalina.tribes.Channel channel) {
+ }
+
+ @Override
+ public org.apache.catalina.tribes.MembershipProvider getMembershipProvider() {
+ return null;
+ }
+ };
+
+ org.apache.catalina.tribes.ChannelSender mockSender =
+ new org.apache.catalina.tribes.ChannelSender() {
+ @Override
+ public void start() {
+ }
+
+ @Override
+ public void stop() {
+ }
+
+ @Override
+ public void heartbeat() {
+ }
+
+ @Override
+ public void add(org.apache.catalina.tribes.Member member) {
+ }
+
+ @Override
+ public void remove(org.apache.catalina.tribes.Member member) {
+ }
+
+ @Override
+ public void sendMessage(org.apache.catalina.tribes.ChannelMessage message,
+ org.apache.catalina.tribes.Member[] destination) {
+ }
+
+ @Override
+ public org.apache.catalina.tribes.Channel getChannel() {
+ return null;
+ }
+
+ @Override
+ public void setChannel(org.apache.catalina.tribes.Channel channel) {
+ }
+ };
+
+ ChannelCoordinator coordinator = new ChannelCoordinator(
+ slowReceiver, mockSender, mockMembership);
+
+ Class> interceptorBase = Class.forName(
+ "org.apache.catalina.tribes.group.ChannelInterceptorBase");
+ Field channelField = interceptorBase.getDeclaredField("channel");
+ channelField.setAccessible(true);
+ channelField.set(coordinator, new org.apache.catalina.tribes.Channel() {
+ @Override
+ public void addInterceptor(org.apache.catalina.tribes.ChannelInterceptor interceptor) {
+ }
+
+ @Override
+ public void start(int svc) {
+ }
+
+ @Override
+ public void stop(int svc) {
+ }
+
+ @Override
+ public org.apache.catalina.tribes.UniqueId send(Member[] destination,
+ java.io.Serializable msg, int options) {
+ return null;
+ }
+
+ @Override
+ public org.apache.catalina.tribes.UniqueId send(Member[] destination,
+ java.io.Serializable msg, int options,
+ org.apache.catalina.tribes.ErrorHandler handler) {
+ return null;
+ }
+
+ @Override
+ public void heartbeat() {
+ }
+
+ @Override
+ public void setHeartbeat(boolean enable) {
+ }
+
+ @Override
+ public void addMembershipListener(org.apache.catalina.tribes.MembershipListener listener) {
+ }
+
+ @Override
+ public void addChannelListener(org.apache.catalina.tribes.ChannelListener listener) {
+ }
+
+ @Override
+ public void removeMembershipListener(org.apache.catalina.tribes.MembershipListener listener) {
+ }
+
+ @Override
+ public void removeChannelListener(org.apache.catalina.tribes.ChannelListener listener) {
+ }
+
+ @Override
+ public boolean hasMembers() {
+ return false;
+ }
+
+ @Override
+ public Member[] getMembers() {
+ return new Member[0];
+ }
+
+ @Override
+ public Member getLocalMember(boolean incAlive) {
+ return mockMembership.getLocalMember(incAlive);
+ }
+
+ @Override
+ public Member getMember(Member mbr) {
+ return null;
+ }
+
+ @Override
+ public String getName() {
+ return "test";
+ }
+
+ @Override
+ public void setName(String name) {
+ }
+
+ @Override
+ public java.util.concurrent.ScheduledExecutorService getUtilityExecutor() {
+ return null;
+ }
+
+ @Override
+ public void setUtilityExecutor(java.util.concurrent.ScheduledExecutorService utilityExecutor) {
+ }
+ });
+
+ Method internalStart = ChannelCoordinator.class.getDeclaredMethod("internalStart", int.class);
+ internalStart.setAccessible(true);
+
+ try {
+ internalStart.invoke(coordinator, Channel.SND_RX_SEQ);
+ Assert.fail("Expected ChannelException when receiver fails to become ready");
+ } catch (java.lang.reflect.InvocationTargetException e) {
+ Throwable cause = e.getCause();
+ Assert.assertTrue("Cause should be ChannelException, was: " + cause.getClass(),
+ cause instanceof org.apache.catalina.tribes.ChannelException);
+ }
+ }
+
+ @Test
+ public void testChannelCoordinatorHandlesInterruptedException() throws Exception {
+ AtomicBoolean interrupted = new AtomicBoolean(false);
+
+ ChannelReceiver interruptingReceiver = new ChannelReceiver() {
+ @Override
+ public void start() {
+ }
+
+ @Override
+ public void stop() {
+ }
+
+ @Override
+ public String getHost() {
+ return "127.0.0.1";
+ }
+
+ @Override
+ public int getPort() {
+ return 4000;
+ }
+
+ @Override
+ public int getSecurePort() {
+ return -1;
+ }
+
+ @Override
+ public int getUdpPort() {
+ return -1;
+ }
+
+ @Override
+ public void setMessageListener(org.apache.catalina.tribes.MessageListener listener) {
+ }
+
+ @Override
+ public org.apache.catalina.tribes.MessageListener getMessageListener() {
+ return null;
+ }
+
+ @Override
+ public org.apache.catalina.tribes.Channel getChannel() {
+ return null;
+ }
+
+ @Override
+ public void setChannel(org.apache.catalina.tribes.Channel channel) {
+ }
+
+ @Override
+ public void heartbeat() {
+ }
+
+ @Override
+ public boolean waitForReady(long timeout, TimeUnit unit) throws InterruptedException {
+ interrupted.set(true);
+ throw new InterruptedException("Simulated interrupt during waitForReady");
+ }
+ };
+
+ MembershipService mockMembership = new MembershipService() {
+ @Override
+ public void setProperties(java.util.Properties properties) {
+ }
+
+ @Override
+ public java.util.Properties getProperties() {
+ return new java.util.Properties();
+ }
+
+ @Override
+ public void start() {
+ }
+
+ @Override
+ public void start(int level) {
+ }
+
+ @Override
+ public void stop(int level) {
+ }
+
+ @Override
+ public boolean hasMembers() {
+ return false;
+ }
+
+ @Override
+ public Member getMember(Member mbr) {
+ return null;
+ }
+
+ @Override
+ public Member[] getMembers() {
+ return new Member[0];
+ }
+
+ @Override
+ public Member getLocalMember(boolean incAliveTime) {
+ return null;
+ }
+
+ @Override
+ public String[] getMembersByName() {
+ return new String[0];
+ }
+
+ @Override
+ public Member findMemberByName(String name) {
+ return null;
+ }
+
+ @Override
+ public void setLocalMemberProperties(String listenHost, int listenPort,
+ int securePort, int udpPort) {
+ }
+
+ @Override
+ public void setMembershipListener(org.apache.catalina.tribes.MembershipListener listener) {
+ }
+
+ @Override
+ public void removeMembershipListener() {
+ }
+
+ @Override
+ public void setPayload(byte[] payload) {
+ }
+
+ @Override
+ public void setDomain(byte[] domain) {
+ }
+
+ @Override
+ public void broadcast(org.apache.catalina.tribes.ChannelMessage message) {
+ }
+
+ @Override
+ public org.apache.catalina.tribes.Channel getChannel() {
+ return null;
+ }
+
+ @Override
+ public void setChannel(org.apache.catalina.tribes.Channel channel) {
+ }
+
+ @Override
+ public org.apache.catalina.tribes.MembershipProvider getMembershipProvider() {
+ return null;
+ }
+ };
+
+ org.apache.catalina.tribes.ChannelSender mockSender =
+ new org.apache.catalina.tribes.ChannelSender() {
+ @Override
+ public void start() {
+ }
+
+ @Override
+ public void stop() {
+ }
+
+ @Override
+ public void heartbeat() {
+ }
+
+ @Override
+ public void add(org.apache.catalina.tribes.Member member) {
+ }
+
+ @Override
+ public void remove(org.apache.catalina.tribes.Member member) {
+ }
+
+ @Override
+ public void sendMessage(org.apache.catalina.tribes.ChannelMessage message,
+ org.apache.catalina.tribes.Member[] destination) {
+ }
+
+ @Override
+ public org.apache.catalina.tribes.Channel getChannel() {
+ return null;
+ }
+
+ @Override
+ public void setChannel(org.apache.catalina.tribes.Channel channel) {
+ }
+ };
+
+ ChannelCoordinator coordinator = new ChannelCoordinator(
+ interruptingReceiver, mockSender, mockMembership);
+
+ Class> interceptorBase = Class.forName(
+ "org.apache.catalina.tribes.group.ChannelInterceptorBase");
+ Field channelField = interceptorBase.getDeclaredField("channel");
+ channelField.setAccessible(true);
+ channelField.set(coordinator, new org.apache.catalina.tribes.Channel() {
+ @Override
+ public void addInterceptor(org.apache.catalina.tribes.ChannelInterceptor interceptor) {
+ }
+
+ @Override
+ public void start(int svc) {
+ }
+
+ @Override
+ public void stop(int svc) {
+ }
+
+ @Override
+ public org.apache.catalina.tribes.UniqueId send(Member[] destination,
+ java.io.Serializable msg, int options) {
+ return null;
+ }
+
+ @Override
+ public org.apache.catalina.tribes.UniqueId send(Member[] destination,
+ java.io.Serializable msg, int options,
+ org.apache.catalina.tribes.ErrorHandler handler) {
+ return null;
+ }
+
+ @Override
+ public void heartbeat() {
+ }
+
+ @Override
+ public void setHeartbeat(boolean enable) {
+ }
+
+ @Override
+ public void addMembershipListener(org.apache.catalina.tribes.MembershipListener listener) {
+ }
+
+ @Override
+ public void addChannelListener(org.apache.catalina.tribes.ChannelListener listener) {
+ }
+
+ @Override
+ public void removeMembershipListener(org.apache.catalina.tribes.MembershipListener listener) {
+ }
+
+ @Override
+ public void removeChannelListener(org.apache.catalina.tribes.ChannelListener listener) {
+ }
+
+ @Override
+ public boolean hasMembers() {
+ return false;
+ }
+
+ @Override
+ public Member[] getMembers() {
+ return new Member[0];
+ }
+
+ @Override
+ public Member getLocalMember(boolean incAlive) {
+ return null;
+ }
+
+ @Override
+ public Member getMember(Member mbr) {
+ return null;
+ }
+
+ @Override
+ public String getName() {
+ return "test";
+ }
+
+ @Override
+ public void setName(String name) {
+ }
+
+ @Override
+ public java.util.concurrent.ScheduledExecutorService getUtilityExecutor() {
+ return null;
+ }
+
+ @Override
+ public void setUtilityExecutor(java.util.concurrent.ScheduledExecutorService utilityExecutor) {
+ }
+ });
+
+ Method internalStart = ChannelCoordinator.class.getDeclaredMethod("internalStart", int.class);
+ internalStart.setAccessible(true);
+
+ try {
+ internalStart.invoke(coordinator, Channel.SND_RX_SEQ);
+ Assert.fail("Expected ChannelException when waitForReady is interrupted");
+ } catch (java.lang.reflect.InvocationTargetException e) {
+ Throwable cause = e.getCause();
+ Assert.assertTrue("Cause should be ChannelException, was: " + cause.getClass(),
+ cause instanceof org.apache.catalina.tribes.ChannelException);
+ Assert.assertTrue("Interrupt should have been triggered", interrupted.get());
+ }
+ }
+}