Changed link monitoring as described by on4kst, Alain

This commit is contained in:
Marc Froehlich
2026-09-13 19:50:33 +02:00
parent f0508f51a2
commit 393c9c49ee
8 changed files with 1028 additions and 90 deletions
@@ -877,10 +877,9 @@ public class MessageBusManagementThread extends Thread {
|| messageToProcess.getMessageText().isEmpty()) {
// No processable data.
} else {
if (On4KstProtocol.isConnectionProbeResponse(
if (On4KstProtocol.isInternalDxqResponse(
messageToProcess.getMessageText())) {
// DXQ is the internal response to the active connection probe.
// Liveness was already recorded by the session manager.
// Preserve the established internal handling of DXQ server data.
return;
}
@@ -2183,7 +2182,7 @@ public class MessageBusManagementThread extends Thread {
// e.printStackTrace();
// }
if (!On4KstProtocol.isConnectionProbeResponse(
if (!On4KstProtocol.isInternalDxqResponse(
messageTextRaw.getMessageText())) {
System.out.println(messageTextRaw.getMessageText() + " <- RXed"); // Stdout at
// Console#######################################################TODO:Wichtig
@@ -47,9 +47,7 @@ final class On4KstConnectionManager {
static final int CONNECT_TIMEOUT_MILLIS = 10_000; //TCP-Connect-Timeout
static final long LOGIN_FALLBACK_MILLIS = 2_000L; //Login-Fallback
static final long HANDSHAKE_TIMEOUT_MILLIS = 45_000L; //Handshake-Timeout
static final long APPLICATION_HEARTBEAT_AFTER_MILLIS = 90_000L; //Application-Heartbeat
/** Idle duration after which the server is asked for current DX data. */
static final long CONNECTION_PROBE_AFTER_MILLIS = 180_000L; //Active connection probe
static final long CLIENT_LIVENESS_PROBE_AFTER_MILLIS = 90_000L;
static final long INBOUND_STALE_AFTER_MILLIS = 210_000L; //Stale-Timeout - time without rxed data
static final List<Long> RECONNECT_DELAYS_MILLIS =
List.of(2_000L, 5_000L, 10_000L, 20_000L, 30_000L); //Reconnect-Backoff if no connection possible
@@ -181,10 +179,10 @@ final class On4KstConnectionManager {
session.lastProgressMillis.set(now);
String opcode = On4KstProtocol.opcode(line);
long probeResponseMillis = session.connectionProbe.acknowledge(now);
long probeResponseMillis = session.clientLivenessProbe.acknowledge(now);
if (probeResponseMillis >= 0L) {
LOGGER.log(Level.INFO,
"ON4KST connection probe confirmed: session {0}, "
"ON4KST client liveness probe confirmed: session {0}, "
+ "received opcode {1}, response time {2} ms",
new Object[] {
sessionId,
@@ -193,8 +191,8 @@ final class On4KstConnectionManager {
});
}
if ("CK".equals(opcode)) {
sendHeartbeat(session);
if (On4KstProtocol.isServerLivenessProbe(line)) {
sendServerLivenessProbeResponse(session);
}
if (!session.loginSent
@@ -265,8 +263,7 @@ final class On4KstConnectionManager {
token,
socket,
receiveQueue,
transmitQueue,
mainCategory);
transmitQueue);
ReadThread readThread = new ReadThread(
token, socket, receiveQueue, this::isActiveSession,
@@ -541,42 +538,35 @@ final class On4KstConnectionManager {
session.transmitQueue.offer(message);
}
private void sendHeartbeat(Session session) {
private void sendServerLivenessProbeResponse(Session session) {
if (session == null || !isActiveSession(session.id)) {
return;
}
long now = System.currentTimeMillis();
session.lastHeartbeatMillis.set(now);
LOGGER.log(Level.FINE,
"Sending application heartbeat for ON4KST session {0}",
"Responding to ON4KST server liveness probe: session {0}, "
+ "opcode CK",
session.id);
ChatMessage heartbeat = new ChatMessage();
heartbeat.setMessageDirectedToServer(true);
heartbeat.setMessageText("");
session.transmitQueue.offer(heartbeat);
sendControl(session, On4KstProtocol.serverLivenessProbeResponse());
}
private void sendConnectionProbe(
private void sendClientLivenessProbe(
Session session,
long now,
long inboundIdle
) {
if (session == null || !isActiveSession(session.id)
|| !session.connectionProbe.tryStart(now)) {
|| !session.clientLivenessProbe.tryStart(now)) {
return;
}
LOGGER.log(Level.INFO,
"Sending ON4KST connection probe: session {0}, main category "
+ "{1}, inbound idle {2} seconds",
"Sending ON4KST client liveness probe: session {0}, "
+ "opcode CK, inbound idle {1} seconds",
new Object[] {
session.id,
session.mainCategory,
inboundIdle / 1_000L
});
sendControl(
session,
On4KstProtocol.connectionProbe(session.mainCategory));
sendControl(session, On4KstProtocol.clientLivenessProbe());
}
private void onConnectionFailure(long sessionId, Throwable failure) {
@@ -682,8 +672,8 @@ final class On4KstConnectionManager {
long inboundIdle = now - lastInboundMillis;
IdleAction idleAction = determineIdleAction(
inboundIdle,
session.lastHeartbeatMillis.get() >= lastInboundMillis,
session.connectionProbe.isOutstanding());
session.online,
session.clientLivenessProbe.isOutstanding());
if (session.lastInboundMillis.get() != lastInboundMillis) {
return;
@@ -695,16 +685,15 @@ final class On4KstConnectionManager {
return;
}
long probeWaitMillis =
session.connectionProbe.responseWaitMillis(now);
session.clientLivenessProbe.responseWaitMillis(now);
if (probeWaitMillis >= 0L) {
LOGGER.log(Level.WARNING,
"ON4KST connection probe timed out: session "
+ "{0}, main category {1}, no response "
+ "for {2} ms, inbound idle {3} seconds; "
"ON4KST client liveness probe timed out: "
+ "session {0}, opcode CK, no response "
+ "for {1} ms, inbound idle {2} seconds; "
+ "reconnecting",
new Object[] {
session.id,
session.mainCategory,
probeWaitMillis,
inboundIdle / 1_000L
});
@@ -713,9 +702,8 @@ final class On4KstConnectionManager {
new SocketException("No ON4KST data received for "
+ inboundIdle / 1_000L + " seconds"));
}
case CONNECTION_PROBE ->
sendConnectionProbe(session, now, inboundIdle);
case HEARTBEAT -> sendHeartbeat(session);
case CLIENT_LIVENESS_PROBE ->
sendClientLivenessProbe(session, now, inboundIdle);
case NONE -> {
// The session is active or already has the required idle action.
}
@@ -731,19 +719,18 @@ final class On4KstConnectionManager {
*/
static IdleAction determineIdleAction(
long inboundIdleMillis,
boolean heartbeatSentForIdlePhase,
boolean online,
boolean probeOutstanding
) {
if (!online) {
return IdleAction.NONE;
}
if (inboundIdleMillis > INBOUND_STALE_AFTER_MILLIS) {
return IdleAction.TIMEOUT;
}
if (inboundIdleMillis >= CONNECTION_PROBE_AFTER_MILLIS
if (inboundIdleMillis > CLIENT_LIVENESS_PROBE_AFTER_MILLIS
&& !probeOutstanding) {
return IdleAction.CONNECTION_PROBE;
}
if (inboundIdleMillis > APPLICATION_HEARTBEAT_AFTER_MILLIS
&& !heartbeatSentForIdlePhase) {
return IdleAction.HEARTBEAT;
return IdleAction.CLIENT_LIVENESS_PROBE;
}
return IdleAction.NONE;
}
@@ -894,15 +881,13 @@ final class On4KstConnectionManager {
private final Socket socket;
private final LinkedBlockingQueue<ChatMessage> receiveQueue;
private final LinkedBlockingQueue<ChatMessage> transmitQueue;
private final int mainCategory;
private final long connectedMillis = System.currentTimeMillis();
private final AtomicLong lastInboundMillis =
new AtomicLong(connectedMillis);
private final AtomicLong lastProgressMillis =
new AtomicLong(connectedMillis);
private final AtomicLong lastHeartbeatMillis = new AtomicLong();
private final ConnectionProbeState connectionProbe =
new ConnectionProbeState();
private final ClientLivenessProbeState clientLivenessProbe =
new ClientLivenessProbeState();
private final Map<Integer, Map<String, ChatMember>> initialMembers =
new ConcurrentHashMap<>();
@@ -920,27 +905,24 @@ final class On4KstConnectionManager {
long id,
Socket socket,
LinkedBlockingQueue<ChatMessage> receiveQueue,
LinkedBlockingQueue<ChatMessage> transmitQueue,
int mainCategory
LinkedBlockingQueue<ChatMessage> transmitQueue
) {
this.id = id;
this.socket = socket;
this.receiveQueue = receiveQueue;
this.transmitQueue = transmitQueue;
this.mainCategory = mainCategory;
}
}
/** Maintenance action selected by the session monitor. */
enum IdleAction {
NONE,
HEARTBEAT,
CONNECTION_PROBE,
CLIENT_LIVENESS_PROBE,
TIMEOUT
}
/** Tracks one outstanding liveness probe for the complete TCP session. */
static final class ConnectionProbeState {
/** Tracks one client-initiated liveness probe for the complete TCP session. */
static final class ClientLivenessProbeState {
private final AtomicLong sentMillis = new AtomicLong();
boolean tryStart(long now) {
@@ -54,13 +54,32 @@ final class On4KstProtocol {
+ "|0|";
}
/** Builds the active liveness probe for the session's main chat. */
static String connectionProbe(int category) {
return "RDXQ|" + category(category) + "|";
/** Builds the session-wide liveness probe initiated by this client. */
static String clientLivenessProbe() {
return "CK|";
}
/** Returns whether a server frame is the expected liveness-probe response. */
static boolean isConnectionProbeResponse(String frame) {
/** Returns whether a server frame acknowledges a client-initiated probe. */
static boolean isClientLivenessProbeResponse(String frame) {
if (frame == null) {
return false;
}
String normalized = frame.trim().toUpperCase(Locale.ROOT);
return "OK".equals(normalized) || "OK|".equals(normalized);
}
/** Returns whether ON4KST initiated its own liveness check. */
static boolean isServerLivenessProbe(String frame) {
return "CK".equals(opcode(frame));
}
/** Builds the established empty response to a server-initiated {@code CK}. */
static String serverLivenessProbeResponse() {
return "";
}
/** Preserves the existing internal handling of DXQ server data. */
static boolean isInternalDxqResponse(String frame) {
return "DXQ".equals(opcode(frame));
}
@@ -86,11 +86,20 @@ public class ReadThread extends Thread {
throw new EOFException("ON4KST closed the TCP connection");
}
if (!sessionIsActive.test(sessionId)) {
break;
}
inboundActivity.accept(response);
if (!sessionIsActive.test(sessionId)) {
break;
}
if (On4KstProtocol.isClientLivenessProbeResponse(response)) {
// OK and OK| acknowledge the client-side CK| probe.
continue;
}
ChatMessage message = new ChatMessage();
message.setMessageText(response);
receiveQueue.put(message);
@@ -127,4 +136,4 @@ public class ReadThread extends Thread {
socket.close();
return true;
}
}
}