Merge branch 'feat/15-session-disconnect-job' into 'main'

Add SessionDisconnectJob to publish DisconnectEvent when Session timeouts

Closes #15

See merge request cs108-fs26/Gruppe-13!51
This commit was merged in pull request #207.
This commit is contained in:
Lars Simon Winzer
2026-03-28 13:17:42 +01:00
5 changed files with 79 additions and 3 deletions
@@ -7,6 +7,7 @@ import ch.unibas.dmi.dbis.cs108.casono.server.network.command.CommandRouter;
import ch.unibas.dmi.dbis.cs108.casono.server.network.events.DisconnectEvent; import ch.unibas.dmi.dbis.cs108.casono.server.network.events.DisconnectEvent;
import ch.unibas.dmi.dbis.cs108.casono.server.network.events.EventBus; import ch.unibas.dmi.dbis.cs108.casono.server.network.events.EventBus;
import ch.unibas.dmi.dbis.cs108.casono.server.network.parser.CommandParserDispatcher; import ch.unibas.dmi.dbis.cs108.casono.server.network.parser.CommandParserDispatcher;
import ch.unibas.dmi.dbis.cs108.casono.server.network.sessions.SessionDisconnectJob;
import ch.unibas.dmi.dbis.cs108.casono.server.network.sessions.SessionManager; import ch.unibas.dmi.dbis.cs108.casono.server.network.sessions.SessionManager;
import java.time.Duration; import java.time.Duration;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
@@ -17,9 +18,12 @@ import org.apache.logging.log4j.Logger;
/** Application class for starting the server. */ /** Application class for starting the server. */
public class ServerApp { public class ServerApp {
public static final int USER_CLEANUP_JOB_DELAY = 0; private static final int USER_CLEANUP_JOB_DELAY = 0;
public static final int USER_CLEANUP_JOB_PERIOD = 10; private static final int USER_CLEANUP_JOB_PERIOD = 10;
public static final int USER_CLEANUP_JOB_RECONNECT_THRESHOLD = 10; private static final int USER_CLEANUP_JOB_RECONNECT_THRESHOLD = 10;
private static final int SESSION_DISCONNECT_JOB_DELAY = 0;
private static final int SESSION_DISCONNECT_JOB_PERIOD = 2;
private static final int SESSION_DISCONNECT_JOB_TIMEOUT = 5;
public static void start(String arg) { public static void start(String arg) {
int port = Integer.parseInt(arg); int port = Integer.parseInt(arg);
@@ -45,6 +49,14 @@ public class ServerApp {
USER_CLEANUP_JOB_DELAY, USER_CLEANUP_JOB_DELAY,
USER_CLEANUP_JOB_PERIOD, USER_CLEANUP_JOB_PERIOD,
TimeUnit.SECONDS); TimeUnit.SECONDS);
scheduler.scheduleAtFixedRate(
new SessionDisconnectJob(
sessionManager,
eventBus,
Duration.ofSeconds(SESSION_DISCONNECT_JOB_TIMEOUT)),
SESSION_DISCONNECT_JOB_DELAY,
SESSION_DISCONNECT_JOB_PERIOD,
TimeUnit.SECONDS);
networkManager.start(); networkManager.start();
} }
@@ -6,12 +6,14 @@ import ch.unibas.dmi.dbis.cs108.casono.server.network.parser.CommandParserDispat
import ch.unibas.dmi.dbis.cs108.casono.server.network.response.PrimitiveResponse; import ch.unibas.dmi.dbis.cs108.casono.server.network.response.PrimitiveResponse;
import ch.unibas.dmi.dbis.cs108.casono.server.network.transport.TransportLayer; import ch.unibas.dmi.dbis.cs108.casono.server.network.transport.TransportLayer;
import java.io.IOException; import java.io.IOException;
import java.time.Instant;
import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue; import java.util.concurrent.BlockingQueue;
/** Represents a client session in the network server. */ /** Represents a client session in the network server. */
public class Session { public class Session {
private final SessionId id; private final SessionId id;
private Instant lastActivity;
private final TransportLayer transport; private final TransportLayer transport;
private final BlockingQueue<PrimitiveResponse> responseQueue; private final BlockingQueue<PrimitiveResponse> responseQueue;
private final CommandParserDispatcher dispatcher; private final CommandParserDispatcher dispatcher;
@@ -31,6 +33,7 @@ public class Session {
CommandParserDispatcher dispatcher, CommandParserDispatcher dispatcher,
CommandRouter router) { CommandRouter router) {
this.id = new SessionId(); this.id = new SessionId();
this.lastActivity = Instant.now();
this.transport = transport; this.transport = transport;
this.dispatcher = dispatcher; this.dispatcher = dispatcher;
this.router = router; this.router = router;
@@ -46,6 +49,20 @@ public class Session {
return this.id; return this.id;
} }
/**
* Gets the timestamp of the last inbound activity for this session.
*
* @return an {@link Instant} representing the time of the last inbound activity
*/
public Instant getLastInboundActivity() {
return lastActivity;
}
/** Updates the timestamp of the last inbound activity for this session. */
public void updateLastInboundActivity() {
this.lastActivity = Instant.now();
}
/** /**
* Returns the TransportLayer of this session * Returns the TransportLayer of this session
* *
@@ -0,0 +1,40 @@
package ch.unibas.dmi.dbis.cs108.casono.server.network.sessions;
import ch.unibas.dmi.dbis.cs108.casono.server.network.events.DisconnectEvent;
import ch.unibas.dmi.dbis.cs108.casono.server.network.events.EventBus;
import java.time.Duration;
import java.time.Instant;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
public class SessionDisconnectJob implements Runnable {
private final Logger logger;
private final SessionManager sessionManager;
private final EventBus eventBus;
private final Duration timeoutThreshold;
public SessionDisconnectJob(
SessionManager sessionManager, EventBus eventBus, Duration timeoutThreshold) {
this.logger = LogManager.getLogger(SessionDisconnectJob.class);
this.sessionManager = sessionManager;
this.eventBus = eventBus;
this.timeoutThreshold = timeoutThreshold;
}
@Override
public void run() {
logger.debug("Job started.");
Instant threshold = Instant.now().minus(timeoutThreshold);
for (Session session : sessionManager.getAllSessions()) {
if (session.getLastInboundActivity().isBefore(threshold)) {
eventBus.publish(new DisconnectEvent(session.getId()));
logger.info(
"Initiated disconnect of {}, as it hasn't been active since a while",
session.getId());
}
}
logger.debug("Job finished.");
}
}
@@ -6,8 +6,10 @@ import ch.unibas.dmi.dbis.cs108.casono.server.network.events.EventBus;
import ch.unibas.dmi.dbis.cs108.casono.server.network.parser.CommandParserDispatcher; import ch.unibas.dmi.dbis.cs108.casono.server.network.parser.CommandParserDispatcher;
import ch.unibas.dmi.dbis.cs108.casono.server.network.transport.TransportLayer; import ch.unibas.dmi.dbis.cs108.casono.server.network.transport.TransportLayer;
import java.io.IOException; import java.io.IOException;
import java.util.Collection;
import java.util.Map; import java.util.Map;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.stream.Collectors;
import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger; import org.apache.logging.log4j.Logger;
@@ -109,4 +111,8 @@ public class SessionManager {
return handle.session(); return handle.session();
} }
public Collection<Session> getAllSessions() {
return sessions.values().stream().map(SessionHandle::session).collect(Collectors.toList());
}
} }
@@ -43,6 +43,7 @@ public class SessionReader implements Runnable {
RawPacket rawPacket = null; RawPacket rawPacket = null;
try { try {
rawPacket = transport.read(); rawPacket = transport.read();
session.updateLastInboundActivity();
logger.debug("Recieved: {}", rawPacket); logger.debug("Recieved: {}", rawPacket);
RawRequest rawRequest = ProtocolParser.parse(rawPacket.payload()); RawRequest rawRequest = ProtocolParser.parse(rawPacket.payload());