const WebSocket = require('ws');
const sqlite3 = require('sqlite3').verbose();
// 1. Inicializar Base de Dados SQLite
const db = new sqlite3.Database('./logs.db');
db.serialize(() => {
db.run(`
CREATE TABLE IF NOT EXISTS desync_triggers (
id INTEGER PRIMARY KEY AUTOINCREMENT,
timestamp DATETIME DEFAULT CURRENT_TIMESTAMP,
expected_seq INTEGER,
received_seq INTEGER,
missed_count INTEGER,
status TEXT DEFAULT 'PENDING'
)
`);
});
// Lista global para monitorizar em tempo real as tarefas assíncronas ativas (soltas)
const activeTasks = [];
// FUNÇÃO ASSÍNCRONA SOLTA
// Executa imediatamente em paralelo e adiciona-se à lista global para tracking
async function runLooseRetroactiveCalculation(triggerId, expectedSeq, receivedSeq) {
const taskInfo = { triggerId, status: 'STARTED', startedAt: new Date() };
activeTasks.push(taskInfo);
console.log(`⚡ [TASKS SOLTAS] Nova tarefa assíncrona disparada em paralelo para o Trigger ID: ${triggerId}`);
console.log(`📊 Total de tarefas soltas a correr em simultâneo: ${activeTasks.length}`);
// 1. Atualizar base de dados para PROCESSING
db.run(`UPDATE desync_triggers SET status = 'PROCESSING' WHERE id = ?`, [triggerId]);
// 2. Simular cálculo pesado (Ex: 2.5 segundos)
// Como correm soltas, se 3 tarefas entrarem juntas, todas vão terminar quase ao mesmo tempo
await new Promise(resolve => setTimeout(resolve, 2500));
for (let i = expectedSeq; i < receivedSeq; i++) {
console.log(` 🔀 [Solta ID ${triggerId}] A processar sequência oculta: ${i}`);
}
// 3. Atualizar base de dados para COMPLETED
db.run(`UPDATE desync_triggers SET status = 'COMPLETED' WHERE id = ?`, [triggerId], (err) => {
if (!err) {
console.log(`✅ [TASKS SOLTAS] Terminou a tarefa do Trigger ID: ${triggerId}`);
// Remover a tarefa da lista de ativas
const index = activeTasks.findIndex(t => t.triggerId === triggerId);
if (index > -1) activeTasks.splice(index, 1);
console.log(`📉 Tarefas soltas restantes em execução: ${activeTasks.length}`);
}
});
}
// 3. Servidor WebSocket
const wss = new WebSocket.Server({ port: 8080 });
console.log("🚀 Servidor WebSocket ativo com Funções Soltas (Paralelo) na porta 8080");
wss.on('connection', (ws) => {
let expectedSequence = 1;
ws.on('message', (message) => {
try {
const data = JSON.parse(message);
const { seq } = data;
// DETEÇÃO DE DESINCRONIZAÇÃO
if (seq !== expectedSequence) {
const missedCount = seq - expectedSequence;
console.error(`❌ [DESYNC] Esperado: ${expectedSequence} | Recebido: ${seq}`);
// Registar na BD
const stmt = db.prepare(`
INSERT INTO desync_triggers (expected_seq, received_seq, missed_count, status)
VALUES (?, ?, ?, 'PENDING')
`);
stmt.run(expectedSequence, seq, missedCount, function(err) {
if (err) return;
const triggerId = this.lastID;
// EXECUÇÃO IMEDIATA E SOLTA: Não há await nem fila. Dispara em paralelo.
runLooseRetroactiveCalculation(triggerId, expectedSequence, seq);
});
stmt.finalize();
ws.send(JSON.stringify({ status: "DISPATCHED_IMMEDIATELY", triggerId }));
expectedSequence = seq + 1;
return;
}
console.log(`✅ [SYNC] Seq ${seq} OK.`);
expectedSequence++;
ws.send(JSON.stringify({ status: "SUCCESS", ack: seq }));
} catch (err) {
console.error("🚨 Erro:", err.message);
}
});
});






No comments:
Post a Comment