diff --git a/docs/systems/federation.md b/docs/systems/federation.md index 7a3a6ea7..45fe5bab 100644 --- a/docs/systems/federation.md +++ b/docs/systems/federation.md @@ -977,16 +977,79 @@ The mutation log entry for reactions stores a simpler payload (no `messageId`/`m ## 12. Initial Sync -### `runInitialSyncForNewPeers()` (`federationWorker.ts:739`) +### `startupBootstrapSync()` (`federationWorker.ts`) -Triggered once at server startup (async, non-blocking). Finds peers with `status = 'active'` and `lastSyncedAt = 0`. +Triggered once at server startup (async, non-blocking). Finds peers with `status = 'active'` and `lastSyncedAt = 0` and calls `onPeerActivated(peerId, 'startup_bootstrap')` for each. This preserves the original startup-sync semantics while unifying the code path with all other activation sites (see "Peer Activation Recovery" below). -**For each unsynced peer:** -1. **DM sync pass:** Paginate through `POST {peerOrigin}/api/federation/sync` with `sinceTimestamp = 0`, `limit = 100` -2. **Direct processing:** Call `processRelayEvents()` to process received events in-process (no HTTP round-trip) -3. **Friend sync pass:** Same pagination with `contextType: 'friend'`, also processed via `processRelayEvents()` -4. Update `lastSyncedAt = Date.now()` after completion -5. On failure: don't update `lastSyncedAt` -- retried on next startup +### Peer Activation Recovery + +Every transition of `federation_peers.status` to `active` invokes `onPeerActivated(peerId, reason)` — one handler wired at all transition sites. Two independent invariants, both unconditional: + +1. **`resetOutboxBackoff`** — sets `nextRetryAt = now` and `attempts = 0` for every outbox entry belonging to the peer. Entries that accumulated exponential backoff before the peer went unreachable are immediately eligible again. Attempts counter is also reset so a freshly-healthy peer's next failure starts at `BACKOFF_SCHEDULE_MS[0]` (30s), not wherever the counter left off. +2. **`syncPeerMutationLog`** — pulls missed events from the peer's `/api/federation/sync` endpoint since `peer.lastSyncedAt`. Three passes: DM, friend, profile (in that order, each paginated). `lastSyncedAt` advances to `Date.now()` on full success; stays put on transient failure so the next activation retries the same window. + +#### Call sites (must remain exhaustive) + +| File | Context | Reason | +|---|---|---| +| `routes/federation.ts` | `/peer/initiate` 200 activation | `initiate_accepted` | +| `routes/federation.ts` | `/peer/accept` existing-rejected override | `accept_rejected_override` | +| `routes/federation.ts` | `/peer/accept` existing-awaiting_approval | `accept_awaiting_approval` | +| `routes/federation.ts` | `/peer/accept` existing-pending | `accept_pending` | +| `routes/federation.ts` | `/peer/accept` new-peer | `accept_new` | +| `routes/federation.ts` | `/approval-requests/:id/approve` success | `approval_handshake` | +| `utils/federationWorker.ts` | Health check unreachable → active | `health_check_recovery` | +| `utils/federationPeering.ts` | `ensurePeered/performHandshake` 200 | `ensure_peered` | +| `utils/federationWorker.ts` | Startup scan (status=active, lastSyncedAt=0) | `startup_bootstrap` | + +HTTP handler sites dispatch fire-and-forget (`.catch(log)`) so the response is not blocked by sync-pull pagination. Worker-internal sites `await` since the worker tick is already async. + +Concurrent activations for the same peer are deduplicated via an in-flight promise map keyed by `peerId`. + +#### Peer-state × outbox-enqueue × recovery matrix + +| Status | `queueOutboxEvent` enqueue | Mutation log captures | Recovery on transition to `active` | +|---|---|---|---| +| `active` | Queue | Yes (for covered event types — see below) | N/A | +| `pending` | Queue | Yes | `onPeerActivated` | +| `unreachable` | Queue | Yes | `onPeerActivated` | +| `awaiting_approval` | **Drop, debug-log** | Yes | `onPeerActivated` | +| `needs_attention` | **Drop, debug-log** | Yes | `onPeerActivated` (fires when the row is re-created via admin Reset + re-peer) | +| `rejected` | Drop, debug-log | Yes | `onPeerActivated` | +| `revoked` | Drop, debug-log | Yes | `onPeerActivated` (fires when the row is re-created via hard-delete + re-initiate) | + +`queueOutboxEvent` uses an exhaustive TypeScript `switch` on the narrowed peer-status union — adding a new status value without handling it fails compile-time typecheck (`const _exhaustive: never = status;`). + +**Mid-call race catch:** when the initial peers SELECT filters out a peer because its status is non-deliverable, but the fallback loop observes the status has since flipped to `active`/`pending`/`unreachable`, the code re-fetches the peer row and appends it to `matchedPeers` so the outer enqueue loop includes it. Silent drops would lose real-time delivery under asymmetric failure (e.g., `/peer/accept` 200 response lost on the wire, health-check transition firing on only one side). + +#### Mutation log coverage + +Event types covered by `appendMutationLog` (replayed on sync-pull): + +| Event type | `contextType` | Source | +|---|---|---| +| DM `create` / `update` / `delete` | `dm` | `federationOutbox.queueDmRelay`, `dm.ts` delete handler | +| `reaction_add` / `reaction_remove` | `dm` | `ws/events.ts` | +| `member_add` / `member_remove` / `ownership_transfer` | `dm` | `dm.ts` | +| `dm_close` / `dm_reopen` | `dm` | `federationOutbox.queueDmCloseRelay` | +| `read_state_update` | `dm` | `federationOutbox.queueReadStateRelay` | +| `file_rejected` | `dm` | `federationWorker.handleSizeRejection` | +| `friend_request_*` / `friend_add` / `friend_remove` | `friend` | `social.ts` | +| `profile_update` | `profile` | `routes/users.ts` (PATCH `/api/users/@me`) | + +Ephemeral events (`dm_typing_*`, `dm_call_*`) are fire-and-forget by design and are NOT captured — missed typing/call-signaling packets are acceptable and carry no durable state. + +#### `/api/federation/sync` contextType filter values + +| Filter | Returns | +|---|---| +| (none) / omitted | DM events (including `dm_close`, `dm_reopen`, `read_state_update`, `file_rejected`) | +| `'friend'` | Friend events | +| `'profile'` | Profile update events | + +#### Known Issues + +- **Poison-pill event in peer's mutation log.** If a peer's mutation log contains a row whose inbound processor throws (e.g., a UNIQUE conflict from a malformed relay payload), `syncPeerMutationLog` catches the error and declines to advance `lastSyncedAt`. Subsequent activations retry the same window and hit the same failure, effectively blocking catch-up for that peer. No automatic poison-pill skip is implemented — recovery requires either: (a) fixing the mutation log on the peer side, or (b) manually advancing `lastSyncedAt` past the offending row via DB admin. Flagged as a follow-up backlog item. ### Sync Endpoint (`POST /api/federation/sync`) @@ -994,7 +1057,7 @@ HMAC-authenticated. Returns events from the `federation_mutation_log`. **Request:** ```typescript -{ sinceTimestamp: number, dmChannelId?: string, federatedId?: string, contextType?: 'dm'|'friend', limit?: 1-500 } +{ sinceTimestamp: number, dmChannelId?: string, federatedId?: string, contextType?: 'dm'|'friend'|'profile', limit?: 1-500 } ``` **Response:** @@ -1125,7 +1188,7 @@ All workers are started by `startFederationWorkers()` on server boot and stopped | File download | 30s | 5 | 60s | `processFileQueueTick` | | Health check | 15min | all unreachable | 10s | `processHealthCheckTick` | | Janitor | 1h | -- | -- | `runFederationJanitor` (sync) | -| Initial sync | Once at startup | -- | 30s per page | `runInitialSyncForNewPeers` | +| Startup bootstrap sync | Once at startup | -- | 30s per page | `startupBootstrapSync` → `onPeerActivated` | ### Janitor Cleanup (`storageJanitor.ts:runFederationJanitor`) diff --git a/docs/systems/social.md b/docs/systems/social.md index 7f44a131..0864c10f 100644 --- a/docs/systems/social.md +++ b/docs/systems/social.md @@ -368,14 +368,16 @@ The sorted join ensures the same pair always produces the same prefix regardless ## 7. Initial Sync: Friend Backfill -When a new peer is established (`federation_peers.lastSyncedAt = 0`), the federation worker runs `runInitialSyncForNewPeers()` on startup. This includes a dedicated friend sync pass. +When a peer transitions to `active` (including at startup for peers with `lastSyncedAt = 0`), the federation worker calls `onPeerActivated(peerId, reason)`. One of its two unconditional invariants is `syncPeerMutationLog`, which pulls missed events from the peer's `/api/federation/sync` endpoint — including a dedicated friend sync pass. -**Flow (`federationWorker.ts:runInitialSyncForNewPeers`):** +**Flow (`federationPeerActivation.ts:syncPeerMutationLog`):** -1. Query all active peers with `lastSyncedAt = 0` -2. **First pass (DM events):** Paginates through `POST /federation/sync` with no `contextType` filter (defaults to DM events), processing each batch via `processRelayEvents()` directly -3. **Second pass (friend events):** Paginates through `POST /federation/sync` with `contextType: 'friend'`, same direct processing -4. After both passes complete, updates `lastSyncedAt = Date.now()` so the sync doesn't repeat +1. **First pass (DM events):** Paginates through `POST /federation/sync` with no `contextType` filter (defaults to DM events), processing each batch via `processRelayEvents()` directly +2. **Second pass (friend events):** Paginates through `POST /federation/sync` with `contextType: 'friend'`, same direct processing +3. **Third pass (profile events):** Paginates through `POST /federation/sync` with `contextType: 'profile'`, same direct processing +4. After all three passes complete, updates `lastSyncedAt = Date.now()` so the window advances on the next activation + +At startup, `startupBootstrapSync()` scans for `status = 'active' AND lastSyncedAt = 0` peers and calls `onPeerActivated(peerId, 'startup_bootstrap')` for each, preserving the original startup-sync semantics while using the unified path. The sync endpoint (`POST /api/federation/sync`) returns events from the `federation_mutation_log` table, which retains entries for 90 days. This means friend relationships established within the last 90 days are backfilled when a new peer connection is created. diff --git a/packages/server/src/routes/federation.ts b/packages/server/src/routes/federation.ts index 4263f95d..9d41f8d6 100644 --- a/packages/server/src/routes/federation.ts +++ b/packages/server/src/routes/federation.ts @@ -16,6 +16,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 { getDmMessageWithUser } from './dm.js'; 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 { .where(eq(schema.federationPeers.id, peerId)) .run(); 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 .select() @@ -560,6 +564,9 @@ export async function federationRoutes(app: FastifyInstance): Promise { } 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 }); } @@ -583,6 +590,9 @@ export async function federationRoutes(app: FastifyInstance): Promise { } 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 }); } @@ -597,6 +607,9 @@ export async function federationRoutes(app: FastifyInstance): Promise { .run(); 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 }); } @@ -613,6 +626,9 @@ export async function federationRoutes(app: FastifyInstance): Promise { }).run(); 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 }); }, @@ -1136,6 +1152,9 @@ export async function federationRoutes(app: FastifyInstance): Promise { .run(); 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 .select() @@ -1615,6 +1634,10 @@ export async function federationRoutes(app: FastifyInstance): Promise { payload: string | null; }>; + // Maps local DM channel ID → federatedId for O(1) lookup in serializers. + // Only populated in the DM branch (friend/profile branches don't need it). + let channelFederatedIdMap = new Map(); + if (contextTypeFilter === 'friend') { // ── Friend event sync: no DM channel logic needed ── mutationRows = rawDb.prepare(` @@ -1624,6 +1647,15 @@ export async function federationRoutes(app: FastifyInstance): Promise { ORDER BY mutated_at ASC LIMIT ? `).all(sinceTimestamp, limit) as typeof mutationRows; + } else if (contextTypeFilter === 'profile') { + // ── Profile event sync: no DM channel logic needed ── + mutationRows = rawDb.prepare(` + SELECT id, entity_id, context_id, context_type, mutation_type, mutated_at, payload + FROM federation_mutation_log + WHERE context_type = 'profile' AND mutated_at > ? + ORDER BY mutated_at ASC + LIMIT ? + `).all(sinceTimestamp, limit) as typeof mutationRows; } else { // ── DM sync path ── // Determine which DM channels to sync. @@ -1631,11 +1663,14 @@ export async function federationRoutes(app: FastifyInstance): Promise { // that should be synced. The peer's relay endpoint will create the channel // if it doesn't exist, or match by federated_id if it does. const sharedChannelRows = rawDb.prepare(` - SELECT id as dm_channel_id FROM dm_channels + SELECT id as dm_channel_id, federated_id FROM dm_channels WHERE federated_id IS NOT NULL AND deleted_at IS NULL - `).all() as Array<{ dm_channel_id: string }>; + `).all() as Array<{ dm_channel_id: string; federated_id: string }>; const sharedChannelIds = sharedChannelRows.map(r => r.dm_channel_id); + channelFederatedIdMap = new Map( + sharedChannelRows.map(r => [r.dm_channel_id, r.federated_id]) + ); // If filtering by federatedId, resolve to local channel ID let effectiveChannelFilter = dmChannelIdFilter; @@ -1679,7 +1714,10 @@ export async function federationRoutes(app: FastifyInstance): Promise { WHERE ml.context_id = ? AND ml.context_type = 'dm' AND ml.mutated_at > ? - AND (dm.id IS NOT NULL OR ml.mutation_type IN ('delete', 'member_add', 'member_remove', 'ownership_transfer')) + AND (dm.id IS NOT NULL OR ml.mutation_type IN ( + 'delete', 'member_add', 'member_remove', 'ownership_transfer', + 'dm_close', 'dm_reopen', 'read_state_update', 'file_rejected' + )) ORDER BY ml.mutated_at ASC LIMIT ? `).all(effectiveChannelFilter, sinceTimestamp, limit) as typeof mutationRows; @@ -1720,7 +1758,10 @@ export async function federationRoutes(app: FastifyInstance): Promise { WHERE ml.context_id IN (${placeholders}) AND ml.context_type = 'dm' AND ml.mutated_at > ? - AND (dm.id IS NOT NULL OR ml.mutation_type IN ('delete', 'member_add', 'member_remove', 'ownership_transfer')) + AND (dm.id IS NOT NULL OR ml.mutation_type IN ( + 'delete', 'member_add', 'member_remove', 'ownership_transfer', + 'dm_close', 'dm_reopen', 'read_state_update', 'file_rejected' + )) ORDER BY ml.mutated_at ASC LIMIT ? `).all(...sharedChannelIds, sinceTimestamp, limit) as typeof mutationRows; @@ -1760,7 +1801,9 @@ export async function federationRoutes(app: FastifyInstance): Promise { const mutationType = mutation.mutation_type as 'create' | 'update' | 'delete' | 'reaction_add' | 'reaction_remove' | 'member_add' | 'member_remove' | 'ownership_transfer' | 'friend_request_create' | 'friend_request_update' | 'friend_request_cancel' - | 'friend_add' | 'friend_remove'; + | 'friend_add' | 'friend_remove' + | 'dm_close' | 'dm_reopen' | 'read_state_update' | 'file_rejected' + | 'profile_update'; if (['member_add', 'member_remove', 'ownership_transfer', 'friend_request_create', 'friend_request_update', 'friend_request_cancel', @@ -1820,6 +1863,86 @@ export async function federationRoutes(app: FastifyInstance): Promise { continue; } + if (mutationType === 'dm_close' || mutationType === 'dm_reopen') { + if (!mutation.payload) continue; + let dmCloseReopenPayload: { homeUserId: string; homeInstance: string } | null = null; + try { dmCloseReopenPayload = JSON.parse(mutation.payload); } catch { continue; } + if (!dmCloseReopenPayload) continue; + const fedIdCloseReopen = channelFederatedIdMap.get(mutation.context_id); + if (!fedIdCloseReopen) continue; + events.push({ + eventType: mutationType, + dmChannelId: mutation.context_id, + messageId: mutation.entity_id, + federatedId: fedIdCloseReopen, + encryptionVersion: 0, + timestamp: mutation.mutated_at, + dmCloseReopen: dmCloseReopenPayload, + }); + continue; + } + + if (mutationType === 'read_state_update') { + if (!mutation.payload) continue; + let readState: NonNullable | null = null; + try { readState = JSON.parse(mutation.payload) as NonNullable; } catch { continue; } + if (!readState) continue; + const fedIdReadState = channelFederatedIdMap.get(mutation.context_id); + if (!fedIdReadState) continue; + events.push({ + eventType: 'read_state_update', + dmChannelId: mutation.context_id, + messageId: mutation.entity_id, + federatedId: fedIdReadState, + encryptionVersion: 0, + timestamp: mutation.mutated_at, + readState, + }); + continue; + } + + if (mutationType === 'file_rejected') { + if (!mutation.payload) continue; + let fileRejectedPayload: { + attachmentId: string; + sourceFilename: string; + rejectionReason: string; + rejectionLimit: number; + affectedUserIds: string[]; + } | null = null; + try { fileRejectedPayload = JSON.parse(mutation.payload); } catch { continue; } + if (!fileRejectedPayload) continue; + events.push({ + eventType: 'file_rejected', + dmChannelId: mutation.context_id, + messageId: mutation.entity_id, + encryptionVersion: 0, + timestamp: mutation.mutated_at, + attachmentId: fileRejectedPayload.attachmentId, + sourceFilename: fileRejectedPayload.sourceFilename, + rejectionReason: fileRejectedPayload.rejectionReason, + rejectionLimit: fileRejectedPayload.rejectionLimit, + affectedUserIds: fileRejectedPayload.affectedUserIds, + }); + continue; + } + + if (mutationType === 'profile_update') { + if (!mutation.payload) continue; + let profileOuter: { profileUpdate?: NonNullable } | null = null; + try { profileOuter = JSON.parse(mutation.payload); } catch { continue; } + if (!profileOuter?.profileUpdate) continue; + events.push({ + eventType: 'profile_update', + contextType: 'profile', + messageId: mutation.entity_id, + encryptionVersion: 0, + timestamp: mutation.mutated_at, + profileUpdate: profileOuter.profileUpdate, + }); + continue; + } + // For create and update: fetch the current message state const message = db .select() diff --git a/packages/server/src/routes/users.ts b/packages/server/src/routes/users.ts index d49ba692..9d82e6a1 100644 --- a/packages/server/src/routes/users.ts +++ b/packages/server/src/routes/users.ts @@ -8,7 +8,7 @@ import { AVATAR_COLORS } from '@backspace/shared'; import { sanitizeUser } from '../utils/sanitize.js'; import { deleteUploadFile, deleteAttachmentByFilename } from '../utils/fileCleanup.js'; import { tombstoneUser, collectDeletionBroadcastTargets, collectProfileBroadcastTargetIds } from '../utils/userDeletion.js'; -import { queueOutboxEvent, isFederationRelayEnabled } from '../utils/federationOutbox.js'; +import { queueOutboxEvent, isFederationRelayEnabled, appendMutationLog } from '../utils/federationOutbox.js'; import { generateSnowflake } from '../utils/snowflake.js'; import { resizeProfileImage } from '../utils/thumbnail.js'; import { config } from '../config.js'; @@ -374,6 +374,13 @@ export async function userRoutes(app: FastifyInstance): Promise { bio: preUpdateUser.bio, }; + appendMutationLog( + preUpdateUser.id, + preUpdateUser.id, + 'profile_update', + JSON.stringify({ profileUpdate: profilePayload }), + 'profile', + ); for (const targetOrigin of newOrigins) { queueOutboxEvent( preUpdateUser.id, @@ -529,6 +536,13 @@ export async function userRoutes(app: FastifyInstance): Promise { bio: updatedUser!.bio, }; + appendMutationLog( + updatedUser!.id, + updatedUser!.id, + 'profile_update', + JSON.stringify({ profileUpdate: profilePayload }), + 'profile', + ); queueOutboxEvent( updatedUser!.id, // entityId — user's ID (coalesces rapid edits) updatedUser!.id, // contextId — user-scoped diff --git a/packages/server/src/utils/federationOutbox.mutationlog.test.ts b/packages/server/src/utils/federationOutbox.mutationlog.test.ts new file mode 100644 index 00000000..19c5242a --- /dev/null +++ b/packages/server/src/utils/federationOutbox.mutationlog.test.ts @@ -0,0 +1,135 @@ +import { describe, it, expect, beforeEach, vi } 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'; +import { eq } from 'drizzle-orm'; + +const __dirname = path.dirname(fileURLToPath(import.meta.url)); +type TestDb = ReturnType>; +let sqlite: Database.Database; +let testDb: TestDb; + +vi.mock('../db/index.js', () => ({ + getDb: () => testDb, + schema, +})); + +vi.mock('../utils/federationAuth.js', () => ({ + getOurOrigin: () => 'https://test.example', + buildFederationHeaders: () => ({}), + generateHmacSecret: () => 'secret', +})); + +let _snowflakeCounter = 1; +vi.mock('../utils/snowflake.js', () => ({ + generateSnowflake: () => String(_snowflakeCounter++), + setWorkerId: vi.fn(), +})); + +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); + } + } +} + +function seedSettings(): void { + testDb.insert(schema.instanceSettings).values({ + id: 1, + federationRelayEnabled: 1, + updatedAt: Date.now(), + }).run(); +} + +function seedChannel(id: string, federatedId: string | null): void { + testDb.insert(schema.dmChannels).values({ + id, federatedId, ownerId: null, createdAt: Date.now(), + }).run(); +} + +function seedUser(id: string, username: string): void { + testDb.insert(schema.users).values({ + id, username, displayName: username, passwordHash: 'x', + createdAt: Date.now(), + }).run(); +} + +function seedDmMember(channelId: string, userId: string): void { + testDb.insert(schema.dmMembers).values({ + dmChannelId: channelId, userId, closed: 0, + }).run(); +} + +describe('queueDmCloseRelay — mutation log capture', () => { + beforeEach(() => { + sqlite = new Database(':memory:'); + testDb = drizzle(sqlite, { schema }); + applyMigrations(sqlite); + seedSettings(); + }); + + it('appends a mutation log row for dm_close', async () => { + const { queueDmCloseRelay } = await import('./federationOutbox.js'); + seedUser('u-1', 'alice'); + seedChannel('ch-1', 'fed-1'); + seedDmMember('ch-1', 'u-1'); + + queueDmCloseRelay('ch-1', 'u-1', 'dm_close'); + + const rows = testDb.select().from(schema.federationMutationLog) + .where(eq(schema.federationMutationLog.mutationType, 'dm_close')).all(); + expect(rows.length).toBe(1); + expect(rows[0]?.contextId).toBe('ch-1'); + expect(rows[0]?.contextType).toBe('dm'); + }); + + it('appends a mutation log row for dm_reopen', async () => { + const { queueDmCloseRelay } = await import('./federationOutbox.js'); + seedUser('u-1', 'alice'); + seedChannel('ch-2', 'fed-2'); + seedDmMember('ch-2', 'u-1'); + + queueDmCloseRelay('ch-2', 'u-1', 'dm_reopen'); + + const rows = testDb.select().from(schema.federationMutationLog) + .where(eq(schema.federationMutationLog.mutationType, 'dm_reopen')).all(); + expect(rows.length).toBe(1); + }); +}); + +describe('queueReadStateRelay — mutation log capture', () => { + beforeEach(() => { + sqlite = new Database(':memory:'); + testDb = drizzle(sqlite, { schema }); + applyMigrations(sqlite); + seedSettings(); + }); + + it('appends a mutation log row for read_state_update', async () => { + const { queueReadStateRelay } = await import('./federationOutbox.js'); + seedUser('u-2', 'bob'); + seedChannel('ch-3', 'fed-3'); + seedDmMember('ch-3', 'u-2'); + testDb.insert(schema.dmMessages).values({ + id: 'm-1', dmChannelId: 'ch-3', userId: 'u-2', content: 'hi', + type: 'user', createdAt: Date.now(), + }).run(); + + queueReadStateRelay('ch-3', 'm-1', 'u-2'); + + const rows = testDb.select().from(schema.federationMutationLog) + .where(eq(schema.federationMutationLog.mutationType, 'read_state_update')).all(); + expect(rows.length).toBe(1); + expect(rows[0]?.contextId).toBe('ch-3'); + expect(rows[0]?.contextType).toBe('dm'); + }); +}); diff --git a/packages/server/src/utils/federationOutbox.state.test.ts b/packages/server/src/utils/federationOutbox.state.test.ts new file mode 100644 index 00000000..68011ffd --- /dev/null +++ b/packages/server/src/utils/federationOutbox.state.test.ts @@ -0,0 +1,92 @@ +import { describe, it, expect, beforeEach, vi } 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'; +import { eq } from 'drizzle-orm'; + +const __dirname = path.dirname(fileURLToPath(import.meta.url)); +type TestDb = ReturnType>; + +// Mutable reference updated in beforeEach — the factory closes over this. +let testDb: TestDb; + +vi.mock('../db/index.js', () => ({ + getDb: () => testDb, + schema, +})); + +// Mock federation-auth helpers to avoid env-var dependency +vi.mock('../utils/federationAuth.js', () => ({ + getOurOrigin: () => 'https://local.example', + buildFederationHeaders: () => ({}), + generateHmacSecret: () => 'test-secret', +})); + +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); + } + } +} + +function seedSettings(): void { + testDb.insert(schema.instanceSettings).values({ + id: 1, + federationRelayEnabled: 1, + federationRelayTtlDays: 30, + updatedAt: Date.now(), + }).run(); +} + +function seedPeer(id: string, origin: string, status: string): void { + testDb.insert(schema.federationPeers).values({ + id, origin, hmacSecret: 'secret', + status, lastSyncedAt: 0, createdAt: Date.now(), + }).run(); +} + +function countOutbox(peerId: string): number { + return testDb.select().from(schema.federationOutbox) + .where(eq(schema.federationOutbox.peerId, peerId)) + .all().length; +} + +// Import once at module level — vi.mock is hoisted and the factory returns the +// live testDb reference, so re-using the cached import is correct. +const { queueOutboxEvent } = await import('./federationOutbox.js'); + +describe('queueOutboxEvent — non-deliverable statuses', () => { + beforeEach(() => { + const sqlite = new Database(':memory:'); + testDb = drizzle(sqlite, { schema }); + applyMigrations(sqlite); + seedSettings(); + vi.restoreAllMocks(); + }); + + it.each([ + ['awaiting_approval'], + ['needs_attention'], + ['rejected'], + ['revoked'], + ])('drops the event and logs a reason for %s peers (no outbox row, no throw)', (status) => { + seedPeer('peer-drop', 'https://drop.example', status); + + const debugSpy = vi.spyOn(console, 'debug').mockImplementation(() => {}); + + queueOutboxEvent('entity-1', 'ctx-1', 'create', '{}', ['https://drop.example'], 'dm'); + + expect(countOutbox('peer-drop')).toBe(0); + expect(debugSpy).toHaveBeenCalled(); + expect(debugSpy.mock.calls[0]![0] as string).toContain(status); + }); +}); diff --git a/packages/server/src/utils/federationOutbox.ts b/packages/server/src/utils/federationOutbox.ts index 36bbb37c..d73ee0df 100644 --- a/packages/server/src/utils/federationOutbox.ts +++ b/packages/server/src/utils/federationOutbox.ts @@ -163,40 +163,84 @@ export function queueOutboxEvent( for (const origin of targetPeerOrigins) { if (matchedOrigins.has(origin)) continue; - // Check if there's a rejected/revoked peer we should skip const existingPeer = db .select({ status: schema.federationPeers.status }) .from(schema.federationPeers) .where(eq(schema.federationPeers.origin, origin)) .get(); - if (existingPeer && (existingPeer.status === 'rejected' || existingPeer.status === 'revoked')) { - console.warn(`[federation] queueOutboxEvent: skipping ${existingPeer.status} peer ${origin}`); - continue; - } - - // No peer record at all — create a pending placeholder - const peerId = generateSnowflake(); - const now = Date.now(); - db.insert(schema.federationPeers) - .values({ + if (!existingPeer) { + // No peer row — create pending placeholder, handshake fires on next tick + const peerId = generateSnowflake(); + const now = Date.now(); + db.insert(schema.federationPeers).values({ id: peerId, origin, hmacSecret: generateHmacSecret(), status: 'pending', createdAt: now, - }) - .run(); + }).run(); + const newPeer = db.select().from(schema.federationPeers) + .where(eq(schema.federationPeers.id, peerId)).get(); + if (newPeer) { + matchedPeers = [...matchedPeers, newPeer]; + console.log(`[federation] queueOutboxEvent: created pending placeholder for ${origin}`); + } + continue; + } - const newPeer = db - .select() - .from(schema.federationPeers) - .where(eq(schema.federationPeers.id, peerId)) - .get(); + // schema.federationPeers.status is plain text — narrow to known union for + // compile-time exhaustiveness check without widening to `string`. + const status = existingPeer.status as + | 'active' + | 'pending' + | 'unreachable' + | 'awaiting_approval' + | 'needs_attention' + | 'rejected' + | 'revoked'; - if (newPeer) { - matchedPeers = [...matchedPeers, newPeer]; - console.log(`[federation] queueOutboxEvent: created pending placeholder for ${origin}`); + switch (status) { + case 'active': + case 'pending': + case 'unreachable': { + // Race: peer transitioned to a deliverable status between the initial + // peers SELECT and this point in the loop. Re-fetch the full row and + // add to matchedPeers so the outer enqueue loop includes this peer. + // Do NOT silently drop — symmetric onPeerActivated on the peer's side + // is not guaranteed to cover asymmetric-failure cases (lost /peer/accept + // 200, health-check-only transition on one side). + const raced = db + .select() + .from(schema.federationPeers) + .where(eq(schema.federationPeers.origin, origin)) + .get(); + if (raced) { + matchedPeers = [...matchedPeers, raced]; + console.log(`[federation] queueOutboxEvent: race-caught ${origin} (now ${status}); enqueueing`); + } + break; + } + case 'awaiting_approval': + console.debug(`[federation] queueOutboxEvent: skipping ${origin} (awaiting_approval); mutation log will replay on activation`); + break; + case 'needs_attention': + console.debug(`[federation] queueOutboxEvent: skipping ${origin} (needs_attention; admin Reset required); mutation log will replay after Reset + re-peer`); + break; + case 'rejected': + console.debug(`[federation] queueOutboxEvent: skipping ${origin} (rejected peering)`); + break; + case 'revoked': + console.debug(`[federation] queueOutboxEvent: skipping ${origin} (revoked by admin)`); + break; + default: { + // Exhaustiveness check — no `as never` cast. TypeScript enforces + // that every status value is handled; adding a new value to the + // union without a case here fails typecheck. + const _exhaustive: never = status; + console.error(`[federation] queueOutboxEvent: unknown peer status for ${origin}: ${String(_exhaustive)}`); + break; + } } } } @@ -656,6 +700,12 @@ export function queueReadStateRelay( }; const targetOrigins = getGroupDmTargetOrigins(channelId); + appendMutationLog( + `read_state:${channel.federatedId}:${userId}`, + channelId, + 'read_state_update', + JSON.stringify({ user: { homeUserId, homeInstance }, messageRef }), + ); queueOutboxEvent( `read_state:${channel.federatedId}:${userId}`, channelId, @@ -711,6 +761,12 @@ export function queueDmCloseRelay( }; const targetOrigins = getGroupDmTargetOrigins(dmChannelId); + appendMutationLog( + `${eventType}:${channel.federatedId}:${userId}`, + dmChannelId, + eventType, + JSON.stringify({ homeUserId, homeInstance }), + ); queueOutboxEvent( `${eventType}:${channel.federatedId}:${userId}`, dmChannelId, diff --git a/packages/server/src/utils/federationPeerActivation.test.ts b/packages/server/src/utils/federationPeerActivation.test.ts new file mode 100644 index 00000000..d1dbc5da --- /dev/null +++ b/packages/server/src/utils/federationPeerActivation.test.ts @@ -0,0 +1,296 @@ +import { describe, it, expect, beforeEach, vi } 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'; +import { eq } from 'drizzle-orm'; + +const __dirname = path.dirname(fileURLToPath(import.meta.url)); + +type TestDb = ReturnType>; +let sqlite: Database.Database; +let testDb: TestDb; + +vi.mock('../db/index.js', () => ({ + getDb: () => testDb, + schema, +})); + +vi.mock('../utils/federationOutbox.js', () => ({ + isFederationRelayEnabled: () => true, +})); + +vi.mock('../utils/federationAuth.js', () => ({ + getOurOrigin: () => 'https://local.example', + buildFederationHeaders: (_body: string, _secret: string, _origin: string) => ({ + 'Content-Type': 'application/json', + 'X-Federation-Origin': _origin, + }), +})); + +vi.mock('../routes/federation.js', () => ({ + processRelayEvents: vi.fn().mockResolvedValue({ accepted: [], rejected: [] }), +})); + +vi.mock('../ws/handler.js', () => ({ + connectionManager: { + sendToAdmins: vi.fn(), + getAllOnlineUserIds: () => [], + sendToUser: vi.fn(), + sendToDmMembers: vi.fn(), + }, +})); + +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); + } + } +} + +function seedPeer(id: string, status: string, lastSyncedAt = 0): void { + testDb.insert(schema.federationPeers).values({ + id, origin: `https://${id}.example`, hmacSecret: 'secret', + status, lastSyncedAt, createdAt: Date.now(), + }).run(); +} + +function seedOutboxEntry(id: string, peerId: string, nextRetryAt: number, attempts: number): void { + testDb.insert(schema.federationOutbox).values({ + id, peerId, contextId: 'ch-1', entityId: `msg-${id}`, + contextType: 'dm', eventType: 'create', payload: '{}', + encryptionVersion: 0, attempts, nextRetryAt, + expiresAt: Date.now() + 30 * 86_400_000, + createdAt: Date.now(), + }).run(); +} + +describe('resetOutboxBackoff', () => { + beforeEach(() => { + sqlite = new Database(':memory:'); + testDb = drizzle(sqlite, { schema }); + applyMigrations(sqlite); + }); + + it('resets nextRetryAt=now and attempts=0 for all peer entries — including past-due ones', async () => { + const { resetOutboxBackoff } = await import('./federationPeerActivation.js'); + seedPeer('peer-a', 'active'); + seedPeer('peer-b', 'active'); + const now = Date.now(); + + // Three entries for peer-a: past-due (already eligible), near-future, far-future + seedOutboxEntry('entry-1', 'peer-a', now - 1000, 5); + seedOutboxEntry('entry-2', 'peer-a', now + 60_000, 3); + seedOutboxEntry('entry-3', 'peer-a', now + 86_400_000, 7); + // Entry for unrelated peer-b (must NOT be touched) + seedOutboxEntry('entry-4', 'peer-b', now + 86_400_000, 9); + + resetOutboxBackoff('peer-a'); + + const a1 = testDb.select().from(schema.federationOutbox).where(eq(schema.federationOutbox.id, 'entry-1')).get(); + const a2 = testDb.select().from(schema.federationOutbox).where(eq(schema.federationOutbox.id, 'entry-2')).get(); + const a3 = testDb.select().from(schema.federationOutbox).where(eq(schema.federationOutbox.id, 'entry-3')).get(); + const b4 = testDb.select().from(schema.federationOutbox).where(eq(schema.federationOutbox.id, 'entry-4')).get(); + + // All peer-a entries reset — including the past-due one (correctness: attempts=0 on those too) + expect(a1?.attempts).toBe(0); + expect(a2?.attempts).toBe(0); + expect(a3?.attempts).toBe(0); + expect(a1?.nextRetryAt).toBeGreaterThanOrEqual(now); + expect(a2?.nextRetryAt).toBeLessThanOrEqual(Date.now()); + expect(a3?.nextRetryAt).toBeLessThanOrEqual(Date.now()); + // peer-b untouched + expect(b4?.attempts).toBe(9); + expect(b4?.nextRetryAt).toBe(now + 86_400_000); + }); + + it('is a no-op when the peer has no outbox entries', async () => { + const { resetOutboxBackoff } = await import('./federationPeerActivation.js'); + seedPeer('peer-empty', 'active'); + expect(() => resetOutboxBackoff('peer-empty')).not.toThrow(); + }); +}); + +describe('syncPeerMutationLog', () => { + beforeEach(() => { + sqlite = new Database(':memory:'); + testDb = drizzle(sqlite, { schema }); + applyMigrations(sqlite); + vi.restoreAllMocks(); + }); + + it('seeds sinceTimestamp from peer.lastSyncedAt for each pass', async () => { + const { syncPeerMutationLog } = await import('./federationPeerActivation.js'); + testDb.insert(schema.federationPeers).values({ + id: 'peer-1', origin: 'https://peer-1.example', hmacSecret: 'secret', + status: 'active', lastSyncedAt: 5000, createdAt: Date.now(), + }).run(); + + const fetchSpy = vi.spyOn(globalThis, 'fetch').mockImplementation(async () => + new Response(JSON.stringify({ events: [], hasMore: false, checkpoint: 5000 }), { status: 200 }) + ); + + await syncPeerMutationLog('peer-1', 'health_check_recovery'); + + // Three passes: dm (no contextType), friend, profile + expect(fetchSpy).toHaveBeenCalledTimes(3); + for (const call of fetchSpy.mock.calls) { + const body = JSON.parse(call[1]?.body as string) as { sinceTimestamp: number }; + expect(body.sinceTimestamp).toBe(5000); + } + const calls = fetchSpy.mock.calls.map(c => JSON.parse(c[1]?.body as string) as { contextType?: string }); + expect(calls[0]!.contextType).toBeUndefined(); // DM pass (no contextType filter) + expect(calls[1]!.contextType).toBe('friend'); + expect(calls[2]!.contextType).toBe('profile'); + }); + + it('advances lastSyncedAt on success', async () => { + const { syncPeerMutationLog } = await import('./federationPeerActivation.js'); + testDb.insert(schema.federationPeers).values({ + id: 'peer-2', origin: 'https://peer-2.example', hmacSecret: 'secret', + status: 'active', lastSyncedAt: 0, createdAt: Date.now(), + }).run(); + + vi.spyOn(globalThis, 'fetch').mockImplementation(async () => + new Response(JSON.stringify({ events: [], hasMore: false, checkpoint: 1000 }), { status: 200 }) + ); + + const before = Date.now(); + await syncPeerMutationLog('peer-2', 'startup_bootstrap'); + const after = Date.now(); + + const row = testDb.select().from(schema.federationPeers) + .where(eq(schema.federationPeers.id, 'peer-2')).get(); + expect(row?.lastSyncedAt).toBeGreaterThanOrEqual(before); + expect(row?.lastSyncedAt).toBeLessThanOrEqual(after); + }); + + it('does NOT update lastSyncedAt on transient failure', async () => { + const { syncPeerMutationLog } = await import('./federationPeerActivation.js'); + testDb.insert(schema.federationPeers).values({ + id: 'peer-3', origin: 'https://peer-3.example', hmacSecret: 'secret', + status: 'active', lastSyncedAt: 42_000, createdAt: Date.now(), + }).run(); + + vi.spyOn(globalThis, 'fetch').mockImplementation(async () => + new Response('internal error', { status: 500 }) + ); + + await syncPeerMutationLog('peer-3', 'ensure_peered'); + const row = testDb.select().from(schema.federationPeers) + .where(eq(schema.federationPeers.id, 'peer-3')).get(); + expect(row?.lastSyncedAt).toBe(42_000); + }); + + it('does nothing when peer is not active', async () => { + const { syncPeerMutationLog } = await import('./federationPeerActivation.js'); + testDb.insert(schema.federationPeers).values({ + id: 'peer-4', origin: 'https://peer-4.example', hmacSecret: 'secret', + status: 'pending', lastSyncedAt: 0, createdAt: Date.now(), + }).run(); + + const fetchSpy = vi.spyOn(globalThis, 'fetch'); + await syncPeerMutationLog('peer-4', 'health_check_recovery'); + expect(fetchSpy).not.toHaveBeenCalled(); + }); + + it('advances sinceTimestamp within a pass using data.checkpoint when hasMore is true', async () => { + const { syncPeerMutationLog } = await import('./federationPeerActivation.js'); + testDb.insert(schema.federationPeers).values({ + id: 'peer-5', origin: 'https://peer-5.example', hmacSecret: 'secret', + status: 'active', lastSyncedAt: 100, createdAt: Date.now(), + }).run(); + + let call = 0; + const fetchSpy = vi.spyOn(globalThis, 'fetch').mockImplementation(async () => { + call++; + // First DM call: one event, hasMore=true, checkpoint advances to 2500 + // Second DM call: empty, hasMore=false, ends the DM pass + // Remaining calls (friend, profile): empty/done immediately + if (call === 1) { + return new Response( + JSON.stringify({ + events: [{ eventType: 'create', messageId: 'm1', timestamp: 200, encryptionVersion: 0 }], + hasMore: true, + checkpoint: 2500, + }), + { status: 200 }, + ); + } + return new Response(JSON.stringify({ events: [], hasMore: false, checkpoint: call === 2 ? 2500 : 100 }), { status: 200 }); + }); + + await syncPeerMutationLog('peer-5', 'health_check_recovery'); + + // Call 1: DM pass, since=100 (peer.lastSyncedAt) + // Call 2: DM pass continuation, since=2500 (advanced by previous checkpoint) + // Call 3: friend pass, since=100 (re-seeded from peer.lastSyncedAt) + // Call 4: profile pass, since=100 (re-seeded from peer.lastSyncedAt) + expect(fetchSpy).toHaveBeenCalledTimes(4); + const bodies = fetchSpy.mock.calls.map(c => JSON.parse(c[1]?.body as string) as { sinceTimestamp: number; contextType?: string }); + expect(bodies[0]?.sinceTimestamp).toBe(100); + expect(bodies[0]?.contextType).toBeUndefined(); + expect(bodies[1]?.sinceTimestamp).toBe(2500); // advanced by checkpoint from call 1 + expect(bodies[1]?.contextType).toBeUndefined(); + expect(bodies[2]?.sinceTimestamp).toBe(100); // friend pass re-seeds from peer.lastSyncedAt + expect(bodies[2]?.contextType).toBe('friend'); + expect(bodies[3]?.sinceTimestamp).toBe(100); // profile pass re-seeds from peer.lastSyncedAt + expect(bodies[3]?.contextType).toBe('profile'); + }); +}); + +describe('onPeerActivated', () => { + beforeEach(() => { + sqlite = new Database(':memory:'); + testDb = drizzle(sqlite, { schema }); + applyMigrations(sqlite); + vi.restoreAllMocks(); + }); + + it('runs resetOutboxBackoff and syncPeerMutationLog once, even under concurrent calls', async () => { + const { onPeerActivated } = await import('./federationPeerActivation.js'); + + testDb.insert(schema.federationPeers).values({ + id: 'peer-x', origin: 'https://peer-x.example', hmacSecret: 'secret', + status: 'active', lastSyncedAt: 0, createdAt: Date.now(), + }).run(); + + let fetchCount = 0; + vi.spyOn(globalThis, 'fetch').mockImplementation(async () => { + fetchCount++; + // Deliberately slow to let the second concurrent call share the in-flight promise. + await new Promise(r => setTimeout(r, 20)); + return new Response(JSON.stringify({ events: [], hasMore: false, checkpoint: 0 }), { status: 200 }); + }); + + const p1 = onPeerActivated('peer-x', 'health_check_recovery'); + const p2 = onPeerActivated('peer-x', 'accept_new'); + await Promise.all([p1, p2]); + + // Three fetch calls for the three sync passes (dm, friend, profile) — not six. + expect(fetchCount).toBe(3); + }); + + it('swallows errors from syncPeerMutationLog so the handler does not throw', async () => { + const { onPeerActivated } = await import('./federationPeerActivation.js'); + + testDb.insert(schema.federationPeers).values({ + id: 'peer-err', origin: 'https://peer-err.example', hmacSecret: 'secret', + status: 'active', lastSyncedAt: 0, createdAt: Date.now(), + }).run(); + + vi.spyOn(globalThis, 'fetch').mockImplementation(async () => { + throw new Error('network down'); + }); + + await expect(onPeerActivated('peer-err', 'ensure_peered')).resolves.toBeUndefined(); + }); +}); diff --git a/packages/server/src/utils/federationPeerActivation.ts b/packages/server/src/utils/federationPeerActivation.ts new file mode 100644 index 00000000..ee34c265 --- /dev/null +++ b/packages/server/src/utils/federationPeerActivation.ts @@ -0,0 +1,179 @@ +import { getDb } from '../db/index.js'; +import * as schema from '../db/schema.js'; +import { and, eq } from 'drizzle-orm'; +import { isFederationRelayEnabled } from './federationOutbox.js'; +import { buildFederationHeaders, getOurOrigin } from './federationAuth.js'; +import type { FederationRelayEvent } from '@backspace/shared'; + +export type PeerActivationReason = + | 'initiate_accepted' + | 'accept_rejected_override' + | 'accept_awaiting_approval' + | 'accept_pending' + | 'accept_new' + | 'approval_handshake' + | 'health_check_recovery' + | 'ensure_peered' + | 'startup_bootstrap'; + +// Dedup: concurrent activations for the same peerId share one promise. +const inFlightActivation = new Map>(); + +/** + * Called whenever federation_peers.status transitions to 'active' for any reason. + * Two independent invariants — both run unconditionally: + * 1. Reset outbox backoff (nextRetryAt = now, attempts = 0) for this peer. + * 2. Pull-sync mutation log from peer's /api/federation/sync since lastSyncedAt. + * + * Call sites (must remain exhaustive — grep `onPeerActivated(` to audit): + * - routes/federation.ts /peer/initiate activation + * - routes/federation.ts /peer/accept existing-rejected override + * - routes/federation.ts /peer/accept existing-awaiting_approval + * - routes/federation.ts /peer/accept existing-pending + * - routes/federation.ts /peer/accept new-peer + * - routes/federation.ts /approval-requests/:id/approve + * - utils/federationWorker.ts health check recovery + * - utils/federationPeering.ts ensurePeered/performHandshake + * - utils/federationWorker.ts startup bootstrap (via startupBootstrapSync) + * + * Deduplicated by peerId — concurrent calls share one promise. + */ +export async function onPeerActivated( + peerId: string, + reason: PeerActivationReason, +): Promise { + const existing = inFlightActivation.get(peerId); + if (existing) return existing; + + const promise = (async () => { + try { + resetOutboxBackoff(peerId); + await syncPeerMutationLog(peerId, reason); + const { connectionManager } = await import('../ws/handler.js'); + connectionManager.sendToAdmins({ type: 'federation_peers_changed' as const }); + } catch (err) { + console.error(`[federation] onPeerActivated(${peerId}, ${reason}) failed:`, err); + } + })(); + + inFlightActivation.set(peerId, promise); + try { + await promise; + } finally { + inFlightActivation.delete(peerId); + } +} + +/** + * Reset all outbox backoff state for a peer (nextRetryAt = now, attempts = 0). + * Unconditional across all entries of the peer — see spec §Invariant 1. + */ +export function resetOutboxBackoff(peerId: string): void { + const db = getDb(); + const now = Date.now(); + const result = db + .update(schema.federationOutbox) + .set({ nextRetryAt: now, attempts: 0 }) + .where(eq(schema.federationOutbox.peerId, peerId)) + .run(); + if (result.changes > 0) { + console.log(`[federation] Reset backoff on ${result.changes} outbox entries for peer ${peerId}`); + } +} + +/** + * Pull-sync mutation log from the peer's /api/federation/sync endpoint. + * Runs three contextType passes (dm, friend, profile), paginating each. + * Updates peer.lastSyncedAt to Date.now() on success; leaves it untouched + * on transient failure so the next activation retries. + */ +export async function syncPeerMutationLog( + peerId: string, + reason: PeerActivationReason, +): Promise { + if (!isFederationRelayEnabled()) return; + + const db = getDb(); + const peer = db.select().from(schema.federationPeers) + .where(eq(schema.federationPeers.id, peerId)).get(); + if (!peer || peer.status !== 'active') return; + + const activePeer = peer; // narrowed by the guard above + + const ourOrigin = getOurOrigin(); + const signingSecret = (activePeer.pendingHmacSecret && activePeer.secretRotationAt) + ? activePeer.pendingHmacSecret + : activePeer.hmacSecret; + + console.log(`[federation] Sync-pull from ${activePeer.origin} (reason=${reason}, since=${activePeer.lastSyncedAt ?? 0})`); + + let totalEvents = 0; + + type SyncRequestBody = { + sinceTimestamp: number; + limit: number; + contextType?: 'friend' | 'profile'; + }; + + async function runPass(contextType?: 'friend' | 'profile'): Promise { + let since = activePeer.lastSyncedAt ?? 0; + while (true) { + const bodyObj: SyncRequestBody = { sinceTimestamp: since, limit: 100 }; + if (contextType) bodyObj.contextType = contextType; + const body = JSON.stringify(bodyObj); + const headers = buildFederationHeaders(body, signingSecret, ourOrigin); + const resp = await fetch(`${activePeer.origin}/api/federation/sync`, { + method: 'POST', headers, body, + signal: AbortSignal.timeout(30_000), + }); + if (!resp.ok) { + console.warn(`[federation] Sync-pull ${contextType ?? 'dm'} pass HTTP ${resp.status} for ${activePeer.origin}`); + return false; + } + const data = await resp.json() as { events: FederationRelayEvent[]; hasMore: boolean; checkpoint: number }; + if (data.events.length === 0) return true; + const { processRelayEvents } = await import('../routes/federation.js'); + await processRelayEvents(data.events, activePeer.origin, activePeer.origin, db); + totalEvents += data.events.length; + since = data.checkpoint; + if (!data.hasMore) return true; + } + } + + try { + if (!(await runPass())) return; + if (!(await runPass('friend'))) return; + if (!(await runPass('profile'))) return; + + db.update(schema.federationPeers) + .set({ lastSyncedAt: Date.now() }) + .where(eq(schema.federationPeers.id, activePeer.id)) + .run(); + + if (totalEvents > 0) { + console.log(`[federation] Sync-pull from ${activePeer.origin} replayed ${totalEvents} events`); + } + } catch (err) { + console.error(`[federation] Sync-pull from ${activePeer.origin} failed:`, err); + } +} + +/** + * Startup bootstrap — scan for freshly-peered rows (status='active', lastSyncedAt=0) + * and run onPeerActivated for each. Replaces runInitialSyncForNewPeers. + * Invoked from startFederationWorkers. + */ +export async function startupBootstrapSync(): Promise { + if (!isFederationRelayEnabled()) return; + + const db = getDb(); + const peers = db.select().from(schema.federationPeers) + .where(and( + eq(schema.federationPeers.status, 'active'), + eq(schema.federationPeers.lastSyncedAt, 0), + )).all(); + + for (const peer of peers) { + await onPeerActivated(peer.id, 'startup_bootstrap'); + } +} diff --git a/packages/server/src/utils/federationPeering.test.ts b/packages/server/src/utils/federationPeering.test.ts index 3961dc0b..2b044795 100644 --- a/packages/server/src/utils/federationPeering.test.ts +++ b/packages/server/src/utils/federationPeering.test.ts @@ -1,4 +1,4 @@ -import { describe, it, expect, vi } from 'vitest'; +import { describe, it, expect, vi, beforeEach } from 'vitest'; import type { EnsurePeeredResult } from './federationPeering.js'; import { racePeering } from './federationPeering.js'; @@ -112,3 +112,59 @@ describe('racePeering', () => { warnSpy.mockRestore(); }); }); + +describe('ensurePeered needs_attention handling', () => { + beforeEach(() => { + vi.resetModules(); + }); + + it('returns rejected without calling performHandshake when peer is in needs_attention', async () => { + const fakeDbGet = vi.fn().mockReturnValue({ + id: 'peer-na', + origin: 'https://remote.example', + status: 'needs_attention', + hmacSecret: 'secret', + createdAt: Date.now(), + lastSyncedAt: 0, + }); + + vi.doMock('../db/index.js', () => ({ + getDb: () => ({ + select: () => ({ + from: () => ({ + where: () => ({ + get: fakeDbGet, + }), + }), + }), + }), + })); + + vi.doMock('../utils/federationAuth.js', () => ({ + getOurOrigin: () => 'https://local.example', + generateHmacSecret: () => 'new-secret', + })); + + vi.doMock('../routes/federation.js', () => ({ + validateOrigin: (o: string) => o, + })); + + vi.doMock('../utils/federationPeerActivation.js', () => ({ + onPeerActivated: vi.fn(), + })); + + const { ensurePeered } = await import('./federationPeering.js'); + const fetchSpy = vi.spyOn(globalThis, 'fetch'); + + const result = await ensurePeered('https://remote.example'); + + expect(result.status).toBe('rejected'); + if (result.status === 'rejected') { + expect(result.error).toContain('needs_attention'); + } + // performHandshake must not have fired — no POST to /peer/accept + expect(fetchSpy).not.toHaveBeenCalled(); + + vi.restoreAllMocks(); + }); +}); diff --git a/packages/server/src/utils/federationPeering.ts b/packages/server/src/utils/federationPeering.ts index 2b582413..fc6b9ae2 100644 --- a/packages/server/src/utils/federationPeering.ts +++ b/packages/server/src/utils/federationPeering.ts @@ -4,6 +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'; // ─── Types ─────────────────────────────────────────────────────────────────── @@ -72,6 +73,9 @@ export async function ensurePeered(origin: string): Promise // Unreachable peers were previously active — treat as active for peering // (the health check will restore them; don't re-handshake) return { status: 'active', peerId: existing.id }; + case 'needs_attention': + // Admin intervention required — do not auto-heal via performHandshake + return { status: 'rejected', error: 'Peer in needs_attention — admin Reset required' }; case 'awaiting_approval': return { status: 'pending', error: 'Awaiting admin approval on remote instance' }; case 'pending': @@ -158,6 +162,9 @@ async function performHandshake( .run(); const { connectionManager } = await import('../ws/handler.js'); 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 }; } diff --git a/packages/server/src/utils/federationWorker.ts b/packages/server/src/utils/federationWorker.ts index d444aa53..e9c218e0 100644 --- a/packages/server/src/utils/federationWorker.ts +++ b/packages/server/src/utils/federationWorker.ts @@ -2,7 +2,7 @@ import { getDb } from '../db/index.js'; import * as schema from '../db/schema.js'; import { eq, and, lte, asc, inArray, sql } from 'drizzle-orm'; import { config } from '../config.js'; -import { isFederationRelayEnabled, queueOutboxEvent } from './federationOutbox.js'; +import { isFederationRelayEnabled, queueOutboxEvent, appendMutationLog } from './federationOutbox.js'; import { runFederationJanitor } from './storageJanitor.js'; import { buildFederationHeaders, getOurOrigin, generateHmacSecret, ROTATION_GRACE_PERIOD_MS } from './federationAuth.js'; import { evaluateAuthFailure, AUTH_FAILURE_THRESHOLD } from './federationAuthFailure.js'; @@ -10,8 +10,8 @@ import { generateSnowflake } from './snowflake.js'; import { getDmMessageWithUser } from '../routes/dm.js'; import { connectionManager } from '../ws/handler.js'; import { generateThumbnail } from './thumbnail.js'; -import { processRelayEvents } from '../routes/federation.js'; import type { FederationRelayRequest, FederationRelayResponse, FederationRelayEvent } from '@backspace/shared'; +import { onPeerActivated, startupBootstrapSync } from './federationPeerActivation.js'; import fs from 'node:fs'; import path from 'node:path'; import crypto from 'node:crypto'; @@ -703,6 +703,18 @@ function handleSizeRejection( affectedUserIds, }; + appendMutationLog( + localMsg.sourceMessageId, + localMsg.dmChannelId, + 'file_rejected', + JSON.stringify({ + attachmentId: att?.id ?? entry.sourceUrl, + sourceFilename, + rejectionReason: 'size_limit_exceeded', + rejectionLimit: maxUploadSize, + affectedUserIds, + }), + ); queueOutboxEvent( localMsg.sourceMessageId, localMsg.dmChannelId, @@ -1054,6 +1066,8 @@ async function processHealthCheckTick(): Promise { console.log( `[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 } catch (err) { @@ -1076,124 +1090,15 @@ function scheduleJanitorTick(): void { // ─── Lifecycle ────────────────────────────────────────────────────────────── -/** - * Trigger checkpoint sync for peers that have never been synced (lastSyncedAt === 0). - * This catches historical messages that existed before the relay was enabled. - */ -async function runInitialSyncForNewPeers(): Promise { - if (!isFederationRelayEnabled()) return; - - const db = getDb(); - const unsyncedPeers = db - .select() - .from(schema.federationPeers) - .where(and( - eq(schema.federationPeers.status, 'active'), - eq(schema.federationPeers.lastSyncedAt, 0), - )) - .all(); - - if (unsyncedPeers.length === 0) return; - - const ourOrigin = getOurOrigin(); - - for (const peer of unsyncedPeers) { - const signingSecret = (peer.pendingHmacSecret && peer.secretRotationAt) - ? peer.pendingHmacSecret - : peer.hmacSecret; - try { - console.log(`[federation-worker] Running initial sync with ${peer.origin}...`); - let sinceTimestamp = 0; - let totalEvents = 0; - - // Paginate through all events from the peer - while (true) { - const body = JSON.stringify({ sinceTimestamp, limit: 100 }); - const headers = buildFederationHeaders(body, signingSecret, ourOrigin); - - const response = await fetch(`${peer.origin}/api/federation/sync`, { - method: 'POST', - headers, - body, - signal: AbortSignal.timeout(30_000), - }); - - if (!response.ok) { - console.error(`[federation-worker] Sync with ${peer.origin} failed: ${response.status}`); - break; - } - - const data = await response.json() as { events: FederationRelayEvent[]; hasMore: boolean; checkpoint: number }; - - if (data.events.length === 0) break; - - // Process events directly — no HTTP round-trip (FED-005) - await processRelayEvents(data.events, peer.origin, peer.origin, db); - - totalEvents += data.events.length; - sinceTimestamp = data.checkpoint; - - if (!data.hasMore) break; - } - - // Second pass: sync friend events - let friendSinceTimestamp = 0; - while (true) { - const friendBody = JSON.stringify({ sinceTimestamp: friendSinceTimestamp, contextType: 'friend', limit: 100 }); - const friendHeaders = buildFederationHeaders(friendBody, signingSecret, ourOrigin); - - const friendResponse = await fetch(`${peer.origin}/api/federation/sync`, { - method: 'POST', - headers: friendHeaders, - body: friendBody, - signal: AbortSignal.timeout(30_000), - }); - - if (!friendResponse.ok) { - console.error(`[federation-worker] Friend sync with ${peer.origin} failed: ${friendResponse.status}`); - break; - } - - const friendData = await friendResponse.json() as { events: FederationRelayEvent[]; hasMore: boolean; checkpoint: number }; - - if (friendData.events.length === 0) break; - - // Process events directly — no HTTP round-trip (FED-005) - await processRelayEvents(friendData.events, peer.origin, peer.origin, db); - - totalEvents += friendData.events.length; - friendSinceTimestamp = friendData.checkpoint; - - if (!friendData.hasMore) break; - } - - // Update lastSyncedAt so this doesn't run again - db.update(schema.federationPeers) - .set({ lastSyncedAt: Date.now() }) - .where(eq(schema.federationPeers.id, peer.id)) - .run(); - - if (totalEvents > 0) { - console.log(`[federation-worker] Initial sync with ${peer.origin}: ${totalEvents} events synced`); - } else { - console.log(`[federation-worker] Initial sync with ${peer.origin}: no events to sync`); - } - } catch (err) { - console.error(`[federation-worker] Initial sync with ${peer.origin} failed:`, err); - // Don't update lastSyncedAt — will retry next startup - } - } -} - export function startFederationWorkers(): void { console.log('[federation-worker] Federation workers started'); scheduleOutboxTick(); scheduleFileQueueTick(); scheduleHealthCheckTick(); scheduleJanitorTick(); - // Run initial sync for newly peered instances (async, non-blocking) - runInitialSyncForNewPeers().catch((err) => { - console.error('[federation-worker] Initial sync error:', err); + // Bootstrap sync for freshly-peered rows (async, non-blocking) + startupBootstrapSync().catch((err) => { + console.error('[federation-worker] Startup bootstrap sync error:', err); }); }