From 0a8949fbfbc2c96d623dba0fca34ba0b754d5435 Mon Sep 17 00:00:00 2001 From: Jannis Braun <151788261+TheZwiss@users.noreply.github.com> Date: Fri, 24 Apr 2026 00:40:42 +0200 Subject: [PATCH] feat(server): onPeerDeactivated utility mirrors onPeerActivated (TDD) --- .../utils/federationPeerActivation.test.ts | 104 ++++++++++++++++++ .../src/utils/federationPeerActivation.ts | 85 ++++++++++++++ 2 files changed, 189 insertions(+) diff --git a/packages/server/src/utils/federationPeerActivation.test.ts b/packages/server/src/utils/federationPeerActivation.test.ts index c684255c..f8127a37 100644 --- a/packages/server/src/utils/federationPeerActivation.test.ts +++ b/packages/server/src/utils/federationPeerActivation.test.ts @@ -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(); + }); +}); diff --git a/packages/server/src/utils/federationPeerActivation.ts b/packages/server/src/utils/federationPeerActivation.ts index a30f22d9..e2f981f2 100644 --- a/packages/server/src/utils/federationPeerActivation.ts +++ b/packages/server/src/utils/federationPeerActivation.ts @@ -191,3 +191,88 @@ export async function startupBootstrapSync(): Promise { 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 + | 'admin_reset'; // admin reset peer clearing to non-active status + +// 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>(); + +/** + * 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 { + 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); + } +}