trails/apps/journal/app/lib/federation-replay.integration.test.ts
Ullrich Schäfer 105659df7c
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>
2026-07-13 22:07:30 +02:00

59 lines
2.5 KiB
TypeScript

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