From 49fd688773329590916a89dc7ae4774635e4b8ba Mon Sep 17 00:00:00 2001 From: Bertrand Martin Date: Fri, 9 Oct 2026 16:58:55 +0200 Subject: [PATCH 1/4] Make a connection safe for its sending, receiving and keep-alive threads StateMachine serializes doTransition() and notifyMessage(), so a late reply cannot interleave with the timeout or close of the request it answers, and its state, address and port are volatile; its observers, the connection listeners and the connector listener lists are copy-on-write lists, so a listener may be registered or unregistered while they are notified. Connection.waitForResponse() rolls the state machine back outside the response lock, which the receiving thread takes while holding the state machine lock; ConnectionManager looks connections up without a lock for the same reason (#96). IntegrityAlgorithm and ConfidentialityAesCbc128 synchronize the use of their Mac and Cipher, which the sending, receiving and keep-alive threads share (#89). Co-Authored-By: Claude Fable 5.1 --- .../core/api/async/IpmiAsyncConnector.java | 36 +++-------- .../security/ConfidentialityAesCbc128.java | 7 +- .../coding/security/IntegrityAlgorithm.java | 7 +- .../ipmi/core/connection/Connection.java | 16 +++-- .../core/connection/ConnectionManager.java | 40 +++++------- .../metricshub/ipmi/core/sm/StateMachine.java | 20 +++--- src/site/markdown/upgrading.md | 6 +- .../api/async/IpmiAsyncConnectorTest.java | 64 +++++++++++++++++++ .../ConfidentialityAesCbc128Test.java | 44 +++++++++++++ .../security/IntegrityAlgorithmTest.java | 49 ++++++++++++++ .../ipmi/core/connection/ConnectionTest.java | 36 +++++++++++ 11 files changed, 254 insertions(+), 71 deletions(-) create mode 100644 src/test/java/org/metricshub/ipmi/core/api/async/IpmiAsyncConnectorTest.java create mode 100644 src/test/java/org/metricshub/ipmi/core/coding/security/ConfidentialityAesCbc128Test.java create mode 100644 src/test/java/org/metricshub/ipmi/core/coding/security/IntegrityAlgorithmTest.java diff --git a/src/main/java/org/metricshub/ipmi/core/api/async/IpmiAsyncConnector.java b/src/main/java/org/metricshub/ipmi/core/api/async/IpmiAsyncConnector.java index 3e9c037..15e7d71 100644 --- a/src/main/java/org/metricshub/ipmi/core/api/async/IpmiAsyncConnector.java +++ b/src/main/java/org/metricshub/ipmi/core/api/async/IpmiAsyncConnector.java @@ -44,8 +44,8 @@ import java.io.FileNotFoundException; import java.io.IOException; import java.net.InetAddress; -import java.util.ArrayList; import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.TimeUnit; import org.slf4j.Logger; @@ -87,8 +87,8 @@ public class IpmiAsyncConnector implements ConnectionListener { private ConnectionManager connectionManager; private SessionManager sessionManager; private int retries; - private final List responseListeners; - private final List inboundMessageListeners; + private final List responseListeners = new CopyOnWriteArrayList(); + private final List inboundMessageListeners = new CopyOnWriteArrayList(); private static Logger logger = LoggerFactory.getLogger(IpmiAsyncConnector.class); @@ -104,8 +104,6 @@ public class IpmiAsyncConnector implements ConnectionListener { * when properties file was not found */ public IpmiAsyncConnector(int port) throws IOException { - responseListeners = new ArrayList(); - inboundMessageListeners = new ArrayList(); connectionManager = new ConnectionManager(port); sessionManager = new SessionManager(); loadProperties(); @@ -126,8 +124,6 @@ public IpmiAsyncConnector(int port) throws IOException { * when properties file was not found */ public IpmiAsyncConnector(int port, InetAddress address) throws IOException { - responseListeners = new ArrayList(); - inboundMessageListeners = new ArrayList(); connectionManager = new ConnectionManager(port, address); sessionManager = new SessionManager(); loadProperties(); @@ -145,8 +141,6 @@ public IpmiAsyncConnector(int port, InetAddress address) throws IOException { * error. */ public IpmiAsyncConnector(int port, long pingPeriod) throws IOException { - responseListeners = new ArrayList<>(); - inboundMessageListeners = new ArrayList<>(); connectionManager = new ConnectionManager(port, pingPeriod); sessionManager = new SessionManager(); loadProperties(); @@ -475,9 +469,7 @@ public int retry(ConnectionHandle connectionHandle, int tag, PayloadType message * {@link IpmiResponseListener} to processResponse */ public void registerListener(IpmiResponseListener listener) { - synchronized (responseListeners) { - responseListeners.add(listener); - } + responseListeners.add(listener); } /** @@ -488,9 +480,7 @@ public void registerListener(IpmiResponseListener listener) { * - the {@link IpmiResponseListener} to unregister */ public void unregisterListener(IpmiResponseListener listener) { - synchronized (responseListeners) { - responseListeners.remove(listener); - } + responseListeners.remove(listener); } /** @@ -500,9 +490,7 @@ public void unregisterListener(IpmiResponseListener listener) { * the {@link InboundMessageListener} to register. */ public void registerIncomingPayloadListener(InboundMessageListener listener) { - synchronized (inboundMessageListeners) { - inboundMessageListeners.add(listener); - } + inboundMessageListeners.add(listener); } /** @@ -512,9 +500,7 @@ public void registerIncomingPayloadListener(InboundMessageListener listener) { * the {@link InboundMessageListener} to unregister. */ public void unregisterIncomingPayloadListener(InboundMessageListener listener) { - synchronized (inboundMessageListeners) { - inboundMessageListeners.remove(listener); - } + inboundMessageListeners.remove(listener); } @Override @@ -539,11 +525,9 @@ public void processResponse(ResponseData responseData, int handle, int tag, Exce new ConnectionHandle(handle, connection.getRemoteMachineAddress(), connection.getRemoteMachinePort())); } - synchronized (responseListeners) { - for (IpmiResponseListener listener : responseListeners) { - if (listener != null) { - listener.notify(response); - } + for (IpmiResponseListener listener : responseListeners) { + if (listener != null) { + listener.notify(response); } } } diff --git a/src/main/java/org/metricshub/ipmi/core/coding/security/ConfidentialityAesCbc128.java b/src/main/java/org/metricshub/ipmi/core/coding/security/ConfidentialityAesCbc128.java index 57cbf77..6f179bb 100644 --- a/src/main/java/org/metricshub/ipmi/core/coding/security/ConfidentialityAesCbc128.java +++ b/src/main/java/org/metricshub/ipmi/core/coding/security/ConfidentialityAesCbc128.java @@ -44,6 +44,7 @@ public class ConfidentialityAesCbc128 extends ConfidentialityAlgorithm { Arrays.fill(CONST2, (byte) 2); } + // The sending and receiving threads share this instance: encrypt() and decrypt() each init the cipher private Cipher cipher; private SecretKeySpec cipherKey; @@ -54,7 +55,7 @@ public byte getCode() { } @Override - public void initialize(byte[] sik, AuthenticationAlgorithm authenticationAlgorithm) + public synchronized void initialize(byte[] sik, AuthenticationAlgorithm authenticationAlgorithm) throws InvalidKeyException, NoSuchAlgorithmException, NoSuchPaddingException { @@ -78,7 +79,7 @@ public void initialize(byte[] sik, AuthenticationAlgorithm authenticationAlgorit } @Override - public byte[] encrypt(byte[] data) throws InvalidKeyException { + public synchronized byte[] encrypt(byte[] data) throws InvalidKeyException { int length = data.length + 17; int pad = 0; if (length % 16 != 0) { @@ -115,7 +116,7 @@ public byte[] encrypt(byte[] data) throws InvalidKeyException { } @Override - public byte[] decrypt(byte[] data) { + public synchronized byte[] decrypt(byte[] data) { byte[] decrypted = null; try { diff --git a/src/main/java/org/metricshub/ipmi/core/coding/security/IntegrityAlgorithm.java b/src/main/java/org/metricshub/ipmi/core/coding/security/IntegrityAlgorithm.java index b739adf..d5684e2 100644 --- a/src/main/java/org/metricshub/ipmi/core/coding/security/IntegrityAlgorithm.java +++ b/src/main/java/org/metricshub/ipmi/core/coding/security/IntegrityAlgorithm.java @@ -67,7 +67,7 @@ private IntegrityAlgorithm(Mac mac) { * @param key - Session Integrity Key calculated during the opening of the * session or user password if 'one-key' logins are enabled. */ - public void initialize(byte[] key) throws InvalidKeyException { + public synchronized void initialize(byte[] key) throws InvalidKeyException { this.sik = key; final String algorithmName = getAlgorithmName(); @@ -103,14 +103,15 @@ protected void setSik(byte[] sik) { public abstract byte getCode(); /** - * Creates AuthCode field for message. + * Creates AuthCode field for message. Synchronized: the sending, receiving and keep-alive threads share this + * instance and its {@link Mac}. * * @param base - data starting with the AuthType/Format field up to and * including the field that immediately precedes the AuthCode field * @return AuthCode field. Might be null if empty AuthCOde field is generated. * @see Rakp1#calculateSik(org.metricshub.ipmi.core.coding.commands.session.Rakp1ResponseData) */ - public byte[] generateAuthCode(final byte[] base) { + public synchronized byte[] generateAuthCode(final byte[] base) { if (sik == null) { throw new NullPointerException("Algorithm not initialized."); diff --git a/src/main/java/org/metricshub/ipmi/core/connection/Connection.java b/src/main/java/org/metricshub/ipmi/core/connection/Connection.java index 862772e..0ffce5d 100644 --- a/src/main/java/org/metricshub/ipmi/core/connection/Connection.java +++ b/src/main/java/org/metricshub/ipmi/core/connection/Connection.java @@ -78,6 +78,7 @@ import java.util.Map; import java.util.Timer; import java.util.TimerTask; +import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; @@ -140,7 +141,7 @@ public void setTimeout(int timeout) { public Connection(Messenger messenger, int handle) { stateMachine = new StateMachine(messenger); this.handle = handle; - listeners = new ArrayList(); + listeners = new CopyOnWriteArrayList(); timeout = Integer.parseInt(PropertiesManager.getInstance().getProperty("timeout")); messageHandlers = new EnumMap(PayloadType.class); currentSessionSequenceNumber = new AtomicInteger(0); @@ -327,6 +328,7 @@ public List getAvailableCipherSuites(int tag) throws Exception { private void waitForResponse() throws Exception { long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(timeout); + InterruptedException interrupted = null; synchronized (responseLock) { try { long remaining = deadline - System.nanoTime(); @@ -335,12 +337,16 @@ private void waitForResponse() throws Exception { remaining = deadline - System.nanoTime(); } } catch (InterruptedException e) { - // The caller gave up on us (Future.cancel): leave the state machine in a state that allows a retry - stateMachine.doTransition(new Timeout()); - Thread.currentThread().interrupt(); - throw e; + interrupted = e; } } + if (interrupted != null) { + // The caller gave up on us (Future.cancel): leave the state machine in a state that allows a retry. + // Outside the response lock: the receiving thread takes it while holding the state machine lock. + stateMachine.doTransition(new Timeout()); + Thread.currentThread().interrupt(); + throw interrupted; + } if (lastAction == null) { stateMachine.doTransition(new Timeout()); diff --git a/src/main/java/org/metricshub/ipmi/core/connection/ConnectionManager.java b/src/main/java/org/metricshub/ipmi/core/connection/ConnectionManager.java index b760af2..a668c80 100644 --- a/src/main/java/org/metricshub/ipmi/core/connection/ConnectionManager.java +++ b/src/main/java/org/metricshub/ipmi/core/connection/ConnectionManager.java @@ -34,13 +34,17 @@ import java.net.InetAddress; import java.util.ArrayList; import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; /** * Manages multiple {@link Connection}s */ public class ConnectionManager { private Messenger messenger; + // Copy-on-write: looked up without a lock by the receiving thread, which holds a state machine lock, while + // connect() takes a state machine lock under the list lock private List connections; + private final Object connectionsLock = new Object(); private static final Object SESSIONLESS_TAG_LOCK = new Object(); private static int sessionlessTag; @@ -104,7 +108,7 @@ public ConnectionManager(Messenger messenger) { } private void initialize() { - connections = new ArrayList(); + connections = new CopyOnWriteArrayList(); if (pingPeriod == -1) { pingPeriod = Long.parseLong(PropertiesManager.getInstance().getProperty("pingPeriod")); } @@ -123,11 +127,9 @@ long getPingPeriod() { * Closes all open connections and disconnects {@link UdpListener}. */ public void close() { - synchronized (connections) { - for (Connection connection : connections) { - if (connection != null && connection.isActive()) { - connection.disconnect(); - } + for (Connection connection : connections) { + if (connection != null && connection.isActive()) { + connection.disconnect(); } } messenger.closeConnection(); @@ -187,10 +189,7 @@ public static void freeTag(int tag) { * - index of the connection to return */ public Connection getConnection(int index) { - Connection connection; - synchronized (connections) { - connection = connections.get(index); - } + Connection connection = connections.get(index); if (connection == null) { throw new IllegalStateException("Connection " + index + " is closed"); } @@ -202,10 +201,7 @@ public Connection getConnection(int index) { * closed connection does nothing. */ public void closeConnection(int index) { - Connection connection; - synchronized (connections) { - connection = connections.set(index, null); - } + Connection connection = connections.set(index, null); if (connection != null) { connection.disconnect(); } @@ -220,14 +216,12 @@ public void closeConnection(int index) { * @return First {@link Connection} to the address or null if none found */ public Connection getConnection(InetAddress address, int port) { - synchronized (connections) { - for (Connection connection : connections) { - if (connection != null - && connection.isActive() - && connection.getRemoteMachineAddress().equals(address) - && connection.getRemoteMachinePort() == port) { - return connection; - } + for (Connection connection : connections) { + if (connection != null + && connection.isActive() + && connection.getRemoteMachineAddress().equals(address) + && connection.getRemoteMachinePort() == port) { + return connection; } } return null; @@ -254,7 +248,7 @@ public int createConnection(InetAddress address, int port, int connectionPingPer private int connect(InetAddress address, int port, long connectionPingPeriod, boolean skipCiphers) throws IOException { - synchronized (connections) { + synchronized (connectionsLock) { Connection connection = new Connection(messenger, connections.size()); connection.connect(address, port, connectionPingPeriod, skipCiphers); connections.add(connection); diff --git a/src/main/java/org/metricshub/ipmi/core/sm/StateMachine.java b/src/main/java/org/metricshub/ipmi/core/sm/StateMachine.java index a9890bb..ec42fe0 100644 --- a/src/main/java/org/metricshub/ipmi/core/sm/StateMachine.java +++ b/src/main/java/org/metricshub/ipmi/core/sm/StateMachine.java @@ -24,8 +24,8 @@ import java.io.IOException; import java.net.InetAddress; -import java.util.ArrayList; import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; import org.metricshub.ipmi.core.coding.rmcp.RmcpDecoder; import org.metricshub.ipmi.core.common.Constants; @@ -44,15 +44,15 @@ */ public class StateMachine implements UdpListener { - private List observers; + private final List observers = new CopyOnWriteArrayList(); - private State current; + private volatile State current; private Messenger messenger; - private InetAddress remoteMachineAddress; - private int remoteMachinePort; + private volatile InetAddress remoteMachineAddress; + private volatile int remoteMachinePort; - private boolean initialized; + private volatile boolean initialized; public State getCurrent() { return current; @@ -72,7 +72,6 @@ public void setCurrent(State current) { */ public StateMachine(Messenger messenger) { this.messenger = messenger; - observers = new ArrayList(); initialized = false; } @@ -154,7 +153,8 @@ public boolean isActive() { /** * Performs a {@link State} transition according to the event and - * {@link #current} state + * {@link #current} state. Transitions and received messages are serialized, so a late reply cannot interleave + * with the timeout or close of the request it answers. * * @param event * - {@link StateMachineEvent} invoking the transition @@ -162,7 +162,7 @@ public boolean isActive() { * - when machine was not yet started * @see #start(InetAddress, int) */ - public void doTransition(StateMachineEvent event) { + public synchronized void doTransition(StateMachineEvent event) { if (!initialized) { throw new NullPointerException("State machine not started"); } @@ -170,7 +170,7 @@ public void doTransition(StateMachineEvent event) { } @Override - public void notifyMessage(UdpMessage message) { + public synchronized void notifyMessage(UdpMessage message) { if (message.getAddress().equals(getRemoteMachineAddress()) && message.getPort() == getRemoteMachinePort()) { current.doAction(this, RmcpDecoder.decode(message.getMessage())); } diff --git a/src/site/markdown/upgrading.md b/src/site/markdown/upgrading.md index af3f1bb..f69505c 100644 --- a/src/site/markdown/upgrading.md +++ b/src/site/markdown/upgrading.md @@ -48,7 +48,11 @@ The `IpmiClient` API is unchanged, and the client is more tolerant of real-world whole connector; * the `PropertiesManager` lookups are logged at `DEBUG` instead of `INFO`, the unused `cleaningFrequency` property is gone from `connection.properties`, and `Constants.TIMEOUT`, - which nothing reads, is deprecated. + which nothing reads, is deprecated; +* the sending, receiving and keep-alive threads of a connection no longer race: the state + machine serializes transitions and received messages, the HMAC and AES objects of a cipher suite + are used by one thread at a time, and the listener lists can be changed while they are being + notified (a listener may unregister itself from its own callback). The decoders follow the IPMI 2.0 and FRU specifications more closely; the visible changes are: diff --git a/src/test/java/org/metricshub/ipmi/core/api/async/IpmiAsyncConnectorTest.java b/src/test/java/org/metricshub/ipmi/core/api/async/IpmiAsyncConnectorTest.java new file mode 100644 index 0000000..16dbfa5 --- /dev/null +++ b/src/test/java/org/metricshub/ipmi/core/api/async/IpmiAsyncConnectorTest.java @@ -0,0 +1,64 @@ +package org.metricshub.ipmi.core.api.async; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.net.InetAddress; +import java.util.concurrent.atomic.AtomicInteger; + +import org.junit.jupiter.api.Test; +import org.metricshub.ipmi.core.api.async.messages.IpmiResponse; +import org.metricshub.ipmi.core.coding.payload.IpmiPayload; +import org.metricshub.ipmi.core.coding.payload.PlainMessage; + +class IpmiAsyncConnectorTest { + + @Test + void aResponseListenerMayUnregisterItselfWhileNotified() throws Exception { + IpmiAsyncConnector connector = new IpmiAsyncConnector(0); + try { + ConnectionHandle handle = connector.createConnection(InetAddress.getLoopbackAddress(), 623); + AtomicInteger notified = new AtomicInteger(); + IpmiResponseListener oneShot = new IpmiResponseListener() { + @Override + public void notify(IpmiResponse response) { + connector.unregisterListener(this); + } + }; + connector.registerListener(oneShot); + connector.registerListener(response -> notified.incrementAndGet()); + + connector.processResponse(null, handle.getHandle(), 1, new Exception("timed out")); + connector.processResponse(null, handle.getHandle(), 2, new Exception("timed out")); + assertEquals(2, notified.get(), "the listener registered after the one-shot one must be notified each time"); + } finally { + connector.tearDown(); + } + } + + @Test + void anInboundListenerMayUnregisterItselfWhileNotified() throws Exception { + IpmiAsyncConnector connector = new IpmiAsyncConnector(0); + try { + AtomicInteger notified = new AtomicInteger(); + InboundMessageListener oneShot = new InboundMessageListener() { + @Override + public boolean isPayloadSupported(IpmiPayload payload) { + return true; + } + + @Override + public void notify(IpmiPayload payload) { + notified.incrementAndGet(); + connector.unregisterIncomingPayloadListener(this); + } + }; + connector.registerIncomingPayloadListener(oneShot); + connector.processRequest(new PlainMessage(new byte[0])); + connector.processRequest(new PlainMessage(new byte[0])); + assertTrue(notified.get() == 1, "the one-shot listener was notified " + notified.get() + " times"); + } finally { + connector.tearDown(); + } + } +} diff --git a/src/test/java/org/metricshub/ipmi/core/coding/security/ConfidentialityAesCbc128Test.java b/src/test/java/org/metricshub/ipmi/core/coding/security/ConfidentialityAesCbc128Test.java new file mode 100644 index 0000000..36f2023 --- /dev/null +++ b/src/test/java/org/metricshub/ipmi/core/coding/security/ConfidentialityAesCbc128Test.java @@ -0,0 +1,44 @@ +package org.metricshub.ipmi.core.coding.security; + +import static org.junit.jupiter.api.Assertions.assertArrayEquals; +import static org.junit.jupiter.api.Assertions.assertNull; + +import java.util.Arrays; +import java.util.concurrent.atomic.AtomicReference; + +import org.junit.jupiter.api.Test; + +class ConfidentialityAesCbc128Test { + + private static final int THREADS = 8; + + private static final int ROUNDS = 2000; + + @Test + void encryptAndDecryptRoundTripWhenSeveralThreadsShareTheAlgorithm() throws Exception { + ConfidentialityAesCbc128 algorithm = new ConfidentialityAesCbc128(); + algorithm.initialize(new byte[20], new AuthenticationRakpHmacSha1()); + + AtomicReference failure = new AtomicReference<>(); + Thread[] threads = new Thread[THREADS]; + for (int i = 0; i < THREADS; i++) { + final byte[] data = new byte[20 + i]; + Arrays.fill(data, (byte) (i + 1)); + threads[i] = new Thread(() -> { + for (int round = 0; round < ROUNDS; round++) { + try { + assertArrayEquals(data, algorithm.decrypt(algorithm.encrypt(data))); + } catch (Throwable e) { + failure.compareAndSet(null, e); + return; + } + } + }); + threads[i].start(); + } + for (Thread thread : threads) { + thread.join(); + } + assertNull(failure.get(), String.valueOf(failure.get())); + } +} diff --git a/src/test/java/org/metricshub/ipmi/core/coding/security/IntegrityAlgorithmTest.java b/src/test/java/org/metricshub/ipmi/core/coding/security/IntegrityAlgorithmTest.java new file mode 100644 index 0000000..667bbae --- /dev/null +++ b/src/test/java/org/metricshub/ipmi/core/coding/security/IntegrityAlgorithmTest.java @@ -0,0 +1,49 @@ +package org.metricshub.ipmi.core.coding.security; + +import static org.junit.jupiter.api.Assertions.assertArrayEquals; +import static org.junit.jupiter.api.Assertions.assertNull; + +import java.util.Arrays; +import java.util.concurrent.atomic.AtomicReference; + +import org.junit.jupiter.api.Test; + +class IntegrityAlgorithmTest { + + private static final int THREADS = 8; + + private static final int ROUNDS = 2000; + + @Test + void generateAuthCodeIsStableWhenSeveralThreadsShareTheAlgorithm() throws Exception { + IntegrityAlgorithm algorithm = new IntegrityHmacSha1_96(); + algorithm.initialize(new byte[20]); + byte[][] bases = new byte[THREADS][40]; + byte[][] expected = new byte[THREADS][]; + for (int i = 0; i < THREADS; i++) { + Arrays.fill(bases[i], (byte) (i + 1)); + expected[i] = algorithm.generateAuthCode(bases[i]); + } + + AtomicReference failure = new AtomicReference<>(); + Thread[] threads = new Thread[THREADS]; + for (int i = 0; i < THREADS; i++) { + final int id = i; + threads[i] = new Thread(() -> { + for (int round = 0; round < ROUNDS; round++) { + try { + assertArrayEquals(expected[id], algorithm.generateAuthCode(bases[id]), "thread " + id); + } catch (AssertionError e) { + failure.compareAndSet(null, e); + return; + } + } + }); + threads[i].start(); + } + for (Thread thread : threads) { + thread.join(); + } + assertNull(failure.get(), String.valueOf(failure.get())); + } +} diff --git a/src/test/java/org/metricshub/ipmi/core/connection/ConnectionTest.java b/src/test/java/org/metricshub/ipmi/core/connection/ConnectionTest.java index 5e2b387..b862b5c 100644 --- a/src/test/java/org/metricshub/ipmi/core/connection/ConnectionTest.java +++ b/src/test/java/org/metricshub/ipmi/core/connection/ConnectionTest.java @@ -199,4 +199,40 @@ public void processRequest(IpmiPayload payload) { connection.disconnect(); } } + + @Test + void aListenerMayUnregisterItselfWhileNotified() throws Exception { + Connection connection = connect(TIMEOUT_MS); + try { + AtomicInteger notified = new AtomicInteger(); + ConnectionListener oneShot = new ConnectionListener() { + @Override + public void processResponse(ResponseData responseData, int handle, int tag, Exception exception) { + connection.unregisterListener(this); + } + + @Override + public void processRequest(IpmiPayload payload) { + connection.unregisterListener(this); + } + }; + connection.registerListener(oneShot); + connection.registerListener(new ConnectionListener() { + @Override + public void processResponse(ResponseData responseData, int handle, int tag, Exception exception) { + notified.incrementAndGet(); + } + + @Override + public void processRequest(IpmiPayload payload) { + notified.incrementAndGet(); + } + }); + connection.notifyResponseListeners(0, 1, null, new Exception("timed out")); + connection.notifyResponseListeners(0, 2, null, new Exception("timed out")); + assertEquals(2, notified.get(), "the second listener must be notified each time"); + } finally { + connection.disconnect(); + } + } } From 870ef5d9466f46c0b8efbd2b34645307f08e3c70 Mon Sep 17 00:00:00 2001 From: Bertrand Martin Date: Fri, 9 Oct 2026 17:19:08 +0200 Subject: [PATCH 2/4] Notify the observers outside the state machine lock; close the manager atomically A state that emits an action during a transition or the processing of a received message queues it; doTransition() and notifyMessage() dispatch the queued actions to the observers once they release the lock, so no listener runs under the state machine monitor. ConnectionManager.close() disconnects the connections under the same lock as connect() and marks the manager closed; connect() then throws IllegalStateException instead of creating a connection on a closed messenger. Co-Authored-By: Claude Fable 5.1 --- .../core/connection/ConnectionManager.java | 13 ++++- .../metricshub/ipmi/core/sm/StateMachine.java | 56 ++++++++++++++++--- .../connection/ConnectionManagerTest.java | 11 ++++ .../ipmi/core/sm/StateMachineTest.java | 47 ++++++++++++++++ 4 files changed, 116 insertions(+), 11 deletions(-) create mode 100644 src/test/java/org/metricshub/ipmi/core/sm/StateMachineTest.java diff --git a/src/main/java/org/metricshub/ipmi/core/connection/ConnectionManager.java b/src/main/java/org/metricshub/ipmi/core/connection/ConnectionManager.java index a668c80..cf39a31 100644 --- a/src/main/java/org/metricshub/ipmi/core/connection/ConnectionManager.java +++ b/src/main/java/org/metricshub/ipmi/core/connection/ConnectionManager.java @@ -45,6 +45,7 @@ public class ConnectionManager { // connect() takes a state machine lock under the list lock private List connections; private final Object connectionsLock = new Object(); + private boolean closed; private static final Object SESSIONLESS_TAG_LOCK = new Object(); private static int sessionlessTag; @@ -127,9 +128,12 @@ long getPingPeriod() { * Closes all open connections and disconnects {@link UdpListener}. */ public void close() { - for (Connection connection : connections) { - if (connection != null && connection.isActive()) { - connection.disconnect(); + synchronized (connectionsLock) { + closed = true; + for (Connection connection : connections) { + if (connection != null && connection.isActive()) { + connection.disconnect(); + } } } messenger.closeConnection(); @@ -249,6 +253,9 @@ public int createConnection(InetAddress address, int port, int connectionPingPer private int connect(InetAddress address, int port, long connectionPingPeriod, boolean skipCiphers) throws IOException { synchronized (connectionsLock) { + if (closed) { + throw new IllegalStateException("The connection manager is closed"); + } Connection connection = new Connection(messenger, connections.size()); connection.connect(address, port, connectionPingPeriod, skipCiphers); connections.add(connection); diff --git a/src/main/java/org/metricshub/ipmi/core/sm/StateMachine.java b/src/main/java/org/metricshub/ipmi/core/sm/StateMachine.java index ec42fe0..49b4699 100644 --- a/src/main/java/org/metricshub/ipmi/core/sm/StateMachine.java +++ b/src/main/java/org/metricshub/ipmi/core/sm/StateMachine.java @@ -24,6 +24,7 @@ import java.io.IOException; import java.net.InetAddress; +import java.util.ArrayList; import java.util.List; import java.util.concurrent.CopyOnWriteArrayList; @@ -46,6 +47,12 @@ public class StateMachine implements UdpListener { private final List observers = new CopyOnWriteArrayList(); + /** + * Actions emitted by the states while the lock is held, dispatched by {@link #doTransition(StateMachineEvent)} + * and {@link #notifyMessage(UdpMessage)} once they release it, so that no observer runs under the lock. + */ + private final List pendingActions = new ArrayList(); + private volatile State current; private Messenger messenger; @@ -100,12 +107,21 @@ public int getRemoteMachinePort() { } /** - * Sends a notification of an action to all {@link MachineObserver}s + * Sends a notification of an action to all {@link MachineObserver}s. Called by a state during a transition, the + * notification is deferred until the transition releases the lock. * * @param action * - a {@link StateMachineAction} to perform */ public void doExternalAction(StateMachineAction action) { + if (Thread.holdsLock(this)) { + pendingActions.add(action); + } else { + notifyObservers(action); + } + } + + private void notifyObservers(StateMachineAction action) { for (MachineObserver observer : observers) { if (observer != null) { observer.notify(action); @@ -113,6 +129,19 @@ public void doExternalAction(StateMachineAction action) { } } + private void dispatch(List actions) { + for (StateMachineAction action : actions) { + notifyObservers(action); + } + } + + /** Returns the actions emitted so far and clears them; called under the lock. */ + private List drainPendingActions() { + List actions = new ArrayList(pendingActions); + pendingActions.clear(); + return actions; + } + /** * Sets the State Machine in the initial state. * @@ -154,7 +183,7 @@ public boolean isActive() { /** * Performs a {@link State} transition according to the event and * {@link #current} state. Transitions and received messages are serialized, so a late reply cannot interleave - * with the timeout or close of the request it answers. + * with the timeout or close of the request it answers; the observers are notified once the lock is released. * * @param event * - {@link StateMachineEvent} invoking the transition @@ -162,18 +191,29 @@ public boolean isActive() { * - when machine was not yet started * @see #start(InetAddress, int) */ - public synchronized void doTransition(StateMachineEvent event) { - if (!initialized) { - throw new NullPointerException("State machine not started"); + public void doTransition(StateMachineEvent event) { + List actions; + synchronized (this) { + if (!initialized) { + throw new NullPointerException("State machine not started"); + } + current.doTransition(this, event); + actions = drainPendingActions(); } - current.doTransition(this, event); + dispatch(actions); } @Override - public synchronized void notifyMessage(UdpMessage message) { - if (message.getAddress().equals(getRemoteMachineAddress()) && message.getPort() == getRemoteMachinePort()) { + public void notifyMessage(UdpMessage message) { + if (!message.getAddress().equals(getRemoteMachineAddress()) || message.getPort() != getRemoteMachinePort()) { + return; + } + List actions; + synchronized (this) { current.doAction(this, RmcpDecoder.decode(message.getMessage())); + actions = drainPendingActions(); } + dispatch(actions); } /** diff --git a/src/test/java/org/metricshub/ipmi/core/connection/ConnectionManagerTest.java b/src/test/java/org/metricshub/ipmi/core/connection/ConnectionManagerTest.java index 4ac6667..0933792 100644 --- a/src/test/java/org/metricshub/ipmi/core/connection/ConnectionManagerTest.java +++ b/src/test/java/org/metricshub/ipmi/core/connection/ConnectionManagerTest.java @@ -136,4 +136,15 @@ void everyOperationOnAReleasedHandleFailsTheSameWay() throws Exception { manager.close(); } } + + @Test + void aClosedManagerCreatesNoConnection() throws Exception { + ConnectionManager manager = new ConnectionManager(new SilentMessenger()); + int handle = manager.createConnection(InetAddress.getLoopbackAddress(), 623); + manager.close(); + assertFalse(manager.getConnection(handle).isActive(), "close() disconnects the connections"); + assertThrows( + IllegalStateException.class, + () -> manager.createConnection(InetAddress.getLoopbackAddress(), 623)); + } } diff --git a/src/test/java/org/metricshub/ipmi/core/sm/StateMachineTest.java b/src/test/java/org/metricshub/ipmi/core/sm/StateMachineTest.java new file mode 100644 index 0000000..4a5e552 --- /dev/null +++ b/src/test/java/org/metricshub/ipmi/core/sm/StateMachineTest.java @@ -0,0 +1,47 @@ +package org.metricshub.ipmi.core.sm; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.net.InetAddress; +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.atomic.AtomicBoolean; + +import org.junit.jupiter.api.Test; +import org.metricshub.ipmi.core.sm.actions.ErrorAction; +import org.metricshub.ipmi.core.sm.actions.StateMachineAction; +import org.metricshub.ipmi.core.sm.events.GetChannelCipherSuitesPending; +import org.metricshub.ipmi.core.sm.states.Uninitialized; +import org.metricshub.ipmi.core.transport.SilentMessenger; +import org.metricshub.ipmi.core.transport.UdpMessage; + +class StateMachineTest { + + @Test + void observersAreNotifiedOutsideTheLockAndTheStateIsRolledBackFirst() throws Exception { + StateMachine machine = new StateMachine(new SilentMessenger() { + @Override + public void send(UdpMessage message) { + throw new IllegalStateException("cable unplugged"); + } + }); + machine.start(InetAddress.getLoopbackAddress(), 623); + List actions = new CopyOnWriteArrayList<>(); + AtomicBoolean lockHeld = new AtomicBoolean(true); + AtomicBoolean rolledBack = new AtomicBoolean(); + machine.register(action -> { + actions.add(action); + lockHeld.set(Thread.holdsLock(machine)); + rolledBack.set(machine.getCurrent() instanceof Uninitialized); + }); + + machine.doTransition(new GetChannelCipherSuitesPending(1)); + + assertEquals(1, actions.size()); + assertTrue(actions.get(0) instanceof ErrorAction, String.valueOf(actions.get(0))); + assertFalse(lockHeld.get(), "the observer must not run under the state machine lock"); + assertTrue(rolledBack.get(), "the transition must be complete when the observer runs"); + } +} From b15047b762f5f367e6d14f3a47b92215d04b8c4f Mon Sep 17 00:00:00 2001 From: Bertrand Martin Date: Fri, 9 Oct 2026 17:34:08 +0200 Subject: [PATCH 3/4] Publish the handshake actions under the state machine lock Only an in-session message, which reaches the application listeners, is dispatched after the lock is released; a handshake reply, an error and the session key are published to the connection under the lock, so the caller sees the reply together with the state it produced and cannot time the request out in between. Co-Authored-By: Claude Fable 5.1 --- .../metricshub/ipmi/core/sm/StateMachine.java | 16 +++-- .../ipmi/core/sm/StateMachineTest.java | 67 +++++++++++++++---- 2 files changed, 64 insertions(+), 19 deletions(-) diff --git a/src/main/java/org/metricshub/ipmi/core/sm/StateMachine.java b/src/main/java/org/metricshub/ipmi/core/sm/StateMachine.java index 49b4699..babbffb 100644 --- a/src/main/java/org/metricshub/ipmi/core/sm/StateMachine.java +++ b/src/main/java/org/metricshub/ipmi/core/sm/StateMachine.java @@ -30,6 +30,7 @@ import org.metricshub.ipmi.core.coding.rmcp.RmcpDecoder; import org.metricshub.ipmi.core.common.Constants; +import org.metricshub.ipmi.core.sm.actions.MessageAction; import org.metricshub.ipmi.core.sm.actions.StateMachineAction; import org.metricshub.ipmi.core.sm.events.StateMachineEvent; import org.metricshub.ipmi.core.sm.states.SessionValid; @@ -48,8 +49,10 @@ public class StateMachine implements UdpListener { private final List observers = new CopyOnWriteArrayList(); /** - * Actions emitted by the states while the lock is held, dispatched by {@link #doTransition(StateMachineEvent)} - * and {@link #notifyMessage(UdpMessage)} once they release it, so that no observer runs under the lock. + * In-session messages received while the lock is held, dispatched by {@link #doTransition(StateMachineEvent)} + * and {@link #notifyMessage(UdpMessage)} once they release it, so that no application listener runs under the + * lock. The other actions (a handshake reply, an error, the session key) are published under the lock: the + * caller then sees the reply together with the state it produced, and cannot time the request out in between. */ private final List pendingActions = new ArrayList(); @@ -107,14 +110,14 @@ public int getRemoteMachinePort() { } /** - * Sends a notification of an action to all {@link MachineObserver}s. Called by a state during a transition, the - * notification is deferred until the transition releases the lock. + * Sends a notification of an action to all {@link MachineObserver}s. A {@link MessageAction} emitted by a state + * while the lock is held is deferred until the transition releases it; the other actions are published at once. * * @param action * - a {@link StateMachineAction} to perform */ public void doExternalAction(StateMachineAction action) { - if (Thread.holdsLock(this)) { + if (action instanceof MessageAction && Thread.holdsLock(this)) { pendingActions.add(action); } else { notifyObservers(action); @@ -183,7 +186,8 @@ public boolean isActive() { /** * Performs a {@link State} transition according to the event and * {@link #current} state. Transitions and received messages are serialized, so a late reply cannot interleave - * with the timeout or close of the request it answers; the observers are notified once the lock is released. + * with the timeout or close of the request it answers; the in-session messages are dispatched once the lock is + * released. * * @param event * - {@link StateMachineEvent} invoking the transition diff --git a/src/test/java/org/metricshub/ipmi/core/sm/StateMachineTest.java b/src/test/java/org/metricshub/ipmi/core/sm/StateMachineTest.java index 4a5e552..6b577fb 100644 --- a/src/test/java/org/metricshub/ipmi/core/sm/StateMachineTest.java +++ b/src/test/java/org/metricshub/ipmi/core/sm/StateMachineTest.java @@ -10,38 +10,79 @@ import java.util.concurrent.atomic.AtomicBoolean; import org.junit.jupiter.api.Test; +import org.metricshub.ipmi.core.coding.Encoder; +import org.metricshub.ipmi.core.coding.commands.IpmiVersion; +import org.metricshub.ipmi.core.coding.commands.session.GetChannelAuthenticationCapabilities; +import org.metricshub.ipmi.core.coding.protocol.encoder.Protocolv20Encoder; +import org.metricshub.ipmi.core.coding.security.CipherSuite; import org.metricshub.ipmi.core.sm.actions.ErrorAction; +import org.metricshub.ipmi.core.sm.actions.MessageAction; import org.metricshub.ipmi.core.sm.actions.StateMachineAction; import org.metricshub.ipmi.core.sm.events.GetChannelCipherSuitesPending; +import org.metricshub.ipmi.core.sm.states.SessionValid; import org.metricshub.ipmi.core.sm.states.Uninitialized; import org.metricshub.ipmi.core.transport.SilentMessenger; import org.metricshub.ipmi.core.transport.UdpMessage; class StateMachineTest { - @Test - void observersAreNotifiedOutsideTheLockAndTheStateIsRolledBackFirst() throws Exception { - StateMachine machine = new StateMachine(new SilentMessenger() { - @Override - public void send(UdpMessage message) { - throw new IllegalStateException("cable unplugged"); - } - }); + private static final int SESSION_ID = 1; + + private final List actions = new CopyOnWriteArrayList<>(); + + private final AtomicBoolean lockHeld = new AtomicBoolean(); + + private final AtomicBoolean rolledBack = new AtomicBoolean(); + + private StateMachine machine; + + private void start(SilentMessenger messenger) { + machine = new StateMachine(messenger); machine.start(InetAddress.getLoopbackAddress(), 623); - List actions = new CopyOnWriteArrayList<>(); - AtomicBoolean lockHeld = new AtomicBoolean(true); - AtomicBoolean rolledBack = new AtomicBoolean(); machine.register(action -> { actions.add(action); lockHeld.set(Thread.holdsLock(machine)); rolledBack.set(machine.getCurrent() instanceof Uninitialized); }); + } + + @Test + void aHandshakeErrorIsPublishedWithTheStateItProduced() { + start(new SilentMessenger() { + @Override + public void send(UdpMessage message) { + throw new IllegalStateException("cable unplugged"); + } + }); machine.doTransition(new GetChannelCipherSuitesPending(1)); assertEquals(1, actions.size()); assertTrue(actions.get(0) instanceof ErrorAction, String.valueOf(actions.get(0))); - assertFalse(lockHeld.get(), "the observer must not run under the state machine lock"); - assertTrue(rolledBack.get(), "the transition must be complete when the observer runs"); + assertTrue(rolledBack.get(), "the state must be rolled back when the error is published"); + assertTrue(lockHeld.get(), "a handshake action is published under the lock, before a timeout can be applied"); + } + + @Test + void anInSessionMessageIsDispatchedOutsideTheLock() throws Exception { + start(new SilentMessenger()); + machine.setCurrent(new SessionValid(CipherSuite.getEmpty(), SESSION_ID)); + byte[] raw = Encoder + .encode( + new Protocolv20Encoder(), + new GetChannelAuthenticationCapabilities(IpmiVersion.V20, IpmiVersion.V20, CipherSuite.getEmpty()), + 1, + 1, + SESSION_ID); + UdpMessage message = new UdpMessage(); + message.setAddress(InetAddress.getLoopbackAddress()); + message.setPort(623); + message.setMessage(raw); + + machine.notifyMessage(message); + + assertEquals(1, actions.size()); + assertTrue(actions.get(0) instanceof MessageAction, String.valueOf(actions.get(0))); + assertFalse(lockHeld.get(), "an in-session message must not be dispatched under the lock"); } } From d828034690a0c1a256673902db1e14fc482b44ed Mon Sep 17 00:00:00 2001 From: Bertrand Martin Date: Fri, 9 Oct 2026 17:55:21 +0200 Subject: [PATCH 4/4] Recheck the reply under the state machine lock before timing out The receiving thread publishes a reply under the state machine lock; waitForResponse() now takes that lock when the wait expires and applies the Timeout only if no reply was published meanwhile, so a reply that arrives as the deadline passes is used instead of being rolled back. Co-Authored-By: Claude Fable 5.1 --- .../ipmi/core/connection/Connection.java | 10 ++++- .../ipmi/core/connection/ConnectionTest.java | 42 +++++++++++++++++++ 2 files changed, 50 insertions(+), 2 deletions(-) diff --git a/src/main/java/org/metricshub/ipmi/core/connection/Connection.java b/src/main/java/org/metricshub/ipmi/core/connection/Connection.java index 84c200d..2491bf2 100644 --- a/src/main/java/org/metricshub/ipmi/core/connection/Connection.java +++ b/src/main/java/org/metricshub/ipmi/core/connection/Connection.java @@ -349,8 +349,14 @@ private void waitForResponse() throws Exception { } if (lastAction == null) { - stateMachine.doTransition(new Timeout()); - throw new ConnectionException("Command timed out"); + // The receiving thread publishes a reply under the state machine lock: once we hold it, a reply that + // is not there yet cannot be processed before the timeout rolls the state back + synchronized (stateMachine) { + if (lastAction == null) { + stateMachine.doTransition(new Timeout()); + throw new ConnectionException("Command timed out"); + } + } } if (!(lastAction instanceof ResponseAction || lastAction instanceof GetSikAction)) { if (lastAction instanceof ErrorAction) { diff --git a/src/test/java/org/metricshub/ipmi/core/connection/ConnectionTest.java b/src/test/java/org/metricshub/ipmi/core/connection/ConnectionTest.java index 1b6fb70..0292982 100644 --- a/src/test/java/org/metricshub/ipmi/core/connection/ConnectionTest.java +++ b/src/test/java/org/metricshub/ipmi/core/connection/ConnectionTest.java @@ -32,6 +32,10 @@ import org.metricshub.ipmi.core.sm.actions.MessageAction; import org.metricshub.ipmi.core.sm.states.SessionValid; import org.metricshub.ipmi.core.transport.UdpMessage; +import java.util.List; +import java.util.concurrent.CountDownLatch; +import org.metricshub.ipmi.core.coding.commands.session.GetChannelCipherSuitesResponseData; +import org.metricshub.ipmi.core.sm.actions.ResponseAction; class ConnectionTest { @@ -322,4 +326,42 @@ public void processRequest(IpmiPayload payload) { connection.disconnect(); } } + + @Test + void aReplyPublishedWhileTheTimeoutIsPendingIsNotLost() throws Exception { + Connection connection = connect(TIMEOUT_MS); + try { + Field field = Connection.class.getDeclaredField("stateMachine"); + field.setAccessible(true); + StateMachine machine = (StateMachine) field.get(connection); + AtomicReference publisherFailure = new AtomicReference<>(); + CountDownLatch published = new CountDownLatch(1); + // Plays the receiving thread: takes the state machine lock once the request is sent, keeps it past the + // deadline (the waiter times out meanwhile and blocks on the lock), then publishes the reply under it + Thread publisher = new Thread(() -> { + try { + Thread.sleep(TIMEOUT_MS / 2); + synchronized (machine) { + Thread.sleep(2 * TIMEOUT_MS); + GetChannelCipherSuitesResponseData data = new GetChannelCipherSuitesResponseData(); + data.setCipherSuiteData(new byte[0]); + connection.notify(new ResponseAction(data)); + published.countDown(); + } + } catch (Throwable t) { + publisherFailure.set(t); + } + }); + publisher.start(); + + List suites = connection.getAvailableCipherSuites(1); + + publisher.join(5000); + assertEquals(null, publisherFailure.get()); + assertEquals(0, published.getCount(), "the reply was published before the waiter could proceed"); + assertTrue(suites.isEmpty(), "the reply, not a timeout, ends the step"); + } finally { + connection.disconnect(); + } + } }