Added ConnectionManager, propagating socketId to processors

This commit is contained in:
Kai S. K. Engelbart 2020-01-03 18:11:38 +02:00
parent 998fc3a91e
commit c42bdffbd7
6 changed files with 105 additions and 9 deletions

View File

@ -0,0 +1,79 @@
package envoy.server;
import java.util.HashMap;
import java.util.HashSet;
import java.util.Map;
import java.util.Set;
import com.jenkov.nioserver.ISocketIdListener;
/**
* Project: <strong>envoy-server-standalone</strong><br>
* File: <strong>ConnectionManager.java</strong><br>
* Created: <strong>03.01.2020</strong><br>
*
* @author Kai S. K. Engelbart
* @since Envoy Server Standalone v0.1-alpha
*/
public class ConnectionManager implements ISocketIdListener {
/**
* Contains all socket IDs that have not yet performed a handshake / acquired
* their corresponding user ID.
*
* @since Envoy Server Standalone v0.1-alpha
*/
private Set<Long> pendingSockets = new HashSet<>();
/**
* Contains all socket IDs that have acquired a user ID as keys to these IDs.
*
* @since Envoy Server Standalone v0.1-alpha
*/
private Map<Long, Long> sockets = new HashMap<>();
private static ConnectionManager connectionManager = new ConnectionManager();
private ConnectionManager() {}
/**
* @return a singleton instance of this object
* @since Envoy Server Standalone v0.1-alpha
*/
public static ConnectionManager getInstance() { return connectionManager; }
@Override
public void socketCancelled(long socketId) {
if (!pendingSockets.remove(socketId))
sockets.entrySet().stream().filter(e -> e.getValue() == socketId).forEach(e -> sockets.remove(e.getValue()));
}
@Override
public void socketRegistered(long socketId) { pendingSockets.add(socketId); }
/**
* Associates a socket ID with a user ID.
*
* @param socketId the socket ID
* @param userId the user ID
* @since Envoy Server Standalone v0.1-alpha
*/
public void registerUser(long socketId, long userId) {
sockets.put(socketId, userId);
pendingSockets.remove(socketId);
}
/**
* @param userId the ID of the user registered at the a socket
* @return the ID of the socket
* @since Envoy Server Standalone v0.1-alpha
*/
public long getSocketId(long userId) { return sockets.get(userId); }
/**
* @param userId the ID of the user to check for
* @return {@code true} if the user is online
* @since Envoy Server Standalone v0.1-alpha
*/
public boolean isOnline(long userId) { return sockets.containsKey(userId); }
}

View File

@ -22,8 +22,10 @@ public class LoginCredentialProcessor implements ObjectProcessor<LoginCredential
public Class<LoginCredentials> getInputClass() { return LoginCredentials.class; }
@Override
public User process(LoginCredentials input) {
System.out.println("Received login credentials " + input);
return new User(currentUserId++, input.getName());
public User process(LoginCredentials input, long socketId) {
System.out.println(String.format("Received login credentials %s from socket ID %d", input, socketId));
User user = new User(currentUserId++, input.getName());
ConnectionManager.getInstance().registerUser(socketId, user.getId());
return user;
}
}

View File

@ -18,5 +18,15 @@ public class MessageProcessor implements ObjectProcessor<Message, Void> {
public Class<Message> getInputClass() { return Message.class; }
@Override
public Void process(Message input) { return null; }
public Void process(Message message, long socketId) {
// TODO: Send message to recipient if online
ConnectionManager connectionManager = ConnectionManager.getInstance();
if (connectionManager.isOnline(message.getRecipientId())) {
}
// TODO: Add message to database
return null;
}
}

View File

@ -22,9 +22,10 @@ public interface ObjectProcessor<T, U> {
Class<T> getInputClass();
/**
* @param input the request object
* @param input the request object
* @param socketId the ID of the socket from which the object was received
* @return the response object
* @since Envoy Server Standalone v0.1-alpha
*/
U process(T input);
U process(T input, long socketId);
}

View File

@ -32,7 +32,9 @@ public class Startup {
Set<ObjectProcessor<?, ?>> processors = new HashSet<>();
processors.add(new LoginCredentialProcessor());
processors.add(new MessageProcessor());
Server server = new Server(8080, () -> new ObjectMessageReader(), new ObjectMessageProcessor(processors));
server.start();
server.getSocketProcessor().registerSocketIdListener(ConnectionManager.getInstance());
}
}

View File

@ -25,7 +25,7 @@ import envoy.server.ObjectProcessor;
*/
public class ObjectMessageProcessor implements IMessageProcessor {
private final Set<ObjectProcessor<?, ?>> processors;
private final Set<ObjectProcessor<?, ?>> processors;
/**
* The constructor to set the {@link ObjectProcessor}s.
@ -33,7 +33,9 @@ public class ObjectMessageProcessor implements IMessageProcessor {
* @param processors the {@link ObjectProcessor} to set
* @since Envoy Server Standalone v0.1-alpha
*/
public ObjectMessageProcessor(Set<ObjectProcessor<?, ?>> processors) { this.processors = processors; }
public ObjectMessageProcessor(Set<ObjectProcessor<?, ?>> processors) {
this.processors = processors;
}
@SuppressWarnings("unchecked")
@Override
@ -44,7 +46,7 @@ public class ObjectMessageProcessor implements IMessageProcessor {
// Process object
processors.stream().filter(p -> p.getInputClass().isInstance(obj)).forEach((@SuppressWarnings("rawtypes") ObjectProcessor p) -> {
Object responseObj = p.process(p.getInputClass().cast(obj));
Object responseObj = p.process(p.getInputClass().cast(obj), message.socketId);
if (responseObj != null) {
// Create message targeted at the client
Message response = writeProxy.getMessage();