From 314df6c5a1d7ec94500f5a892c7e00f2e2e15725 Mon Sep 17 00:00:00 2001 From: Jannis Braun <151788261+TheZwiss@users.noreply.github.com> Date: Fri, 24 Apr 2026 00:49:04 +0200 Subject: [PATCH] feat(server): 30s federated-call sentinel worker (TDD) --- .../server/src/utils/federationWorker.test.ts | 133 +++++++++++++++++- packages/server/src/utils/federationWorker.ts | 58 ++++++++ 2 files changed, 190 insertions(+), 1 deletion(-) diff --git a/packages/server/src/utils/federationWorker.test.ts b/packages/server/src/utils/federationWorker.test.ts index 3e8e4d56..11d58850 100644 --- a/packages/server/src/utils/federationWorker.test.ts +++ b/packages/server/src/utils/federationWorker.test.ts @@ -1,4 +1,4 @@ -import { describe, it, expect, beforeEach, vi } from 'vitest'; +import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest'; import Database from 'better-sqlite3'; import { drizzle } from 'drizzle-orm/better-sqlite3'; import fs from 'node:fs'; @@ -23,6 +23,8 @@ vi.mock('../ws/handler.js', () => ({ getAllOnlineUserIds: () => [], sendToUser: vi.fn(), sendToDmMembers: vi.fn(), + evictFederatedCallsForHost: vi.fn().mockReturnValue(0), + getAllFederatedCalls: vi.fn(() => new Map()), }, })); @@ -41,6 +43,7 @@ vi.mock('../utils/federationOutbox.js', () => ({ vi.mock('../utils/federationPeerActivation.js', () => ({ onPeerActivated: vi.fn(), + onPeerDeactivated: vi.fn().mockResolvedValue(undefined), startupBootstrapSync: vi.fn(), })); @@ -157,3 +160,131 @@ describe('outbox worker — duplicate rejection is terminal', () => { expect(remaining).toBeDefined(); }); }); + +// ─── Sentinel test helpers ─────────────────────────────────────────────────── + +function seedSentinelPeer( + id: string, + origin: string, + status: string, + instanceName?: string, +): void { + testDb.insert(schema.federationPeers).values({ + id, + origin, + hmacSecret: 'secret', + status, + instanceName: instanceName ?? null, + lastSyncedAt: 0, + createdAt: Date.now(), + }).run(); +} + +type FedCallEntry = import('../ws/handler.js').FederatedCallEntry; + +function makeFedCall(partial: Partial): FedCallEntry { + return { + dmChannelId: null, + federatedId: `fed-${Math.random().toString(36).slice(2, 8)}`, + callerId: 'caller', + callerHomeUserId: 'caller-home', + federatedCallHost: 'https://hostA.example', + livekitUrl: 'wss://lk.example', + tokens: new Map(), + ringedUserIds: [], + state: 'active', + startedAt: Date.now(), + ...partial, + }; +} + +// ─── Sentinel describe block ───────────────────────────────────────────────── + +describe('federatedCallSentinel', () => { + beforeEach(() => { + sqlite = new Database(':memory:'); + testDb = drizzle(sqlite, { schema }); + applyMigrations(sqlite); + }); + + afterEach(() => { + sqlite.close(); + }); + + it('early-exits when there are no federated calls (no DB work, no eviction calls)', async () => { + const { connectionManager } = await import('../ws/handler.js'); + vi.mocked(connectionManager.getAllFederatedCalls).mockReturnValue(new Map()); + vi.mocked(connectionManager.evictFederatedCallsForHost).mockClear(); + + const { runFederatedCallSentinelTick } = await import('./federationWorker.js'); + await runFederatedCallSentinelTick(); + + expect(connectionManager.evictFederatedCallsForHost).not.toHaveBeenCalled(); + }); + + it('evicts entries whose host peer is non-active', async () => { + seedSentinelPeer('p-A', 'https://hostA.example', 'unreachable', 'HostA'); + seedSentinelPeer('p-B', 'https://hostB.example', 'active', 'HostB'); + + const calls = new Map([ + ['fed-1', makeFedCall({ federatedId: 'fed-1', federatedCallHost: 'https://hostA.example' })], + ['fed-2', makeFedCall({ federatedId: 'fed-2', federatedCallHost: 'https://hostB.example' })], + ]); + + const { connectionManager } = await import('../ws/handler.js'); + vi.mocked(connectionManager.getAllFederatedCalls).mockReturnValue(calls); + vi.mocked(connectionManager.evictFederatedCallsForHost).mockClear(); + + const { runFederatedCallSentinelTick } = await import('./federationWorker.js'); + await runFederatedCallSentinelTick(); + + expect(connectionManager.evictFederatedCallsForHost).toHaveBeenCalledTimes(1); + expect(connectionManager.evictFederatedCallsForHost).toHaveBeenCalledWith( + 'https://hostA.example', + { reason: 'peer_transient_failure', peerLabel: 'HostA' }, + ); + }); + + it('maps rejected/revoked to peer_rejected', async () => { + seedSentinelPeer('p-rej', 'https://hostRej.example', 'rejected', 'Rej'); + seedSentinelPeer('p-rev', 'https://hostRev.example', 'revoked', 'Rev'); + + const calls = new Map([ + ['fed-rej', makeFedCall({ federatedId: 'fed-rej', federatedCallHost: 'https://hostRej.example' })], + ['fed-rev', makeFedCall({ federatedId: 'fed-rev', federatedCallHost: 'https://hostRev.example' })], + ]); + + const { connectionManager } = await import('../ws/handler.js'); + vi.mocked(connectionManager.getAllFederatedCalls).mockReturnValue(calls); + vi.mocked(connectionManager.evictFederatedCallsForHost).mockClear(); + + const { runFederatedCallSentinelTick } = await import('./federationWorker.js'); + await runFederatedCallSentinelTick(); + + expect(connectionManager.evictFederatedCallsForHost).toHaveBeenCalledWith( + 'https://hostRej.example', + { reason: 'peer_rejected', peerLabel: 'Rej' }, + ); + expect(connectionManager.evictFederatedCallsForHost).toHaveBeenCalledWith( + 'https://hostRev.example', + { reason: 'peer_rejected', peerLabel: 'Rev' }, + ); + }); + + it('treats a missing peer row as transient failure', async () => { + const calls = new Map([ + ['fed-missing', makeFedCall({ federatedId: 'fed-missing', federatedCallHost: 'https://ghost.example' })], + ]); + const { connectionManager } = await import('../ws/handler.js'); + vi.mocked(connectionManager.getAllFederatedCalls).mockReturnValue(calls); + vi.mocked(connectionManager.evictFederatedCallsForHost).mockClear(); + + const { runFederatedCallSentinelTick } = await import('./federationWorker.js'); + await runFederatedCallSentinelTick(); + + expect(connectionManager.evictFederatedCallsForHost).toHaveBeenCalledWith( + 'https://ghost.example', + { reason: 'peer_transient_failure', peerLabel: undefined }, + ); + }); +}); diff --git a/packages/server/src/utils/federationWorker.ts b/packages/server/src/utils/federationWorker.ts index 8d3ad240..3ddbf3bf 100644 --- a/packages/server/src/utils/federationWorker.ts +++ b/packages/server/src/utils/federationWorker.ts @@ -1109,6 +1109,54 @@ async function processHealthCheckTick(): Promise { } } +// ─── Federated Call Health Sweep ──────────────────────────────────────────── +// +// Periodic backstop: iterate active FederatedCallEntry objects, look up each +// distinct host's peer status, and evict entries whose host is non-active. +// Covers the gap where a peer transitioned to non-active outside of any hook +// site, or was already non-active when the entry was created. +// +// Latency note: real eviction = peer-status-update-lag + tick-period (≤30s). +// Worst case 15.5min for idle instances with no outbox traffic (health-check +// worker is the only status source). Documented in the design spec. + +export const FEDERATED_CALL_SENTINEL_MS = 30_000; + +export async function runFederatedCallSentinelTick(): Promise { + const calls = connectionManager.getAllFederatedCalls(); + if (calls.size === 0) return; + + const distinctHosts = new Set(); + for (const entry of calls.values()) { + distinctHosts.add(entry.federatedCallHost); + } + + const db = getDb(); + for (const host of distinctHosts) { + const row = db + .select({ + status: schema.federationPeers.status, + instanceName: schema.federationPeers.instanceName, + }) + .from(schema.federationPeers) + .where(eq(schema.federationPeers.origin, host)) + .get(); + + if (row && row.status === 'active') continue; + + const isRejectedLike = row?.status === 'rejected' || row?.status === 'revoked'; + const reason: 'peer_rejected' | 'peer_transient_failure' = + isRejectedLike ? 'peer_rejected' : 'peer_transient_failure'; + + connectionManager.evictFederatedCallsForHost(host, { + reason, + peerLabel: row?.instanceName ?? undefined, + }); + } +} + +let federatedCallSentinelTimer: ReturnType | null = null; + // ─── Janitor Worker ────────────────────────────────────────────────────────── function scheduleJanitorTick(): void { @@ -1126,6 +1174,11 @@ export function startFederationWorkers(): void { scheduleFileQueueTick(); scheduleHealthCheckTick(); scheduleJanitorTick(); + federatedCallSentinelTimer = setInterval(() => { + runFederatedCallSentinelTick().catch(err => + console.error('[federation-worker] federatedCallSentinel tick failed:', err) + ); + }, FEDERATED_CALL_SENTINEL_MS); // Bootstrap sync for freshly-peered rows (async, non-blocking) startupBootstrapSync().catch((err) => { console.error('[federation-worker] Startup bootstrap sync error:', err); @@ -1151,6 +1204,11 @@ export function stopFederationWorkers(): void { janitorTimer = null; } + if (federatedCallSentinelTimer) { + clearInterval(federatedCallSentinelTimer); + federatedCallSentinelTimer = null; + } + outboxAbortController?.abort(); outboxAbortController = null;