Add SessionDisconnectJob to publish DisconnectEvent when Session timeouts #207
@@ -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
|
||||||
*
|
*
|
||||||
|
|||||||
+40
@@ -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
@@ -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());
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+1
@@ -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());
|
||||||
|
|||||||
Reference in New Issue
Block a user