Bem-vindo à Apostila Completa!
Esta é sua jornada completa através do mundo dos serviços em tempo real com Python.
WebSocket Nativo
Protocolo fundamental para comunicação bidirecional
Socket.IO
Framework poderoso com fallbacks e recursos avançados
Aplicações Práticas
Exemplos reais e casos de uso em produção
O que são Serviços em Tempo Real?
- Comunicação instantânea entre cliente e servidor
- Dados sincronizados em tempo real
- Experiência de usuário interativa
- Notificações push instantâneas
Tecnologias que Vamos Usar
Python
Linguagem principal
Socket.IO
Framework completo
WebSocket
Protocolo nativo
Casos de Uso Populares
Sistemas de Chat
Jogos Multiplayer
Dashboards em Tempo Real
Notificações Push
Jogos Multiplayer
Dashboards em Tempo Real
Notificações Push
Colaboração Online
Trading de Criptomoedas
Rastreamento de Localização
Streaming de Vídeo
Trading de Criptomoedas
Rastreamento de Localização
Streaming de Vídeo
WebSocket Básico
1. Servidor WebSocket Simples
# servidor_websocket.py
import asyncio
import websockets
import json
from datetime import datetime
# Lista de clientes conectados
connected_clients = set()
async def register_client(websocket):
"""Registra um novo cliente"""
connected_clients.add(websocket)
print(f"Cliente conectado. Total: {len(connected_clients)}")
async def unregister_client(websocket):
"""Remove um cliente"""
connected_clients.remove(websocket)
print(f"Cliente desconectado. Total: {len(connected_clients)}")
async def broadcast_message(message, sender=None):
"""Envia mensagem para todos os clientes conectados"""
if connected_clients:
# Remove o remetente da lista se especificado
recipients = connected_clients.copy()
if sender in recipients:
recipients.remove(sender)
# Envia para todos os destinatários
await asyncio.gather(
*[client.send(message) for client in recipients],
return_exceptions=True
)
async def handle_client(websocket, path):
"""Manipula conexões de clientes"""
await register_client(websocket)
try:
async for message in websocket:
try:
# Parse da mensagem JSON
data = json.loads(message)
# Adiciona timestamp
data['timestamp'] = datetime.now().isoformat()
# Log no servidor
print(f"Mensagem recebida: {data}")
# Retransmite para outros clientes
response = json.dumps(data)
await broadcast_message(response, sender=websocket)
except json.JSONDecodeError:
# Mensagem inválida
error_response = json.dumps({
'type': 'error',
'message': 'Formato de mensagem inválido'
})
await websocket.send(error_response)
except websockets.exceptions.ConnectionClosed:
pass
finally:
await unregister_client(websocket)
# Configuração do servidor
async def start_server():
"""Inicia o servidor WebSocket"""
print("Iniciando servidor WebSocket na porta 8765...")
# Servidor básico
server = await websockets.serve(
handle_client,
"localhost",
8765,
# Configurações opcionais
ping_interval=20, # Ping a cada 20s
ping_timeout=10, # Timeout de 10s
max_size=1024*1024 # Max 1MB por mensagem
)
print("Servidor rodando! Acesse ws://localhost:8765")
return server
if __name__ == "__main__":
# Executa o servidor
asyncio.run(start_server())
2. Cliente WebSocket Python
# cliente_websocket.py
import asyncio
import websockets
import json
import threading
from datetime import datetime
class WebSocketClient:
def __init__(self, uri):
self.uri = uri
self.websocket = None
self.connected = False
async def connect(self):
"""Conecta ao servidor"""
try:
self.websocket = await websockets.connect(self.uri)
self.connected = True
print(f"Conectado ao servidor: {self.uri}")
# Inicia o loop de recebimento
await self.listen_for_messages()
except Exception as e:
print(f"Erro na conexão: {e}")
self.connected = False
async def listen_for_messages(self):
"""Escuta mensagens do servidor"""
try:
async for message in self.websocket:
data = json.loads(message)
await self.handle_message(data)
except websockets.exceptions.ConnectionClosed:
print("Conexão fechada pelo servidor")
self.connected = False
except Exception as e:
print(f"Erro ao receber mensagem: {e}")
async def handle_message(self, data):
"""Processa mensagens recebidas"""
msg_type = data.get('type', 'unknown')
timestamp = data.get('timestamp', '')
if msg_type == 'chat':
user = data.get('user', 'Anônimo')
message = data.get('message', '')
print(f"[{timestamp}] {user}: {message}")
elif msg_type == 'notification':
message = data.get('message', '')
print(f"🔔 Notificação: {message}")
elif msg_type == 'error':
message = data.get('message', '')
print(f"❌ Erro: {message}")
else:
print(f"Mensagem recebida: {data}")
async def send_message(self, message_type, content):
"""Envia mensagem para o servidor"""
if not self.connected or not self.websocket:
print("Não conectado ao servidor")
return
try:
message = {
'type': message_type,
'timestamp': datetime.now().isoformat(),
**content
}
await self.websocket.send(json.dumps(message))
except Exception as e:
print(f"Erro ao enviar mensagem: {e}")
async def send_chat_message(self, user, message):
"""Envia mensagem de chat"""
await self.send_message('chat', {
'user': user,
'message': message
})
async def close(self):
"""Fecha a conexão"""
if self.websocket:
await self.websocket.close()
self.connected = False
print("Conexão fechada")
# Exemplo de uso interativo
async def interactive_client():
"""Cliente interativo via terminal"""
client = WebSocketClient("ws://localhost:8765")
# Conecta ao servidor em uma task separada
connect_task = asyncio.create_task(client.connect())
# Aguarda um pouco para estabelecer conexão
await asyncio.sleep(1)
if not client.connected:
print("Falha na conexão!")
return
# Input do usuário
user_name = input("Digite seu nome: ")
print("Digite suas mensagens (ou 'quit' para sair):")
def get_input():
"""Thread para capturar input do usuário"""
while client.connected:
try:
message = input()
if message.lower() == 'quit':
break
# Envia mensagem (precisa ser executado no loop principal)
asyncio.run_coroutine_threadsafe(
client.send_chat_message(user_name, message),
asyncio.get_event_loop()
)
except EOFError:
break
# Inicia thread de input
input_thread = threading.Thread(target=get_input)
input_thread.daemon = True
input_thread.start()
# Aguarda a conexão terminar
try:
await connect_task
except KeyboardInterrupt:
pass
finally:
await client.close()
if __name__ == "__main__":
asyncio.run(interactive_client())
Como Executar
- Instale o websockets:
pip install websockets - Execute o servidor:
python servidor_websocket.py - Em outro terminal, execute o cliente:
python cliente_websocket.py - Abra vários clientes para testar o chat!
3. Cliente JavaScript (Frontend)
// websocket_client.js
class WebSocketClient {
constructor(url) {
this.url = url;
this.ws = null;
this.connected = false;
this.reconnectAttempts = 0;
this.maxReconnectAttempts = 5;
this.reconnectDelay = 1000;
}
connect() {
try {
this.ws = new WebSocket(this.url);
this.ws.onopen = (event) => {
console.log('Conectado ao WebSocket');
this.connected = true;
this.reconnectAttempts = 0;
this.onConnected(event);
};
this.ws.onmessage = (event) => {
try {
const data = JSON.parse(event.data);
this.onMessage(data);
} catch (e) {
console.error('Erro ao parsear mensagem:', e);
}
};
this.ws.onclose = (event) => {
console.log('Conexão fechada');
this.connected = false;
this.onDisconnected(event);
// Tentativa de reconexão automática
if (this.reconnectAttempts < this.maxReconnectAttempts) {
setTimeout(() => {
this.reconnectAttempts++;
console.log(`Tentativa de reconexão ${this.reconnectAttempts}/${this.maxReconnectAttempts}`);
this.connect();
}, this.reconnectDelay * this.reconnectAttempts);
}
};
this.ws.onerror = (error) => {
console.error('Erro no WebSocket:', error);
this.onError(error);
};
} catch (error) {
console.error('Erro ao conectar:', error);
}
}
send(type, data) {
if (this.connected && this.ws.readyState === WebSocket.OPEN) {
const message = {
type: type,
timestamp: new Date().toISOString(),
...data
};
this.ws.send(JSON.stringify(message));
} else {
console.warn('WebSocket não está conectado');
}
}
sendChatMessage(user, message) {
this.send('chat', { user, message });
}
close() {
if (this.ws) {
this.ws.close();
}
}
// Métodos a serem sobrescritos
onConnected(event) { }
onDisconnected(event) { }
onMessage(data) { }
onError(error) { }
}
// Exemplo de uso
const client = new WebSocketClient('ws://localhost:8765');
client.onConnected = () => {
document.getElementById('status').textContent = 'Conectado ✅';
};
client.onDisconnected = () => {
document.getElementById('status').textContent = 'Desconectado ❌';
};
client.onMessage = (data) => {
if (data.type === 'chat') {
addMessageToChat(data.user, data.message, data.timestamp);
}
};
function addMessageToChat(user, message, timestamp) {
const chatContainer = document.getElementById('chat-messages');
const messageElement = document.createElement('div');
messageElement.className = 'message';
const time = new Date(timestamp).toLocaleTimeString();
messageElement.innerHTML = `
[${time}]
${user}:
${message}
`;
chatContainer.appendChild(messageElement);
chatContainer.scrollTop = chatContainer.scrollHeight;
}
// Conecta automaticamente
client.connect();
Socket.IO - Framework Completo
Por que Socket.IO?
Fallback automático para polling
Reconexão automática
Salas e namespaces
Broadcast seletivo
Reconexão automática
Salas e namespaces
Broadcast seletivo
Middleware avançado
Autenticação integrada
Multiplexing
Compatibilidade ampla
Autenticação integrada
Multiplexing
Compatibilidade ampla
1. Servidor Socket.IO Completo
# servidor_socketio.py
import socketio
import asyncio
import uvicorn
from datetime import datetime
import json
import logging
# Configuração de logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
# Cria servidor Socket.IO com configurações avançadas
sio = socketio.AsyncServer(
cors_allowed_origins="*", # Em produção, especifique domínios
logger=True,
engineio_logger=True,
async_mode='asgi'
)
# Dicionário para armazenar dados dos usuários
connected_users = {}
chat_rooms = {}
class ChatRoom:
def __init__(self, name):
self.name = name
self.users = set()
self.messages = []
self.created_at = datetime.now()
def add_user(self, user_id, username):
self.users.add(user_id)
logger.info(f"Usuário {username} entrou na sala {self.name}")
def remove_user(self, user_id):
self.users.discard(user_id)
if user_id in connected_users:
username = connected_users[user_id]['username']
logger.info(f"Usuário {username} saiu da sala {self.name}")
def add_message(self, user_id, message):
username = connected_users.get(user_id, {}).get('username', 'Anônimo')
msg_data = {
'id': len(self.messages),
'user_id': user_id,
'username': username,
'message': message,
'timestamp': datetime.now().isoformat()
}
self.messages.append(msg_data)
return msg_data
# Middleware para autenticação
@sio.event
async def connect(sid, environ, auth):
"""Evento de conexão"""
try:
# Extrai dados de autenticação
username = auth.get('username', f'User_{sid[:8]}') if auth else f'User_{sid[:8]}'
user_id = auth.get('user_id', sid) if auth else sid
# Armazena dados do usuário
connected_users[sid] = {
'user_id': user_id,
'username': username,
'connected_at': datetime.now(),
'rooms': set()
}
logger.info(f"Usuário conectado: {username} (SID: {sid})")
# Envia dados de conexão para o cliente
await sio.emit('connected', {
'user_id': user_id,
'username': username,
'timestamp': datetime.now().isoformat()
}, room=sid)
# Atualiza lista de usuários online
await broadcast_user_list()
except Exception as e:
logger.error(f"Erro na conexão: {e}")
return False # Rejeita a conexão
@sio.event
async def disconnect(sid):
"""Evento de desconexão"""
if sid in connected_users:
user_data = connected_users[sid]
username = user_data['username']
# Remove usuário de todas as salas
for room_name in user_data['rooms'].copy():
await leave_room_handler(sid, room_name)
# Remove usuário da lista
del connected_users[sid]
logger.info(f"Usuário desconectado: {username} (SID: {sid})")
# Atualiza lista de usuários
await broadcast_user_list()
@sio.event
async def join_room(sid, data):
"""Usuário entra em uma sala"""
try:
room_name = data.get('room', 'general')
if sid not in connected_users:
return {'error': 'Usuário não autenticado'}
user_data = connected_users[sid]
username = user_data['username']
# Cria sala se não existir
if room_name not in chat_rooms:
chat_rooms[room_name] = ChatRoom(room_name)
# Adiciona usuário à sala
await sio.enter_room(sid, room_name)
chat_rooms[room_name].add_user(sid, username)
user_data['rooms'].add(room_name)
# Notifica outros usuários da sala
await sio.emit('user_joined', {
'username': username,
'room': room_name,
'timestamp': datetime.now().isoformat()
}, room=room_name, skip_sid=sid)
# Envia histórico de mensagens para o usuário
recent_messages = chat_rooms[room_name].messages[-50:] # Últimas 50
await sio.emit('room_history', {
'room': room_name,
'messages': recent_messages
}, room=sid)
# Confirma entrada na sala
await sio.emit('joined_room', {
'room': room_name,
'users_count': len(chat_rooms[room_name].users)
}, room=sid)
return {'success': True, 'room': room_name}
except Exception as e:
logger.error(f"Erro ao entrar na sala: {e}")
return {'error': str(e)}
@sio.event
async def leave_room_handler(sid, room_name):
"""Usuário sai de uma sala"""
try:
if sid in connected_users and room_name in chat_rooms:
user_data = connected_users[sid]
username = user_data['username']
# Remove da sala
await sio.leave_room(sid, room_name)
chat_rooms[room_name].remove_user(sid)
user_data['rooms'].discard(room_name)
# Notifica outros usuários
await sio.emit('user_left', {
'username': username,
'room': room_name,
'timestamp': datetime.now().isoformat()
}, room=room_name)
except Exception as e:
logger.error(f"Erro ao sair da sala: {e}")
@sio.event
async def send_message(sid, data):
"""Envia mensagem para uma sala"""
try:
if sid not in connected_users:
return {'error': 'Usuário não autenticado'}
room_name = data.get('room', 'general')
message = data.get('message', '').strip()
if not message:
return {'error': 'Mensagem vazia'}
if room_name not in chat_rooms:
return {'error': 'Sala não encontrada'}
# Adiciona mensagem à sala
msg_data = chat_rooms[room_name].add_message(sid, message)
# Envia para todos na sala
await sio.emit('new_message', msg_data, room=room_name)
return {'success': True, 'message_id': msg_data['id']}
except Exception as e:
logger.error(f"Erro ao enviar mensagem: {e}")
return {'error': str(e)}
@sio.event
async def private_message(sid, data):
"""Envia mensagem privada"""
try:
if sid not in connected_users:
return {'error': 'Usuário não autenticado'}
target_user = data.get('target_user')
message = data.get('message', '').strip()
if not message:
return {'error': 'Mensagem vazia'}
sender = connected_users[sid]
# Encontra SID do usuário alvo
target_sid = None
for user_sid, user_data in connected_users.items():
if user_data['username'] == target_user:
target_sid = user_sid
break
if not target_sid:
return {'error': 'Usuário não encontrado'}
msg_data = {
'from': sender['username'],
'to': target_user,
'message': message,
'timestamp': datetime.now().isoformat(),
'type': 'private'
}
# Envia para ambos os usuários
await sio.emit('private_message', msg_data, room=target_sid)
await sio.emit('private_message', msg_data, room=sid)
return {'success': True}
except Exception as e:
logger.error(f"Erro ao enviar mensagem privada: {e}")
return {'error': str(e)}
@sio.event
async def get_rooms(sid):
"""Retorna lista de salas disponíveis"""
try:
rooms_info = []
for room_name, room in chat_rooms.items():
rooms_info.append({
'name': room_name,
'users_count': len(room.users),
'created_at': room.created_at.isoformat()
})
return {'rooms': rooms_info}
except Exception as e:
logger.error(f"Erro ao buscar salas: {e}")
return {'error': str(e)}
async def broadcast_user_list():
"""Envia lista de usuários online para todos"""
try:
users_online = []
for sid, user_data in connected_users.items():
users_online.append({
'username': user_data['username'],
'rooms': list(user_data['rooms'])
})
await sio.emit('users_online', {'users': users_online})
except Exception as e:
logger.error(f"Erro ao broadcast da lista de usuários: {e}")
# Cria aplicação ASGI
app = socketio.ASGIApp(sio)
if __name__ == "__main__":
# Configuração do servidor
config = uvicorn.Config(
app=app,
host="0.0.0.0",
port=8000,
log_level="info",
reload=True # Remove em produção
)
server = uvicorn.Server(config)
print("Servidor Socket.IO iniciando na porta 8000...")
print("Acesse: http://localhost:8000")
# Inicia o servidor
asyncio.run(server.serve())
2. Cliente Python Socket.IO
# cliente_socketio.py
import socketio
import asyncio
import threading
import json
from datetime import datetime
class SocketIOClient:
def __init__(self, server_url, username):
self.server_url = server_url
self.username = username
self.sio = socketio.AsyncClient()
self.connected = False
self.current_room = None
# Registra eventos
self.setup_events()
def setup_events(self):
"""Configura eventos do Socket.IO"""
@self.sio.event
async def connect():
print(f"✅ Conectado como {self.username}")
self.connected = True
@self.sio.event
async def disconnect():
print("❌ Desconectado do servidor")
self.connected = False
@self.sio.event
async def connected(data):
print(f"Dados de conexão: {data}")
@self.sio.event
async def new_message(data):
timestamp = datetime.fromisoformat(data['timestamp'].replace('Z', '+00:00'))
time_str = timestamp.strftime('%H:%M:%S')
print(f"[{time_str}] {data['username']}: {data['message']}")
@self.sio.event
async def private_message(data):
timestamp = datetime.fromisoformat(data['timestamp'].replace('Z', '+00:00'))
time_str = timestamp.strftime('%H:%M:%S')
print(f"🔒 [{time_str}] Privado de {data['from']}: {data['message']}")
@self.sio.event
async def user_joined(data):
print(f"👋 {data['username']} entrou na sala {data['room']}")
@self.sio.event
async def user_left(data):
print(f"👋 {data['username']} saiu da sala {data['room']}")
@self.sio.event
async def joined_room(data):
self.current_room = data['room']
print(f"🏠 Você entrou na sala '{data['room']}' ({data['users_count']} usuários)")
@self.sio.event
async def room_history(data):
print(f"\n📜 Histórico da sala '{data['room']}':")
for msg in data['messages'][-10:]: # Últimas 10 mensagens
timestamp = datetime.fromisoformat(msg['timestamp'])
time_str = timestamp.strftime('%H:%M:%S')
print(f" [{time_str}] {msg['username']}: {msg['message']}")
print("---")
@self.sio.event
async def users_online(data):
if data['users']:
print(f"\n👥 Usuários online ({len(data['users'])}):")
for user in data['users']:
rooms = ', '.join(user['rooms']) if user['rooms'] else 'nenhuma sala'
print(f" • {user['username']} ({rooms})")
async def connect_to_server(self):
"""Conecta ao servidor"""
try:
auth_data = {
'username': self.username,
'user_id': f"user_{self.username}"
}
await self.sio.connect(
self.server_url,
auth=auth_data,
wait_timeout=10
)
except Exception as e:
print(f"Erro na conexão: {e}")
return False
return True
async def join_room(self, room_name):
"""Entra em uma sala"""
if not self.connected:
print("Não conectado ao servidor")
return
result = await self.sio.call('join_room', {'room': room_name})
if result.get('error'):
print(f"Erro ao entrar na sala: {result['error']}")
async def send_message(self, message):
"""Envia mensagem para a sala atual"""
if not self.current_room:
print("Você não está em nenhuma sala")
return
result = await self.sio.call('send_message', {
'room': self.current_room,
'message': message
})
if result.get('error'):
print(f"Erro ao enviar mensagem: {result['error']}")
async def send_private_message(self, target_user, message):
"""Envia mensagem privada"""
result = await self.sio.call('private_message', {
'target_user': target_user,
'message': message
})
if result.get('error'):
print(f"Erro ao enviar mensagem privada: {result['error']}")
async def get_rooms(self):
"""Lista salas disponíveis"""
result = await self.sio.call('get_rooms')
if result.get('rooms'):
print("\n🏠 Salas disponíveis:")
for room in result['rooms']:
print(f" • {room['name']} ({room['users_count']} usuários)")
else:
print("Nenhuma sala disponível")
async def disconnect(self):
"""Desconecta do servidor"""
if self.connected:
await self.sio.disconnect()
# Interface interativa
async def interactive_client():
"""Cliente interativo via terminal"""
# Configuração inicial
print("=== Cliente Socket.IO ===")
username = input("Digite seu nome de usuário: ")
server_url = input("URL do servidor (padrão: http://localhost:8000): ") or "http://localhost:8000"
client = SocketIOClient(server_url, username)
# Conecta ao servidor
print("Conectando...")
if not await client.connect_to_server():
print("Falha na conexão!")
return
# Aguarda estabelecer conexão
await asyncio.sleep(1)
print(f"""
✅ Conectado com sucesso!
Comandos disponíveis:
/join - Entrar em uma sala
/rooms - Listar salas
/private - Mensagem privada
/help - Mostrar ajuda
/quit - Sair
Para enviar mensagem na sala: digite diretamente
""")
# Loop principal
try:
while client.connected:
message = await asyncio.get_event_loop().run_in_executor(
None, input, f"[{client.current_room or 'sem sala'}] > "
)
if message.startswith('/'):
# Processa comandos
parts = message.split(' ', 2)
command = parts[0]
if command == '/quit':
break
elif command == '/join' and len(parts) > 1:
await client.join_room(parts[1])
elif command == '/rooms':
await client.get_rooms()
elif command == '/private' and len(parts) > 2:
await client.send_private_message(parts[1], parts[2])
elif command == '/help':
print("Comandos: /join, /rooms, /private, /help, /quit")
else:
print("Comando inválido. Use /help para ver comandos")
elif message.strip():
# Envia mensagem normal
await client.send_message(message)
except KeyboardInterrupt:
pass
finally:
await client.disconnect()
print("Desconectado!")
if __name__ == "__main__":
# Instala dependências necessárias
# pip install python-socketio[asyncio_client] uvicorn
asyncio.run(interactive_client())
3. Cliente JavaScript/HTML Completo
Chat Socket.IO
💬 Chat Socket.IO
Desconectado
Sala: - | Usuários online: 0
Usuários Online:
Demo de Chat em Tempo Real
Simulador de Chat (Demonstração)
Sala: Geral
3 usuários online
Chat em Flask com Socket.IO
# app.py - Chat completo com Flask-SocketIO
from flask import Flask, render_template, request
from flask_socketio import SocketIO, emit, join_room, leave_room, rooms
import json
from datetime import datetime
import uuid
app = Flask(__name__)
app.config['SECRET_KEY'] = 'sua_chave_secreta_aqui'
# Configuração do SocketIO
socketio = SocketIO(
app,
cors_allowed_origins="*",
async_mode='threading'
)
# Armazenamento em memória (use Redis em produção)
users_online = {}
chat_rooms = {
'geral': {'users': set(), 'messages': []},
'python': {'users': set(), 'messages': []},
'tecnologia': {'users': set(), 'messages': []}
}
@app.route('/')
def index():
"""Página principal do chat"""
return render_template('chat.html')
@app.route('/api/rooms')
def get_rooms():
"""API para listar salas"""
rooms_info = []
for room_name, room_data in chat_rooms.items():
rooms_info.append({
'name': room_name,
'users_count': len(room_data['users']),
'last_activity': datetime.now().isoformat()
})
return {'rooms': rooms_info}
@socketio.on('connect')
def handle_connect(auth):
"""Usuário conectou"""
user_id = request.sid
username = auth.get('username', f'Usuario_{user_id[:8]}') if auth else f'Usuario_{user_id[:8]}'
# Registra usuário
users_online[user_id] = {
'username': username,
'connected_at': datetime.now(),
'current_room': None
}
print(f"✅ {username} conectou (ID: {user_id})")
# Notifica cliente sobre conexão bem-sucedida
emit('connected', {
'user_id': user_id,
'username': username,
'timestamp': datetime.now().isoformat()
})
# Atualiza contadores globais
update_users_count()
@socketio.on('disconnect')
def handle_disconnect():
"""Usuário desconectou"""
user_id = request.sid
if user_id in users_online:
username = users_online[user_id]['username']
current_room = users_online[user_id]['current_room']
# Remove das salas
if current_room and current_room in chat_rooms:
chat_rooms[current_room]['users'].discard(user_id)
emit('user_left', {
'username': username,
'room': current_room,
'timestamp': datetime.now().isoformat()
}, room=current_room)
# Remove do registro
del users_online[user_id]
print(f"❌ {username} desconectou")
update_users_count()
@socketio.on('join_room')
def handle_join_room(data):
"""Usuário entra em uma sala"""
user_id = request.sid
room_name = data.get('room', 'geral')
if user_id not in users_online:
emit('error', {'message': 'Usuário não autenticado'})
return
username = users_online[user_id]['username']
# Cria sala se não existir
if room_name not in chat_rooms:
chat_rooms[room_name] = {'users': set(), 'messages': []}
# Sai da sala anterior
previous_room = users_online[user_id]['current_room']
if previous_room and previous_room in chat_rooms:
leave_room(previous_room)
chat_rooms[previous_room]['users'].discard(user_id)
emit('user_left', {
'username': username,
'room': previous_room,
'timestamp': datetime.now().isoformat()
}, room=previous_room)
# Entra na nova sala
join_room(room_name)
chat_rooms[room_name]['users'].add(user_id)
users_online[user_id]['current_room'] = room_name
# Notifica outros usuários
emit('user_joined', {
'username': username,
'room': room_name,
'timestamp': datetime.now().isoformat()
}, room=room_name, include_self=False)
# Envia histórico para o usuário
recent_messages = chat_rooms[room_name]['messages'][-20:]
emit('room_history', {
'room': room_name,
'messages': recent_messages
})
# Confirma entrada na sala
emit('joined_room', {
'room': room_name,
'users_count': len(chat_rooms[room_name]['users'])
})
print(f"🏠 {username} entrou na sala '{room_name}'")
@socketio.on('send_message')
def handle_send_message(data):
"""Usuário envia mensagem"""
user_id = request.sid
message = data.get('message', '').strip()
if user_id not in users_online:
emit('error', {'message': 'Usuário não autenticado'})
return
if not message:
emit('error', {'message': 'Mensagem vazia'})
return
user_data = users_online[user_id]
username = user_data['username']
room_name = user_data['current_room']
if not room_name:
emit('error', {'message': 'Não está em nenhuma sala'})
return
# Cria dados da mensagem
message_data = {
'id': str(uuid.uuid4()),
'user_id': user_id,
'username': username,
'message': message,
'timestamp': datetime.now().isoformat(),
'room': room_name
}
# Armazena mensagem
chat_rooms[room_name]['messages'].append(message_data)
# Mantém apenas últimas 100 mensagens por sala
if len(chat_rooms[room_name]['messages']) > 100:
chat_rooms[room_name]['messages'] = chat_rooms[room_name]['messages'][-100:]
# Envia para todos na sala
emit('new_message', message_data, room=room_name)
print(f"💬 [{room_name}] {username}: {message}")
@socketio.on('private_message')
def handle_private_message(data):
"""Mensagem privada entre usuários"""
sender_id = request.sid
target_username = data.get('target_user')
message = data.get('message', '').strip()
if sender_id not in users_online:
emit('error', {'message': 'Usuário não autenticado'})
return
if not message:
emit('error', {'message': 'Mensagem vazia'})
return
# Encontra usuário alvo
target_id = None
for uid, user_data in users_online.items():
if user_data['username'] == target_username:
target_id = uid
break
if not target_id:
emit('error', {'message': 'Usuário não encontrado'})
return
sender_username = users_online[sender_id]['username']
message_data = {
'id': str(uuid.uuid4()),
'from': sender_username,
'to': target_username,
'message': message,
'timestamp': datetime.now().isoformat(),
'type': 'private'
}
# Envia para ambos os usuários
emit('private_message', message_data, room=target_id)
emit('private_message', message_data, room=sender_id)
print(f"🔒 {sender_username} → {target_username}: {message}")
@socketio.on('typing_start')
def handle_typing_start(data):
"""Usuário começou a digitar"""
user_id = request.sid
if user_id in users_online:
username = users_online[user_id]['username']
room_name = users_online[user_id]['current_room']
if room_name:
emit('user_typing', {
'username': username,
'typing': True
}, room=room_name, include_self=False)
@socketio.on('typing_stop')
def handle_typing_stop(data):
"""Usuário parou de digitar"""
user_id = request.sid
if user_id in users_online:
username = users_online[user_id]['username']
room_name = users_online[user_id]['current_room']
if room_name:
emit('user_typing', {
'username': username,
'typing': False
}, room=room_name, include_self=False)
def update_users_count():
"""Atualiza contador de usuários para todos"""
users_by_room = {}
for user_id, user_data in users_online.items():
room = user_data.get('current_room')
if room:
if room not in users_by_room:
users_by_room[room] = []
users_by_room[room].append({
'username': user_data['username'],
'connected_at': user_data['connected_at'].isoformat()
})
socketio.emit('users_update', {
'total_online': len(users_online),
'users_by_room': users_by_room
})
if __name__ == '__main__':
print("🚀 Iniciando servidor Flask-SocketIO...")
print("📡 Acesse: http://localhost:5000")
socketio.run(
app,
host='0.0.0.0',
port=5000,
debug=True,
allow_unsafe_werkzeug=True
)
Instalação e Execução
# Instala dependências
pip install flask flask-socketio
# Cria estrutura de pastas
mkdir templates static
# Executa o servidor
python app.py
# Acesse http://localhost:5000
Aplicações em Tempo Real
Dashboard em Tempo Real
0
Usuários Online
0
Mensagens/min
Aguardando atividades...
Sistema de Notificações
Nenhuma notificação
Sistema de Monitoramento em Tempo Real
# monitor_sistema.py - Monitoramento de sistema em tempo real
import asyncio
import psutil
import json
import time
from datetime import datetime
import socketio
from typing import Dict, List
class SystemMonitor:
def __init__(self):
self.sio = socketio.AsyncServer(cors_allowed_origins="*")
self.app = socketio.ASGIApp(self.sio)
self.monitoring = False
self.clients = set()
# Histórico de métricas
self.metrics_history = {
'cpu': [],
'memory': [],
'disk': [],
'network': []
}
self.setup_events()
def setup_events(self):
@self.sio.event
async def connect(sid, environ):
self.clients.add(sid)
print(f"Cliente conectado: {sid}")
# Envia estado inicial
await self.sio.emit('monitoring_status', {
'active': self.monitoring,
'clients_count': len(self.clients)
}, room=sid)
# Envia histórico recente
await self.send_metrics_history(sid)
@self.sio.event
async def disconnect(sid):
self.clients.discard(sid)
print(f"Cliente desconectado: {sid}")
@self.sio.event
async def start_monitoring(sid):
if not self.monitoring:
self.monitoring = True
asyncio.create_task(self.monitor_loop())
await self.sio.emit('monitoring_started', {'timestamp': datetime.now().isoformat()})
@self.sio.event
async def stop_monitoring(sid):
self.monitoring = False
await self.sio.emit('monitoring_stopped', {'timestamp': datetime.now().isoformat()})
@self.sio.event
async def get_process_list(sid):
processes = self.get_top_processes()
await self.sio.emit('process_list', {'processes': processes}, room=sid)
def get_system_metrics(self) -> Dict:
"""Coleta métricas do sistema"""
# CPU
cpu_percent = psutil.cpu_percent(interval=1)
cpu_count = psutil.cpu_count()
cpu_freq = psutil.cpu_freq()
# Memória
memory = psutil.virtual_memory()
swap = psutil.swap_memory()
# Disco
disk = psutil.disk_usage('/')
disk_io = psutil.disk_io_counters()
# Rede
network_io = psutil.net_io_counters()
network_connections = len(psutil.net_connections())
# Processos
process_count = len(psutil.pids())
# Temperatura (se disponível)
try:
temperatures = psutil.sensors_temperatures()
temp_avg = 0
if temperatures:
temp_values = []
for name, entries in temperatures.items():
for entry in entries:
if entry.current:
temp_values.append(entry.current)
temp_avg = sum(temp_values) / len(temp_values) if temp_values else 0
except:
temp_avg = 0
metrics = {
'timestamp': datetime.now().isoformat(),
'cpu': {
'percent': cpu_percent,
'count': cpu_count,
'freq_current': cpu_freq.current if cpu_freq else 0,
'freq_max': cpu_freq.max if cpu_freq else 0
},
'memory': {
'total': memory.total,
'available': memory.available,
'percent': memory.percent,
'used': memory.used,
'free': memory.free,
'swap_total': swap.total,
'swap_used': swap.used,
'swap_percent': swap.percent
},
'disk': {
'total': disk.total,
'used': disk.used,
'free': disk.free,
'percent': disk.percent,
'read_bytes': disk_io.read_bytes if disk_io else 0,
'write_bytes': disk_io.write_bytes if disk_io else 0
},
'network': {
'bytes_sent': network_io.bytes_sent,
'bytes_recv': network_io.bytes_recv,
'packets_sent': network_io.packets_sent,
'packets_recv': network_io.packets_recv,
'connections': network_connections
},
'system': {
'process_count': process_count,
'temperature': temp_avg,
'uptime': time.time() - psutil.boot_time()
}
}
return metrics
def get_top_processes(self, limit: int = 10) -> List[Dict]:
"""Retorna os processos que mais consomem recursos"""
processes = []
for proc in psutil.process_iter(['pid', 'name', 'cpu_percent', 'memory_percent', 'status']):
try:
proc_info = proc.info
proc_info['cpu_percent'] = proc.cpu_percent()
processes.append(proc_info)
except (psutil.NoSuchProcess, psutil.AccessDenied):
pass
# Ordena por CPU
processes.sort(key=lambda x: x.get('cpu_percent', 0), reverse=True)
return processes[:limit]
async def monitor_loop(self):
"""Loop principal de monitoramento"""
while self.monitoring and self.clients:
try:
# Coleta métricas
metrics = self.get_system_metrics()
# Armazena no histórico (mantém últimas 100 entradas)
for key in self.metrics_history:
if key in metrics:
self.metrics_history[key].append({
'timestamp': metrics['timestamp'],
'value': metrics[key]
})
if len(self.metrics_history[key]) > 100:
self.metrics_history[key] = self.metrics_history[key][-100:]
# Envia para todos os clientes
await self.sio.emit('system_metrics', metrics)
# Verifica alertas
alerts = self.check_alerts(metrics)
if alerts:
await self.sio.emit('system_alerts', {'alerts': alerts})
# Aguarda próxima coleta
await asyncio.sleep(2)
except Exception as e:
print(f"Erro no monitoramento: {e}")
await asyncio.sleep(5)
self.monitoring = False
def check_alerts(self, metrics: Dict) -> List[Dict]:
"""Verifica condições de alerta"""
alerts = []
# CPU alto
if metrics['cpu']['percent'] > 80:
alerts.append({
'type': 'warning',
'title': 'CPU Alto',
'message': f"Uso de CPU em {metrics['cpu']['percent']:.1f}%",
'timestamp': datetime.now().isoformat()
})
# Memória alta
if metrics['memory']['percent'] > 85:
alerts.append({
'type': 'warning',
'title': 'Memória Alta',
'message': f"Uso de memória em {metrics['memory']['percent']:.1f}%",
'timestamp': datetime.now().isoformat()
})
# Disco cheio
if metrics['disk']['percent'] > 90:
alerts.append({
'type': 'error',
'title': 'Disco Cheio',
'message': f"Uso de disco em {metrics['disk']['percent']:.1f}%",
'timestamp': datetime.now().isoformat()
})
# Temperatura alta
if metrics['system']['temperature'] > 75:
alerts.append({
'type': 'warning',
'title': 'Temperatura Alta',
'message': f"Temperatura do sistema: {metrics['system']['temperature']:.1f}°C",
'timestamp': datetime.now().isoformat()
})
return alerts
async def send_metrics_history(self, sid: str):
"""Envia histórico de métricas para um cliente"""
await self.sio.emit('metrics_history', {
'history': self.metrics_history
}, room=sid)
# Servidor de exemplo
if __name__ == "__main__":
import uvicorn
monitor = SystemMonitor()
print("🖥️ Iniciando monitor de sistema...")
print("📊 Acesse: http://localhost:8001")
config = uvicorn.Config(
app=monitor.app,
host="0.0.0.0",
port=8001,
log_level="info"
)
server = uvicorn.Server(config)
asyncio.run(server.serve())
Trading/Criptomoedas em Tempo Real
# crypto_monitor.py - Monitor de criptomoedas em tempo real
import asyncio
import aiohttp
import json
from datetime import datetime
import socketio
from typing import Dict, List
import logging
class CryptoMonitor:
def __init__(self):
self.sio = socketio.AsyncServer(cors_allowed_origins="*")
self.app = socketio.ASGIApp(self.sio)
# Configuração
self.api_url = "https://api.coingecko.com/api/v3"
self.coins_to_track = [
'bitcoin', 'ethereum', 'binancecoin', 'cardano',
'solana', 'polkadot', 'dogecoin', 'avalanche-2'
]
# Estado
self.monitoring = False
self.clients = {}
self.price_history = {}
self.alerts = {}
self.setup_events()
def setup_events(self):
@self.sio.event
async def connect(sid, environ, auth):
user_id = auth.get('user_id', sid) if auth else sid
self.clients[sid] = {
'user_id': user_id,
'connected_at': datetime.now(),
'subscriptions': set(self.coins_to_track) # Por padrão, todas as moedas
}
print(f"📱 Cliente conectado: {user_id}")
# Envia dados iniciais
await self.send_initial_data(sid)
@self.sio.event
async def disconnect(sid):
if sid in self.clients:
user_id = self.clients[sid]['user_id']
del self.clients[sid]
print(f"📱 Cliente desconectado: {user_id}")
@self.sio.event
async def subscribe_coin(sid, data):
"""Cliente se inscreve para uma moeda específica"""
coin_id = data.get('coin_id')
if sid in self.clients and coin_id:
self.clients[sid]['subscriptions'].add(coin_id)
await self.sio.emit('subscription_updated', {
'coin_id': coin_id,
'subscribed': True
}, room=sid)
@self.sio.event
async def unsubscribe_coin(sid, data):
"""Cliente cancela inscrição de uma moeda"""
coin_id = data.get('coin_id')
if sid in self.clients and coin_id:
self.clients[sid]['subscriptions'].discard(coin_id)
await self.sio.emit('subscription_updated', {
'coin_id': coin_id,
'subscribed': False
}, room=sid)
@self.sio.event
async def set_price_alert(sid, data):
"""Define alerta de preço"""
coin_id = data.get('coin_id')
price_target = data.get('price_target')
alert_type = data.get('type', 'above') # 'above' ou 'below'
if sid in self.clients and coin_id and price_target:
user_id = self.clients[sid]['user_id']
if user_id not in self.alerts:
self.alerts[user_id] = {}
alert_id = f"{coin_id}_{alert_type}_{price_target}"
self.alerts[user_id][alert_id] = {
'coin_id': coin_id,
'price_target': float(price_target),
'type': alert_type,
'created_at': datetime.now().isoformat(),
'triggered': False
}
await self.sio.emit('alert_set', {
'alert_id': alert_id,
'coin_id': coin_id,
'price_target': price_target,
'type': alert_type
}, room=sid)
@self.sio.event
async def get_coin_history(sid, data):
"""Retorna histórico de uma moeda"""
coin_id = data.get('coin_id')
if coin_id in self.price_history:
await self.sio.emit('coin_history', {
'coin_id': coin_id,
'history': self.price_history[coin_id][-100:] # Últimos 100 pontos
}, room=sid)
async def fetch_crypto_prices(self) -> Dict:
"""Busca preços das criptomoedas"""
try:
coins_str = ','.join(self.coins_to_track)
url = f"{self.api_url}/simple/price"
params = {
'ids': coins_str,
'vs_currencies': 'usd,brl',
'include_24hr_change': 'true',
'include_24hr_vol': 'true',
'include_market_cap': 'true'
}
async with aiohttp.ClientSession() as session:
async with session.get(url, params=params) as response:
if response.status == 200:
data = await response.json()
# Processa dados
processed_data = {}
timestamp = datetime.now().isoformat()
for coin_id, coin_data in data.items():
processed_data[coin_id] = {
'id': coin_id,
'name': coin_id.replace('-', ' ').title(),
'price_usd': coin_data.get('usd', 0),
'price_brl': coin_data.get('brl', 0),
'change_24h': coin_data.get('usd_24h_change', 0),
'volume_24h': coin_data.get('usd_24h_vol', 0),
'market_cap': coin_data.get('usd_market_cap', 0),
'timestamp': timestamp
}
# Armazena no histórico
if coin_id not in self.price_history:
self.price_history[coin_id] = []
self.price_history[coin_id].append({
'timestamp': timestamp,
'price_usd': coin_data.get('usd', 0),
'price_brl': coin_data.get('brl', 0),
'volume': coin_data.get('usd_24h_vol', 0)
})
# Mantém apenas últimos 500 pontos
if len(self.price_history[coin_id]) > 500:
self.price_history[coin_id] = self.price_history[coin_id][-500:]
return processed_data
else:
print(f"Erro na API: {response.status}")
return {}
except Exception as e:
print(f"Erro ao buscar preços: {e}")
return {}
async def check_price_alerts(self, prices: Dict):
"""Verifica alertas de preço"""
for user_id, user_alerts in self.alerts.items():
for alert_id, alert_data in user_alerts.items():
if alert_data['triggered']:
continue
coin_id = alert_data['coin_id']
if coin_id not in prices:
continue
current_price = prices[coin_id]['price_usd']
target_price = alert_data['price_target']
alert_type = alert_data['type']
triggered = False
if alert_type == 'above' and current_price >= target_price:
triggered = True
elif alert_type == 'below' and current_price <= target_price:
triggered = True
if triggered:
alert_data['triggered'] = True
alert_data['triggered_at'] = datetime.now().isoformat()
alert_data['triggered_price'] = current_price
# Envia alerta para o usuário
for sid, client_data in self.clients.items():
if client_data['user_id'] == user_id:
await self.sio.emit('price_alert_triggered', {
'alert_id': alert_id,
'coin_id': coin_id,
'coin_name': prices[coin_id]['name'],
'current_price': current_price,
'target_price': target_price,
'type': alert_type,
'timestamp': alert_data['triggered_at']
}, room=sid)
async def send_initial_data(self, sid: str):
"""Envia dados iniciais para um cliente"""
# Dados das moedas
if self.price_history:
latest_prices = {}
for coin_id in self.coins_to_track:
if coin_id in self.price_history and self.price_history[coin_id]:
latest_prices[coin_id] = self.price_history[coin_id][-1]
await self.sio.emit('initial_data', {
'prices': latest_prices,
'monitoring': self.monitoring
}, room=sid)
async def monitor_loop(self):
"""Loop principal de monitoramento"""
self.monitoring = True
while self.monitoring and self.clients:
try:
# Busca preços atualizados
prices = await self.fetch_crypto_prices()
if prices:
# Verifica alertas
await self.check_price_alerts(prices)
# Envia atualizações para clientes
for sid, client_data in self.clients.items():
subscriptions = client_data['subscriptions']
# Filtra apenas moedas inscritas
filtered_prices = {
coin_id: price_data
for coin_id, price_data in prices.items()
if coin_id in subscriptions
}
if filtered_prices:
await self.sio.emit('price_update', {
'prices': filtered_prices,
'timestamp': datetime.now().isoformat()
}, room=sid)
# Aguarda próxima atualização (30 segundos)
await asyncio.sleep(30)
except Exception as e:
print(f"Erro no loop de monitoramento: {e}")
await asyncio.sleep(60) # Aguarda mais tempo em caso de erro
self.monitoring = False
async def start_monitoring(self):
"""Inicia o monitoramento"""
if not self.monitoring:
asyncio.create_task(self.monitor_loop())
# Execução
if __name__ == "__main__":
import uvicorn
crypto_monitor = CryptoMonitor()
# Inicia monitoramento
asyncio.create_task(crypto_monitor.start_monitoring())
print("💰 Iniciando monitor de criptomoedas...")
print("📊 Acesse: http://localhost:8002")
config = uvicorn.Config(
app=crypto_monitor.app,
host="0.0.0.0",
port=8002,
log_level="info"
)
server = uvicorn.Server(config)
asyncio.run(server.serve())
Deploy em Produção
Considerações Críticas para Produção
Segurança e autenticação
Escalabilidade horizontal
Persistência de dados
Monitoramento e métricas
Escalabilidade horizontal
Persistência de dados
Monitoramento e métricas
Load balancing
Recuperação de falhas
Rate limiting
CI/CD pipeline
Recuperação de falhas
Rate limiting
CI/CD pipeline
1. Configuração Docker Completa
# Dockerfile
FROM python:3.11-slim
# Instala dependências do sistema
RUN apt-get update && apt-get install -y \
gcc \
redis-tools \
curl \
&& rm -rf /var/lib/apt/lists/*
# Configura usuário não-root
RUN useradd --create-home --shell /bin/bash app
WORKDIR /home/app
# Instala dependências Python
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# Copia código da aplicação
COPY --chown=app:app . .
# Muda para usuário não-root
USER app
# Porta da aplicação
EXPOSE 8000
# Comando de execução
CMD ["gunicorn", "--worker-class", "eventlet", "-w", "1", "--bind", "0.0.0.0:8000", "app:app"]
# docker-compose.yml
version: '3.8'
services:
# Aplicação principal
app:
build: .
ports:
- "8000:8000"
environment:
- REDIS_URL=redis://redis:6379
- DATABASE_URL=postgresql://postgres:password@postgres:5432/realtime_db
- SECRET_KEY=${SECRET_KEY}
- DEBUG=false
depends_on:
- redis
- postgres
restart: unless-stopped
volumes:
- ./logs:/home/app/logs
networks:
- app-network
# Redis para cache e pub/sub
redis:
image: redis:7-alpine
ports:
- "6379:6379"
command: redis-server --appendonly yes
volumes:
- redis_data:/data
restart: unless-stopped
networks:
- app-network
# PostgreSQL para dados persistentes
postgres:
image: postgres:15-alpine
environment:
- POSTGRES_DB=realtime_db
- POSTGRES_USER=postgres
- POSTGRES_PASSWORD=password
ports:
- "5432:5432"
volumes:
- postgres_data:/var/lib/postgresql/data
- ./init.sql:/docker-entrypoint-initdb.d/init.sql
restart: unless-stopped
networks:
- app-network
# Nginx como proxy reverso
nginx:
image: nginx:alpine
ports:
- "80:80"
- "443:443"
volumes:
- ./nginx.conf:/etc/nginx/nginx.conf
- ./ssl:/etc/nginx/ssl
depends_on:
- app
restart: unless-stopped
networks:
- app-network
# Monitoramento com Prometheus
prometheus:
image: prom/prometheus
ports:
- "9090:9090"
volumes:
- ./prometheus.yml:/etc/prometheus/prometheus.yml
- prometheus_data:/prometheus
command:
- '--config.file=/etc/prometheus/prometheus.yml'
- '--storage.tsdb.path=/prometheus'
- '--web.console.libraries=/etc/prometheus/console_libraries'
- '--web.console.templates=/etc/prometheus/consoles'
networks:
- app-network
# Grafana para dashboards
grafana:
image: grafana/grafana
ports:
- "3000:3000"
environment:
- GF_SECURITY_ADMIN_PASSWORD=admin
volumes:
- grafana_data:/var/lib/grafana
networks:
- app-network
volumes:
redis_data:
postgres_data:
prometheus_data:
grafana_data:
networks:
app-network:
driver: bridge
2. Aplicação Escalável com Redis
# app.py - Aplicação escalável para produção
import os
import redis
import asyncio
import logging
from datetime import datetime, timedelta
from typing import Dict, List, Optional
import json
import jwt
from functools import wraps
import socketio
from aioredis import Redis
import asyncpg
from fastapi import FastAPI, HTTPException, Depends, Security
from fastapi.security import HTTPBearer, HTTPAuthorizationCredentials
from fastapi.middleware.cors import CORSMiddleware
import uvicorn
# Configuração de logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
# Configurações de ambiente
REDIS_URL = os.getenv('REDIS_URL', 'redis://localhost:6379')
DATABASE_URL = os.getenv('DATABASE_URL', 'postgresql://postgres:password@localhost:5432/realtime_db')
SECRET_KEY = os.getenv('SECRET_KEY', 'sua-chave-secreta-super-segura')
DEBUG = os.getenv('DEBUG', 'false').lower() == 'true'
class ProductionSocketIOServer:
def __init__(self):
# Redis para pub/sub entre instâncias
self.redis_client = None
self.db_pool = None
# Socket.IO com Redis adapter
self.sio = socketio.AsyncServer(
cors_allowed_origins="*" if DEBUG else ["https://seudominio.com"],
async_mode='asgi',
redis_url=REDIS_URL,
logger=DEBUG,
engineio_logger=DEBUG
)
# FastAPI para APIs REST
self.app = FastAPI(title="Realtime API", debug=DEBUG)
self.security = HTTPBearer()
# Configurações de CORS
self.app.add_middleware(
CORSMiddleware,
allow_origins=["*"] if DEBUG else ["https://seudominio.com"],
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
)
# Métricas e rate limiting
self.metrics = {
'connections': 0,
'messages_sent': 0,
'messages_received': 0,
'errors': 0
}
self.rate_limits = {} # {user_id: {action: [timestamps]}}
self.setup_events()
self.setup_routes()
async def init_services(self):
"""Inicializa serviços externos"""
try:
# Redis
self.redis_client = redis.from_url(REDIS_URL)
await self.redis_client.ping()
logger.info("✅ Redis conectado")
# PostgreSQL
self.db_pool = await asyncpg.create_pool(DATABASE_URL)
logger.info("✅ PostgreSQL conectado")
except Exception as e:
logger.error(f"❌ Erro ao conectar serviços: {e}")
raise
def verify_jwt_token(self, credentials: HTTPAuthorizationCredentials = Security(HTTPBearer())):
"""Verifica token JWT"""
try:
token = credentials.credentials
payload = jwt.decode(token, SECRET_KEY, algorithms=['HS256'])
return payload
except jwt.ExpiredSignatureError:
raise HTTPException(status_code=401, detail="Token expirado")
except jwt.InvalidTokenError:
raise HTTPException(status_code=401, detail="Token inválido")
def check_rate_limit(self, user_id: str, action: str, limit: int = 60, window: int = 60) -> bool:
"""Verifica rate limiting"""
now = datetime.now()
window_start = now - timedelta(seconds=window)
if user_id not in self.rate_limits:
self.rate_limits[user_id] = {}
if action not in self.rate_limits[user_id]:
self.rate_limits[user_id][action] = []
# Remove timestamps antigos
self.rate_limits[user_id][action] = [
ts for ts in self.rate_limits[user_id][action]
if ts > window_start
]
# Verifica limite
if len(self.rate_limits[user_id][action]) >= limit:
return False
# Adiciona timestamp atual
self.rate_limits[user_id][action].append(now)
return True
def setup_events(self):
@self.sio.event
async def connect(sid, environ, auth):
"""Evento de conexão com autenticação"""
try:
# Verifica autenticação
if not auth or 'token' not in auth:
logger.warning(f"Conexão rejeitada - sem token: {sid}")
return False
# Valida token
try:
payload = jwt.decode(auth['token'], SECRET_KEY, algorithms=['HS256'])
user_id = payload.get('user_id')
username = payload.get('username')
if not user_id or not username:
logger.warning(f"Token inválido - dados faltando: {sid}")
return False
except jwt.InvalidTokenError as e:
logger.warning(f"Token JWT inválido: {e}")
return False
# Rate limiting para conexões
if not self.check_rate_limit(user_id, 'connect', limit=10, window=60):
logger.warning(f"Rate limit excedido para conexão: {user_id}")
return False
# Armazena dados do usuário na sessão
async with self.sio.session(sid) as session:
session['user_id'] = user_id
session['username'] = username
session['connected_at'] = datetime.now().isoformat()
# Armazena no Redis para compartilhamento entre instâncias
user_data = {
'sid': sid,
'user_id': user_id,
'username': username,
'connected_at': datetime.now().isoformat(),
'instance': os.getpid() # ID da instância
}
await self.redis_client.hset(
f"user_sessions:{user_id}",
sid,
json.dumps(user_data)
)
# Métricas
self.metrics['connections'] += 1
logger.info(f"✅ Usuário conectado: {username} ({user_id}) - SID: {sid}")
# Envia confirmação
await self.sio.emit('authenticated', {
'user_id': user_id,
'username': username,
'server_time': datetime.now().isoformat()
}, room=sid)
return True
except Exception as e:
logger.error(f"Erro na conexão: {e}")
self.metrics['errors'] += 1
return False
@self.sio.event
async def disconnect(sid):
"""Evento de desconexão"""
try:
# Recupera dados da sessão
async with self.sio.session(sid) as session:
user_id = session.get('user_id')
username = session.get('username')
if user_id:
# Remove do Redis
await self.redis_client.hdel(f"user_sessions:{user_id}", sid)
# Se não há mais sessões ativas, remove usuário completamente
sessions = await self.redis_client.hgetall(f"user_sessions:{user_id}")
if not sessions:
await self.redis_client.delete(f"user_sessions:{user_id}")
logger.info(f"❌ Usuário desconectado: {username} ({user_id})")
self.metrics['connections'] -= 1
except Exception as e:
logger.error(f"Erro na desconexão: {e}")
@self.sio.event
async def join_room(sid, data):
"""Usuário entra em uma sala"""
try:
async with self.sio.session(sid) as session:
user_id = session.get('user_id')
username = session.get('username')
if not user_id:
await self.sio.emit('error', {'message': 'Não autenticado'}, room=sid)
return
room_name = data.get('room', '').strip()
if not room_name:
await self.sio.emit('error', {'message': 'Nome da sala inválido'}, room=sid)
return
# Rate limiting
if not self.check_rate_limit(user_id, 'join_room', limit=20, window=60):
await self.sio.emit('error', {'message': 'Rate limit excedido'}, room=sid)
return
# Valida permissão (implementar lógica de negócio)
if not await self.user_can_join_room(user_id, room_name):
await self.sio.emit('error', {'message': 'Sem permissão para esta sala'}, room=sid)
return
# Entra na sala
await self.sio.enter_room(sid, room_name)
# Armazena no banco de dados
async with self.db_pool.acquire() as conn:
await conn.execute("""
INSERT INTO room_memberships (user_id, room_name, joined_at)
VALUES ($1, $2, $3)
ON CONFLICT (user_id, room_name)
DO UPDATE SET joined_at = $3
""", user_id, room_name, datetime.now())
# Notifica outros usuários
await self.sio.emit('user_joined', {
'user_id': user_id,
'username': username,
'room': room_name,
'timestamp': datetime.now().isoformat()
}, room=room_name, skip_sid=sid)
# Confirma para o usuário
await self.sio.emit('joined_room', {
'room': room_name,
'timestamp': datetime.now().isoformat()
}, room=sid)
logger.info(f"🏠 {username} entrou na sala {room_name}")
except Exception as e:
logger.error(f"Erro ao entrar na sala: {e}")
await self.sio.emit('error', {'message': 'Erro interno'}, room=sid)
self.metrics['errors'] += 1
@self.sio.event
async def send_message(sid, data):
"""Envia mensagem para uma sala"""
try:
async with self.sio.session(sid) as session:
user_id = session.get('user_id')
username = session.get('username')
if not user_id:
await self.sio.emit('error', {'message': 'Não autenticado'}, room=sid)
return
room_name = data.get('room', '').strip()
message = data.get('message', '').strip()
if not room_name or not message:
await self.sio.emit('error', {'message': 'Dados inválidos'}, room=sid)
return
# Rate limiting para mensagens
if not self.check_rate_limit(user_id, 'send_message', limit=30, window=60):
await self.sio.emit('error', {'message': 'Muitas mensagens. Aguarde.'}, room=sid)
return
# Validação de conteúdo (implementar filtros)
if not await self.validate_message_content(message):
await self.sio.emit('error', {'message': 'Mensagem contém conteúdo proibido'}, room=sid)
return
# Verifica se usuário está na sala
if not await self.user_in_room(user_id, room_name):
await self.sio.emit('error', {'message': 'Você não está nesta sala'}, room=sid)
return
# Cria mensagem
message_id = await self.save_message(user_id, room_name, message)
message_data = {
'id': message_id,
'user_id': user_id,
'username': username,
'message': message,
'room': room_name,
'timestamp': datetime.now().isoformat()
}
# Envia para todos na sala
await self.sio.emit('new_message', message_data, room=room_name)
# Métricas
self.metrics['messages_sent'] += 1
self.metrics['messages_received'] += 1
logger.info(f"💬 [{room_name}] {username}: {message[:50]}...")
except Exception as e:
logger.error(f"Erro ao enviar mensagem: {e}")
await self.sio.emit('error', {'message': 'Erro interno'}, room=sid)
self.metrics['errors'] += 1
async def user_can_join_room(self, user_id: str, room_name: str) -> bool:
"""Verifica se usuário pode entrar na sala"""
# Implementar lógica de negócio específica
# Ex: salas privadas, banimentos, etc.
return True
async def user_in_room(self, user_id: str, room_name: str) -> bool:
"""Verifica se usuário está na sala"""
async with self.db_pool.acquire() as conn:
result = await conn.fetchval("""
SELECT EXISTS(
SELECT 1 FROM room_memberships
WHERE user_id = $1 AND room_name = $2
)
""", user_id, room_name)
return result
async def validate_message_content(self, message: str) -> bool:
"""Valida conteúdo da mensagem"""
# Implementar filtros de spam, palavrões, etc.
if len(message) > 1000: # Limite de caracteres
return False
# Lista de palavras proibidas (exemplo simplificado)
forbidden_words = ['spam', 'hack', 'phishing']
message_lower = message.lower()
for word in forbidden_words:
if word in message_lower:
return False
return True
async def save_message(self, user_id: str, room_name: str, message: str) -> str:
"""Salva mensagem no banco de dados"""
async with self.db_pool.acquire() as conn:
message_id = await conn.fetchval("""
INSERT INTO messages (user_id, room_name, content, created_at)
VALUES ($1, $2, $3, $4)
RETURNING id
""", user_id, room_name, message, datetime.now())
return str(message_id)
def setup_routes(self):
"""Configura rotas da API REST"""
@self.app.get("/health")
async def health_check():
"""Health check para load balancer"""
return {
"status": "healthy",
"timestamp": datetime.now().isoformat(),
"metrics": self.metrics
}
@self.app.get("/metrics")
async def get_metrics():
"""Métricas da aplicação"""
return {
"metrics": self.metrics,
"redis_info": await self.get_redis_info(),
"db_pool_info": {
"size": self.db_pool.get_size(),
"max_size": self.db_pool.get_max_size(),
"min_size": self.db_pool.get_min_size()
}
}
@self.app.post("/auth/token")
async def create_token(user_data: dict):
"""Cria token JWT (implementar autenticação real)"""
# Em produção, verificar credenciais no banco
payload = {
'user_id': user_data.get('user_id'),
'username': user_data.get('username'),
'exp': datetime.utcnow() + timedelta(hours=24)
}
token = jwt.encode(payload, SECRET_KEY, algorithm='HS256')
return {"access_token": token, "token_type": "bearer"}
async def get_redis_info(self):
"""Informações do Redis"""
try:
info = await self.redis_client.info()
return {
"connected_clients": info.get('connected_clients'),
"used_memory_human": info.get('used_memory_human'),
"total_commands_processed": info.get('total_commands_processed')
}
except:
return {"error": "Redis não disponível"}
# Inicialização da aplicação
server = ProductionSocketIOServer()
app = socketio.ASGIApp(server.sio, server.app)
@app.on_event("startup")
async def startup_event():
await server.init_services()
if __name__ == "__main__":
config = uvicorn.Config(
app=app,
host="0.0.0.0",
port=8000,
log_level="info",
access_log=True
)
server = uvicorn.Server(config)
asyncio.run(server.serve())
3. Scripts de Deploy e Monitoramento
#!/bin/bash
# deploy.sh - Script de deploy automatizado
set -e
echo "🚀 Iniciando deploy da aplicação..."
# Variáveis
APP_NAME="realtime-app"
DOCKER_REGISTRY="seu-registry.com"
VERSION=$(git rev-parse --short HEAD)
# Build da imagem
echo "📦 Construindo imagem Docker..."
docker build -t $DOCKER_REGISTRY/$APP_NAME:$VERSION .
docker tag $DOCKER_REGISTRY/$APP_NAME:$VERSION $DOCKER_REGISTRY/$APP_NAME:latest
# Push para registry
echo "📤 Enviando para registry..."
docker push $DOCKER_REGISTRY/$APP_NAME:$VERSION
docker push $DOCKER_REGISTRY/$APP_NAME:latest
# Deploy usando docker-compose
echo "🔄 Fazendo deploy..."
export IMAGE_VERSION=$VERSION
docker-compose -f docker-compose.prod.yml down
docker-compose -f docker-compose.prod.yml up -d
# Aguarda aplicação ficar saudável
echo "🏥 Verificando saúde da aplicação..."
for i in {1..30}; do
if curl -f http://localhost:8000/health > /dev/null 2>&1; then
echo "✅ Aplicação está saudável!"
break
fi
echo "Aguardando... ($i/30)"
sleep 10
done
# Testes de smoke
echo "🧪 Executando testes de smoke..."
curl -f http://localhost:8000/health
curl -f http://localhost:8000/metrics
echo "🎉 Deploy concluído com sucesso!"
# prometheus.yml - Configuração do Prometheus
global:
scrape_interval: 15s
evaluation_interval: 15s
rule_files:
- "alert_rules.yml"
alerting:
alertmanagers:
- static_configs:
- targets:
- alertmanager:9093
scrape_configs:
- job_name: 'realtime-app'
static_configs:
- targets: ['app:8000']
metrics_path: '/metrics'
scrape_interval: 5s
- job_name: 'redis'
static_configs:
- targets: ['redis:6379']
- job_name: 'postgres'
static_configs:
- targets: ['postgres:5432']
- job_name: 'nginx'
static_configs:
- targets: ['nginx:80']
# kubernetes.yml - Deploy no Kubernetes
apiVersion: apps/v1
kind: Deployment
metadata:
name: realtime-app
labels:
app: realtime-app
spec:
replicas: 3
selector:
matchLabels:
app: realtime-app
template:
metadata:
labels:
app: realtime-app
spec:
containers:
- name: app
image: seu-registry.com/realtime-app:latest
ports:
- containerPort: 8000
env:
- name: REDIS_URL
value: "redis://redis-service:6379"
- name: DATABASE_URL
valueFrom:
secretKeyRef:
name: db-secret
key: url
- name: SECRET_KEY
valueFrom:
secretKeyRef:
name: app-secret
key: secret-key
resources:
limits:
cpu: 500m
memory: 512Mi
requests:
cpu: 200m
memory: 256Mi
livenessProbe:
httpGet:
path: /health
port: 8000
initialDelaySeconds: 30
periodSeconds: 10
readinessProbe:
httpGet:
path: /health
port: 8000
initialDelaySeconds: 5
periodSeconds: 5
---
apiVersion: v1
kind: Service
metadata:
name: realtime-app-service
spec:
selector:
app: realtime-app
ports:
- protocol: TCP
port: 80
targetPort: 8000
type: LoadBalancer
---
apiVersion: autoscaling/v1
kind: HorizontalPodAutoscaler
metadata:
name: realtime-app-hpa
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: realtime-app
minReplicas: 2
maxReplicas: 10
targetCPUUtilizationPercentage: 70