Unicast implementado, mas falta testes

This commit is contained in:
GustavoHMDS 2026-08-17 14:07:12 -03:00
parent 81c37225fa
commit 45e3bf6876
8 changed files with 834 additions and 119 deletions

6
.idea/vcs.xml Normal file
View File

@ -0,0 +1,6 @@
<?xml version="1.0" encoding="UTF-8"?>
<project version="4">
<component name="VcsDirectoryMappings">
<mapping directory="$PROJECT_DIR$" vcs="Git" />
</component>
</project>

View File

@ -0,0 +1,6 @@
<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,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,
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.
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);
```
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.

View File

@ -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<Plugin> 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();
}
this.socket = new Socket(this);
this.socket.start();
if (this.socketTCP != null) {
this.socketTCP.close();
}
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() {
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);
}

View File

@ -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.
*
* <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

@ -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.
*
* <p>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.</p>
* <p>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.</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 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.
*
* <p>O socket é associado à porta da aplicação e ingressa no grupo
* multicast configurado.</p>
*
* @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.
*
* <p>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.</p>
* <p>Mensagens de heartbeat são encaminhadas ao gerenciador de heartbeat.
* As demais mensagens são encaminhadas aos plugins carregados pela
* aplicação.</p>
*
* @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);
}
}
}

View File

@ -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.
*
* <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();
}
}