feat(journal): inbound federation replay defense

Task group 2 of federation-hardening. The narrow inbox
(Follow/Undo/Accept/Reject) had no replay protection — only Create(Note)
did, via the activities.remote_origin_iri unique constraint. A remote
redelivering a signed follow-graph activity would re-run its side
effects.

- New `federation_processed_activities` table (activity IRI PK,
  received_at + index). Additive, so drizzle-kit push creates it; no
  hand-written migration needed.
- `federation-replay.server.ts`: `markInboundActivityProcessed` does an
  insert-or-drop (ON CONFLICT DO NOTHING RETURNING) and reports whether
  the IRI is fresh; `sweepProcessedActivities` deletes rows > 30 days old
  (signature date-freshness already rejects older replays).
- Each inbox listener drops a duplicate before side effects. The
  follow-graph handlers are idempotent, so a handler failure whose retry
  is later dropped as a duplicate can't corrupt state.
- `federation-dedup-sweep` job (daily 04:30 UTC) runs the TTL sweep.

Verified: db + journal typecheck + lint clean; drizzle-kit push creates
the table; replay integration test (fresh-vs-duplicate + 30-day sweep)
green against real Postgres; journal unit suite 355 passing.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Ullrich Schäfer 2026-07-13 21:55:44 +02:00
parent 71d92306db
commit 105659df7c
No known key found for this signature in database
GPG key ID: A32FF691A0F752D9
8 changed files with 157 additions and 4 deletions

View file

@ -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 };
},
});

View file

@ -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;

View file

@ -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);
});
});

View file

@ -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<number> {
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;
}

View file

@ -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<void> {
// 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<void> {
// 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<void> {
// 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<void> {
.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<ReturnType<typeof findLocalPublicUserByIri>> = null;

View file

@ -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());

View file

@ -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

View file

@ -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