Mantém conexões persistentes com os peers e reutiliza essas conexões + * para o envio de mensagens. As mensagens recebidas são encaminhadas aos + * plugins, que são responsáveis por interpretá-las.
+ * + *Novas conexões são aceitas pelo servidor e cada conexão possui uma + * tarefa própria para receber mensagens, permitindo que múltiplas conexões + * sejam mantidas simultaneamente.
+ * + *As conexões são associadas aos endereços dos peers e podem ser + * encerradas quando um peer deixa de estar ativo.
+ * + * @author Gustavo + */ +public class SocketTCP extends Thread implements PeerListener { + private final ExecutorService connectionExecutor = Executors.newCachedThreadPool(); + private final MapO servidor utiliza a porta definida pela aplicação e inicia sem + * estabelecer conexões com os peers. As conexões são criadas conforme + * necessário durante o envio ou recebidas de outros peers.
+ * + * @param main janela principal da aplicação. + */ + public SocketTCP(MainWindow main) { + this.main = main; + connections = new ConcurrentHashMap<>(); + try { + serverSocket = new ServerSocket(main.getPort()); + } catch (IOException ex) { + System.out.println("There is no socket connection. Sorry."); + System.out.println(ex); + } + } + + /** + * Envia uma mensagem TCP para um peer. + * + *Uma conexão existente com o peer é reutilizada. Caso não exista uma + * conexão válida, uma nova conexão é criada e adicionada ao conjunto de + * conexões ativas.
+ * + * @param msg mensagem serializada a ser enviada. + * @param destinationAddress endereço IP do peer destinatário. + */ + public void send(byte[] msg, InetAddress destinationAddress) { + try { + TcpConnection connection; + + synchronized (connections) { + connection= connections.get(destinationAddress); + + if (connection == null || connection.isClosed()) { + Socket socket = new Socket(destinationAddress, main.getPort()); + + connection = new TcpConnection(socket); + + connections.put(destinationAddress, connection); + + startReceiver(destinationAddress, connection); + } + } + + connection.send(msg); + + } catch (IOException ex) { + System.out.println("Could not send TCP message."); + System.out.println(ex); + } + } + + /** + * Aguarda e aceita novas conexões TCP. + * + *Cada conexão aceita é associada ao endereço do peer e registrada para + * que possa ser utilizada tanto para recepção quanto para envio de + * mensagens.
+ */ + private void receive() { + while (running) { + + try { + Socket socket = serverSocket.accept(); + + InetAddress address = socket.getInetAddress(); + + TcpConnection connection = new TcpConnection(socket); + + TcpConnection oldConnection = connections.put(address, connection); + + if (oldConnection != null) { + oldConnection.close(); + } + + startReceiver(address, connection); + + } catch (IOException ex) { + + if (!running) { + break; + } + + System.out.println("Error accepting TCP connection."); + System.out.println(ex); + } + } + } + + /** + * Inicia a recepção de mensagens de uma conexão TCP. + * + *A recepção é executada de forma independente para que uma conexão + * não impeça a aceitação ou o processamento de outras conexões.
+ * + *Cada mensagem recebida é encaminhada aos plugins da aplicação.
+ * + * @param address endereço do peer associado à conexão. + * @param connection conexão TCP utilizada para a comunicação. + */ + private void startReceiver(InetAddress address, TcpConnection connection) { + connectionExecutor.submit(() -> { + try { + while (running) { + byte[] message = connection.receive(); + + for (Plugin plugin : main.getPlugins()) { + plugin.receiveMessage(message); + } + } + + } catch (IOException ex) { + System.out.println( + "Connection with " + address.getHostAddress() + " closed." + ); + + } finally { + connections.remove(address, connection); + + try { + connection.close(); + } catch (IOException ignored) { + } + } + }); + } + + /** + * Inicia o processo de recepção de conexões TCP. + */ + @Override + public void run() { + receive(); + } + + /** + * Encerra o servidor TCP e as conexões associadas. + * + *As conexões existentes são encerradas e o servidor deixa + * de aceitar novas conexões.
+ */ + public void close() { + running = false; + + if (serverSocket != null) { + try { + serverSocket.close(); + } catch (IOException _) { + } + } + + for (TcpConnection connection : connections.values()) { + try { + connection.close(); + } catch (IOException _) { + } + } + + connections.clear(); + + connectionExecutor.shutdownNow(); + } + + /** + * Atualiza as conexões de acordo com os peers atualmente ativos. + * + *Conexões associadas a peers que não estão mais ativos são encerradas + * e removidas.
+ * + * @param peers lista atualizada de peers ativos. + */ + @Override + public void onPeersChanged(ListEsta classe gerencia o envio e o recebimento de mensagens, além de - * manter a conexão com o grupo multicast utilizado pela aplicação.
+ *Mantém a conexão com o grupo multicast da aplicação e permite o envio + * de mensagens tanto para o grupo quanto diretamente para um peer.
+ * + *As mensagens recebidas são verificadas para identificar mensagens de + * heartbeat. As demais mensagens são encaminhadas aos plugins, que são + * responsáveis por interpretá-las.
* * @author flavio * @author tony + * @author Gustavo */ -public class Socket extends Thread { +public class SocketUDP extends Thread { private MulticastSocket multicastSocket; + private InetAddress address; private MainWindow main; private volatile boolean running = true; public final static String INET_ADDR = "224.0.0.3"; /** - * Cria e inicializa a conexão multicast utilizada pela aplicação. + * Cria e inicializa o socket UDP utilizado pela aplicação. + * + *O socket é associado à porta da aplicação e ingressa no grupo + * multicast configurado.
* * @param main janela principal da aplicação. */ - public Socket(MainWindow main) { + public SocketUDP(MainWindow main) { this.main = main; try { - this.address = InetAddress.getByName(Socket.INET_ADDR); + this.address = InetAddress.getByName(SocketUDP.INET_ADDR); } catch (UnknownHostException ex) { - System.getLogger(Socket.class.getName()).log(System.Logger.Level.ERROR, (String) null, ex); + System.getLogger(SocketUDP.class.getName()).log(System.Logger.Level.ERROR, (String) null, ex); } try { multicastSocket = new MulticastSocket(this.main.getPort()); @@ -56,7 +61,7 @@ public class Socket extends Thread { multicastSocket.joinGroup(address); } catch (IOException ex) { System.out.println("There is no socket connection. Sorry."); - System.out.println(ex.toString()); + System.out.println(ex); } } @@ -64,11 +69,33 @@ public class Socket extends Thread { * Envia uma mensagem para o grupo multicast da aplicação. * * @param msg mensagem serializada a ser enviada. - * @throws IOException caso ocorra um erro durante o envio da mensagem. + * @throws IOException caso ocorra um erro durante o envio. */ - public void send(byte[] msg) throws IOException { + public void sendMulticast(byte[] msg) throws IOException { + send(msg, this.address); + } + + /** + * Envia uma mensagem diretamente para um peer utilizando UDP. + * + * @param msg mensagem serializada a ser enviada. + * @param destinationAddress endereço IP do peer destinatário. + * @throws IOException caso ocorra um erro durante o envio. + */ + public void sendUnicast(byte[] msg, InetAddress destinationAddress) throws IOException { + send(msg, destinationAddress); + } + + /** + * Envia uma mensagem UDP para o endereço especificado. + * + * @param msg mensagem serializada a ser enviada. + * @param destinationAddress endereço de destino. + * @throws IOException caso ocorra um erro durante o envio. + */ + private void send(byte[] msg, InetAddress destinationAddress) throws IOException { DatagramPacket msgPacket; - msgPacket = new DatagramPacket(msg, msg.length, this.address, this.main.getPort()); + msgPacket = new DatagramPacket(msg, msg.length, destinationAddress, this.main.getPort()); multicastSocket.send(msgPacket); } @@ -81,7 +108,7 @@ public class Socket extends Thread { try { multicastSocket.leaveGroup(address); } catch (IOException ex) { - + System.out.println(ex); } multicastSocket.close(); } @@ -89,15 +116,17 @@ public class Socket extends Thread { /** - * Inicia o loop de recepção de mensagens da aplicação. + * Recebe e processa mensagens UDP. * - *As mensagens recebidas são verificadas inicialmente para identificar - * mensagens de heartbeat. Caso não sejam heartbeats, seu conteúdo é - * encaminhado para todos os plugins carregados, que decidem se devem ou - * não processá-lo.
+ *Mensagens de heartbeat são encaminhadas ao gerenciador de heartbeat. + * As demais mensagens são encaminhadas aos plugins carregados pela + * aplicação.
* - * @throws IOException caso ocorra um erro durante a recepção da mensagem. - * @throws ClassNotFoundException caso a desserialização de uma mensagem falhe. + * @throws UnknownHostException caso não seja possível resolver o endereço + * utilizado pela comunicação. + * @throws IOException caso ocorra um erro durante a recepção. + * @throws ClassNotFoundException caso ocorra um erro ao verificar uma + * mensagem de heartbeat. */ public void receive() throws UnknownHostException, IOException, ClassNotFoundException { byte[] buf = new byte[256000]; @@ -124,11 +153,11 @@ public class Socket extends Thread { } /** - * Tenta desserializar os dados recebidos como uma {@code HeartbeatMessage}. + * Tenta identificar os dados recebidos como uma mensagem de heartbeat. * - * @param data dados serializados da mensagem. - * @return a mensagem desserializada caso os dados representem uma - * {@code HeartbeatMessage}; caso contrário, {@code null}. + * @param data dados serializados recebidos. + * @return a mensagem de heartbeat caso os dados correspondam a uma; + * {@code null} caso contrário. */ private HeartbeatMessage tryParseHeartbeat(byte[] data) { try { @@ -141,12 +170,15 @@ public class Socket extends Thread { } } + /** + * Inicia a recepção de mensagens UDP. + */ @Override public void run() { try { this.receive(); } catch (IOException | ClassNotFoundException ex) { - Logger.getLogger(Socket.class.getName()).log(Level.SEVERE, null, ex); + Logger.getLogger(SocketUDP.class.getName()).log(Level.SEVERE, null, ex); } } } diff --git a/pitiupi/src/pitiupi/net/TcpConnection.java b/pitiupi/src/pitiupi/net/TcpConnection.java new file mode 100644 index 0000000..4921ad8 --- /dev/null +++ b/pitiupi/src/pitiupi/net/TcpConnection.java @@ -0,0 +1,88 @@ +package pitiupi.net; + +import java.io.DataInputStream; +import java.io.DataOutputStream; +import java.io.IOException; +import java.net.Socket; + +/** + * Representa uma conexão TCP persistente com um peer. + * + *Encapsula o socket e os streams utilizados para enviar e receber + * mensagens. Também é responsável pelo framing das mensagens, permitindo + * que várias mensagens sejam transmitidas pela mesma conexão.
+ * + * @author Gustavo + */ +class TcpConnection { + + private final Socket socket; + private final DataInputStream input; + private final DataOutputStream output; + + /** + * Cria uma conexão a partir de um socket existente. + * + * @param socket socket TCP utilizado pela conexão. + * @throws IOException caso não seja possível obter os streams do socket. + */ + public TcpConnection(Socket socket) throws IOException { + this.socket = socket; + this.input = new DataInputStream(socket.getInputStream()); + this.output = new DataOutputStream(socket.getOutputStream()); + } + + /** + * Envia uma mensagem pela conexão. + * + *A mensagem é precedida por seu tamanho para permitir que o receptor + * identifique o limite entre mensagens transmitidas pela mesma conexão.
+ * + * @param message mensagem serializada a ser enviada. + * @throws IOException caso ocorra um erro durante o envio. + */ + public void send(byte[] message) throws IOException { + synchronized (output) { + output.writeInt(message.length); + output.write(message); + output.flush(); + } + } + + /** + * Recebe uma mensagem completa da conexão. + * + *O tamanho da mensagem é lido primeiro e, em seguida, os dados + * correspondentes são recebidos.
+ * + * @return mensagem recebida. + * @throws IOException caso ocorra um erro durante a recepção. + */ + public byte[] receive() throws IOException { + int length = input.readInt(); + + if (length < 0) { + throw new IOException("Invalid message length."); + } + + return input.readNBytes(length); + } + + /** + * Verifica se o socket da conexão está fechado. + * + * @return {@code true} caso o socket esteja fechado. + */ + public boolean isClosed() { + return socket.isClosed(); + } + + /** + * Encerra a conexão TCP. + * + * @throws IOException caso ocorra um erro ao fechar o socket. + */ + public void close() throws IOException { + socket.close(); + } +}