feat(federation): queue S2S presence_update on auth/disconnect/status/activity changes

New FederationPresenceUpdatePayload + queuePresenceRelay() helper. Five WS
sites now project the native user's status (and optional activities) to all
active peers via the outbox: WS auth-success, finalizeDisconnect,
manual presence_update, activity_update, showActivity-toggle clear.

Outbox-only (no mutation-log entry) — presence is ephemeral; the upcoming
peer-activation hook re-emits a fresh snapshot so peers recovering from
unreachable converge without history replay. No-op for replicated users.
This commit is contained in:
Jannis Braun
2026-05-05 16:01:28 +02:00
parent bdbc90ebd2
commit 613424e1c7
6 changed files with 230 additions and 1 deletions
+10
View File
@@ -511,6 +511,11 @@ function handlePresenceUpdate(event: Record<string, unknown>, userId: string): v
// Also send to self (other tabs)
connectionManager.sendToUser(userId, payload);
// S2S: project to all active peers
void import('../utils/federationPresence.js').then(({ queuePresenceRelay }) => {
try { queuePresenceRelay(userId, status as 'online' | 'idle' | 'dnd', activities); } catch (e) { console.warn('[ws] queuePresenceRelay(manual) failed', e); }
});
}
function handleActivityUpdate(event: Record<string, unknown>, userId: string): void {
@@ -535,6 +540,11 @@ function handleActivityUpdate(event: Record<string, unknown>, userId: string): v
connectionManager.sendToSpace(spaceId, payload, userId);
}
connectionManager.sendToUser(userId, payload);
// S2S: project to all active peers (activities + current status).
void import('../utils/federationPresence.js').then(({ queuePresenceRelay }) => {
try { queuePresenceRelay(userId, status as 'online' | 'idle' | 'dnd' | 'offline', activities); } catch (e) { console.warn('[ws] queuePresenceRelay(activity) failed', e); }
});
}
// ─── Voice Handlers (Unified Room API) ─────────────────────────────────────
+12
View File
@@ -307,6 +307,12 @@ class ConnectionManager {
});
}
// S2S: project offline to all active peers (mirrors profile_update fanout).
// Imported lazily to avoid circular import (federationPresence → db → ws/handler).
void import('../utils/federationPresence.js').then(({ queuePresenceRelay }) => {
try { queuePresenceRelay(userId, 'offline', []); } catch (e) { console.warn('[ws] queuePresenceRelay(offline) failed', e); }
});
// Clean up userSpaces (re-populated on next connect via setUserSpaces)
this.userSpaces.delete(userId);
@@ -1672,6 +1678,12 @@ export async function registerWebSocket(app: FastifyInstance): Promise<void> {
status: 'online',
}, userId);
}
// S2S: project online to all active peers (mirrors profile_update fanout).
const _uid = userId;
void import('../utils/federationPresence.js').then(({ queuePresenceRelay }) => {
try { queuePresenceRelay(_uid, 'online', []); } catch (e) { console.warn('[ws] queuePresenceRelay(online) failed', e); }
});
} catch {
ws.send(JSON.stringify({ type: 'error', message: 'Invalid token' }));
ws.close();