// ws-server.js require("dotenv").config(); const WebSocket = require("ws"); const jwt = require("jsonwebtoken"); const { randomUUID } = require("crypto"); const dbWs = require("./db/dbWs"); const PORT = process.env.WS_PORT || 4002; const HOST = process.env.WS_HOST || "0.0.0.0"; const JWT_SECRET = process.env.JWT_SECRET; if (!JWT_SECRET) { console.error("❌ ERRORE: JWT_SECRET non definito nel .env"); process.exit(1); } const wss = new WebSocket.Server({ port: PORT, host: HOST }); console.log(`🚀 WebSocket server attivo su ws://${HOST}:${PORT}`); // ===================================================== // PULIZIA SESSIONI E PENDING AL RIAVVIO SERVER // ===================================================== // Nota: // questa pulizia è lasciata ATTIVA perché la usi per i test. // La mancanza di sessione NON viene interpretata come full sync, // ma solo come need_recovery=true. (async () => { try { console.log("🧹 [WS CLEANUP] Pulizia sessioni e pending al riavvio server..."); await dbWs.deleteAllSessions(); await dbWs.deleteAllPendingEvents(); console.log("🧹 [WS CLEANUP] Completata: tutte le sessioni e pending rimossi"); } catch (err) { console.error("❌ [WS CLEANUP ERROR]:", err); } })(); // Mappa session_id -> Map(device_id -> ws) const wsBySession = new Map(); // ===================================================== // COSTANTI // ===================================================== const WS_DORMANT_MS = 2 * 60 * 1000; // 2 minuti // ===================================================== // HELPER: invio con log // ===================================================== function sendToClient(ws, msg) { const json = JSON.stringify(msg); console.log(`📤 [WS SEND] → user=${ws.user}: ${json}`); ws.send(json); } // ===================================================== // HELPER: decide recovery/full sync // ===================================================== async function decideAuthRecovery({ session_id, user, device_id, }) { const sessionRow = await dbWs.getSession(session_id, device_id); let needFull = false; let needRecovery = false; let reason = "none"; if (!sessionRow) { // Se cancelli le sessioni al riavvio per test, è normale non trovarla. // Questo NON deve significare full bootstrap. // Significa solo: fai recovery/delta. needRecovery = true; reason = "session not found after cleanup"; } else { const lastAck = sessionRow.last_ack || 0; const dormancy = lastAck > 0 ? Date.now() - lastAck : 0; if (sessionRow.need_full_sync) { // Manteniamo compatibilità con la vecchia colonna, // ma semanticamente la trattiamo come "need recovery". needRecovery = true; reason = "db recovery flag"; } else if (lastAck > 0 && dormancy >= WS_DORMANT_MS) { // Dormant non è full sync. // È recovery progressiva. needRecovery = true; reason = "dormant"; } else { needRecovery = false; reason = "recent or clean auth"; } } // Per compatibilità con lo schema attuale: // salviamo needRecovery nella colonna need_full_sync, // ma nel protocollo WS mandiamo need_recovery separato. await dbWs.setSessionNeedFullSync(session_id, device_id, needRecovery); console.log("🧪 [WS AUTH DECISION]", { user, session_id, device_id, needFull, needRecovery, reason, last_ack: sessionRow?.last_ack, last_sync: sessionRow?.last_sync, db_need_full_sync: sessionRow?.need_full_sync, }); return { needFull, needRecovery, reason, }; } // ===================================================== // CONNESSIONE CLIENT // ===================================================== wss.on("connection", (ws) => { console.log("🔌 [WS CONNECT] Nuovo client"); ws.authenticated = false; ws.user = null; ws.session_id = null; ws.device_id = null; ws.request_id = null; ws.on("message", async (msg) => { console.log(`📩 [WS RAW] msg=${msg.toString()}`); let data; try { data = JSON.parse(msg.toString()); } catch (err) { console.error("❌ Errore parsing JSON:", err); return; } // ===================================================== // AUTENTICAZIONE // ===================================================== if (data.type === "auth") { try { const payload = jwt.verify(data.token, JWT_SECRET); const user = payload.name; const session_id = data.session_id || randomUUID(); // IMPORTANTE: // Non generare randomUUID() a ogni connessione se il client // non manda device_id, altrimenti il server vede sempre // un device nuovo. // // Per test singolo dispositivo va bene default-device. // In futuro meglio far mandare un device_id persistente dal client. const device_id = data.device_id || "default-device"; const request_id = data.request_id || null; console.log("🧪 AUTH DEBUG:", { data_device_id: data.device_id, ws_device_id: device_id, session_id, user, }); ws.user = user; ws.session_id = session_id; ws.device_id = device_id; ws.request_id = request_id; ws.authenticated = true; if (!wsBySession.has(session_id)) { wsBySession.set(session_id, new Map()); } wsBySession.get(session_id).set(device_id, ws); console.log( `🔐 [WS AUTH OK] user=${user} session=${session_id} device=${device_id}` ); await dbWs.upsertSession(session_id, user, device_id); const decision = await decideAuthRecovery({ session_id, user, device_id, }); sendToClient(ws, { type: "auth_ok", user, session_id, // Full bootstrap vero. // Per dormancy/session cleanup deve restare false. need_full_sync: decision.needFull, // Flag corretto per delta/progressive recovery. need_recovery: decision.needRecovery, reason: decision.reason, }); } catch (err) { console.error("❌ [WS AUTH ERROR]:", err); sendToClient(ws, { type: "auth_error", error: err.message, }); ws.close(); } return; } // ===================================================== // BLOCCO NON AUTENTICATO // ===================================================== if (!ws.authenticated) { sendToClient(ws, { type: "error", message: "Not authenticated", }); return; } // ===================================================== // ACK EVENTO // ===================================================== if (data.type === "ack" && data.event_id) { console.log(`🟢 [WS ACK] user=${ws.user} event=${data.event_id}`); await dbWs.updateSessionAck(ws.session_id, ws.device_id); await dbWs.deletePendingEvent({ user_id: ws.user, session_id: ws.session_id, device_id: ws.device_id, event_id: data.event_id, }); return; } // ===================================================== // CLIENT RECOVERED // ===================================================== if (data.type === "client_recovered") { console.log(`🔄 [WS CLIENT RECOVERED] user=${ws.user}`); const session = await dbWs.getSession(ws.session_id, ws.device_id); if (!session) { console.log("⚠️ [WS CLIENT RECOVERED] session missing -> need_recovery"); sendToClient(ws, { type: "need_recovery", reason: "session missing", }); return; } const lastAck = session.last_ack || 0; const dormancy = lastAck > 0 ? Date.now() - lastAck : 0; if (dormancy >= WS_DORMANT_MS || session.need_full_sync) { sendToClient(ws, { type: "need_recovery", reason: dormancy >= WS_DORMANT_MS ? "dormant" : "db recovery flag", }); return; } const pending = await dbWs.getPendingEvents( ws.user, ws.session_id, ws.device_id ); for (const ev of pending) { ws.send(ev.payload); } return; } // ===================================================== // PONG // ===================================================== if (data.type === "pong") { await dbWs.updateSessionAck(ws.session_id, ws.device_id); return; } // ===================================================== // RECOVERY DONE // ===================================================== if (data.type === "recovery_done") { console.log( `✅ [WS RECOVERY DONE] user=${ws.user} session=${ws.session_id} device=${ws.device_id}` ); await dbWs.clearNeedFullSync(ws.session_id, ws.device_id); await dbWs.updateSessionAck(ws.session_id, ws.device_id); return; } // ===================================================== // ALTRI MESSAGGI // ===================================================== console.log(`📨 [WS MESSAGE] user=${ws.user} data=${JSON.stringify(data)}`); }); ws.on("close", () => { console.log(`❌ [WS DISCONNECT] user=${ws.user} session=${ws.session_id}`); if (ws.session_id && ws.device_id) { const devices = wsBySession.get(ws.session_id); if (devices) { devices.delete(ws.device_id); if (devices.size === 0) { wsBySession.delete(ws.session_id); } } } }); }); // ===================================================== // BROADCAST AFFIDABILE // ===================================================== wss.broadcastToUserReliable = async function (user, payload) { console.log(`📣 [WS BROADCAST] user=${user}`); const sessions = await dbWs.getSessionsByUser(user); if (!sessions || sessions.length === 0) return; const event_id = randomUUID(); const now = Date.now(); // 1. Salva pending per TUTTI i device registrati nel DB, // anche se offline. for (const s of sessions) { const deviceRows = await dbWs.getDevicesForSession(s.session_id); for (const device_id of deviceRows) { await dbWs.insertPendingEvent({ user_id: user, session_id: s.session_id, device_id, event_id, payload: JSON.stringify(payload), created_at: now, }); } } // 2. Invia WS ai device attivi. for (const s of sessions) { const devices = wsBySession.get(s.session_id); if (!devices) continue; for (const [device_id, ws] of devices) { if (ws.readyState === WebSocket.OPEN && ws.authenticated) { ws.send( JSON.stringify({ ...payload, event_id, }) ); } } } }; wss.broadcastToUser = wss.broadcastToUserReliable; module.exports = wss;