diff --git a/apps/journal/app/jobs/federation-dedup-sweep.ts b/apps/journal/app/jobs/federation-dedup-sweep.ts new file mode 100644 index 0000000..8fe3980 --- /dev/null +++ b/apps/journal/app/jobs/federation-dedup-sweep.ts @@ -0,0 +1,22 @@ +import { defineJournalJob } from "./payloads.ts"; +import { sweepProcessedActivities } from "../lib/federation-replay.server.ts"; +import { logger } from "../lib/logger.server.ts"; + +/** + * Daily cleanup of federation_processed_activities rows older than 30 + * days (spec: federation-operations "Inbound replay defense"). Replays of + * activities that old are already rejected by HTTP-signature date + * freshness, so the dedup record is no longer needed; this keeps the + * table bounded. + */ +export const federationDedupSweepJob = defineJournalJob({ + name: "federation-dedup-sweep", + cron: "30 4 * * *", // daily at 04:30 UTC (offset from federation-kv-sweep at 04:15) + retryLimit: 1, + expireInSeconds: 60, + async handler() { + const purged = await sweepProcessedActivities(); + logger.info({ purged }, "federation-dedup-sweep"); + return { purged }; + }, +}); diff --git a/apps/journal/app/jobs/payloads.ts b/apps/journal/app/jobs/payloads.ts index 3fa63f6..e3ba254 100644 --- a/apps/journal/app/jobs/payloads.ts +++ b/apps/journal/app/jobs/payloads.ts @@ -17,6 +17,7 @@ export interface JobPayloads { "deliver-activity": DeliveryPayload; "demo-bot-generate": void; "demo-bot-prune": void; + "federation-dedup-sweep": void; "federation-kv-sweep": void; "garmin-import-activity": GarminImportData; "import-batches-sweep": void; diff --git a/apps/journal/app/lib/federation-replay.integration.test.ts b/apps/journal/app/lib/federation-replay.integration.test.ts new file mode 100644 index 0000000..d55bd87 --- /dev/null +++ b/apps/journal/app/lib/federation-replay.integration.test.ts @@ -0,0 +1,59 @@ +import { describe, it, expect, afterAll } from "vitest"; +import { like, eq } from "drizzle-orm"; +import { getDb } from "./db.ts"; +import { federationProcessedActivities } from "@trails-cool/db/schema/journal"; +import { markInboundActivityProcessed, sweepProcessedActivities } from "./federation-replay.server.ts"; + +// Opt-in: talks to real Postgres (same convention as the other +// *.integration.test.ts files, e.g. federation-kv). +const runIntegration = process.env.FEDERATION_INTEGRATION === "1"; + +// Namespace IRIs per run so concurrent runs can't collide. +const NS = `https://replay.test/${Date.now()}`; + +describe.runIf(runIntegration)("federation replay defense (integration)", () => { + afterAll(async () => { + const db = getDb(); + await db.delete(federationProcessedActivities).where(like(federationProcessedActivities.activityIri, `${NS}%`)); + }); + + it("first delivery is fresh, redelivery is dropped as a duplicate", async () => { + const iri = `${NS}/like/1`; + expect((await markInboundActivityProcessed(iri)).fresh).toBe(true); + // The same signed activity delivered again: no-op, reported not fresh. + expect((await markInboundActivityProcessed(iri)).fresh).toBe(false); + expect((await markInboundActivityProcessed(iri)).fresh).toBe(false); + }); + + it("distinct activity IRIs are each fresh", async () => { + expect((await markInboundActivityProcessed(`${NS}/delete/1`)).fresh).toBe(true); + expect((await markInboundActivityProcessed(`${NS}/update/1`)).fresh).toBe(true); + }); + + it("sweeps rows older than 30 days, keeps recent ones", async () => { + const db = getDb(); + const oldIri = `${NS}/old`; + const freshIri = `${NS}/recent`; + // 31 days ago vs now. + await db.insert(federationProcessedActivities).values({ + activityIri: oldIri, + receivedAt: new Date(Date.now() - 31 * 24 * 60 * 60 * 1000), + }); + await markInboundActivityProcessed(freshIri); + + const purged = await sweepProcessedActivities(); + expect(purged).toBeGreaterThanOrEqual(1); + + const remaining = await db + .select({ iri: federationProcessedActivities.activityIri }) + .from(federationProcessedActivities) + .where(eq(federationProcessedActivities.activityIri, oldIri)); + expect(remaining).toHaveLength(0); + + const kept = await db + .select({ iri: federationProcessedActivities.activityIri }) + .from(federationProcessedActivities) + .where(eq(federationProcessedActivities.activityIri, freshIri)); + expect(kept).toHaveLength(1); + }); +}); diff --git a/apps/journal/app/lib/federation-replay.server.ts b/apps/journal/app/lib/federation-replay.server.ts new file mode 100644 index 0000000..f4e4958 --- /dev/null +++ b/apps/journal/app/lib/federation-replay.server.ts @@ -0,0 +1,50 @@ +// Inbound-activity replay defense (spec: federation-operations "Inbound +// replay defense"). The narrow inbox processes Follow/Undo/Accept/Reject; +// none of those carried replay protection before (only Create(Note) did, +// via the activities.remote_origin_iri unique constraint). A hostile or +// buggy remote redelivering a signed activity would re-run its side +// effects. This records each processed activity IRI and lets the inbox +// handlers drop a duplicate before doing anything. + +import { lt } from "drizzle-orm"; +import { federationProcessedActivities } from "@trails-cool/db/schema/journal"; +import { getDb } from "./db.ts"; + +/** + * Record an inbound activity IRI, returning whether this is the first + * time we've seen it. Insert-or-drop: on primary-key conflict no row is + * inserted and `fresh` is false, so the caller drops the activity as a + * replay before running side effects. + * + * Callers must be idempotent regardless: if a handler fails after this + * records the IRI, Fedify's retry would be dropped here as a duplicate, + * so the follow-graph handlers (recordRemoteFollow/removeRemoteFollow/ + * settle/reject) are all safe to under-run. + */ +export async function markInboundActivityProcessed( + activityIri: string, +): Promise<{ fresh: boolean }> { + const db = getDb(); + const inserted = await db + .insert(federationProcessedActivities) + .values({ activityIri }) + .onConflictDoNothing() + .returning({ iri: federationProcessedActivities.activityIri }); + return { fresh: inserted.length > 0 }; +} + +const THIRTY_DAYS_MS = 30 * 24 * 60 * 60 * 1000; + +/** + * Delete processed-activity rows older than 30 days. Called by the + * `federation-dedup-sweep` job. Returns the number of rows removed. + */ +export async function sweepProcessedActivities(now: Date = new Date()): Promise { + const db = getDb(); + const cutoff = new Date(now.getTime() - THIRTY_DAYS_MS); + const deleted = await db + .delete(federationProcessedActivities) + .where(lt(federationProcessedActivities.receivedAt, cutoff)) + .returning({ iri: federationProcessedActivities.activityIri }); + return deleted.length; +} diff --git a/apps/journal/app/lib/federation.server.ts b/apps/journal/app/lib/federation.server.ts index 190f37a..b461834 100644 --- a/apps/journal/app/lib/federation.server.ts +++ b/apps/journal/app/lib/federation.server.ts @@ -36,6 +36,7 @@ import { getOrigin } from "./config.server.ts"; import { localActorIri } from "./actor-iri.ts"; import { PostgresKvStore } from "./federation-kv.server.ts"; import { PgBossMessageQueue } from "./federation-queue.server.ts"; +import { markInboundActivityProcessed } from "./federation-replay.server.ts"; import { ensureUserKeypair, loadUserKeypair } from "./federation-keys.server.ts"; import { activityToCreate, activityToNote } from "./federation-objects.server.ts"; import { @@ -271,6 +272,7 @@ function buildFederation(): Federation { // when the local target is public; otherwise drop (the actor // already 404s for private users). if (follow.id == null || follow.actorId == null || follow.objectId == null) return; + if (!(await markInboundActivityProcessed(follow.id.href)).fresh) return; // replay: drop const parsed = ctx.parseUri(follow.objectId); if (parsed?.type !== "actor") return; const { outcome } = await recordRemoteFollow(follow.actorId.href, parsed.identifier); @@ -291,6 +293,7 @@ function buildFederation(): Federation { // Spec 4.3: Undo(Follow) removes the follow row. Other Undos are // acknowledged and dropped. if (undo.actorId == null) return; + if (undo.id != null && !(await markInboundActivityProcessed(undo.id.href)).fresh) return; // replay: drop const undoObjectId = undo.objectId; // capture before dereference (see Accept) const object = await undo.getObject(ctx); if (object instanceof Follow && object.objectId != null) { @@ -319,6 +322,7 @@ function buildFederation(): Federation { // Spec 4.4: a remote accepted our outgoing Follow — settle the // Pending row and trigger the first outbox poll for that actor. if (accept.actorId == null) return; + if (accept.id != null && !(await markInboundActivityProcessed(accept.id.href)).fresh) return; // replay: drop // Capture the raw object reference BEFORE dereferencing: // getObject() memoizes the fetched document, after which objectId // reports the fetched object's id (fragment stripped) instead of @@ -358,6 +362,7 @@ function buildFederation(): Federation { .on(Reject, async (ctx, reject) => { // Spec 4.5: remote refused our Follow — drop the Pending row. if (reject.actorId == null) return; + if (reject.id != null && !(await markInboundActivityProcessed(reject.id.href)).fresh) return; // replay: drop const objectId = reject.objectId; // capture before dereference (see Accept) const object = await reject.getObject(ctx); let localUser: Awaited> = null; diff --git a/apps/journal/server.ts b/apps/journal/server.ts index cda9678..3e0d8f4 100644 --- a/apps/journal/server.ts +++ b/apps/journal/server.ts @@ -147,10 +147,11 @@ server.listen(port, async () => { if (process.env.FEDERATION_ENABLED === "true") { const { backfillUserKeypairsJob } = await import("./app/jobs/backfill-user-keypairs.ts"); const { federationKvSweepJob } = await import("./app/jobs/federation-kv-sweep.ts"); + const { federationDedupSweepJob } = await import("./app/jobs/federation-dedup-sweep.ts"); const { deliverActivityJob } = await import("./app/jobs/deliver-activity.ts"); const { pollRemoteActorJob } = await import("./app/jobs/poll-remote-actor.ts"); const { pollRemoteOutboxesJob } = await import("./app/jobs/poll-remote-outboxes.ts"); - jobs.push(backfillUserKeypairsJob, federationKvSweepJob, deliverActivityJob, pollRemoteActorJob, pollRemoteOutboxesJob); + jobs.push(backfillUserKeypairsJob, federationKvSweepJob, federationDedupSweepJob, deliverActivityJob, pollRemoteActorJob, pollRemoteOutboxesJob); } const boss = createBoss(getDatabaseUrl()); diff --git a/openspec/changes/federation-hardening/tasks.md b/openspec/changes/federation-hardening/tasks.md index fe3ebbd..182789a 100644 --- a/openspec/changes/federation-hardening/tasks.md +++ b/openspec/changes/federation-hardening/tasks.md @@ -6,9 +6,9 @@ ## 2. Replay defense -- [ ] 2.1 Add `federation_processed_activities` table (activity IRI PK, received_at) + migration; insert-or-drop check in inbox handlers before side effects -- [ ] 2.2 30-day TTL sweep in the jobs worker -- [ ] 2.3 Tests: duplicate Like/Delete/Update dropped as no-ops; Create double-delivery still covered by remoteOriginIri constraint +- [x] 2.1 Add `federation_processed_activities` table (activity IRI PK, received_at) + migration; insert-or-drop check in inbox handlers before side effects +- [x] 2.2 30-day TTL sweep in the jobs worker +- [x] 2.3 Tests: duplicate Like/Delete/Update dropped as no-ops; Create double-delivery still covered by remoteOriginIri constraint ## 3. Blocklist diff --git a/packages/db/src/schema/journal.ts b/packages/db/src/schema/journal.ts index b13c46c..00b5605 100644 --- a/packages/db/src/schema/journal.ts +++ b/packages/db/src/schema/journal.ts @@ -424,6 +424,21 @@ export const federationKv = journalSchema.table("federation_kv", { expiresAtIdx: index("federation_kv_expires_at_idx").on(t.expiresAt), })); +// Inbound-activity replay defense (spec: federation-operations "Inbound +// replay defense"). One row per processed inbound activity IRI; inbox +// handlers insert-or-drop before side effects so a redelivered activity +// is a no-op. Rows older than 30 days are swept by the +// `federation-dedup-sweep` job — HTTP-signature date freshness already +// rejects older replays — so this never grows unbounded. Create(Note) +// keeps its own idempotency via activities.remote_origin_iri; this table +// covers the follow-graph activities (Follow/Undo/Accept/Reject). +export const federationProcessedActivities = journalSchema.table("federation_processed_activities", { + activityIri: text("activity_iri").primaryKey(), + receivedAt: timestamp("received_at", { withTimezone: true }).notNull().defaultNow(), +}, (t) => ({ + receivedAtIdx: index("federation_processed_activities_received_at_idx").on(t.receivedAt), +})); + // Cache of remote ActivityPub actors we interact with (spec: // social-federation). One row per actor IRI: display fields for feed // cards, inbox/outbox URLs for delivery and polling, the public key for