Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -87,8 +87,8 @@ public class IpmiAsyncConnector implements ConnectionListener {
private ConnectionManager connectionManager;
private SessionManager sessionManager;
private int retries;
private final List<IpmiResponseListener> responseListeners;
private final List<InboundMessageListener> inboundMessageListeners;
private final List<IpmiResponseListener> responseListeners = new CopyOnWriteArrayList<IpmiResponseListener>();
private final List<InboundMessageListener> inboundMessageListeners = new CopyOnWriteArrayList<InboundMessageListener>();

private static Logger logger = LoggerFactory.getLogger(IpmiAsyncConnector.class);

Expand All @@ -104,8 +104,6 @@ public class IpmiAsyncConnector implements ConnectionListener {
* when properties file was not found
*/
public IpmiAsyncConnector(int port) throws IOException {
responseListeners = new ArrayList<IpmiResponseListener>();
inboundMessageListeners = new ArrayList<InboundMessageListener>();
connectionManager = new ConnectionManager(port);
sessionManager = new SessionManager();
loadProperties();
Expand All @@ -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<IpmiResponseListener>();
inboundMessageListeners = new ArrayList<InboundMessageListener>();
connectionManager = new ConnectionManager(port, address);
sessionManager = new SessionManager();
loadProperties();
Expand All @@ -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();
Expand Down Expand Up @@ -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);
}

/**
Expand All @@ -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);
}

/**
Expand All @@ -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);
}

/**
Expand All @@ -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
Expand All @@ -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);
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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 {
Expand All @@ -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) {
Expand Down Expand Up @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();

Expand Down Expand Up @@ -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.");
Expand Down
26 changes: 19 additions & 7 deletions src/main/java/org/metricshub/ipmi/core/connection/Connection.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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<ConnectionListener>();
listeners = new CopyOnWriteArrayList<ConnectionListener>();
timeout = Integer.parseInt(PropertiesManager.getInstance().getProperty("timeout"));
messageHandlers = new EnumMap<PayloadType, MessageHandler>(PayloadType.class);
currentSessionSequenceNumber = new AtomicInteger(0);
Expand Down Expand Up @@ -327,6 +328,7 @@ public List<CipherSuite> 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();
Expand All @@ -335,16 +337,26 @@ 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());
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) {
Expand Down
68 changes: 56 additions & 12 deletions src/main/java/org/metricshub/ipmi/core/sm/StateMachine.java
Original file line number Diff line number Diff line change
Expand Up @@ -26,9 +26,11 @@
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;
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;
Expand All @@ -44,15 +46,23 @@
*/
public class StateMachine implements UdpListener {

private List<MachineObserver> observers;
private final List<MachineObserver> observers = new CopyOnWriteArrayList<MachineObserver>();

private State current;
/**
* 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<StateMachineAction> pendingActions = new ArrayList<StateMachineAction>();

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;
Expand All @@ -72,7 +82,6 @@ public void setCurrent(State current) {
*/
public StateMachine(Messenger messenger) {
this.messenger = messenger;
observers = new ArrayList<MachineObserver>();
initialized = false;
}

Expand Down Expand Up @@ -101,19 +110,41 @@ 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. 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 (action instanceof MessageAction && Thread.holdsLock(this)) {
pendingActions.add(action);
} else {
notifyObservers(action);
}
}

private void notifyObservers(StateMachineAction action) {
for (MachineObserver observer : observers) {
if (observer != null) {
observer.notify(action);
}
}
}

private void dispatch(List<StateMachineAction> actions) {
for (StateMachineAction action : actions) {
notifyObservers(action);
}
}

/** Returns the actions emitted so far and clears them; called under the lock. */
private List<StateMachineAction> drainPendingActions() {
List<StateMachineAction> actions = new ArrayList<StateMachineAction>(pendingActions);
pendingActions.clear();
return actions;
}

/**
* Sets the State Machine in the initial state.
*
Expand Down Expand Up @@ -154,7 +185,9 @@ 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; the in-session messages are dispatched once the lock is
* released.
*
* @param event
* - {@link StateMachineEvent} invoking the transition
Expand All @@ -163,17 +196,28 @@ public boolean isActive() {
* @see #start(InetAddress, int)
*/
public void doTransition(StateMachineEvent event) {
if (!initialized) {
throw new NullPointerException("State machine not started");
List<StateMachineAction> actions;
synchronized (this) {
if (!initialized) {
throw new NullPointerException("State machine not started");
}
current.doTransition(this, event);
Comment thread
bertysentry marked this conversation as resolved.
actions = drainPendingActions();
}
current.doTransition(this, event);
dispatch(actions);
}

@Override
public void notifyMessage(UdpMessage message) {
if (message.getAddress().equals(getRemoteMachineAddress()) && message.getPort() == getRemoteMachinePort()) {
if (!message.getAddress().equals(getRemoteMachineAddress()) || message.getPort() != getRemoteMachinePort()) {
return;
}
List<StateMachineAction> actions;
synchronized (this) {
current.doAction(this, RmcpDecoder.decode(message.getMessage()));
actions = drainPendingActions();
}
dispatch(actions);
Comment thread
bertysentry marked this conversation as resolved.
}

/**
Expand Down
6 changes: 5 additions & 1 deletion src/site/markdown/upgrading.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:

Expand Down
Loading