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 @@ -52,7 +52,8 @@ public class IpmiClientConfiguration {
* @param password Password used to establish the connection with the host via the IPMI protocol.
* @param bmcKey The key that should be provided if the two-key authentication is enabled, null otherwise.
* @param skipAuth Whether the client should skip authentication
* @param timeout Timeout used for each IPMI request.
* @param timeout Overall deadline of each {@code IpmiClient} call, in seconds. It also caps the timeout of each
* message.
*/
public IpmiClientConfiguration(String hostname, String username, char[] password,
byte[] bmcKey, boolean skipAuth, long timeout) {
Expand All @@ -73,7 +74,8 @@ public IpmiClientConfiguration(String hostname, String username, char[] password
* @param password Password used to establish the connection with the host via the IPMI protocol.
* @param bmcKey The key that should be provided if the two-key authentication is enabled, null otherwise.
* @param skipAuth Whether the client should skip authentication
* @param timeout Timeout used for each IPMI request.
* @param timeout Overall deadline of each {@code IpmiClient} call, in seconds. It also caps the timeout of each
* message.
*/
public IpmiClientConfiguration(String hostname, int port, String username, char[] password,
byte[] bmcKey, boolean skipAuth, long timeout) {
Expand All @@ -89,7 +91,8 @@ public IpmiClientConfiguration(String hostname, int port, String username, char[
* @param password Password used to establish the connection with the host via the IPMI protocol.
* @param bmcKey The key that should be provided if the two-key authentication is enabled, null otherwise.
* @param skipAuth Whether the client should skip authentication
* @param timeout Timeout used for each IPMI request.
* @param timeout Overall deadline of each {@code IpmiClient} call, in seconds. It also caps the timeout of each
* message.
* @param pingPeriod The period in milliseconds used to send the keep alive messages.<br>
* Set pingPeriod to 0 to turn off keep-alive messages sent to the remote host.
*/
Expand Down Expand Up @@ -217,18 +220,18 @@ public void setSkipAuth(boolean skipAuth) {
}

/**
* Returns the timeout used for each IPMI request.
* Returns the overall deadline of each {@code IpmiClient} call, in seconds.
*
* @return The timeout used for each IPMI request.
* @return The overall deadline of each {@code IpmiClient} call, in seconds.
*/
public long getTimeout() {
return timeout;
}

/**
* Sets the timeout used for each IPMI request.
* Sets the overall deadline of each {@code IpmiClient} call, in seconds. It also caps the timeout of each message.
*
* @param timeout The timeout used for each IPMI request.
* @param timeout The overall deadline of each {@code IpmiClient} call, in seconds.
*/
public void setTimeout(long timeout) {
this.timeout = timeout;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -147,6 +147,7 @@ protected void startSession() throws Exception {
ipmiConfiguration.getPort(),
Connection.getDefaultCipherSuite(),
PrivilegeLevel.User);
capMessageTimeout();
}

// Start the session, provide user name and password, and optionally the
Expand Down Expand Up @@ -175,6 +176,7 @@ public void authenticate() throws Exception {
.createConnection(
InetAddress.getByName(ipmiConfiguration.getHostname()),
ipmiConfiguration.getPort());
capMessageTimeout();

// Get available cipher suites list via getAvailableCipherSuites and
// pick one of them that will be used further in the session.
Expand Down Expand Up @@ -212,8 +214,24 @@ protected CipherSuite getAvailableCipherSuite() throws Exception {
return suites.get(0);
}

/**
* Caps the timeout of each message by the overall deadline of the call, so that a lost reply is retried
* within that deadline rather than reported after it.
*/
private void capMessageTimeout() {
long deadlineMs = ipmiConfiguration.getTimeout() * 1000;
if (deadlineMs > 0 && deadlineMs < connector.getTimeout(handle)) {
connector.setTimeout(handle, (int) deadlineMs);
}
}

@Override
public void close() {
// startSession() may have failed before the connector or the handle existed
if (connector == null) {
return;
}

if (handle != null) {
// Close the session
try {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@
import java.net.InetAddress;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.TimeUnit;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
Expand Down Expand Up @@ -225,6 +226,8 @@ public List<CipherSuite> getAvailableCipherSuites(
++tries;
result = connectionManager
.getAvailableCipherSuites(connectionHandle.getHandle());
} catch (InterruptedException e) {
throw e;
} catch (Exception e) {
logger.warn(FAILED_TO_RECEIVE_ANSWER_CAUSE_MESSAGE, e);
if (tries > retries) {
Expand Down Expand Up @@ -269,6 +272,8 @@ public GetChannelAuthenticationCapabilitiesResponseData getChannelAuthentication
requestedPrivilegeLevel);
connectionHandle.setCipherSuite(cipherSuite);
connectionHandle.setPrivilegeLevel(requestedPrivilegeLevel);
} catch (InterruptedException e) {
throw e;
} catch (Exception e) {
logger.warn(FAILED_TO_RECEIVE_ANSWER_CAUSE_MESSAGE, e);
if (tries > retries) {
Expand Down Expand Up @@ -326,6 +331,8 @@ public Session openSession(
session = sessionManager.registerSession(sessionId, connectionHandle);

succeded = true;
} catch (InterruptedException e) {
throw e;
} catch (Exception e) {
logger.warn(FAILED_TO_RECEIVE_ANSWER_CAUSE_MESSAGE, e);
if (tries > retries) {
Expand Down Expand Up @@ -416,26 +423,26 @@ public int sendMessage(
throws Exception {
int tries = 0;
int tag = -1;
Connection connection = connectionManager.getConnection(connectionHandle.getHandle());
while (tries <= retries && tag < 0) {
try {
++tries;
// tag < 0 means that the MessageQueue is full: wait for a slot, at most one message timeout
long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(connection.getTimeout());
while (tag < 0) {
tag = connectionManager
.getConnection(
connectionHandle.getHandle())
.sendMessage(
request,
isOneWay);
tag = connection.sendMessage(request, isOneWay);
if (tag < 0) {
Thread.sleep(10); // tag < 0 means that MessageQueue is
// full so we need to wait and retry
if (System.nanoTime() >= deadline) {
throw new ConnectionException("Message queue is full");
}
Thread.sleep(10);
}
}
logger
.debug(
"Sending message with tag " + tag + ", try "
+ tries);
} catch (IllegalArgumentException e) {
} catch (IllegalArgumentException | InterruptedException e) {
throw e;
} catch (Exception e) {
logger.warn("Failed to send message, cause:", e);
Expand Down Expand Up @@ -585,4 +592,14 @@ public void setTimeout(ConnectionHandle handle, int timeout) {
connectionManager.getConnection(handle.getHandle()).setTimeout(timeout);
}

/**
* Returns the timeout of a single message on the given connection.
*
* @param handle {@link ConnectionHandle} of the connection
* @return the timeout in milliseconds after which a message without a reply is reported as timed out
*/
public int getTimeout(ConnectionHandle handle) {
return connectionManager.getConnection(handle.getHandle()).getTimeout();
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -425,7 +425,7 @@ private ResponseData sendThroughAsyncConnector(
}

messageSent = true;
} catch (IllegalArgumentException e) {
} catch (IllegalArgumentException | InterruptedException e) {
throw e;
Comment on lines +428 to 429

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Unregister the listener when interruption propagates

When a thread using the synchronous low-level IpmiConnector.sendMessage is interrupted inside waitForAnswer, this new branch immediately rethrows, but sendMessage() unregisters its MessageListener only after a normal return. The listener therefore remains in IpmiAsyncConnector; repeated cancellations leak listeners, and every later response must traverse these stale registrations. Move listener removal into a finally block.

Useful? React with 馃憤聽/ 馃憥.

} catch (IPMIException e) {
handleErrorResponse(tries, e);
Expand Down Expand Up @@ -496,6 +496,17 @@ public void setTimeout(ConnectionHandle handle, int timeout) {
asyncConnector.setTimeout(handle, timeout);
}

/**
* Returns the timeout of a single message on the connection with the given handle.
*
* @param handle
* - {@link ConnectionHandle} associated with the remote host.
* @return the timeout in ms after which a message without a reply is reported as timed out
*/
public int getTimeout(ConnectionHandle handle) {
return asyncConnector.getTimeout(handle);
}

/**
* Returns configured number of retries.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -86,24 +86,25 @@ public ResponseData waitForAnswer(int messageTag) throws Exception {
if (messageTag < 0 || messageTag > 63) {
throw new IllegalArgumentException("Corrupted message tag");
}
IpmiResponse answer;
synchronized (this) {
// Forget the outcome of the previous try: a retry must wait for the reply of the resent message
response = null;
this.tag = messageTag;
for (IpmiResponse quickResponse : quickMessages) {
this.notify(quickResponse);
}
}

while (response == null) {
Thread.sleep(1);
}
if (response instanceof IpmiResponseData) {
synchronized (this) {
this.tag = -1;
quickMessages.clear();
while (response == null) {
wait();
}
return ((IpmiResponseData) response).getResponseData();
} else /* response instanceof IpmiError */ {
throw ((IpmiError) response).getException();
answer = response;
this.tag = -1;
quickMessages.clear();
}
if (answer instanceof IpmiResponseData) {
return ((IpmiResponseData) answer).getResponseData();
} else /* answer instanceof IpmiError */ {
throw ((IpmiError) answer).getException();
}
}

Expand All @@ -114,6 +115,7 @@ public synchronized void notify(IpmiResponse ipmiResponse) {
quickMessages.add(ipmiResponse);
} else if (ipmiResponse.getTag() == tag) {
this.response = ipmiResponse;
notifyAll();
}
}
}
Expand Down
25 changes: 18 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 @@ -76,6 +76,7 @@
import java.util.Map;
import java.util.Timer;
import java.util.TimerTask;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;

/**
Expand All @@ -96,6 +97,7 @@ public class Connection extends TimerTask implements MachineObserver {
*/
private volatile int timeout = -1;
private volatile StateMachineAction lastAction;
private final Object responseLock = new Object();
private volatile int sessionId;
private volatile int managedSystemSessionId;
private volatile byte[] sik;
Expand Down Expand Up @@ -208,7 +210,7 @@ public void connect(InetAddress address, int port, long pingPeriod, boolean skip
// If the pingPeriod greater than 0, start the timer otherwise don't start it
// means that the connection won't be kept alive by sending no-op messages
if (pingPeriod > 0) {
timer = new Timer();
timer = new Timer(true);
timer.schedule(this, pingPeriod, pingPeriod);
}

Expand Down Expand Up @@ -321,15 +323,21 @@ public List<CipherSuite> getAvailableCipherSuites(int tag) throws Exception {
}

private void waitForResponse() throws Exception {
int time = 0;
long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(timeout);

while (time < timeout && lastAction == null) {
synchronized (responseLock) {
try {
Thread.sleep(1);
long remaining = deadline - System.nanoTime();
while (lastAction == null && remaining > 0) {
TimeUnit.NANOSECONDS.timedWait(responseLock, remaining);
remaining = deadline - System.nanoTime();
}
} catch (InterruptedException e) {
LOGGER.error(e.getMessage(), 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;
}
++time;
}

if (lastAction == null) {
Expand Down Expand Up @@ -644,7 +652,10 @@ public void notify(StateMachineAction action) {
if (action instanceof GetSikAction) {
sik = ((GetSikAction) action).getSik();
} else if (!(action instanceof MessageAction)) {
lastAction = action;
synchronized (responseLock) {
lastAction = action;
responseLock.notifyAll();
}
if (action instanceof ErrorAction) {
ErrorAction errorAction = (ErrorAction) action;
LOGGER.error(errorAction.getException().getMessage(), errorAction.getException());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,7 @@ public MessageQueue(Connection connection, int timeout, int minSequenceNumber, i
this.connection = connection;
queue = new ArrayList<QueueElement>();
setTimeout(timeout);
timer = new Timer();
timer = new Timer(true);
timer.schedule(this, cleaningFrequency, cleaningFrequency);
}

Expand Down Expand Up @@ -117,7 +117,7 @@ private synchronized boolean isReserved(int tag) {
* @return true if tag was reserved successfully, false otherwise
*/
private synchronized boolean reserveTag(int tag) {
if (isReserved(tag)) {
if (!isReserved(tag)) {
reservedTags.add(tag);
return true;
}
Expand Down Expand Up @@ -349,23 +349,29 @@ private boolean messageJustTimedOut(QueueElement oldestQueueElement) {
return now.getTime() - oldestQueueElement.getTimestamp().getTime() > (long) timeout;
}

/**
* Removes the oldest message from the queue; when it timed out (rather than being answered), the response
* listeners are told so, which lets the sender retry it with a fresh tag.
*/
private void processObsoleteMessage(QueueElement message, boolean done) {
int tag = message.getId();
boolean previouslyTimedOut = message.isTimedOut();

if (previouslyTimedOut || done) {
queue.remove(0);
logger.info("Removing message after timeout, tag: " + tag);
releaseTag(tag);
} else {
message.makeTimedOut();
message.refreshTimestamp();
connection
.notifyResponseListeners(
connection.getHandle(),
tag,
null,
new ConnectionException("Message timed out"));

queue.remove(0);
releaseTag(tag);

if (!done) {
logger.debug("Message timed out, tag: {}", tag);
try {
connection
.notifyResponseListeners(
connection.getHandle(),
tag,
null,
new ConnectionException("Message timed out"));
} catch (RuntimeException e) {
// A failing listener must not kill the timer thread that expires the other messages
logger.warn("Response listener failed while handling the timeout of tag {}", tag, e);
}
}
}

Expand Down
Loading
Loading