Compare commits

..

No commits in common. "e9ce17abb3c30bd44af17c37b53fe6d1347d5af1" and "81c37225fac91d24e8651da729d0282cc09af6a2" have entirely different histories.

13 changed files with 207 additions and 828 deletions

5
.idea/misc.xml Normal file
View File

@ -0,0 +1,5 @@
<project version="4">
<component name="ProjectRootManager">
<output url="file://$PROJECT_DIR$/out" />
</component>
</project>

83
.idea/workspace.xml Normal file
View File

@ -0,0 +1,83 @@
<?xml version="1.0" encoding="UTF-8"?>
<project version="4">
<component name="AutoImportSettings">
<option name="autoReloadType" value="SELECTIVE" />
</component>
<component name="ChangeListManager">
<list default="true" id="62118a1b-608f-4144-b29c-cf46419ca676" name="Changes" comment="" />
<option name="SHOW_DIALOG" value="false" />
<option name="HIGHLIGHT_CONFLICTS" value="true" />
<option name="HIGHLIGHT_NON_ACTIVE_CHANGELIST" value="false" />
<option name="LAST_RESOLUTION" value="IGNORE" />
</component>
<component name="Git.Settings">
<option name="RECENT_GIT_ROOT_PATH" value="$PROJECT_DIR$" />
</component>
<component name="ProjectColorInfo">{
&quot;associatedIndex&quot;: 0
}</component>
<component name="ProjectId" id="3GbOdQKFlKqBjAjqjAv5TRPxr6Z" />
<component name="ProjectViewState">
<option name="hideEmptyMiddlePackages" value="true" />
<option name="showLibraryContents" value="true" />
</component>
<component name="PropertiesComponent">{
&quot;keyToString&quot;: {
&quot;Application.Main.executor&quot;: &quot;Run&quot;,
&quot;ModuleVcsDetector.initialDetectionPerformed&quot;: &quot;true&quot;,
&quot;RunOnceActivity.ShowReadmeOnStart&quot;: &quot;true&quot;,
&quot;RunOnceActivity.TerminalTabsStorage.copyFrom.TerminalArrangementManager.252&quot;: &quot;true&quot;,
&quot;RunOnceActivity.git.unshallow&quot;: &quot;true&quot;,
&quot;git-widget-placeholder&quot;: &quot;feature/help-and-documentation&quot;,
&quot;ignore.virus.scanning.warn.message&quot;: &quot;true&quot;,
&quot;kotlin-language-version-configured&quot;: &quot;true&quot;,
&quot;last_opened_file_path&quot;: &quot;C:/Users/extre/OneDrive/Documentos/Faculdade/2026_02/Monitoria/Pitiupi&quot;,
&quot;settings.editor.selected.configurable&quot;: &quot;preferences.keymap&quot;
}
}</component>
<component name="RecentsManager">
<key name="CopyFile.RECENT_KEYS">
<recent name="C:\Users\extre\OneDrive\Documentos\Faculdade\2026_02\Monitoria\Pitiupi\pitiupi\plugins" />
</key>
<key name="MoveFile.RECENT_KEYS">
<recent name="C:\Users\extre\OneDrive\Documentos\Faculdade\2026_02\Monitoria\Pitiupi\pitiupi\plugins" />
</key>
</component>
<component name="RunManager">
<configuration name="Main" type="Application" factoryName="Application" temporary="true" nameIsGenerated="true">
<option name="MAIN_CLASS_NAME" value="pitiupi.Main" />
<module name="pitiupi" />
<extension name="coverage">
<pattern>
<option name="PATTERN" value="pitiupi.*" />
<option name="ENABLED" value="true" />
</pattern>
</extension>
<method v="2">
<option name="Make" enabled="true" />
</method>
</configuration>
<recent_temporary>
<list>
<item itemvalue="Application.Main" />
</list>
</recent_temporary>
</component>
<component name="SharedIndexes">
<attachedChunks>
<set>
<option value="bundled-jdk-9823dce3aa75-bf35d07a577b-intellij.indexing.shared.core-IU-252.28539.54" />
</set>
</attachedChunks>
</component>
<component name="TaskManager">
<task active="true" id="Default" summary="Default task">
<changelist id="62118a1b-608f-4144-b29c-cf46419ca676" name="Changes" comment="" />
<created>1784236888334</created>
<option name="number" value="Default" />
<option name="presentableId" value="Default" />
<updated>1784236888334</updated>
</task>
<servers />
</component>
</project>

BIN
pitiupi.chat/dist/pitiupi.chat.jar vendored Normal file

Binary file not shown.

Binary file not shown.

View File

@ -1,6 +0,0 @@
<component name="InspectionProjectProfileManager">
<profile version="1.0">
<option name="myName" value="Project Default" />
<inspection_tool class="LanguageDetectionInspection" enabled="false" level="WEAK WARNING" enabled_by_default="false" />
</profile>
</component>

View File

@ -1,23 +1,21 @@
# Visão geral # Arquitetura da aplicação
O `Pitiupi` é uma plataforma P2P baseada em uma arquitetura de plugins, ## 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 desenvolvida para permitir a criação de aplicações distribuídas sobre uma
infraestrutura comum de comunicação em rede. infraestrutura comum de comunicação em rede.
A aplicação fornece mecanismos de descoberta de usuários, comunicação entre A aplicação fornece mecanismos de descoberta de usuários, comunicação
peers, troca de mensagens serializadas e integração de extensões por meio de multicast, troca de mensagens serializadas e integração de extensões por meio
plugins. A comunicação de rede utiliza dois protocolos de transporte: UDP, de plugins. Dessa forma, novos comportamentos podem ser adicionados sem
empregado principalmente na descoberta e comunicação multicast entre as modificar o núcleo da aplicação.
instâncias, e TCP, utilizado para comunicação direta e persistente entre peers.
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 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 Computadores, servindo também como base para experimentação de aplicações
distribuídas e desenvolvimento de projetos acadêmicos. distribuídas e desenvolvimento de projetos acadêmicos.
# Arquitetura geral ## Arquitetura geral
A aplicação é dividida em três grupos principais: A aplicação é dividida em três grupos principais:
@ -25,396 +23,103 @@ A aplicação é dividida em três grupos principais:
sendo responsável por inicializar e integrar esses componentes durante a sendo responsável por inicializar e integrar esses componentes durante a
execução. execução.
- **Sistema de plugins**: responsável por adicionar funcionalidades sem - **Sistema de plugins**: responsável por adicionar funcionalidades sem
alterar o núcleo da aplicação. alterar o núcleo.
- **Comunicação de rede**: responsável pela troca de mensagens entre - **Comunicação de rede**: responsável pela troca de mensagens entre
instâncias da aplicação. instâncias da aplicação.
## Inicialização da aplicação ## Inicialização da aplicação
1. A MainWindow é criada. 1. A MainWindow é criada.
2. O SocketUDP é inicializado, associa-se à porta configurada e ingressa no 2. O Socket é inicializado e inicia a comunicação multicast.
grupo multicast. 3. O HeartbeatManager é criado e inicia o envio periódico de heartbeats.
3. O SocketTCP é inicializado e cria um ServerSocket na porta da aplicação 4. O painel de usuários online é registrado como listener do HeartbeatManager.
para aceitar conexões TCP. 5. O PluginLoader procura e carrega os plugins disponíveis.
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
A comunicação entre as instâncias da aplicação é realizada por meio de dois ### Socket
mecanismos complementares:
- **UDP multicast/unicast**, implementado por `SocketUDP` A classe `Socket` é responsável pela comunicação de rede da aplicação.
- **TCP**, implementado por `SocketTCP` e `TcpConnection` Ela utiliza comunicação multicast UDP para permitir que diferentes
instâncias da aplicação troquem mensagens sem a necessidade de um servidor
central.
Ambos os mecanismos utilizam mensagens serializadas em `byte[]`, permitindo que A classe estende `Thread`, pois mantém um processo contínuo de recepção de
os plugins permaneçam independentes dos detalhes do protocolo de transporte. mensagens enquanto a aplicação continua executando suas demais atividades.
## Comunicação UDP 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.
A classe `SocketUDP` é responsável pela comunicação UDP da aplicação. #### Envio de mensagens
Ela utiliza um `MulticastSocket` associado à porta configurada pela aplicação e O método `send()` recebe uma mensagem já serializada em bytes e cria um
ingressa no grupo multicast `224.0.0.3`. Por ser baseada em UDP, a comunicação não `DatagramPacket`, enviando-o para o grupo multicast.
estabelece uma conexão permanente entre os peers.
A classe estende `Thread`, mantendo um processo contínuo de recepção de datagramas De forma geral, o socket atua apenas como mecanismo de transporte das mensagens.
enquanto a aplicação está em execução. A única exceção é o tratamento de HeartbeatMessage, utilizado pela infraestrutura
da aplicação para manter a lista de peers ativos.
O SocketUDP suporta dois modos de envio: #### Recepção de mensagens
- **multicast**, destinado ao grupo de peers O método `receive()` mantém um loop aguardando novas mensagens multicast.
- **unicast**, destinado diretamente ao endereço IP de um peer específico
### Inicialização 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()`.
Durante sua inicialização, o `SocketUDP`: Esse comportamento permite que a aplicação principal trate mensagens de
infraestrutura (como heartbeat), enquanto mensagens específicas de plugins
1. obtém o endereço do grupo multicast; são processadas pelas extensões correspondentes.
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 ### 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 Todas as mensagens utilizadas pelo sistema devem herdar dessa classe. Ela
`Serializable`, permitindo que objetos de mensagem sejam convertidos em vetores de implementa `Serializable`, permitindo que objetos de mensagem sejam
bytes para transmissão pela rede. convertidos em vetores de bytes para transmissão pela rede.
A serialização é realizada pelo método `toByteArray()`, que transforma uma instância A serialização é realizada pelo método `toByteArray()`, que transforma uma
da mensagem em uma representação binária enviada pelos mecanismos de comunicação. instância da mensagem em uma representação binária enviada pelo `Socket`.
As mensagens específicas devem estender essa classe, adicionando os atributos e Mensagens específicas devem estender essa classe adicionando os atributos e
comportamentos necessários para cada funcionalidade. 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 ### Heartbeat
O mecanismo de heartbeat é utilizado para identificar quais usuários estão ativos O mecanismo de heartbeat é utilizado para identificar quais usuários estão
na rede. ativos na rede.
A aplicação envia periodicamente uma `HeartbeatMessage` contendo informações do usuário, A aplicação envia periodicamente uma `HeartbeatMessage` contendo
um identificador único da instância da aplicação e um timestamp. 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.
Essas mensagens são transmitidas utilizando a comunicação UDP multicast. 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.
Quando outra instância recebe uma mensagem de heartbeat, o `SocketUDP` identifica O fluxo de heartbeat é independente das mensagens dos plugins:
o tipo da mensagem e encaminha o objeto ao `HeartbeatManager`, juntamente com o endereço
IP de origem.
O `HeartbeatManager` registra ou atualiza o peer correspondente, utilizando o 1. `HeartbeatManager` cria uma `HeartbeatMessage`.
identificador da instância para diferenciá-lo das demais. 2. A mensagem é serializada utilizando `Message.toByteArray()`.
3. O `Socket` transmite a mensagem pela rede multicast.
Um peer é considerado inativo quando permanece mais de 15 segundos sem receber um novo 4. O `Socket` identifica mensagens de heartbeat recebidas e encaminha ao
heartbeat. O `HeartbeatManager` então remove o peer da lista e notifica todos os `HeartbeatManager`.
componentes registrados como listeners de mudanças de peers. 5. O `HeartbeatManager` atualiza a lista de peers ativos e notifica os
listeners.
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 ## 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 ### Interface Plugin
A interface `Plugin` define o contrato que deve ser implementado por qualquer A interface `Plugin` define o contrato que deve ser implementado por qualquer
@ -481,7 +186,7 @@ MeusPlugins.plugins.NomeDaClasse
Para serem carregados, os plugins devem estar empacotados como arquivos JAR Para serem carregados, os plugins devem estar empacotados como arquivos JAR
e colocados no diretório plugins da aplicação. e colocados no diretório plugins da aplicação.
### Carregamento de plugins ### Carregamento
Ao iniciar a aplicação, o `PluginLoader` executa os seguintes passos: Ao iniciar a aplicação, o `PluginLoader` executa os seguintes passos:
@ -504,44 +209,6 @@ Durante a inicialização, cada plugin recebe uma referência para a
Essa referência permite integrar componentes gráficos à aplicação, sendo 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: principal forma de extensão a adição de menus à barra de menus:
``` ```java
JMenu menu = new JMenu("Meu Plugin"); JMenu menu = new JMenu("Meu Plugin");
window.addMenu(menu); 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.

Binary file not shown.

View File

@ -5,8 +5,7 @@
*/ */
package pitiupi.GUI; package pitiupi.GUI;
import pitiupi.net.SocketTCP; import pitiupi.net.Socket;
import pitiupi.net.SocketUDP;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; import java.util.List;
import javax.swing.JMenu; import javax.swing.JMenu;
@ -40,8 +39,7 @@ public class MainWindow extends javax.swing.JFrame {
private OnlineUsersPanel onlineUsersPanel; private OnlineUsersPanel onlineUsersPanel;
private SocketUDP socketUDP; private Socket socket;
private SocketTCP socketTCP;
private String userName; private String userName;
private int port; private int port;
@ -93,11 +91,8 @@ public class MainWindow extends javax.swing.JFrame {
} }
private void initSocket() { private void initSocket() {
this.socketUDP = new SocketUDP(this); this.socket = new Socket(this);
this.socketTCP = new SocketTCP(this); this.socket.start();
this.socketUDP.start();
this.socketTCP.start();
} }
private void initHeartbeat() { private void initHeartbeat() {
@ -135,8 +130,8 @@ public class MainWindow extends javax.swing.JFrame {
return this.port; return this.port;
} }
public SocketUDP getSocket() { public Socket getSocket() {
return this.socketUDP; return this.socket;
} }
public List<Plugin> getPlugins() { public List<Plugin> getPlugins() {
@ -159,17 +154,11 @@ public class MainWindow extends javax.swing.JFrame {
} }
private void reconnect() { private void reconnect() {
if (this.socketUDP != null) { if (this.socket != null) {
this.socketUDP.close(); this.socket.close();
} }
if (this.socketTCP != null) { this.socket = new Socket(this);
this.socketTCP.close(); this.socket.start();
}
this.socketUDP = new SocketUDP(this);
this.socketTCP = new SocketTCP(this);
this.socketUDP.start();
this.socketTCP.start();
} }
} }

View File

@ -66,7 +66,7 @@ public class HeartbeatManager {
private void announce() { private void announce() {
try { try {
HeartbeatMessage msg = new HeartbeatMessage(mainWindow.getUserName(), instanceId); HeartbeatMessage msg = new HeartbeatMessage(mainWindow.getUserName(), instanceId);
mainWindow.getSocket().sendMulticast(msg.toByteArray()); mainWindow.getSocket().send(msg.toByteArray());
} catch (IOException ex) { } catch (IOException ex) {
Logger.getLogger(HeartbeatManager.class.getName()).log(Level.WARNING, "Falha ao enviar heartbeat", ex); Logger.getLogger(HeartbeatManager.class.getName()).log(Level.WARNING, "Falha ao enviar heartbeat", ex);
} }

View File

@ -10,48 +10,43 @@ import java.io.ByteArrayInputStream;
import java.io.IOException; import java.io.IOException;
import java.io.ObjectInput; import java.io.ObjectInput;
import java.io.ObjectInputStream; import java.io.ObjectInputStream;
import java.net.*; import java.net.DatagramPacket;
import java.net.InetAddress;
import java.net.MulticastSocket;
import java.net.UnknownHostException;
import java.util.logging.Level; import java.util.logging.Level;
import java.util.logging.Logger; import java.util.logging.Logger;
import pitiupi.plugin.Plugin; import pitiupi.plugin.Plugin;
import pitiupi.net.HeartbeatMessage;
/** /**
* Gerencia a comunicação UDP multicast e unicast da aplicação. * "Responsável pela comunicação multicast da aplicação e pelo encaminhamento de mensagens aos plugins."
* *
* <p>Mantém a conexão com o grupo multicast da aplicação e permite o envio * <p>Esta classe gerencia o envio e o recebimento de mensagens, além de
* de mensagens tanto para o grupo quanto diretamente para um peer.</p> * manter a conexão com o grupo multicast utilizado pela aplicação.</p>
*
* <p>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.</p>
* *
* @author flavio * @author flavio
* @author tony * @author tony
* @author Gustavo
*/ */
public class SocketUDP extends Thread { public class Socket extends Thread {
private MulticastSocket multicastSocket; private MulticastSocket multicastSocket;
private InetAddress address; private InetAddress address;
private MainWindow main; private MainWindow main;
private volatile boolean running = true; private volatile boolean running = true;
public final static String INET_ADDR = "224.0.0.3"; public final static String INET_ADDR = "224.0.0.3";
/** /**
* Cria e inicializa o socket UDP utilizado pela aplicação. * Cria e inicializa a conexão multicast utilizada pela aplicação.
*
* <p>O socket é associado à porta da aplicação e ingressa no grupo
* multicast configurado.</p>
* *
* @param main janela principal da aplicação. * @param main janela principal da aplicação.
*/ */
public SocketUDP(MainWindow main) { public Socket(MainWindow main) {
this.main = main; this.main = main;
try { try {
this.address = InetAddress.getByName(SocketUDP.INET_ADDR); this.address = InetAddress.getByName(Socket.INET_ADDR);
} catch (UnknownHostException ex) { } catch (UnknownHostException ex) {
System.getLogger(SocketUDP.class.getName()).log(System.Logger.Level.ERROR, (String) null, ex); System.getLogger(Socket.class.getName()).log(System.Logger.Level.ERROR, (String) null, ex);
} }
try { try {
multicastSocket = new MulticastSocket(this.main.getPort()); multicastSocket = new MulticastSocket(this.main.getPort());
@ -61,7 +56,7 @@ public class SocketUDP extends Thread {
multicastSocket.joinGroup(address); multicastSocket.joinGroup(address);
} catch (IOException ex) { } catch (IOException ex) {
System.out.println("There is no socket connection. Sorry."); System.out.println("There is no socket connection. Sorry.");
System.out.println(ex); System.out.println(ex.toString());
} }
} }
@ -69,33 +64,11 @@ public class SocketUDP extends Thread {
* Envia uma mensagem para o grupo multicast da aplicação. * Envia uma mensagem para o grupo multicast da aplicação.
* *
* @param msg mensagem serializada a ser enviada. * @param msg mensagem serializada a ser enviada.
* @throws IOException caso ocorra um erro durante o envio. * @throws IOException caso ocorra um erro durante o envio da mensagem.
*/ */
public void sendMulticast(byte[] msg) throws IOException { public void send(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; DatagramPacket msgPacket;
msgPacket = new DatagramPacket(msg, msg.length, destinationAddress, this.main.getPort()); msgPacket = new DatagramPacket(msg, msg.length, this.address, this.main.getPort());
multicastSocket.send(msgPacket); multicastSocket.send(msgPacket);
} }
@ -108,7 +81,7 @@ public class SocketUDP extends Thread {
try { try {
multicastSocket.leaveGroup(address); multicastSocket.leaveGroup(address);
} catch (IOException ex) { } catch (IOException ex) {
System.out.println(ex);
} }
multicastSocket.close(); multicastSocket.close();
} }
@ -116,17 +89,15 @@ public class SocketUDP extends Thread {
/** /**
* Recebe e processa mensagens UDP. * Inicia o loop de recepção de mensagens da aplicação.
* *
* <p>Mensagens de heartbeat são encaminhadas ao gerenciador de heartbeat. * <p>As mensagens recebidas são verificadas inicialmente para identificar
* As demais mensagens são encaminhadas aos plugins carregados pela * mensagens de heartbeat. Caso não sejam heartbeats, seu conteúdo é
* aplicação.</p> * encaminhado para todos os plugins carregados, que decidem se devem ou
* não processá-lo.</p>
* *
* @throws UnknownHostException caso não seja possível resolver o endereço * @throws IOException caso ocorra um erro durante a recepção da mensagem.
* utilizado pela comunicação. * @throws ClassNotFoundException caso a desserialização de uma mensagem falhe.
* @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 { public void receive() throws UnknownHostException, IOException, ClassNotFoundException {
byte[] buf = new byte[256000]; byte[] buf = new byte[256000];
@ -153,11 +124,11 @@ public class SocketUDP extends Thread {
} }
/** /**
* Tenta identificar os dados recebidos como uma mensagem de heartbeat. * Tenta desserializar os dados recebidos como uma {@code HeartbeatMessage}.
* *
* @param data dados serializados recebidos. * @param data dados serializados da mensagem.
* @return a mensagem de heartbeat caso os dados correspondam a uma; * @return a mensagem desserializada caso os dados representem uma
* {@code null} caso contrário. * {@code HeartbeatMessage}; caso contrário, {@code null}.
*/ */
private HeartbeatMessage tryParseHeartbeat(byte[] data) { private HeartbeatMessage tryParseHeartbeat(byte[] data) {
try { try {
@ -170,15 +141,12 @@ public class SocketUDP extends Thread {
} }
} }
/**
* Inicia a recepção de mensagens UDP.
*/
@Override @Override
public void run() { public void run() {
try { try {
this.receive(); this.receive();
} catch (IOException | ClassNotFoundException ex) { } catch (IOException | ClassNotFoundException ex) {
Logger.getLogger(SocketUDP.class.getName()).log(Level.SEVERE, null, ex); Logger.getLogger(Socket.class.getName()).log(Level.SEVERE, null, ex);
} }
} }
} }

View File

@ -1,239 +0,0 @@
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.
*
* <p>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.</p>
*
* <p>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.</p>
*
* <p>As conexões são associadas aos endereços dos peers e podem ser
* encerradas quando um peer deixa de estar ativo.</p>
*
* @author Gustavo
*/
public class SocketTCP extends Thread implements PeerListener {
private final ExecutorService connectionExecutor = Executors.newCachedThreadPool();
private final Map<InetAddress, TcpConnection> connections;
private ServerSocket serverSocket;
private final MainWindow main;
private volatile boolean running = true;
/**
* Cria o servidor TCP da aplicação.
*
* <p>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.</p>
*
* @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.
*
* <p>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.</p>
*
* @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.
*
* <p>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.</p>
*/
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.
*
* <p>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.</p>
*
* <p>Cada mensagem recebida é encaminhada aos plugins da aplicação.</p>
*
* @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.
*
* <p>As conexões existentes são encerradas e o servidor deixa
* de aceitar novas conexões.</p>
*/
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.
*
* <p>Conexões associadas a peers que não estão mais ativos são encerradas
* e removidas.</p>
*
* @param peers lista atualizada de peers ativos.
*/
@Override
public void onPeersChanged(List<PeerInfo> peers) {
Set<InetAddress> 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) {
}
}
}
}
}
}

View File

@ -1,88 +0,0 @@
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.
*
* <p>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.</p>
*
* @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.
*
* <p>A mensagem é precedida por seu tamanho para permitir que o receptor
* identifique o limite entre mensagens transmitidas pela mesma conexão.</p>
*
* @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.
*
* <p>O tamanho da mensagem é lido primeiro e, em seguida, os dados
* correspondentes são recebidos.</p>
*
* @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();
}
}