feat(federation): wire onPeerActivated into 8 transition sites

Every code location that sets federation_peers.status='active'
now invokes onPeerActivated(peerId, reason). HTTP handler sites
use fire-and-forget (.catch(log)) so the response isn't blocked
by sync-pull pagination. The worker-internal health-check site
awaits the handler since the tick is already async.

Sites: /peer/initiate, /peer/accept (4 branches), /approval-
requests/:id/approve, health check recovery, ensurePeered/
performHandshake.
This commit is contained in:
Jannis Braun
2026-04-22 00:39:58 +02:00
parent 57d7ca66d3
commit 250596c0f6
3 changed files with 26 additions and 1 deletions
+19
View File
@@ -16,6 +16,7 @@ import { sanitizeUser } from '../utils/sanitize.js';
import { deleteAttachmentFiles, deleteUploadFile } from '../utils/fileCleanup.js'; import { deleteAttachmentFiles, deleteUploadFile } from '../utils/fileCleanup.js';
import { tombstoneUser, collectDeletionBroadcastTargets, collectProfileBroadcastTargetIds } from '../utils/userDeletion.js'; import { tombstoneUser, collectDeletionBroadcastTargets, collectProfileBroadcastTargetIds } from '../utils/userDeletion.js';
import { computeFederatedId, getDmParticipants, sendCallRelay } from '../utils/federationOutbox.js'; import { computeFederatedId, getDmParticipants, sendCallRelay } from '../utils/federationOutbox.js';
import { onPeerActivated } from '../utils/federationPeerActivation.js';
import { getDmMessageWithUser } from './dm.js'; import { getDmMessageWithUser } from './dm.js';
import type { FederationRelayRequest, FederationRelayResponse, FederationRelayEvent, FederationRelayAttachment, FederationSyncRequest, FederationSyncResponse, DmMessageWithUser, DmChannel, FederationRelayProfileSnapshot, FederationIdentityDeleteS2SRequest, FederationProfileUpdatePayload, ServerEvent } from '@backspace/shared'; import type { FederationRelayRequest, FederationRelayResponse, FederationRelayEvent, FederationRelayAttachment, FederationSyncRequest, FederationSyncResponse, DmMessageWithUser, DmChannel, FederationRelayProfileSnapshot, FederationIdentityDeleteS2SRequest, FederationProfileUpdatePayload, ServerEvent } from '@backspace/shared';
@@ -354,6 +355,9 @@ export async function federationRoutes(app: FastifyInstance): Promise<void> {
.where(eq(schema.federationPeers.id, peerId)) .where(eq(schema.federationPeers.id, peerId))
.run(); .run();
connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const }); connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const });
onPeerActivated(peerId, 'initiate_accepted').catch(err =>
console.error('[federation] onPeerActivated from /peer/initiate failed:', err)
);
const peer = db const peer = db
.select() .select()
@@ -560,6 +564,9 @@ export async function federationRoutes(app: FastifyInstance): Promise<void> {
} }
connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const }); connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const });
onPeerActivated(existing.id, 'accept_rejected_override').catch(err =>
console.error('[federation] onPeerActivated from /peer/accept (rejected override) failed:', err)
);
return reply.code(200).send({ accepted: true }); return reply.code(200).send({ accepted: true });
} }
@@ -583,6 +590,9 @@ export async function federationRoutes(app: FastifyInstance): Promise<void> {
} }
connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const }); connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const });
onPeerActivated(existing.id, 'accept_awaiting_approval').catch(err =>
console.error('[federation] onPeerActivated from /peer/accept (awaiting_approval) failed:', err)
);
return reply.code(200).send({ accepted: true }); return reply.code(200).send({ accepted: true });
} }
@@ -597,6 +607,9 @@ export async function federationRoutes(app: FastifyInstance): Promise<void> {
.run(); .run();
connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const }); connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const });
onPeerActivated(existing.id, 'accept_pending').catch(err =>
console.error('[federation] onPeerActivated from /peer/accept (pending) failed:', err)
);
return reply.code(200).send({ accepted: true }); return reply.code(200).send({ accepted: true });
} }
@@ -613,6 +626,9 @@ export async function federationRoutes(app: FastifyInstance): Promise<void> {
}).run(); }).run();
connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const }); connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const });
onPeerActivated(peerId, 'accept_new').catch(err =>
console.error('[federation] onPeerActivated from /peer/accept (new) failed:', err)
);
return reply.code(200).send({ accepted: true }); return reply.code(200).send({ accepted: true });
}, },
@@ -1136,6 +1152,9 @@ export async function federationRoutes(app: FastifyInstance): Promise<void> {
.run(); .run();
connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const }); connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const });
onPeerActivated(peerId, 'approval_handshake').catch(err =>
console.error('[federation] onPeerActivated from /approval-requests/:id/approve failed:', err)
);
const peer = db const peer = db
.select() .select()
@@ -4,6 +4,7 @@ import { eq } from 'drizzle-orm';
import { generateSnowflake } from './snowflake.js'; import { generateSnowflake } from './snowflake.js';
import { getOurOrigin, generateHmacSecret } from './federationAuth.js'; import { getOurOrigin, generateHmacSecret } from './federationAuth.js';
import { validateOrigin } from '../routes/federation.js'; import { validateOrigin } from '../routes/federation.js';
import { onPeerActivated } from './federationPeerActivation.js';
// ─── Types ─────────────────────────────────────────────────────────────────── // ─── Types ───────────────────────────────────────────────────────────────────
@@ -158,6 +159,9 @@ async function performHandshake(
.run(); .run();
const { connectionManager } = await import('../ws/handler.js'); const { connectionManager } = await import('../ws/handler.js');
connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const }); connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const });
onPeerActivated(peerId, 'ensure_peered').catch(err =>
console.error('[federation] onPeerActivated from ensurePeered failed:', err)
);
return { status: 'active', peerId }; return { status: 'active', peerId };
} }
@@ -11,7 +11,7 @@ import { getDmMessageWithUser } from '../routes/dm.js';
import { connectionManager } from '../ws/handler.js'; import { connectionManager } from '../ws/handler.js';
import { generateThumbnail } from './thumbnail.js'; import { generateThumbnail } from './thumbnail.js';
import type { FederationRelayRequest, FederationRelayResponse, FederationRelayEvent } from '@backspace/shared'; import type { FederationRelayRequest, FederationRelayResponse, FederationRelayEvent } from '@backspace/shared';
import { startupBootstrapSync } from './federationPeerActivation.js'; import { onPeerActivated, startupBootstrapSync } from './federationPeerActivation.js';
import fs from 'node:fs'; import fs from 'node:fs';
import path from 'node:path'; import path from 'node:path';
import crypto from 'node:crypto'; import crypto from 'node:crypto';
@@ -1054,6 +1054,8 @@ async function processHealthCheckTick(): Promise<void> {
console.log( console.log(
`[federation-worker] Peer ${peer.origin} recovered — marked active`, `[federation-worker] Peer ${peer.origin} recovered — marked active`,
); );
await onPeerActivated(peer.id, 'health_check_recovery');
} }
// If not ok, leave as unreachable — will check again next cycle // If not ok, leave as unreachable — will check again next cycle
} catch (err) { } catch (err) {