From 613424e1c704dd5746e7a09d2c7e1e9d0dc4b052 Mon Sep 17 00:00:00 2001 From: Jannis Braun <151788261+TheZwiss@users.noreply.github.com> Date: Tue, 5 May 2026 16:01:28 +0200 Subject: [PATCH] feat(federation): queue S2S presence_update on auth/disconnect/status/activity changes MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- packages/server/src/routes/users.ts | 11 ++ .../src/utils/federationPresence.test.ts | 120 ++++++++++++++++++ .../server/src/utils/federationPresence.ts | 61 +++++++++ packages/server/src/ws/events.ts | 10 ++ packages/server/src/ws/handler.ts | 12 ++ packages/shared/src/types.ts | 17 ++- 6 files changed, 230 insertions(+), 1 deletion(-) create mode 100644 packages/server/src/utils/federationPresence.test.ts create mode 100644 packages/server/src/utils/federationPresence.ts diff --git a/packages/server/src/routes/users.ts b/packages/server/src/routes/users.ts index 5a1f4b18..654600a9 100644 --- a/packages/server/src/routes/users.ts +++ b/packages/server/src/routes/users.ts @@ -447,6 +447,17 @@ export async function userRoutes(app: FastifyInstance): Promise { connectionManager.sendToSpace(spaceId, clearPayload, request.userId); } connectionManager.sendToUser(request.userId, clearPayload); + + // S2S: project the cleared-activities snapshot to all active peers. + void import('../utils/federationPresence.js').then(({ queuePresenceRelay }) => { + try { + queuePresenceRelay( + request.userId, + (connectionManager.getUserStatus(request.userId) ?? 'online') as 'online' | 'idle' | 'dnd' | 'offline', + [], + ); + } catch (e) { console.warn('[users] queuePresenceRelay(showActivity-clear) failed', e); } + }); } } diff --git a/packages/server/src/utils/federationPresence.test.ts b/packages/server/src/utils/federationPresence.test.ts new file mode 100644 index 00000000..6dd2b234 --- /dev/null +++ b/packages/server/src/utils/federationPresence.test.ts @@ -0,0 +1,120 @@ +import { describe, it, expect, beforeEach, vi } from 'vitest'; +import Database from 'better-sqlite3'; +import { drizzle } from 'drizzle-orm/better-sqlite3'; +import fs from 'node:fs'; +import path from 'node:path'; +import { fileURLToPath } from 'node:url'; +import * as schema from '../db/schema.js'; + +const __dirname = path.dirname(fileURLToPath(import.meta.url)); +type TestDb = ReturnType>; +let sqlite: Database.Database; +let testDb: TestDb; + +const queueCalls: Array<{ + entityId: string; + contextId: string; + eventType: string; + payload: string; + targetPeerOrigins: string[] | undefined; + contextType: string; +}> = []; +const mutationLogCalls: Array<{ entityId: string; eventType: string }> = []; + +vi.mock('../db/index.js', () => ({ + getDb: () => testDb, + getRawDb: () => sqlite, + schema, +})); + +vi.mock('./federationAuth.js', () => ({ + getOurOrigin: () => 'https://nova.ddns.net', +})); + +vi.mock('./federationOutbox.js', () => ({ + isFederationRelayEnabled: () => true, + queueOutboxEvent: vi.fn((entityId, contextId, eventType, payload, targetPeerOrigins, contextType) => { + queueCalls.push({ entityId, contextId, eventType, payload, targetPeerOrigins, contextType }); + }), + appendMutationLog: vi.fn((entityId, _ctxId, eventType) => { + mutationLogCalls.push({ entityId, eventType }); + }), +})); + +function applyMigrations(db: Database.Database): void { + const migrationsDir = path.resolve(__dirname, '../../drizzle'); + const files = fs.readdirSync(migrationsDir).filter(f => f.endsWith('.sql')).sort(); + for (const f of files) { + const sql = fs.readFileSync(path.join(migrationsDir, f), 'utf8'); + for (const stmt of sql.split(/-->\s*statement-breakpoint/)) { + const clean = stmt.trim(); + if (clean) db.exec(clean); + } + } +} + +beforeEach(() => { + sqlite = new Database(':memory:'); + testDb = drizzle(sqlite, { schema }); + applyMigrations(sqlite); + queueCalls.length = 0; + mutationLogCalls.length = 0; + // Native local user + testDb.insert(schema.users).values({ + id: 'native-1', + username: 'youruser', + passwordHash: 'x', + status: 'online', + isAdmin: 0, + homeUserId: 'native-1', + createdAt: Date.now(), + }).run(); +}); + +describe('queuePresenceRelay', () => { + it('queues an outbox event with status + activities for a native user', async () => { + const { queuePresenceRelay } = await import('./federationPresence.js'); + queuePresenceRelay('native-1', 'online', [{ type: 'playing', name: 'Test' }]); + + expect(queueCalls.length).toBe(1); + const call = queueCalls[0]!; + expect(call.eventType).toBe('presence_update'); + expect(call.contextType).toBe('profile'); + expect(call.targetPeerOrigins).toBeUndefined(); // broadcast to all active peers + const event = JSON.parse(call.payload); + expect(event.eventType).toBe('presence_update'); + expect(event.presenceUpdate.status).toBe('online'); + expect(event.presenceUpdate.activities).toEqual([{ type: 'playing', name: 'Test' }]); + // appendMutationLog NOT called — presence is outbox-only + expect(mutationLogCalls).toEqual([]); + }); + + it('omits activities field when none are passed', async () => { + const { queuePresenceRelay } = await import('./federationPresence.js'); + queuePresenceRelay('native-1', 'offline', []); + const event = JSON.parse(queueCalls[0]!.payload); + expect(event.presenceUpdate.activities).toBeUndefined(); + }); + + it('is a no-op for replicated users (homeInstance set)', async () => { + testDb.insert(schema.users).values({ + id: 'stub-1', + username: 'pbtest3@orbit.ddns.net', + passwordHash: '!federation-replicated', + status: 'online', + isAdmin: 0, + homeInstance: 'orbit.ddns.net', + homeUserId: 'remote-1', + createdAt: Date.now(), + }).run(); + const { queuePresenceRelay } = await import('./federationPresence.js'); + queuePresenceRelay('stub-1', 'online', []); + expect(queueCalls).toEqual([]); + }); + + it('is a no-op for unknown user IDs', async () => { + const { queuePresenceRelay } = await import('./federationPresence.js'); + queuePresenceRelay('does-not-exist', 'online', []); + expect(queueCalls).toEqual([]); + }); +}); diff --git a/packages/server/src/utils/federationPresence.ts b/packages/server/src/utils/federationPresence.ts new file mode 100644 index 00000000..83e3b2d9 --- /dev/null +++ b/packages/server/src/utils/federationPresence.ts @@ -0,0 +1,61 @@ +import { eq } from 'drizzle-orm'; +import type { Activity, FederationRelayEvent, FederationPresenceUpdatePayload } from '@backspace/shared'; +import { getDb, schema } from '../db/index.js'; +import { getOurOrigin } from './federationAuth.js'; +import { isFederationRelayEnabled, queueOutboxEvent } from './federationOutbox.js'; + +export type PresenceStatus = 'online' | 'idle' | 'dnd' | 'offline'; + +/** + * Queue a presence_update event for the given native user. Broadcast to all + * active peers (mirrors profile_update). Outbox-only — presence is ephemeral; + * stale replays from a mutation log are wrong, so we never call + * appendMutationLog. The peer-activation hook re-emits a fresh snapshot for + * peer-related online natives, so a peer recovering from unreachable converges + * without history replay. + * + * No-op for replicated users (their home instance owns presence projection). + */ +export function queuePresenceRelay( + userId: string, + status: PresenceStatus, + activities: Activity[], +): void { + if (!isFederationRelayEnabled()) return; + + const db = getDb(); + const user = db.select().from(schema.users).where(eq(schema.users.id, userId)).get(); + if (!user) return; + if (user.homeInstance) return; // replicated — not our authority + + const ts = Date.now(); + const payload: FederationPresenceUpdatePayload = { + homeUserId: user.id, + homeInstance: getOurOrigin(), + status, + ts, + ...(activities.length > 0 ? { activities } : {}), + }; + + const event: FederationRelayEvent = { + eventType: 'presence_update', + contextType: 'profile', + messageId: `presence:${user.id}:${ts}`, + encryptionVersion: 0, + timestamp: ts, + presenceUpdate: payload, + }; + + // entityId = userId so the outbox coalesces rapid status flaps into the latest. + // contextId = userId, contextType = 'profile' (reuses existing routing). + // targetPeerOrigins = undefined → broadcast to all active peers. + // NO appendMutationLog — presence must not be replayed from history. + queueOutboxEvent( + user.id, + user.id, + 'presence_update', + JSON.stringify(event), + undefined, + 'profile', + ); +} diff --git a/packages/server/src/ws/events.ts b/packages/server/src/ws/events.ts index 6aaa859c..7fb191c2 100644 --- a/packages/server/src/ws/events.ts +++ b/packages/server/src/ws/events.ts @@ -511,6 +511,11 @@ function handlePresenceUpdate(event: Record, 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, userId: string): void { @@ -535,6 +540,11 @@ function handleActivityUpdate(event: Record, 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) ───────────────────────────────────── diff --git a/packages/server/src/ws/handler.ts b/packages/server/src/ws/handler.ts index c730af2e..59618974 100644 --- a/packages/server/src/ws/handler.ts +++ b/packages/server/src/ws/handler.ts @@ -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 { 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(); diff --git a/packages/shared/src/types.ts b/packages/shared/src/types.ts index 64b87c3a..00e4666a 100644 --- a/packages/shared/src/types.ts +++ b/packages/shared/src/types.ts @@ -882,7 +882,7 @@ export interface FederationRelayEvent { | 'friend_add' | 'friend_remove' | 'file_rejected' | 'dm_call_start' | 'dm_call_accept' | 'dm_call_reject' | 'dm_call_end' | 'dm_typing_start' | 'dm_typing_stop' - | 'profile_update' + | 'profile_update' | 'presence_update' | 'read_state_update' | 'dm_close' | 'dm_reopen'; contextType?: 'dm' | 'friend' | 'profile'; @@ -922,6 +922,7 @@ export interface FederationRelayEvent { username: string; }; profileUpdate?: FederationProfileUpdatePayload; + presenceUpdate?: FederationPresenceUpdatePayload; readState?: { user: { homeUserId: string; homeInstance: string }; messageRef: { sourceInstance: string; sourceMessageId: string }; @@ -986,6 +987,20 @@ export interface FederationProfileUpdatePayload { bio: string | null; } +/** + * Presence projection from a home instance to peers. Carries the user's current + * online status and (optionally) rich activities. Outbox-only on the wire — never + * written to federation_mutation_log; presence is ephemeral and stale replays on + * peer activation are wrong (the activation hook re-emits a fresh snapshot). + */ +export interface FederationPresenceUpdatePayload { + homeUserId: string; + homeInstance: string; + status: 'online' | 'idle' | 'dnd' | 'offline'; + activities?: Activity[]; + ts: number; // emitter clock; receiver may use for last-write-wins +} + export interface FederationFriendshipPayload { from: FederationRelayParticipant; to: FederationRelayParticipant;