feat: persistent queue, pairing sessions, read receipts
This commit is contained in:
@@ -1,10 +1,8 @@
|
||||
# المرحلة الأولى: تنزيل الـ dependencies فقط
|
||||
FROM node:20-alpine AS deps
|
||||
WORKDIR /app
|
||||
COPY package.json ./
|
||||
RUN npm install --omit=dev
|
||||
|
||||
# المرحلة الثانية: الـ image النهائي
|
||||
FROM node:20-alpine
|
||||
WORKDIR /app
|
||||
|
||||
|
||||
@@ -6,20 +6,22 @@ services:
|
||||
env_file: .env
|
||||
volumes:
|
||||
- fchati-files:/data/files
|
||||
- fchati-queue:/data/queue
|
||||
- fchati-pairing:/data/pairing
|
||||
networks:
|
||||
- traefik-net
|
||||
labels:
|
||||
- "traefik.enable=true"
|
||||
- "traefik.http.routers.fchati.rule=Host(`fchati.diyaa.de`)"
|
||||
- "traefik.http.routers.fchati.entrypoints=websecure"
|
||||
# غيّر letsencrypt إذا كان اسم certresolver عندك مختلف
|
||||
- "traefik.http.routers.fchati.tls.certresolver=letsencrypt"
|
||||
- "traefik.http.services.fchati.loadbalancer.server.port=3000"
|
||||
|
||||
volumes:
|
||||
fchati-files:
|
||||
fchati-queue:
|
||||
fchati-pairing:
|
||||
|
||||
networks:
|
||||
traefik-net:
|
||||
# غيّر هذا ليطابق اسم network Traefik عندك
|
||||
external: true
|
||||
|
||||
+218
-53
@@ -3,22 +3,28 @@ const { WebSocketServer } = require('ws');
|
||||
const { createServer } = require('http');
|
||||
const { randomBytes, randomUUID } = require('crypto');
|
||||
const multer = require('multer');
|
||||
const path = require('path');
|
||||
const fs = require('fs');
|
||||
const path = require('path');
|
||||
|
||||
// ─── إعداد ───────────────────────────────────────────────────────────────────
|
||||
// ─── Config ───────────────────────────────────────────────────────────────────
|
||||
|
||||
const PORT = process.env.PORT || 3000;
|
||||
const MAX_FILE_BYTES = parseInt(process.env.MAX_FILE_SIZE_MB ?? '25') * 1024 * 1024;
|
||||
const UPLOADS_DIR = process.env.UPLOADS_DIR ?? '/data/files';
|
||||
const PAIRING_TTL_MS = 5 * 60 * 1000; // 5 دقائق
|
||||
const QUEUE_DIR = process.env.QUEUE_DIR ?? '/data/queue';
|
||||
const PAIRING_DIR = process.env.PAIRING_DIR ?? '/data/pairing';
|
||||
const PAIRING_TTL_MS = 5 * 60 * 1000;
|
||||
const QUEUE_TTL_MS = 7 * 24 * 60 * 60 * 1000;
|
||||
const QUEUE_MAX = 500;
|
||||
|
||||
fs.mkdirSync(UPLOADS_DIR, { recursive: true });
|
||||
fs.mkdirSync(QUEUE_DIR, { recursive: true });
|
||||
fs.mkdirSync(PAIRING_DIR, { recursive: true });
|
||||
|
||||
const app = express();
|
||||
app.use(express.json());
|
||||
|
||||
// ─── رفع الملفات ─────────────────────────────────────────────────────────────
|
||||
// ─── File storage ─────────────────────────────────────────────────────────────
|
||||
|
||||
const storage = multer.diskStorage({
|
||||
destination: UPLOADS_DIR,
|
||||
@@ -26,18 +32,135 @@ const storage = multer.diskStorage({
|
||||
});
|
||||
const upload = multer({ storage, limits: { fileSize: MAX_FILE_BYTES } });
|
||||
|
||||
// ─── الذاكرة الداخلية ─────────────────────────────────────────────────────────
|
||||
// ─── In-memory state ──────────────────────────────────────────────────────────
|
||||
|
||||
// code → { code, creatorID, creatorName, creatorToken, joinerID, joinerName, joinerToken, expiresAt }
|
||||
const pairingSessions = new Map();
|
||||
const pairingSessions = new Map(); // code -> session
|
||||
const connections = new Map(); // peerID -> { ws, name }
|
||||
const fileRegistry = new Map(); // fileID -> { diskPath, originalName, size }
|
||||
|
||||
// peerID → { ws, name }
|
||||
const connections = new Map();
|
||||
// ─── Persistent queue (disk-backed) ──────────────────────────────────────────
|
||||
//
|
||||
// Layout on disk:
|
||||
// /data/queue/{peerID}/{messageID}.json
|
||||
//
|
||||
// A message is written to disk the moment it is queued.
|
||||
// It is deleted from disk the moment it is delivered.
|
||||
// On server startup all existing files are loaded back into memory.
|
||||
|
||||
// fileID → { diskPath, originalName, size }
|
||||
const fileRegistry = new Map();
|
||||
function queueDir(peerID) {
|
||||
return path.join(QUEUE_DIR, peerID);
|
||||
}
|
||||
|
||||
// ─── أدوات مساعدة ─────────────────────────────────────────────────────────────
|
||||
function queuePath(peerID, messageID) {
|
||||
return path.join(queueDir(peerID), `${messageID}.json`);
|
||||
}
|
||||
|
||||
function persistMessage(peerID, messageID, envelope) {
|
||||
fs.mkdirSync(queueDir(peerID), { recursive: true });
|
||||
fs.writeFileSync(
|
||||
queuePath(peerID, messageID),
|
||||
JSON.stringify({ envelope, queuedAt: Date.now() })
|
||||
);
|
||||
}
|
||||
|
||||
function deletePersistedMessage(peerID, messageID) {
|
||||
try { fs.unlinkSync(queuePath(peerID, messageID)); } catch (_) {}
|
||||
}
|
||||
|
||||
function loadQueueFromDisk() {
|
||||
const queues = new Map();
|
||||
if (!fs.existsSync(QUEUE_DIR)) return queues;
|
||||
|
||||
for (const peerID of fs.readdirSync(QUEUE_DIR)) {
|
||||
const dir = queueDir(peerID);
|
||||
if (!fs.statSync(dir).isDirectory()) continue;
|
||||
const entries = [];
|
||||
for (const file of fs.readdirSync(dir)) {
|
||||
if (!file.endsWith('.json')) continue;
|
||||
try {
|
||||
const raw = fs.readFileSync(path.join(dir, file), 'utf8');
|
||||
const { envelope, queuedAt } = JSON.parse(raw);
|
||||
// Drop messages older than TTL
|
||||
if (Date.now() - queuedAt > QUEUE_TTL_MS) {
|
||||
fs.unlinkSync(path.join(dir, file));
|
||||
continue;
|
||||
}
|
||||
const messageID = file.replace('.json', '');
|
||||
entries.push({ messageID, envelope, queuedAt });
|
||||
} catch (_) {}
|
||||
}
|
||||
if (entries.length > 0) {
|
||||
entries.sort((a, b) => a.queuedAt - b.queuedAt);
|
||||
queues.set(peerID, entries);
|
||||
}
|
||||
}
|
||||
return queues;
|
||||
}
|
||||
|
||||
// In-memory queue mirrors disk — both are always in sync
|
||||
const offlineQueues = loadQueueFromDisk();
|
||||
console.log(`[queue] loaded ${[...offlineQueues.values()].reduce((s, q) => s + q.length, 0)} queued messages from disk`);
|
||||
|
||||
function enqueue(peerID, messageID, envelope) {
|
||||
if (!offlineQueues.has(peerID)) offlineQueues.set(peerID, []);
|
||||
const queue = offlineQueues.get(peerID);
|
||||
if (queue.length >= QUEUE_MAX) {
|
||||
const dropped = queue.shift();
|
||||
deletePersistedMessage(peerID, dropped.messageID);
|
||||
}
|
||||
queue.push({ messageID, envelope, queuedAt: Date.now() });
|
||||
persistMessage(peerID, messageID, envelope);
|
||||
}
|
||||
|
||||
function flushQueue(peerID, ws) {
|
||||
const queue = offlineQueues.get(peerID);
|
||||
if (!queue || queue.length === 0) return;
|
||||
const now = Date.now();
|
||||
for (const { messageID, envelope, queuedAt } of queue) {
|
||||
if (now - queuedAt < QUEUE_TTL_MS) send(ws, envelope);
|
||||
deletePersistedMessage(peerID, messageID);
|
||||
}
|
||||
offlineQueues.delete(peerID);
|
||||
}
|
||||
|
||||
// ─── Persistent pairing sessions (disk-backed) ────────────────────────────────
|
||||
//
|
||||
// Pairing sessions are written to disk so tokens survive a server restart.
|
||||
// Without this, everyone would need to re-pair after every server update.
|
||||
|
||||
function pairingPath(code) {
|
||||
return path.join(PAIRING_DIR, `${code}.json`);
|
||||
}
|
||||
|
||||
function persistSession(session) {
|
||||
fs.writeFileSync(pairingPath(session.code), JSON.stringify(session));
|
||||
}
|
||||
|
||||
function deleteSession(code) {
|
||||
try { fs.unlinkSync(pairingPath(code)); } catch (_) {}
|
||||
}
|
||||
|
||||
function loadSessionsFromDisk() {
|
||||
if (!fs.existsSync(PAIRING_DIR)) return;
|
||||
const now = Date.now();
|
||||
for (const file of fs.readdirSync(PAIRING_DIR)) {
|
||||
if (!file.endsWith('.json')) continue;
|
||||
try {
|
||||
const session = JSON.parse(fs.readFileSync(path.join(PAIRING_DIR, file), 'utf8'));
|
||||
if (session.expiresAt && session.expiresAt < now) {
|
||||
fs.unlinkSync(path.join(PAIRING_DIR, file));
|
||||
continue;
|
||||
}
|
||||
// Permanent sessions (joined pairs) have no expiry — keep them forever
|
||||
pairingSessions.set(session.code, session);
|
||||
} catch (_) {}
|
||||
}
|
||||
}
|
||||
|
||||
loadSessionsFromDisk();
|
||||
console.log(`[pairing] loaded ${pairingSessions.size} sessions from disk`);
|
||||
|
||||
// ─── Helpers ──────────────────────────────────────────────────────────────────
|
||||
|
||||
function generateCode() {
|
||||
const alphabet = 'ABCDEFGHJKLMNPQRSTUVWXYZ23456789';
|
||||
@@ -62,39 +185,50 @@ function peerIDFromToken(session, token) {
|
||||
return session.creatorToken === token ? session.creatorID : session.joinerID;
|
||||
}
|
||||
|
||||
function cleanExpiredSessions() {
|
||||
const now = Date.now();
|
||||
for (const [code, s] of pairingSessions) {
|
||||
if (s.expiresAt < now) pairingSessions.delete(code);
|
||||
}
|
||||
function peerPartnerID(session, peerID) {
|
||||
return session.creatorID === peerID ? session.joinerID : session.creatorID;
|
||||
}
|
||||
setInterval(cleanExpiredSessions, 60_000);
|
||||
|
||||
function send(ws, obj) {
|
||||
if (ws.readyState === 1) ws.send(JSON.stringify(obj));
|
||||
}
|
||||
|
||||
// ─── Middleware مصادقة HTTP ───────────────────────────────────────────────────
|
||||
function cleanExpiredSessions() {
|
||||
const now = Date.now();
|
||||
for (const [code, s] of pairingSessions) {
|
||||
// Only remove pending (not yet joined) sessions that expired
|
||||
if (!s.joinerID && s.expiresAt < now) {
|
||||
pairingSessions.delete(code);
|
||||
deleteSession(code);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
setInterval(cleanExpiredSessions, 60_000);
|
||||
|
||||
// ─── Auth middleware (HTTP) ───────────────────────────────────────────────────
|
||||
|
||||
function requireToken(req, res, next) {
|
||||
const auth = req.headers.authorization ?? '';
|
||||
if (!auth.startsWith('Bearer ')) return res.status(401).json({ error: 'Unauthorized' });
|
||||
const session = sessionByToken(auth.slice(7));
|
||||
const token = auth.slice(7);
|
||||
const session = sessionByToken(token);
|
||||
if (!session) return res.status(401).json({ error: 'Invalid token' });
|
||||
req.peerID = peerIDFromToken(session, auth.slice(7));
|
||||
req.peerID = peerIDFromToken(session, token);
|
||||
req.session = session;
|
||||
next();
|
||||
}
|
||||
|
||||
// ─── مسارات الـ Pairing ───────────────────────────────────────────────────────
|
||||
// ─── Pairing routes ───────────────────────────────────────────────────────────
|
||||
|
||||
app.post('/pairing/create', (req, res) => {
|
||||
const name = req.body?.displayName?.trim();
|
||||
if (!name) return res.status(400).json({ error: 'displayName مطلوب' });
|
||||
if (!name) return res.status(400).json({ error: 'displayName is required' });
|
||||
|
||||
let code, tries = 0;
|
||||
do {
|
||||
code = generateCode();
|
||||
if (++tries > 200) return res.status(503).json({ error: 'حاول مجدداً' });
|
||||
if (++tries > 200) return res.status(503).json({ error: 'Please try again' });
|
||||
} while (pairingSessions.has(code));
|
||||
|
||||
const session = {
|
||||
@@ -108,6 +242,7 @@ app.post('/pairing/create', (req, res) => {
|
||||
expiresAt: Date.now() + PAIRING_TTL_MS,
|
||||
};
|
||||
pairingSessions.set(code, session);
|
||||
persistSession(session);
|
||||
|
||||
res.json({
|
||||
code,
|
||||
@@ -120,55 +255,56 @@ app.post('/pairing/create', (req, res) => {
|
||||
app.post('/pairing/join', (req, res) => {
|
||||
const code = req.body?.code?.trim().toUpperCase();
|
||||
const name = req.body?.displayName?.trim();
|
||||
if (!code || !name) return res.status(400).json({ error: 'code و displayName مطلوبان' });
|
||||
if (!code || !name) return res.status(400).json({ error: 'code and displayName are required' });
|
||||
|
||||
const session = pairingSessions.get(code);
|
||||
if (!session || Date.now() > session.expiresAt) {
|
||||
if (!session || (!session.joinerID && Date.now() > session.expiresAt)) {
|
||||
pairingSessions.delete(code);
|
||||
return res.status(404).json({ error: 'الكود غير موجود أو انتهت صلاحيته' });
|
||||
deleteSession(code);
|
||||
return res.status(404).json({ error: 'Code not found or expired' });
|
||||
}
|
||||
if (session.joinerID) return res.status(409).json({ error: 'الكود استُخدم مسبقاً' });
|
||||
if (session.joinerID) return res.status(409).json({ error: 'Code already used' });
|
||||
|
||||
session.joinerID = randomUUID();
|
||||
session.joinerName = name;
|
||||
session.joinerToken = generateToken();
|
||||
delete session.expiresAt; // paired sessions never expire
|
||||
persistSession(session);
|
||||
|
||||
res.json({
|
||||
token: session.joinerToken,
|
||||
peerID: session.joinerID,
|
||||
peer: {
|
||||
id: session.creatorID,
|
||||
displayName: session.creatorName,
|
||||
},
|
||||
peer: { id: session.creatorID, displayName: session.creatorName },
|
||||
});
|
||||
});
|
||||
|
||||
// ─── مسارات الملفات ───────────────────────────────────────────────────────────
|
||||
// ─── File routes ──────────────────────────────────────────────────────────────
|
||||
|
||||
app.post('/files', requireToken, upload.single('file'), (req, res) => {
|
||||
if (!req.file) return res.status(400).json({ error: 'لم يُرسَل أي ملف' });
|
||||
|
||||
if (!req.file) return res.status(400).json({ error: 'No file provided' });
|
||||
const fileID = req.file.filename;
|
||||
fileRegistry.set(fileID, {
|
||||
diskPath: req.file.path,
|
||||
originalName: req.file.originalname,
|
||||
size: req.file.size,
|
||||
});
|
||||
|
||||
res.json({ id: fileID, name: req.file.originalname, size: req.file.size });
|
||||
});
|
||||
|
||||
app.get('/files/:id', requireToken, (req, res) => {
|
||||
const meta = fileRegistry.get(req.params.id);
|
||||
if (!meta) return res.status(404).json({ error: 'الملف غير موجود' });
|
||||
|
||||
if (!meta) return res.status(404).json({ error: 'File not found' });
|
||||
res.setHeader('Content-Disposition', `attachment; filename="${encodeURIComponent(meta.originalName)}"`);
|
||||
res.sendFile(meta.diskPath);
|
||||
});
|
||||
|
||||
// ─── فحص الصحة ───────────────────────────────────────────────────────────────
|
||||
// ─── Health check ─────────────────────────────────────────────────────────────
|
||||
|
||||
app.get('/health', (_req, res) => res.json({ ok: true, connections: connections.size }));
|
||||
app.get('/health', (_req, res) => res.json({
|
||||
ok: true,
|
||||
connections: connections.size,
|
||||
queued: [...offlineQueues.values()].reduce((sum, q) => sum + q.length, 0),
|
||||
}));
|
||||
|
||||
// ─── WebSocket ────────────────────────────────────────────────────────────────
|
||||
|
||||
@@ -178,6 +314,7 @@ const wss = new WebSocketServer({ server, path: '/ws' });
|
||||
wss.on('connection', (ws) => {
|
||||
let peerID = null;
|
||||
let peerName = null;
|
||||
let session = null;
|
||||
let authed = false;
|
||||
|
||||
const authTimeout = setTimeout(() => {
|
||||
@@ -187,14 +324,15 @@ wss.on('connection', (ws) => {
|
||||
ws.on('message', (raw) => {
|
||||
let msg;
|
||||
try { msg = JSON.parse(raw.toString()); }
|
||||
catch { return ws.close(1003, 'JSON غير صالح'); }
|
||||
catch { return ws.close(1003, 'Invalid JSON'); }
|
||||
|
||||
// ─── Auth ─────────────────────────────────────────────────────────────────
|
||||
if (!authed) {
|
||||
if (msg.type !== 'auth' || !msg.token) {
|
||||
return ws.close(1008, 'أرسل { type: "auth", token: "..." } أولاً');
|
||||
return ws.close(1008, 'Send { type: "auth", token: "..." } first');
|
||||
}
|
||||
const session = sessionByToken(msg.token);
|
||||
if (!session) return ws.close(1008, 'Token غير صالح');
|
||||
session = sessionByToken(msg.token);
|
||||
if (!session) return ws.close(1008, 'Invalid token');
|
||||
|
||||
clearTimeout(authTimeout);
|
||||
authed = true;
|
||||
@@ -203,21 +341,48 @@ wss.on('connection', (ws) => {
|
||||
connections.set(peerID, { ws, name: peerName });
|
||||
|
||||
send(ws, { type: 'auth.ok', peerID });
|
||||
console.log(`[ws] متصل: ${peerName} (${peerID})`);
|
||||
console.log(`[ws] connected: ${peerName} (${peerID})`);
|
||||
|
||||
// Deliver messages that arrived while this peer was offline
|
||||
flushQueue(peerID, ws);
|
||||
return;
|
||||
}
|
||||
|
||||
// ─── Chat message ─────────────────────────────────────────────────────────
|
||||
if (msg.type === 'chat.message') {
|
||||
const target = connections.get(msg.to);
|
||||
const partnerID = peerPartnerID(session, peerID);
|
||||
const messageID = msg.id ?? randomUUID();
|
||||
const outgoing = { ...msg, id: messageID, from: peerID, fromName: peerName };
|
||||
const target = connections.get(partnerID);
|
||||
|
||||
if (target?.ws.readyState === 1) {
|
||||
send(target.ws, { ...msg, from: peerID, fromName: peerName });
|
||||
send(ws, { type: 'delivered', id: msg.id });
|
||||
send(target.ws, outgoing);
|
||||
send(ws, { type: 'delivered', id: messageID });
|
||||
} else {
|
||||
send(ws, { type: 'not_delivered', id: msg.id, reason: 'peer_offline' });
|
||||
enqueue(partnerID, messageID, outgoing);
|
||||
send(ws, { type: 'queued', id: messageID });
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
// ─── Read receipt ─────────────────────────────────────────────────────────
|
||||
// Client sends: { type: "read", messageIDs: ["id1", "id2", ...] }
|
||||
// Server forwards to the sender of those messages so they see "read" ticks
|
||||
if (msg.type === 'read') {
|
||||
const partnerID = peerPartnerID(session, peerID);
|
||||
const receipt = { type: 'read', messageIDs: msg.messageIDs, by: peerID };
|
||||
const target = connections.get(partnerID);
|
||||
|
||||
if (target?.ws.readyState === 1) {
|
||||
send(target.ws, receipt);
|
||||
} else {
|
||||
// Queue the read receipt too — partner deserves to know
|
||||
enqueue(partnerID, `read-${randomUUID()}`, receipt);
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
// ─── Ping ─────────────────────────────────────────────────────────────────
|
||||
if (msg.type === 'ping') {
|
||||
send(ws, { type: 'pong' });
|
||||
}
|
||||
@@ -226,18 +391,18 @@ wss.on('connection', (ws) => {
|
||||
ws.on('close', () => {
|
||||
if (peerID) {
|
||||
connections.delete(peerID);
|
||||
console.log(`[ws] انقطع: ${peerName} (${peerID})`);
|
||||
console.log(`[ws] disconnected: ${peerName} (${peerID})`);
|
||||
}
|
||||
clearTimeout(authTimeout);
|
||||
});
|
||||
|
||||
ws.on('error', (err) => {
|
||||
console.error(`[ws] خطأ:`, err.message);
|
||||
console.error('[ws] error:', err.message);
|
||||
});
|
||||
});
|
||||
|
||||
// ─── تشغيل ────────────────────────────────────────────────────────────────────
|
||||
// ─── Start ────────────────────────────────────────────────────────────────────
|
||||
|
||||
server.listen(PORT, () => {
|
||||
console.log(`Fchati Relay يعمل على المنفذ ${PORT}`);
|
||||
console.log(`Fchati Relay listening on port ${PORT}`);
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user