feat(server): 30s federated-call sentinel worker (TDD)
This commit is contained in:
@@ -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 Database from 'better-sqlite3';
|
||||||
import { drizzle } from 'drizzle-orm/better-sqlite3';
|
import { drizzle } from 'drizzle-orm/better-sqlite3';
|
||||||
import fs from 'node:fs';
|
import fs from 'node:fs';
|
||||||
@@ -23,6 +23,8 @@ vi.mock('../ws/handler.js', () => ({
|
|||||||
getAllOnlineUserIds: () => [],
|
getAllOnlineUserIds: () => [],
|
||||||
sendToUser: vi.fn(),
|
sendToUser: vi.fn(),
|
||||||
sendToDmMembers: 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', () => ({
|
vi.mock('../utils/federationPeerActivation.js', () => ({
|
||||||
onPeerActivated: vi.fn(),
|
onPeerActivated: vi.fn(),
|
||||||
|
onPeerDeactivated: vi.fn().mockResolvedValue(undefined),
|
||||||
startupBootstrapSync: vi.fn(),
|
startupBootstrapSync: vi.fn(),
|
||||||
}));
|
}));
|
||||||
|
|
||||||
@@ -157,3 +160,131 @@ describe('outbox worker — duplicate rejection is terminal', () => {
|
|||||||
expect(remaining).toBeDefined();
|
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>): 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<string, FedCallEntry>([
|
||||||
|
['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<string, FedCallEntry>([
|
||||||
|
['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<string, FedCallEntry>([
|
||||||
|
['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 },
|
||||||
|
);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|||||||
@@ -1109,6 +1109,54 @@ async function processHealthCheckTick(): Promise<void> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ─── 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<void> {
|
||||||
|
const calls = connectionManager.getAllFederatedCalls();
|
||||||
|
if (calls.size === 0) return;
|
||||||
|
|
||||||
|
const distinctHosts = new Set<string>();
|
||||||
|
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<typeof setInterval> | null = null;
|
||||||
|
|
||||||
// ─── Janitor Worker ──────────────────────────────────────────────────────────
|
// ─── Janitor Worker ──────────────────────────────────────────────────────────
|
||||||
|
|
||||||
function scheduleJanitorTick(): void {
|
function scheduleJanitorTick(): void {
|
||||||
@@ -1126,6 +1174,11 @@ export function startFederationWorkers(): void {
|
|||||||
scheduleFileQueueTick();
|
scheduleFileQueueTick();
|
||||||
scheduleHealthCheckTick();
|
scheduleHealthCheckTick();
|
||||||
scheduleJanitorTick();
|
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)
|
// Bootstrap sync for freshly-peered rows (async, non-blocking)
|
||||||
startupBootstrapSync().catch((err) => {
|
startupBootstrapSync().catch((err) => {
|
||||||
console.error('[federation-worker] Startup bootstrap sync error:', err);
|
console.error('[federation-worker] Startup bootstrap sync error:', err);
|
||||||
@@ -1151,6 +1204,11 @@ export function stopFederationWorkers(): void {
|
|||||||
janitorTimer = null;
|
janitorTimer = null;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (federatedCallSentinelTimer) {
|
||||||
|
clearInterval(federatedCallSentinelTimer);
|
||||||
|
federatedCallSentinelTimer = null;
|
||||||
|
}
|
||||||
|
|
||||||
outboxAbortController?.abort();
|
outboxAbortController?.abort();
|
||||||
outboxAbortController = null;
|
outboxAbortController = null;
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user