diff --git a/docs/systems/federation.md b/docs/systems/federation.md index 1a3521b6..b845672d 100644 --- a/docs/systems/federation.md +++ b/docs/systems/federation.md @@ -466,8 +466,9 @@ Trigger (API/WS handler) 5. Sign with `buildFederationHeaders(body, peerHmacSecret, ourOrigin)` 6. POST to `{peerOrigin}/api/federation/relay` 7. On success (200): - - Delete accepted entries from outbox (matched by `entityId` -> `outboxId`) - - Log rejected entries (remain in outbox for retry) + - Compute the **terminal entity set** = accepted entries ∪ duplicate-rejected entries + - Delete all terminal entries from outbox (matched by `entityId` -> `outboxId`) + - Log remaining (non-duplicate) rejected entries at `console.warn` (they stay in outbox for retry) - Store `result.maxUploadSize` on peer record - Update peer: `lastSeenAt = now`, `consecutiveFailures = 0` 8. On failure (non-200 or network error): @@ -476,6 +477,14 @@ Trigger (API/WS handler) - Increment peer `consecutiveFailures`, set `lastFailureAt` - If `consecutiveFailures >= PEER_UNREACHABLE_THRESHOLD (10)` -> mark peer `unreachable` +#### Terminal rejection: `duplicate` + +The outbox delivery worker treats a relay response of `{ rejected: [{ reason: 'duplicate', ... }] }` as effectively-accepted — the outbox entry is deleted rather than retained for retry. The `duplicate` reason is emitted by the receiving instance's inbound processors when a row with the same `(sourceInstance, sourceMessageId)` already exists; retrying will fail identically until TTL (30 days). Since the peer already has the message, terminal removal is the correct outcome. + +Logged at `console.log` ("outbox entry removed (terminal)") to distinguish from retained-for-retry `console.warn` messages. + +Other rejection reasons (`attribution_mismatch`, `missing_*_payload`, `unknown_event_type`, `unauthorized_source`, `channel_not_found`, `participant_not_found`, `processing_error`, …) remain on the retry path. Some are arguably terminal too; treating them as such is deferred until they are observed accumulating in practice. + ### Retry Backoff Schedule | Attempt | Delay | diff --git a/packages/server/src/utils/federationWorker.test.ts b/packages/server/src/utils/federationWorker.test.ts new file mode 100644 index 00000000..3e8e4d56 --- /dev/null +++ b/packages/server/src/utils/federationWorker.test.ts @@ -0,0 +1,159 @@ +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('../ws/handler.js', () => ({ + connectionManager: { + sendToAdmins: vi.fn(), + getAllOnlineUserIds: () => [], + sendToUser: vi.fn(), + sendToDmMembers: vi.fn(), + }, +})); + +vi.mock('../utils/federationAuth.js', () => ({ + getOurOrigin: () => 'https://test.example', + buildFederationHeaders: () => ({ 'Content-Type': 'application/json' }), + generateHmacSecret: () => 'secret', + ROTATION_GRACE_PERIOD_MS: 15 * 60 * 1000, +})); + +vi.mock('../utils/federationOutbox.js', () => ({ + isFederationRelayEnabled: () => true, + queueOutboxEvent: vi.fn(), + appendMutationLog: vi.fn(), +})); + +vi.mock('../utils/federationPeerActivation.js', () => ({ + onPeerActivated: vi.fn(), + startupBootstrapSync: vi.fn(), +})); + +vi.mock('../utils/storageJanitor.js', () => ({ + runFederationJanitor: vi.fn(), +})); + +vi.mock('../utils/thumbnail.js', () => ({ + generateThumbnail: vi.fn(), +})); + +vi.mock('../routes/dm.js', () => ({ + getDmMessageWithUser: vi.fn(), +})); + +vi.mock('../utils/federationAuthFailure.js', () => ({ + evaluateAuthFailure: vi.fn().mockReturnValue({ kind: 'increment', newAuthFailures: 1 }), + AUTH_FAILURE_THRESHOLD: 5, +})); + +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): void { + testDb.insert(schema.federationPeers).values({ + id, origin: 'https://peer.example', hmacSecret: 'secret', + status: 'active', lastSyncedAt: Date.now(), createdAt: Date.now(), + }).run(); +} + +function seedOutboxEntry(id: string, peerId: string, entityId: string): void { + testDb.insert(schema.federationOutbox).values({ + id, peerId, contextId: 'ch-1', entityId, + contextType: 'dm', eventType: 'create', payload: JSON.stringify({ + message: { userId: 'u', homeUserId: 'u', homeInstance: 'test.example', content: 'hi', replyToId: null, editedAt: null, createdAt: Date.now() }, + }), + encryptionVersion: 0, attempts: 0, nextRetryAt: Date.now() - 1000, + expiresAt: Date.now() + 30 * 86_400_000, + createdAt: Date.now(), + }).run(); +} + +describe('outbox worker — duplicate rejection is terminal', () => { + beforeEach(() => { + sqlite = new Database(':memory:'); + testDb = drizzle(sqlite, { schema }); + applyMigrations(sqlite); + vi.restoreAllMocks(); + // Re-apply the static mocks that vi.restoreAllMocks() would undo. + // isFederationRelayEnabled is mocked at module level via vi.mock (hoisted), + // so it survives restoreAllMocks — spies created with vi.spyOn are the ones + // that get restored. The fetch spy is re-created per test via mockImplementation. + }); + + it('deletes the outbox entry when the peer responds with duplicate rejection', async () => { + seedPeer('peer-dup'); + seedOutboxEntry('entry-dup', 'peer-dup', 'msg-already-there'); + + // Mock the fetch to return a relay response with duplicate rejection + vi.spyOn(globalThis, 'fetch').mockImplementation(async () => + new Response(JSON.stringify({ + accepted: [], + rejected: [{ messageId: 'msg-already-there', reason: 'duplicate' }], + }), { status: 200, headers: { 'Content-Type': 'application/json' } }), + ); + + const workerModule = await import('./federationWorker.js'); + const processOutboxTick = (workerModule as { processOutboxTick?: () => Promise }).processOutboxTick; + if (!processOutboxTick) { + throw new Error('processOutboxTick must be exported for this test. If not yet exported, export it.'); + } + await processOutboxTick(); + + // Verify the outbox entry was deleted (terminal treatment) + const remaining = testDb.select().from(schema.federationOutbox) + .where(eq(schema.federationOutbox.id, 'entry-dup')).get(); + expect(remaining).toBeUndefined(); + }); + + it('retains outbox entries for non-duplicate rejection reasons (e.g., processing_error)', async () => { + seedPeer('peer-transient'); + seedOutboxEntry('entry-transient', 'peer-transient', 'msg-transient'); + + vi.spyOn(globalThis, 'fetch').mockImplementation(async () => + new Response(JSON.stringify({ + accepted: [], + rejected: [{ messageId: 'msg-transient', reason: 'processing_error' }], + }), { status: 200, headers: { 'Content-Type': 'application/json' } }), + ); + + const workerModule = await import('./federationWorker.js'); + const processOutboxTick = (workerModule as { processOutboxTick?: () => Promise }).processOutboxTick; + if (!processOutboxTick) { + throw new Error('processOutboxTick must be exported for this test.'); + } + await processOutboxTick(); + + const remaining = testDb.select().from(schema.federationOutbox) + .where(eq(schema.federationOutbox.id, 'entry-transient')).get(); + // The entry must be retained — a 200 OK with a transient rejection reason + // does not delete the outbox entry. (Backoff is only applied on non-OK HTTP + // responses; a 200 with a rejection means the peer processed the batch but + // declined this particular message — the entry stays for the next tick.) + expect(remaining).toBeDefined(); + }); +}); diff --git a/packages/server/src/utils/federationWorker.ts b/packages/server/src/utils/federationWorker.ts index e9c218e0..233561fc 100644 --- a/packages/server/src/utils/federationWorker.ts +++ b/packages/server/src/utils/federationWorker.ts @@ -100,7 +100,7 @@ function scheduleOutboxTick(): void { }, OUTBOX_INTERVAL_MS); } -async function processOutboxTick(): Promise { +export async function processOutboxTick(): Promise { if (!isFederationRelayEnabled()) { return; } @@ -230,26 +230,43 @@ async function processOutboxTick(): Promise { if (response.ok) { const result = await response.json() as FederationRelayResponse; - // Delete accepted entries - if (result.accepted.length > 0) { - // Map accepted messageIds to outbox IDs - const acceptedSet = new Set(result.accepted); - const acceptedOutboxIds = peerEntries - .filter((e) => acceptedSet.has(e.entityId)) + // Terminal outcomes = accepted + duplicate-rejected. + // `duplicate` means the peer already has the message (e.g., delivered + // earlier via outbox or pulled via sync). Retrying will fail with + // `duplicate` forever until TTL expires — treat it as effectively- + // accepted and remove the outbox entry. + const terminalEntityIds = new Set(result.accepted); + for (const rejection of result.rejected) { + if (rejection.reason === 'duplicate') { + terminalEntityIds.add(rejection.messageId); + } + } + + if (terminalEntityIds.size > 0) { + const terminalOutboxIds = peerEntries + .filter((e) => terminalEntityIds.has(e.entityId)) .map((e) => e.outboxId); - if (acceptedOutboxIds.length > 0) { + if (terminalOutboxIds.length > 0) { db.delete(schema.federationOutbox) - .where(inArray(schema.federationOutbox.id, acceptedOutboxIds)) + .where(inArray(schema.federationOutbox.id, terminalOutboxIds)) .run(); } } - // Log rejected entries (they remain in outbox for retry) + // Log rejected entries. Duplicate is terminal (outbox entry already + // removed above) — log at info level. Other reasons are transient / + // retained for retry — log at warn level. for (const rejection of result.rejected) { - console.warn( - `[federation-worker] Peer ${peerOrigin} rejected message ${rejection.messageId}: ${rejection.reason}`, - ); + if (rejection.reason === 'duplicate') { + console.log( + `[federation-worker] Peer ${peerOrigin} rejected message ${rejection.messageId} as duplicate — outbox entry removed (terminal)`, + ); + } else { + console.warn( + `[federation-worker] Peer ${peerOrigin} rejected message ${rejection.messageId}: ${rejection.reason}`, + ); + } } // Store the peer's max upload size for informational display