From 45e3bf68762bc269122e6ab6be9b1e55f50239b0 Mon Sep 17 00:00:00 2001 From: GustavoHMDS Date: Mon, 17 Aug 2026 14:07:12 -0300 Subject: [PATCH] Unicast implementado, mas falta testes --- .idea/vcs.xml | 6 + .../inspectionProfiles/Project_Default.xml | 6 + pitiupi/docs/Documentation.md | 487 +++++++++++++++--- pitiupi/src/pitiupi/GUI/MainWindow.java | 33 +- .../src/pitiupi/control/HeartbeatManager.java | 2 +- pitiupi/src/pitiupi/net/SocketTCP.java | 239 +++++++++ .../net/{Socket.java => SocketUDP.java} | 92 ++-- pitiupi/src/pitiupi/net/TcpConnection.java | 88 ++++ 8 files changed, 834 insertions(+), 119 deletions(-) create mode 100644 .idea/vcs.xml create mode 100644 pitiupi/.idea/inspectionProfiles/Project_Default.xml create mode 100644 pitiupi/src/pitiupi/net/SocketTCP.java rename pitiupi/src/pitiupi/net/{Socket.java => SocketUDP.java} (55%) create mode 100644 pitiupi/src/pitiupi/net/TcpConnection.java diff --git a/.idea/vcs.xml b/.idea/vcs.xml new file mode 100644 index 0000000..94a25f7 --- /dev/null +++ b/.idea/vcs.xml @@ -0,0 +1,6 @@ + + + + + + \ No newline at end of file diff --git a/pitiupi/.idea/inspectionProfiles/Project_Default.xml b/pitiupi/.idea/inspectionProfiles/Project_Default.xml new file mode 100644 index 0000000..a05c602 --- /dev/null +++ b/pitiupi/.idea/inspectionProfiles/Project_Default.xml @@ -0,0 +1,6 @@ + + + + \ No newline at end of file diff --git a/pitiupi/docs/Documentation.md b/pitiupi/docs/Documentation.md index 2b9fee1..fb7024a 100644 --- a/pitiupi/docs/Documentation.md +++ b/pitiupi/docs/Documentation.md @@ -1,21 +1,23 @@ -# Arquitetura da aplicação +# Visão geral -## Visão geral - -O Share My Sheet é uma plataforma P2P baseada em uma arquitetura de plugins, -desenvolvida para permitir a criação de aplicações distribuídas sobre uma +O `Pitiupi` é uma plataforma P2P baseada em uma arquitetura de plugins, +desenvolvida para permitir a criação de aplicações distribuídas sobre uma infraestrutura comum de comunicação em rede. -A aplicação fornece mecanismos de descoberta de usuários, comunicação -multicast, troca de mensagens serializadas e integração de extensões por meio -de plugins. Dessa forma, novos comportamentos podem ser adicionados sem -modificar o núcleo da aplicação. +A aplicação fornece mecanismos de descoberta de usuários, comunicação entre +peers, troca de mensagens serializadas e integração de extensões por meio de +plugins. A comunicação de rede utiliza dois protocolos de transporte: UDP, +empregado principalmente na descoberta e comunicação multicast entre as +instâncias, e TCP, utilizado para comunicação direta e persistente entre peers. -O projeto foi desenvolvido como suporte às práticas da disciplina de Redes de -Computadores, servindo também como base para experimentação de aplicações +Dessa forma, novos comportamentos podem ser adicionados por meio de plugins +sem modificar o núcleo da aplicação. + +O projeto foi desenvolvido como suporte às práticas da disciplina de Redes de +Computadores, servindo também como base para experimentação de aplicações distribuídas e desenvolvimento de projetos acadêmicos. -## Arquitetura geral +# Arquitetura geral A aplicação é dividida em três grupos principais: @@ -23,103 +25,396 @@ A aplicação é dividida em três grupos principais: sendo responsável por inicializar e integrar esses componentes durante a execução. - **Sistema de plugins**: responsável por adicionar funcionalidades sem - alterar o núcleo. + alterar o núcleo da aplicação. - **Comunicação de rede**: responsável pela troca de mensagens entre instâncias da aplicação. ## Inicialização da aplicação 1. A MainWindow é criada. -2. O Socket é inicializado e inicia a comunicação multicast. -3. O HeartbeatManager é criado e inicia o envio periódico de heartbeats. -4. O painel de usuários online é registrado como listener do HeartbeatManager. -5. O PluginLoader procura e carrega os plugins disponíveis. +2. O SocketUDP é inicializado, associa-se à porta configurada e ingressa no + grupo multicast. +3. O SocketTCP é inicializado e cria um ServerSocket na porta da aplicação + para aceitar conexões TCP. +4. Os mecanismos de comunicação iniciam seus respectivos processos de recepção. +5. O HeartbeatManager é criado e inicia o envio periódico de heartbeats. +6. O painel de usuários online é registrado como listener do HeartbeatManager. +7. O PluginLoader procura e carrega os plugins disponíveis. -## Sistema de comunicação +# Sistema de comunicação -### Socket +A comunicação entre as instâncias da aplicação é realizada por meio de dois +mecanismos complementares: -A classe `Socket` é responsável pela comunicação de rede da aplicação. -Ela utiliza comunicação multicast UDP para permitir que diferentes -instâncias da aplicação troquem mensagens sem a necessidade de um servidor -central. +- **UDP multicast/unicast**, implementado por `SocketUDP` +- **TCP**, implementado por `SocketTCP` e `TcpConnection` -A classe estende `Thread`, pois mantém um processo contínuo de recepção de -mensagens enquanto a aplicação continua executando suas demais atividades. +Ambos os mecanismos utilizam mensagens serializadas em `byte[]`, permitindo que +os plugins permaneçam independentes dos detalhes do protocolo de transporte. -Durante sua inicialização, o socket: -- obtém o endereço do grupo multicast; -- cria um `MulticastSocket` utilizando a porta configurada; -- configura os buffers de comunicação; -- entra no grupo multicast. +## Comunicação UDP -#### Envio de mensagens +A classe `SocketUDP` é responsável pela comunicação UDP da aplicação. -O método `send()` recebe uma mensagem já serializada em bytes e cria um -`DatagramPacket`, enviando-o para o grupo multicast. +Ela utiliza um `MulticastSocket` associado à porta configurada pela aplicação e +ingressa no grupo multicast `224.0.0.3`. Por ser baseada em UDP, a comunicação não +estabelece uma conexão permanente entre os peers. -De forma geral, o socket atua apenas como mecanismo de transporte das mensagens. -A única exceção é o tratamento de HeartbeatMessage, utilizado pela infraestrutura -da aplicação para manter a lista de peers ativos. +A classe estende `Thread`, mantendo um processo contínuo de recepção de datagramas +enquanto a aplicação está em execução. -#### Recepção de mensagens +O SocketUDP suporta dois modos de envio: -O método `receive()` mantém um loop aguardando novas mensagens multicast. +- **multicast**, destinado ao grupo de peers +- **unicast**, destinado diretamente ao endereço IP de um peer específico -Ao receber uma mensagem: -1. Os dados são analisados para verificar se correspondem a um - `HeartbeatMessage`. -2. Caso seja um heartbeat, a mensagem é encaminhada ao `HeartbeatManager` - para atualização dos usuários ativos. -3. Caso contrário, a mensagem é enviada para todos os plugins carregados - através do método `receiveMessage()`. +### Inicialização -Esse comportamento permite que a aplicação principal trate mensagens de -infraestrutura (como heartbeat), enquanto mensagens específicas de plugins -são processadas pelas extensões correspondentes. +Durante sua inicialização, o `SocketUDP`: + +1. obtém o endereço do grupo multicast; +2. cria um `MulticastSocket` utilizando a porta configurada na `MainWindow`; +3. configura os buffers de envio e recepção; +4. habilita o reuso do endereço; +5. ingressa no grupo multicast. + +O tamanho dos buffers de comunicação é configurado para `256000` bytes. + +### Envio de mensagens UDP + +O `SocketUDP` suporta dois modos de envio de mensagens: multicast e unicast, ambos +baseados no mesmo mecanismo interno de criação e envio de DatagramPacket. + +O método `sendMulticast()` envia uma mensagem para o grupo multicast da aplicação. +Nesse caso, o `DatagramPacket` é direcionado ao endereço multicast configurado, +permitindo que todas as instâncias participantes do grupo recebam a mensagem +simultaneamente. + +Já o método `sendUnicast()` permite enviar uma mensagem diretamente para um peer +específico. Nesse caso, o `DatagramPacket` é direcionado ao endereço IP informado pelo +chamador, mantendo a mesma porta utilizada pela aplicação. + +Ambos os métodos são, na prática, abstrações de um único mecanismo interno de envio, +que utiliza o mesmo processo de criação e envio de `DatagramPacket` por baixo dos +panos. A diferença entre eles está apenas no endereço de destino: enquanto o multicast +utiliza o grupo compartilhado, o unicast direciona o pacote a um IP específico. + +### Recepção de mensagens + +O método `receive()` mantém um loop aguardando novos datagramas. + +Ao receber uma mensagem, o `SocketUDP` inicialmente tenta identificar se os dados +correspondem a uma `HeartbeatMessage`. + +O processamento segue o seguinte fluxo: + +1. O datagrama é recebido pelo `MulticastSocket`. +2. Os dados recebidos são analisados pelo método `tryParseHeartbeat()`. +3. Caso seja identificada uma `HeartbeatMessage`, ela é encaminhada ao + `HeartbeatManager`, juntamente com o endereço IP de origem. +4. Caso não seja um heartbeat, a mensagem é encaminhada para todos os plugins + carregados. +5. Cada plugin decide se a mensagem recebida pertence à sua funcionalidade. + +Essa separação mantém as mensagens de infraestrutura, como heartbeat, sob +responsabilidade do núcleo da aplicação, enquanto as mensagens específicas das +funcionalidades são processadas pelos plugins. + +### Identificação de mensagens de heartbeat + +O método `tryParseHeartbeat()` tenta desserializar os dados recebidos e verificar se o +objeto resultante é uma instância de `HeartbeatMessage`. + +Caso a mensagem não possa ser desserializada ou não corresponda ao tipo esperado, o +método retorna `null`. Nesse caso, a mensagem é tratada como uma mensagem destinada aos +plugins. + +### Encerramento + +O método `close()` interrompe a execução do processo de recepção, remove a aplicação +do grupo multicast e fecha o `MulticastSocket`. + +A variável `running` é utilizada para sinalizar que o processo de recepção deve ser +encerrado. + +## Comunicação TCP + +A comunicação TCP é implementada pela classe `SocketTCP`. + +Diferentemente do UDP, o TCP estabelece conexões entre dois peers. Essas conexões +são mantidas enquanto forem necessárias e podem ser reutilizadas para o envio de +múltiplas mensagens. + +A arquitetura utiliza uma conexão TCP persistente por endereço de peer. + +O SocketTCP possui um mapa de conexões: + +`InetAddress` -> `TcpConnection` + +Esse mapa permite localizar rapidamente uma conexão existente para determinado peer e +reutilizá-la durante novos envios. + +### SocketTCP + +A classe `SocketTCP` estende `Thread` e implementa `PeerListener`. + +Ao ser criada, ela inicializa um `ServerSocket` utilizando a porta configurada na +aplicação. + +O servidor TCP permanece aguardando novas conexões enquanto a aplicação estiver +em execução. + +### Envio de mensagens TCP + +O método `send()` recebe uma mensagem serializada e o endereço IP do peer destinatário. + +Antes de enviar a mensagem, o `SocketTCP` verifica se já existe uma conexão válida com +o peer. + +O comportamento é: + +1. Procurar uma conexão existente no mapa de conexões. +2. Verificar se a conexão está fechada. +3. Caso não exista uma conexão válida, criar um novo Socket TCP para o endereço do peer. +4. Encapsular o socket em uma `TcpConnection`. +5. Registrar a conexão no mapa. +6. Iniciar uma tarefa de recepção para essa conexão. +7. Enviar a mensagem pela conexão. + +Assim, uma conexão TCP existente pode ser reutilizada para várias mensagens, evitando +a criação de um novo socket para cada envio. + +### Recepção de conexões + +O método `receive()` permanece aguardando novas conexões através de +`ServerSocket.accept()`. + +Quando uma conexão é aceita: + +O endereço IP do peer é obtido. +Um objeto `TcpConnection` é criado para encapsular o socket. +A conexão é registrada no mapa de conexões. +Caso já exista uma conexão associada ao mesmo endereço, a conexão anterior é encerrada. +Uma tarefa independente de recepção é criada para a nova conexão. + +Esse modelo permite que o servidor aceite múltiplas conexões sem bloquear o +processamento das conexões já existentes. + +### Recepção concorrente + +A recepção das mensagens TCP é realizada de maneira concorrente por meio de um +`ExecutorService`. + +Cada conexão possui uma tarefa própria responsável por receber continuamente as +mensagens. + +O fluxo de recepção é: + +``` +SocketTCP +| ++-- Peer A → TcpConnection → tarefa de recepção +| ++-- Peer B → TcpConnection → tarefa de recepção +| ++-- Peer C → TcpConnection → tarefa de recepção +``` + +Dessa forma, uma conexão bloqueada ou aguardando dados não impede que outras conexões +continuem sendo processadas. + +Quando uma mensagem é recebida por uma `TcpConnection`, ela é encaminhada para todos os +plugins registrados na `MainWindow`. + +O plugin responsável pela funcionalidade deve identificar o tipo da mensagem e +realizar o processamento correspondente. + +### TcpConnection + +A classe TcpConnection encapsula uma conexão TCP individual entre dois peers. + +Ela mantém: + +o Socket utilizado na conexão; +um DataInputStream para recepção; +um DataOutputStream para envio. + +Além de encapsular esses recursos, a classe é responsável pelo framing das mensagens. + +### Framing das mensagens TCP + +A classe `TcpConnection` encapsula uma conexão TCP individual entre dois peers. + +Ela mantém: + +1. o `Socket` utilizado na conexão; +2. um `DataInputStream` para recepção; +3. um `DataOutputStream` para envio. + +Além de encapsular esses recursos, a classe é responsável pelo framing das mensagens. + +Framing das mensagens TCP + +O protocolo TCP fornece um fluxo contínuo de bytes e não preserva os limites das +mensagens enviadas pela aplicação. + +Por esse motivo, é adicionado o tamanho da mensagem antes do conteúdo. + +O formato utilizado é: + +``` ++----------------------+----------------------+ +| tamanho da mensagem | dados da mensagem | ++----------------------+----------------------+ +4 bytes N bytes +``` + +O tamanho é enviado como um `int`, seguido pelos bytes da mensagem. + +No envio, o método `TcpConnection.send()`: + +1. obtém o tamanho do vetor de bytes; +2. escreve o tamanho no `DataOutputStream`; +3. escreve os dados da mensagem; +4. realiza `flush()` no stream. + +O acesso ao stream de saída é sincronizado para evitar que envios concorrentes +misturem seus dados. + +Na recepção, o método `receive()`: + +1. lê o tamanho da mensagem; +2. verifica se o tamanho é válido; +3. lê a quantidade correspondente de bytes; +4. retorna o vetor de bytes completo. + +Esse mecanismo permite que várias mensagens sejam transmitidas pela mesma conexão +TCP sem que o receptor perca a delimitação entre elas. + +### Integração entre TCP e gerenciamento de peers + +O `SocketTCP` implementa a interface PeerListener para acompanhar as alterações na +lista de peers ativos. + +O `HeartbeatManager` mantém a lista de peers atualmente detectados na rede. Quando +essa lista é alterada, o SocketTCP recebe a notificação por meio do método +`onPeersChanged()`. + +O método obtém os endereços IP dos peers atualmente ativos e compara esses endereços +com as conexões TCP existentes. + +Quando uma conexão pertence a um peer que não está mais ativo: + +a conexão é removida do mapa; +a conexão TCP é encerrada; +os recursos associados são liberados. + +Esse mecanismo impede que conexões TCP permaneçam abertas indefinidamente para peers +que já deixaram de participar da rede. + +O fluxo de integração pode ser representado da seguinte forma: +``` +HeartbeatManager +| +| peers ativos alterados +| +v +SocketTCP.onPeersChanged() +| +| compara peers ativos +| +v ++---- peer ativo ------> mantém conexão +| ++---- peer inativo ----> fecha conexão +``` ### Message -A classe `Message` representa a estrutura base das mensagens trocadas pela -aplicação. +A classe `Message` representa a estrutura base das mensagens trocadas pela aplicação. -Todas as mensagens utilizadas pelo sistema devem herdar dessa classe. Ela -implementa `Serializable`, permitindo que objetos de mensagem sejam -convertidos em vetores de bytes para transmissão pela rede. +Todas as mensagens utilizadas pelo sistema devem herdar dessa classe. Ela implementa +`Serializable`, permitindo que objetos de mensagem sejam convertidos em vetores de +bytes para transmissão pela rede. -A serialização é realizada pelo método `toByteArray()`, que transforma uma -instância da mensagem em uma representação binária enviada pelo `Socket`. +A serialização é realizada pelo método `toByteArray()`, que transforma uma instância +da mensagem em uma representação binária enviada pelos mecanismos de comunicação. -Mensagens específicas devem estender essa classe adicionando os atributos e +As mensagens específicas devem estender essa classe, adicionando os atributos e comportamentos necessários para cada funcionalidade. +A camada de transporte não depende do conteúdo específico das mensagens. Tanto o +UDP quanto o TCP recebem mensagens serializadas como `byte[]`. + ### Heartbeat -O mecanismo de heartbeat é utilizado para identificar quais usuários estão -ativos na rede. +O mecanismo de heartbeat é utilizado para identificar quais usuários estão ativos +na rede. -A aplicação envia periodicamente uma `HeartbeatMessage` contendo -informações do usuário, um identificador único da instância da aplicação -e um timestamp. Quando outra instância recebe essa mensagem, o `HeartbeatManager` -registra ou atualiza o peer correspondente, utilizando o identificador -da instância para diferenciá-lo das demais. +A aplicação envia periodicamente uma `HeartbeatMessage` contendo informações do usuário, +um identificador único da instância da aplicação e um timestamp. -Um peer é considerado inativo quando permanece mais de 15 segundos -sem receber um novo heartbeat. O `HeartbeatManager` então remove o peer -da lista e notifica os listeners sobre a alteração. +Essas mensagens são transmitidas utilizando a comunicação UDP multicast. -O fluxo de heartbeat é independente das mensagens dos plugins: +Quando outra instância recebe uma mensagem de heartbeat, o `SocketUDP` identifica +o tipo da mensagem e encaminha o objeto ao `HeartbeatManager`, juntamente com o endereço +IP de origem. -1. `HeartbeatManager` cria uma `HeartbeatMessage`. -2. A mensagem é serializada utilizando `Message.toByteArray()`. -3. O `Socket` transmite a mensagem pela rede multicast. -4. O `Socket` identifica mensagens de heartbeat recebidas e encaminha ao - `HeartbeatManager`. -5. O `HeartbeatManager` atualiza a lista de peers ativos e notifica os - listeners. +O `HeartbeatManager` registra ou atualiza o peer correspondente, utilizando o +identificador da instância para diferenciá-lo das demais. + +Um peer é considerado inativo quando permanece mais de 15 segundos sem receber um novo +heartbeat. O `HeartbeatManager` então remove o peer da lista e notifica todos os +componentes registrados como listeners de mudanças de peers. + +Esses listeners podem incluir diferentes partes da aplicação, como componentes +da interface gráfica, mecanismos de comunicação (como o `SocketTCP`) e também plugins +que desejem reagir à presença ou ausência de usuários na rede. + +Dessa forma, o sistema permite que tanto o núcleo da aplicação quanto extensões +externas sejam notificados sobre alterações no estado da rede, mantendo a arquitetura +flexível e extensível. + +O fluxo completo é: + +``` +HeartbeatManager +| +| cria HeartbeatMessage +| +v +Message.toByteArray() +| +v +SocketUDP.sendMulticast() +| +v +UDP Multicast +| +v +SocketUDP.receive() +| +v +HeartbeatManager.receiveHeartbeat() +| +v +Lista de peers ativos +| ++----> Interface de usuários online +| ++----> SocketTCP (PeerListener) +| ++----> Plugins (PeerListener) +``` ## Sistema de plugins +O sistema de plugins é responsável por permitir a extensão das funcionalidades da +aplicação de forma modular, sem a necessidade de alterações no núcleo do sistema. +Ele define um contrato comum que deve ser seguido por todas as extensões, garantindo +que possam ser carregadas dinamicamente e integradas à comunicação e à interface da +aplicação. + ### Interface Plugin A interface `Plugin` define o contrato que deve ser implementado por qualquer @@ -186,7 +481,7 @@ MeusPlugins.plugins.NomeDaClasse Para serem carregados, os plugins devem estar empacotados como arquivos JAR e colocados no diretório plugins da aplicação. -### Carregamento +### Carregamento de plugins Ao iniciar a aplicação, o `PluginLoader` executa os seguintes passos: @@ -209,6 +504,44 @@ Durante a inicialização, cada plugin recebe uma referência para a Essa referência permite integrar componentes gráficos à aplicação, sendo a principal forma de extensão a adição de menus à barra de menus: -```java +``` JMenu menu = new JMenu("Meu Plugin"); -window.addMenu(menu); \ No newline at end of file +window.addMenu(menu); +``` + +Visão geral da comunicação + +A arquitetura de comunicação pode ser resumida da seguinte maneira: +``` + +------------------+ + | MainWindow | + +--------+---------+ + | + +--------------------+--------------------+ + | | | + v v v + HeartbeatManager SocketUDP SocketTCP + | | | + | | | + v v v + Heartbeat UDP Multicast TCP Connections + | | + | | + +---------+----------+ + | + v + Plugins + | + v + Message +``` + +O UDP multicast fornece a infraestrutura de descoberta e manutenção da presença dos +peers, além de permitir comunicação multicast e unicast. + +O TCP fornece canais diretos e persistentes entre peers, com uma conexão reutilizável +por endereço. + +Os dois mecanismos encaminham as mensagens destinadas às funcionalidades da aplicação +para os plugins, mantendo separadas as responsabilidades de transporte, infraestrutura +e lógica específica de cada extensão. \ No newline at end of file diff --git a/pitiupi/src/pitiupi/GUI/MainWindow.java b/pitiupi/src/pitiupi/GUI/MainWindow.java index f1276aa..e3f4ffd 100644 --- a/pitiupi/src/pitiupi/GUI/MainWindow.java +++ b/pitiupi/src/pitiupi/GUI/MainWindow.java @@ -5,7 +5,8 @@ */ package pitiupi.GUI; -import pitiupi.net.Socket; +import pitiupi.net.SocketTCP; +import pitiupi.net.SocketUDP; import java.util.ArrayList; import java.util.List; import javax.swing.JMenu; @@ -39,7 +40,8 @@ public class MainWindow extends javax.swing.JFrame { private OnlineUsersPanel onlineUsersPanel; - private Socket socket; + private SocketUDP socketUDP; + private SocketTCP socketTCP; private String userName; private int port; @@ -91,8 +93,11 @@ public class MainWindow extends javax.swing.JFrame { } private void initSocket() { - this.socket = new Socket(this); - this.socket.start(); + this.socketUDP = new SocketUDP(this); + this.socketTCP = new SocketTCP(this); + + this.socketUDP.start(); + this.socketTCP.start(); } private void initHeartbeat() { @@ -130,8 +135,8 @@ public class MainWindow extends javax.swing.JFrame { return this.port; } - public Socket getSocket() { - return this.socket; + public SocketUDP getSocket() { + return this.socketUDP; } public List getPlugins() { @@ -154,11 +159,17 @@ public class MainWindow extends javax.swing.JFrame { } private void reconnect() { - if (this.socket != null) { - this.socket.close(); + if (this.socketUDP != null) { + this.socketUDP.close(); + } + if (this.socketTCP != null) { + this.socketTCP.close(); } - this.socket = new Socket(this); - this.socket.start(); - } + this.socketUDP = new SocketUDP(this); + this.socketTCP = new SocketTCP(this); + + this.socketUDP.start(); + this.socketTCP.start(); + } } diff --git a/pitiupi/src/pitiupi/control/HeartbeatManager.java b/pitiupi/src/pitiupi/control/HeartbeatManager.java index 725cede..6fb380e 100644 --- a/pitiupi/src/pitiupi/control/HeartbeatManager.java +++ b/pitiupi/src/pitiupi/control/HeartbeatManager.java @@ -66,7 +66,7 @@ public class HeartbeatManager { private void announce() { try { HeartbeatMessage msg = new HeartbeatMessage(mainWindow.getUserName(), instanceId); - mainWindow.getSocket().send(msg.toByteArray()); + mainWindow.getSocket().sendMulticast(msg.toByteArray()); } catch (IOException ex) { Logger.getLogger(HeartbeatManager.class.getName()).log(Level.WARNING, "Falha ao enviar heartbeat", ex); } diff --git a/pitiupi/src/pitiupi/net/SocketTCP.java b/pitiupi/src/pitiupi/net/SocketTCP.java new file mode 100644 index 0000000..a798af6 --- /dev/null +++ b/pitiupi/src/pitiupi/net/SocketTCP.java @@ -0,0 +1,239 @@ +package pitiupi.net; + +import pitiupi.GUI.MainWindow; +import pitiupi.control.PeerInfo; +import pitiupi.control.PeerListener; +import pitiupi.plugin.Plugin; + +import java.io.*; +import java.net.*; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.stream.Collectors; + +/** + * Gerencia a comunicação TCP da aplicação. + * + *

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 Map connections; + private ServerSocket serverSocket; + + private final MainWindow main; + + private volatile boolean running = true; + + /** + * Cria o servidor TCP da aplicação. + * + *

O 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(List peers) { + Set activePeers = peers.stream() + .map(PeerInfo::getAddress) + .collect(Collectors.toSet()); + + for (InetAddress address : connections.keySet()) { + if (!activePeers.contains(address)) { + + TcpConnection connection = connections.remove(address); + if (connection != null) { + try { + connection.close(); + } catch (IOException e) { + + } + } + } + } + } +} + diff --git a/pitiupi/src/pitiupi/net/Socket.java b/pitiupi/src/pitiupi/net/SocketUDP.java similarity index 55% rename from pitiupi/src/pitiupi/net/Socket.java rename to pitiupi/src/pitiupi/net/SocketUDP.java index 91c2c89..285e314 100644 --- a/pitiupi/src/pitiupi/net/Socket.java +++ b/pitiupi/src/pitiupi/net/SocketUDP.java @@ -10,43 +10,48 @@ import java.io.ByteArrayInputStream; import java.io.IOException; import java.io.ObjectInput; import java.io.ObjectInputStream; -import java.net.DatagramPacket; -import java.net.InetAddress; -import java.net.MulticastSocket; -import java.net.UnknownHostException; +import java.net.*; import java.util.logging.Level; import java.util.logging.Logger; import pitiupi.plugin.Plugin; -import pitiupi.net.HeartbeatMessage; /** - * "Responsável pela comunicação multicast da aplicação e pelo encaminhamento de mensagens aos plugins." + * Gerencia a comunicação UDP multicast e unicast da aplicação. * - *

Esta 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(); + } +}