Merge branch 'feat/remote-participant-host-unreachable'

This commit is contained in:
Jannis Braun
2026-04-24 01:09:35 +02:00
13 changed files with 715 additions and 7 deletions
+15
View File
@@ -1015,6 +1015,21 @@ HTTP handler sites dispatch fire-and-forget (`.catch(log)`) so the response is n
Concurrent activations for the same peer are deduplicated via an in-flight promise map keyed by `peerId`.
### onPeerDeactivated
Mirror of `onPeerActivated` for the transition *out* of `active`. Invoked wherever `federation_peers.status` is written to `unreachable`, `needs_attention`, `rejected`, or `revoked`. Responsibility: sweep `ConnectionManager.federatedCalls` for entries whose `federatedCallHost` matches the deactivated peer and evict them — emitting `dm_call_undeliverable { phase: 'host_unreachable', terminal: true }` to each entry's `ringedUserIds`. See `docs/systems/voice.md` for the client teardown contract.
**Call sites (exhaustive — grep `onPeerDeactivated(` to audit):**
- `utils/federationWorker.ts` handleOutboxDeliveryFailure when status flips to `unreachable`
- `utils/federationWorker.ts` auth-failure path when status flips to `needs_attention`
- `utils/federationWorker.ts` resolvePendingPeers case `'rejected'`
- `routes/federation.ts` admin revoke endpoint
- `utils/federationPeering.ts` performHandshake 403 `PEERING_REQUIRES_APPROVAL` path
Deduplicated by peerId using a **separate** `inFlightDeactivation` map (not shared with activation) so flapping peers retain clean activate-then-deactivate ordering.
A 30s periodic sentinel in `federationWorker.ts` (`runFederatedCallSentinelTick`) is the backstop — it scans active FederatedCallEntries, compares each host's current peer status against reality, and catches transitions missed by the hook sites.
#### Peer-state × outbox-enqueue × recovery matrix
| Status | `queueOutboxEvent` enqueue | Mutation log captures | Recovery on transition to `active` |
+13 -1
View File
@@ -70,6 +70,7 @@ All `dm_call_*` signaling events (`start`, `accept`, `reject`, `end`) are relaye
| `accept` | false | Host → peer fan-out of accept failed; local host call continues. | No state change; info toast. |
| `reject` | false | Rejector's relay to host failed OR host's fan-out after a local reject failed; state already cleared. | No state change; info toast. |
| `end` | false | Ender's relay to host failed OR host's fan-out after a local end failed; state already cleared. | No state change; info toast. |
| `host_unreachable` | true | A FederatedCallEntry's `federatedCallHost` peer transitions out of `active`, OR the 30s sentinel detects a non-active host for an existing entry. | Clear `activeDmCall` + `incomingCall`, disconnect LK, warning toast (*"Call ended — {label} became unreachable."*). |
**Accept-rollback semantics.** `handleDmCallAccept` Path 2 transitions the `FederatedCallEntry` to active and broadcasts `dm_call_accepted` optimistically so the acceptor's UI flips immediately. If the B→host relay fails, the server clears the entry, fans `dm_call_undeliverable { phase: 'accept', terminal: true }` out to all ringed users on B (via `sendToFederatedCallUsers`), and the client tears its call state back down.
@@ -77,7 +78,18 @@ All `dm_call_*` signaling events (`start`, `accept`, `reject`, `end`) are relaye
**Ring-timeout fan-out.** When the host's 60 s ringing timeout fires without an accept, `dm_call_end` is fanned out to all remote peers so stranded Path-A/B ringees on other instances exit their ring state instead of lingering. Registered via `connectionManager.setRingTimeoutFanoutHook` from the WS events module.
**Remaining edge.** When a non-host participant (Bob on B) ends an active call and the relay to the host (Alice on A) fails, Alice's `activeDmCall` marker lingers until she manually ends — LK `ParticipantDisconnected` tears down her voice UI but does not clear the DM-call marker. This is a host-side cleanup concern, tracked separately.
**Remaining edge.** When a non-host participant ends an active call and the relay to the host fails, the host's `activeDmCall` marker lingers until manual end — LK `ParticipantDisconnected` tears down the voice UI but does not clear the DM-call marker on the host side. This is the caller-side mirror of the remote-participant problem and is not covered by the Remote-Participant Host Unreachable Eviction mechanism above (which only reasons about FederatedCallEntry state). Tracked separately.
### Remote-Participant Host Unreachable Eviction
When a FederatedCallEntry's `federatedCallHost` becomes unreachable (peer status transitions to `unreachable`, `needs_attention`, `rejected`, or `revoked`), the entry owner evicts the stranded state and notifies its local ringed users with `dm_call_undeliverable { phase: 'host_unreachable', terminal: true }`. Two signals drive the eviction:
1. **Fast path (`onPeerDeactivated` hook):** every peer-status transition out of `active` invokes `ConnectionManager.evictFederatedCallsForHost(peerOrigin, ...)`. Call sites are listed in the `onPeerDeactivated` docstring (audit via `grep onPeerDeactivated(`).
2. **Backstop (30s sentinel):** `runFederatedCallSentinelTick` in `federationWorker.ts` iterates active entries, looks up each distinct `federatedCallHost`'s current peer status, and evicts non-active matches.
Typical eviction latency is ~90s (time for outbox traffic to fail the unreachable threshold + one sentinel tick). Worst case on an idle instance with no outbox traffic is ~15.5min (health-check cadence + sentinel).
Covers the ringing and active states on the remote-participant side. The caller-side mirror — host's own `activeDmCall` lingering when its LK room empties silently — is a separate, documented out-of-scope edge.
### Dual-Path Processing
+5 -1
View File
@@ -18,7 +18,7 @@ import { sanitizeUser } from '../utils/sanitize.js';
import { deleteAttachmentFiles, deleteUploadFile } from '../utils/fileCleanup.js';
import { tombstoneUser, collectDeletionBroadcastTargets, collectProfileBroadcastTargetIds } from '../utils/userDeletion.js';
import { computeFederatedId, getDmParticipants, sendCallRelay } from '../utils/federationOutbox.js';
import { onPeerActivated } from '../utils/federationPeerActivation.js';
import { onPeerActivated, onPeerDeactivated } from '../utils/federationPeerActivation.js';
import { getDmMessageWithUser } from './dm.js';
import type { FederationRelayRequest, FederationRelayResponse, FederationRelayEvent, FederationRelayAttachment, FederationSyncRequest, FederationSyncResponse, DmMessageWithUser, DmChannel, FederationRelayProfileSnapshot, FederationIdentityDeleteS2SRequest, FederationProfileUpdatePayload, ServerEvent } from '@backspace/shared';
@@ -883,6 +883,10 @@ export async function federationRoutes(app: FastifyInstance): Promise<void> {
.where(eq(schema.federationPeers.id, id))
.run();
onPeerDeactivated(id, 'admin_revoked').catch(err =>
console.error('[federation] onPeerDeactivated from admin revoke failed:', err),
);
// Delete all outbox entries for this peer
db.delete(schema.federationOutbox)
.where(eq(schema.federationOutbox.peerId, id))
@@ -41,6 +41,7 @@ vi.mock('../ws/handler.js', () => ({
getAllOnlineUserIds: () => [],
sendToUser: vi.fn(),
sendToDmMembers: vi.fn(),
evictFederatedCallsForHost: vi.fn().mockReturnValue(0),
},
}));
@@ -348,3 +349,106 @@ describe('onPeerActivated', () => {
await expect(onPeerActivated('peer-err', 'ensure_peered')).resolves.toBeUndefined();
});
});
describe('onPeerDeactivated', () => {
beforeEach(() => {
sqlite = new Database(':memory:');
testDb = drizzle(sqlite, { schema });
applyMigrations(sqlite);
vi.clearAllMocks();
});
it('evicts federated calls for the peer origin with peer_transient_failure on network_threshold', async () => {
seedPeer('peer-net', 'unreachable');
testDb.update(schema.federationPeers)
.set({ instanceName: 'NetPeer' })
.where(eq(schema.federationPeers.id, 'peer-net'))
.run();
const { onPeerDeactivated } = await import('./federationPeerActivation.js');
await onPeerDeactivated('peer-net', 'network_threshold');
const { connectionManager } = await import('../ws/handler.js');
expect(connectionManager.evictFederatedCallsForHost).toHaveBeenCalledWith(
'https://peer-net.example',
{ reason: 'peer_transient_failure', peerLabel: 'NetPeer' },
);
});
it('maps rejected status to peer_rejected reason', async () => {
seedPeer('peer-rej', 'rejected');
const { onPeerDeactivated } = await import('./federationPeerActivation.js');
await onPeerDeactivated('peer-rej', 'remote_rejected');
const { connectionManager } = await import('../ws/handler.js');
expect(connectionManager.evictFederatedCallsForHost).toHaveBeenCalledWith(
'https://peer-rej.example',
{ reason: 'peer_rejected', peerLabel: undefined },
);
});
it('maps revoked status to peer_rejected reason', async () => {
seedPeer('peer-rev', 'revoked');
const { onPeerDeactivated } = await import('./federationPeerActivation.js');
await onPeerDeactivated('peer-rev', 'admin_revoked');
const { connectionManager } = await import('../ws/handler.js');
expect(connectionManager.evictFederatedCallsForHost).toHaveBeenCalledWith(
'https://peer-rev.example',
{ reason: 'peer_rejected', peerLabel: undefined },
);
});
it('broadcasts federation_peers_changed to admins', async () => {
seedPeer('peer-broad', 'unreachable');
const { onPeerDeactivated } = await import('./federationPeerActivation.js');
await onPeerDeactivated('peer-broad', 'network_threshold');
const { connectionManager } = await import('../ws/handler.js');
expect(connectionManager.sendToAdmins).toHaveBeenCalledWith({ type: 'federation_peers_changed' });
});
it('aborts silently when the peer row is missing', async () => {
const { onPeerDeactivated } = await import('./federationPeerActivation.js');
await expect(onPeerDeactivated('peer-missing', 'network_threshold')).resolves.toBeUndefined();
const { connectionManager } = await import('../ws/handler.js');
expect(connectionManager.evictFederatedCallsForHost).not.toHaveBeenCalled();
});
it('deduplicates concurrent calls for the same peerId', async () => {
seedPeer('peer-x', 'unreachable');
const { onPeerDeactivated } = await import('./federationPeerActivation.js');
const { connectionManager } = await import('../ws/handler.js');
const p1 = onPeerDeactivated('peer-x', 'network_threshold');
const p2 = onPeerDeactivated('peer-x', 'network_threshold');
await Promise.all([p1, p2]);
// Exactly one eviction call, not two
expect(connectionManager.evictFederatedCallsForHost).toHaveBeenCalledTimes(1);
});
it('uses a dedup map SEPARATE from onPeerActivated', async () => {
seedPeer('peer-flap', 'active');
const { onPeerActivated, onPeerDeactivated } = await import('./federationPeerActivation.js');
// Simulate an activation already in flight — spawn onPeerActivated then
// immediately kick off a deactivation for the same peer id. The latter
// must NOT be swallowed as a dedup hit against the activation.
const actPromise = onPeerActivated('peer-flap', 'ensure_peered');
// Mark peer non-active now — otherwise the deactivation utility's
// own status guard would skip it.
testDb.update(schema.federationPeers)
.set({ status: 'unreachable' })
.where(eq(schema.federationPeers.id, 'peer-flap'))
.run();
const deactPromise = onPeerDeactivated('peer-flap', 'network_threshold');
await Promise.all([actPromise, deactPromise]);
const { connectionManager } = await import('../ws/handler.js');
expect(connectionManager.evictFederatedCallsForHost).toHaveBeenCalled();
});
});
@@ -191,3 +191,87 @@ export async function startupBootstrapSync(): Promise<void> {
await onPeerActivated(peer.id, 'startup_bootstrap');
}
}
export type PeerDeactivationReason =
| 'network_threshold' // outbox worker hit PEER_UNREACHABLE_THRESHOLD
| 'auth_threshold' // outbox worker hit AUTH_FAILURE_THRESHOLD
| 'remote_rejected' // auto-peer handshake got 403 PEERING_REQUIRES_APPROVAL
| 'admin_revoked'; // admin revoked peering from this side
// Dedup: concurrent deactivations for the same peerId share one promise.
// SEPARATE from inFlightActivation — a flapping peer's activate-then-deactivate
// sequence must not collapse into one slot.
const inFlightDeactivation = new Map<string, Promise<void>>();
/**
* Called whenever federation_peers.status transitions OUT OF 'active' for any reason.
* Sweeps connectionManager.federatedCalls for entries whose federatedCallHost matches
* the peer origin, emitting dm_call_undeliverable { phase: 'host_unreachable', terminal: true }
* to stranded ringed users and clearing the entries.
*
* Call sites (must remain exhaustive — grep `onPeerDeactivated(` to audit):
* - utils/federationWorker.ts handleOutboxDeliveryFailure when status flips to 'unreachable'
* - utils/federationWorker.ts auth-failure path when status flips to 'needs_attention'
* - utils/federationWorker.ts resolvePendingPeers case 'rejected'
* - routes/federation.ts admin revoke endpoint
* - routes/federation.ts admin reset endpoint (when it transitions to a non-active status)
* - utils/federationPeering.ts performHandshake 403 PEERING_REQUIRES_APPROVAL path
*
* Deduplicated by peerId — concurrent calls share one promise. Separate map from
* onPeerActivated so flapping peers don't collapse transitions.
*/
export async function onPeerDeactivated(
peerId: string,
reason: PeerDeactivationReason,
): Promise<void> {
const existing = inFlightDeactivation.get(peerId);
if (existing) return existing;
const promise = (async () => {
try {
const db = getDb();
const peer = db.select({
origin: schema.federationPeers.origin,
status: schema.federationPeers.status,
instanceName: schema.federationPeers.instanceName,
})
.from(schema.federationPeers)
.where(eq(schema.federationPeers.id, peerId))
.get();
if (!peer) {
// Peer row gone — nothing to sweep against.
return;
}
const { connectionManager } = await import('../ws/handler.js');
// Map status to user-facing reason.
const isRejectedLike = peer.status === 'rejected' || peer.status === 'revoked';
const mappedReason: 'peer_rejected' | 'peer_transient_failure' =
isRejectedLike ? 'peer_rejected' : 'peer_transient_failure';
const evicted = connectionManager.evictFederatedCallsForHost(peer.origin, {
reason: mappedReason,
peerLabel: peer.instanceName ?? undefined,
});
if (evicted > 0) {
console.log(
`[federation] onPeerDeactivated(${peerId}, ${reason}) evicted ${evicted} FederatedCallEntry object${evicted === 1 ? '' : 's'} for ${peer.origin}`,
);
}
connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const });
} catch (err) {
console.error(`[federation] onPeerDeactivated(${peerId}, ${reason}) failed:`, err);
}
})();
inFlightDeactivation.set(peerId, promise);
try {
await promise;
} finally {
inFlightDeactivation.delete(peerId);
}
}
@@ -4,7 +4,7 @@ import { eq } from 'drizzle-orm';
import { generateSnowflake } from './snowflake.js';
import { getOurOrigin, generateHmacSecret } from './federationAuth.js';
import { validateOrigin } from '../routes/federation.js';
import { onPeerActivated } from './federationPeerActivation.js';
import { onPeerActivated, onPeerDeactivated } from './federationPeerActivation.js';
// ─── Types ───────────────────────────────────────────────────────────────────
@@ -187,6 +187,9 @@ async function performHandshake(
.run();
const { connectionManager } = await import('../ws/handler.js');
connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const });
onPeerDeactivated(peerId, 'remote_rejected').catch(err =>
console.error('[federation] onPeerDeactivated from performHandshake rejected failed:', err)
);
return { status: 'rejected', error: errorMessage };
}
@@ -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>): 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 },
);
});
});
+72 -1
View File
@@ -11,7 +11,7 @@ import { getDmMessageWithUser } from '../routes/dm.js';
import { connectionManager } from '../ws/handler.js';
import { generateThumbnail } from './thumbnail.js';
import type { FederationRelayRequest, FederationRelayResponse, FederationRelayEvent } from '@backspace/shared';
import { onPeerActivated, startupBootstrapSync } from './federationPeerActivation.js';
import { onPeerActivated, startupBootstrapSync, onPeerDeactivated } from './federationPeerActivation.js';
import fs from 'node:fs';
import path from 'node:path';
import crypto from 'node:crypto';
@@ -310,6 +310,9 @@ export async function processOutboxTick(): Promise<void> {
})
.where(eq(schema.federationPeers.id, peerId))
.run();
onPeerDeactivated(peerId, 'auth_threshold').catch(err =>
console.error('[federation-worker] onPeerDeactivated from auth threshold failed:', err)
);
console.warn(
`[federation-worker] Peer ${peerOrigin} transitioned to needs_attention after ${decision.newAuthFailures} consecutive ${response.status} responses`,
);
@@ -416,6 +419,12 @@ function handleOutboxDeliveryFailure(
.set(updates)
.where(eq(schema.federationPeers.id, peerId))
.run();
if (newFailures >= PEER_UNREACHABLE_THRESHOLD) {
onPeerDeactivated(peerId, 'network_threshold').catch(err =>
console.error('[federation-worker] onPeerDeactivated from unreachable threshold failed:', err)
);
}
}
// ─── Pending Peer Resolution ───────────────────────────────────────────────
@@ -468,6 +477,10 @@ async function resolvePendingPeers(): Promise<void> {
.where(eq(schema.federationOutbox.peerId, peerId))
.run();
onPeerDeactivated(peerId, 'remote_rejected').catch(err =>
console.error('[federation-worker] onPeerDeactivated from resolvePendingPeers rejected failed:', err)
);
// Push federation_peer_rejected WS event to affected users
pushPeerRejectedEvent(peerOrigin, contextMap);
// Notify admins of peer state change
@@ -1096,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 ──────────────────────────────────────────────────────────
function scheduleJanitorTick(): void {
@@ -1113,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);
@@ -1138,6 +1204,11 @@ export function stopFederationWorkers(): void {
janitorTimer = null;
}
if (federatedCallSentinelTimer) {
clearInterval(federatedCallSentinelTimer);
federatedCallSentinelTimer = null;
}
outboxAbortController?.abort();
outboxAbortController = null;
@@ -0,0 +1,196 @@
import { describe, it, expect, vi, beforeEach, afterEach } 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<typeof drizzle<typeof schema>>;
let testDb: TestDb;
vi.mock('../db/index.js', () => ({
getDb: () => testDb,
schema,
}));
vi.mock('../utils/federationAuth.js', () => ({
getOurOrigin: () => 'https://local.example',
}));
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');
const statements = sql.split(/-->\s*statement-breakpoint/);
for (const stmt of statements) {
const clean = stmt.trim();
if (clean) db.exec(clean);
}
}
}
async function importManager() {
const mod = await import('./handler.js');
return mod.connectionManager;
}
type FedCallEntry = import('./handler.js').FederatedCallEntry;
function makeFedCall(partial: Partial<FedCallEntry> = {}): FedCallEntry {
return {
dmChannelId: 'dm-1',
federatedId: `fed-${Math.random().toString(36).slice(2, 10)}`,
callerId: 'caller-user',
callerHomeUserId: 'caller-home',
federatedCallHost: 'https://hostA.example',
livekitUrl: 'wss://lk.example',
tokens: new Map([['caller-home', 'tok']]),
ringedUserIds: ['user-1', 'user-2'],
state: 'active',
startedAt: Date.now(),
...partial,
};
}
let sqlite: Database.Database;
describe('ConnectionManager.evictFederatedCallsForHost', () => {
beforeEach(async () => {
sqlite = new Database(':memory:');
testDb = drizzle(sqlite, { schema });
applyMigrations(sqlite);
const cm = await importManager();
// Reset federatedCalls between tests — use the public API.
for (const [fedId] of cm.getAllFederatedCalls()) {
cm.clearFederatedCall(fedId);
}
vi.useFakeTimers();
});
afterEach(() => {
vi.restoreAllMocks();
vi.useRealTimers();
sqlite.close();
});
it('emits dm_call_undeliverable with host_unreachable to every ringed user and clears entries', async () => {
const cm = await importManager();
const sendSpy = vi.spyOn(cm, 'sendToUser').mockImplementation(() => undefined);
const ringing = makeFedCall({
federatedId: 'fed-ringing',
state: 'ringing',
ringedUserIds: ['alice', 'bob'],
federatedCallHost: 'https://hostA.example',
});
const active = makeFedCall({
federatedId: 'fed-active',
state: 'active',
ringedUserIds: ['carol'],
federatedCallHost: 'https://hostA.example',
});
cm.createFederatedCall(ringing);
cm.createFederatedCall(active);
const count = cm.evictFederatedCallsForHost('https://hostA.example', {
reason: 'peer_transient_failure',
peerLabel: 'Host A',
});
expect(count).toBe(2);
// 3 users total × 1 event each
expect(sendSpy).toHaveBeenCalledTimes(3);
const userIds = sendSpy.mock.calls.map(c => c[0]);
expect(userIds.sort()).toEqual(['alice', 'bob', 'carol']);
for (const call of sendSpy.mock.calls) {
const ev = call[1] as Record<string, unknown>;
expect(ev.type).toBe('dm_call_undeliverable');
expect(ev.phase).toBe('host_unreachable');
expect(ev.terminal).toBe(true);
const failures = ev.failures as Array<{ reason: string; peerOrigin: string; peerLabel?: string }>;
expect(failures).toHaveLength(1);
expect(failures[0]).toMatchObject({
reason: 'peer_transient_failure',
peerOrigin: 'https://hostA.example',
peerLabel: 'Host A',
});
}
expect(cm.getFederatedCall('fed-ringing')).toBeUndefined();
expect(cm.getFederatedCall('fed-active')).toBeUndefined();
});
it('leaves entries pointing at a different host untouched', async () => {
const cm = await importManager();
vi.spyOn(cm, 'sendToUser').mockImplementation(() => undefined);
cm.createFederatedCall(makeFedCall({
federatedId: 'fed-A',
federatedCallHost: 'https://hostA.example',
}));
cm.createFederatedCall(makeFedCall({
federatedId: 'fed-B',
federatedCallHost: 'https://hostB.example',
}));
const count = cm.evictFederatedCallsForHost('https://hostA.example', {
reason: 'peer_rejected',
});
expect(count).toBe(1);
expect(cm.getFederatedCall('fed-A')).toBeUndefined();
expect(cm.getFederatedCall('fed-B')).toBeDefined();
});
it('is idempotent — second call for the same host returns 0 and broadcasts nothing new', async () => {
const cm = await importManager();
const sendSpy = vi.spyOn(cm, 'sendToUser').mockImplementation(() => undefined);
cm.createFederatedCall(makeFedCall({
federatedId: 'fed-once',
ringedUserIds: ['u-1'],
federatedCallHost: 'https://hostA.example',
}));
const first = cm.evictFederatedCallsForHost('https://hostA.example', {
reason: 'peer_transient_failure',
});
const second = cm.evictFederatedCallsForHost('https://hostA.example', {
reason: 'peer_transient_failure',
});
expect(first).toBe(1);
expect(second).toBe(0);
expect(sendSpy).toHaveBeenCalledTimes(1);
});
it('cancels the 60s ring timer on eviction so no late dm_call_ended fires', async () => {
const cm = await importManager();
const sendSpy = vi.spyOn(cm, 'sendToUser').mockImplementation(() => undefined);
cm.createFederatedCall(makeFedCall({
federatedId: 'fed-ringing',
state: 'ringing',
ringedUserIds: ['u-ring'],
federatedCallHost: 'https://hostA.example',
}));
cm.evictFederatedCallsForHost('https://hostA.example', { reason: 'peer_transient_failure' });
// Advance past the 60s ring-timeout; if the timer is still armed we'd see
// a late 'dm_call_ended' broadcast.
vi.advanceTimersByTime(61_000);
const dmCallEndedCalls = sendSpy.mock.calls.filter(c => {
const ev = c[1] as Record<string, unknown>;
return ev.type === 'dm_call_ended';
});
expect(dmCallEndedCalls).toHaveLength(0);
});
});
+50
View File
@@ -519,6 +519,56 @@ class ConnectionManager {
}
}
/**
* Evict all FederatedCallEntry objects whose federatedCallHost matches the given peer origin.
* Emits dm_call_undeliverable { phase: 'host_unreachable', terminal: true } to each entry's
* ringedUserIds, then clears the entry (and its 60s ring timer if still armed).
*
* Idempotent: re-invocation with an already-evicted host returns 0.
* Called from onPeerDeactivated (signal 1) and the 30s sentinel (signal 2 / backstop).
*/
evictFederatedCallsForHost(
peerOrigin: string,
ctx: {
reason: 'peer_transient_failure' | 'peer_rejected';
peerLabel?: string;
},
): number {
const matches: FederatedCallEntry[] = [];
for (const entry of this.federatedCalls.values()) {
if (entry.federatedCallHost === peerOrigin) matches.push(entry);
}
if (matches.length === 0) return 0;
let evicted = 0;
for (const entry of matches) {
// Re-check — concurrent teardown may have removed it between collect and broadcast.
if (!this.federatedCalls.has(entry.federatedId)) continue;
const event: ServerEvent = {
type: 'dm_call_undeliverable',
dmChannelId: entry.dmChannelId,
federatedCallId: entry.federatedId,
terminal: true,
phase: 'host_unreachable',
failures: [{
reason: ctx.reason,
peerOrigin,
peerLabel: ctx.peerLabel,
}],
};
for (const uid of entry.ringedUserIds) {
this.sendToUser(uid, event);
}
this.clearFederatedCall(entry.federatedId);
evicted += 1;
}
return evicted;
}
/** Late-bind a dmChannelId onto a Path B FederatedCallEntry. */
lateBindFederatedCall(federatedId: string, dmChannelId: string): void {
const call = this.federatedCalls.get(federatedId);
+1 -1
View File
@@ -364,7 +364,7 @@ export type DmCallUndeliverableReason =
| 'peer_transient_failure'
| 'livekit_unavailable';
export type DmCallPhase = 'start' | 'accept' | 'reject' | 'end';
export type DmCallPhase = 'start' | 'accept' | 'reject' | 'end' | 'host_unreachable';
export interface DmCallUndeliverableFailure {
reason: DmCallUndeliverableReason;
@@ -34,4 +34,24 @@ describe('buildCallUndeliverableToast', () => {
it('legacy two-arg signature still works', () => {
expect(buildCallUndeliverableToast([fail()], true)).toMatch(/Could not reach nova/);
});
it('builds warning copy for host_unreachable (peer_transient_failure)', () => {
const msg = buildCallUndeliverableToast(
[{ reason: 'peer_transient_failure', peerOrigin: 'https://orbit.local', peerLabel: 'Orbit' }],
true,
'host_unreachable',
);
expect(msg.toLowerCase()).toContain('orbit');
expect(msg.toLowerCase()).toContain('unreachable');
});
it('builds warning copy for host_unreachable (peer_rejected)', () => {
const msg = buildCallUndeliverableToast(
[{ reason: 'peer_rejected', peerOrigin: 'https://orbit.local', peerLabel: 'Orbit' }],
true,
'host_unreachable',
);
expect(msg.toLowerCase()).toContain('orbit');
expect(msg.toLowerCase()).toContain('peered');
});
});
@@ -9,6 +9,9 @@
* - `reject`: the rejector's relay to the host failed; state was already cleared
* locally, so non-terminal info toast only.
* - `end`: the ender's relay to the host failed; state was already cleared locally.
* - `host_unreachable`: the call was terminated by the sentinel worker because the
* host peer became permanently unreachable. Always terminal. A single failure entry
* is expected; multiple fall back to a generic line.
*
* Extracted from `useWebSocket.ts` so it can be unit-tested without pulling in
* the full WS handler graph (livekit / audio deps).
@@ -16,7 +19,7 @@
export function buildCallUndeliverableToast(
failures: Array<{ reason: string; peerOrigin?: string; peerLabel?: string }>,
terminal: boolean,
phase: 'start' | 'accept' | 'reject' | 'end' = 'start',
phase: 'start' | 'accept' | 'reject' | 'end' | 'host_unreachable' = 'start',
): string {
const primary = failures[0];
const labelFor = (f: { peerLabel?: string; peerOrigin?: string }) =>
@@ -37,6 +40,21 @@ export function buildCallUndeliverableToast(
return `Couldn't notify ${labels} that you hung up. Remote participants may see the call for up to 60 seconds.`;
}
// host_unreachable: call terminated because the host peer became unreachable.
// Terminal is always true in this phase. A single failure entry is expected;
// zero or multiple fall back to a generic line.
if (phase === 'host_unreachable') {
const [f] = failures;
if (!f || failures.length !== 1) {
return 'Call ended — host instance became unreachable.';
}
const label = f.peerLabel || f.peerOrigin?.replace(/^https?:\/\//, '') || 'the host instance';
if (f.reason === 'peer_rejected') {
return `Call ended — this instance is no longer peered with ${label}.`;
}
return `Call ended — ${label} became unreachable.`;
}
// phase === 'start' (default + legacy)
if (!terminal) {
const labels = failures.map(labelFor).join(', ');