diff --git a/src/types.ts b/src/types.ts index c075418..9f744b4 100644 --- a/src/types.ts +++ b/src/types.ts @@ -34,4 +34,8 @@ export interface Env { ROOKERY_KNOT_ADMIN_ADD_MEMBER_URL?: string; /** Secret paired with the knot admin endpoint; configured as a Worker secret. */ ROOKERY_KNOT_ADMIN_SECRET?: string; + /** The operator's alert endpoint for security events; configured as a Worker secret. */ + HUB_WEBHOOK_URL?: string; + /** Shared secret sent as X-Hub-Secret with each alert; configured as a Worker secret. */ + HUB_WEBHOOK_SECRET?: string; } diff --git a/src/worker.ts b/src/worker.ts index c5e5747..426c57e 100644 --- a/src/worker.ts +++ b/src/worker.ts @@ -428,6 +428,62 @@ function syncKnotMember(c: { env: Env; executionCtx: ExecutionContext }, subject ); } +type TakedownAlert = { + actor: string; + did: string; + handle: string; + recordsDeleted: number; + blobsDeleted: number; + collections: string[]; +}; + +async function postTakedownAlert(env: Env, alert: TakedownAlert): Promise { + const endpoint = env.HUB_WEBHOOK_URL?.trim(); + if (!endpoint) { + return; + } + + const response = await fetch(endpoint, { + method: "POST", + headers: { + "Content-Type": "application/json", + "X-Hub-Secret": env.HUB_WEBHOOK_SECRET ?? "", + }, + body: JSON.stringify({ + office: "cso", + ts: new Date().toISOString(), + type: "account_takedown", + tier: "T4", + actor: alert.actor, + did: alert.did, + handle: alert.handle, + records_deleted: alert.recordsDeleted, + blobs_deleted: alert.blobsDeleted, + collections: alert.collections, + }), + }); + if (!response.ok) { + const body = await response.text().catch(() => ""); + throw new Error( + `takedown security alert delivery failed: ${response.status} ${body.slice(0, 200)}`, + ); + } +} + +function alertTakedown( + c: { env: Env; executionCtx: ExecutionContext }, + alert: TakedownAlert, +): void { + c.executionCtx.waitUntil( + postTakedownAlert(c.env, alert).catch((err) => { + console.error("takedown security alert delivery failed", { + did: alert.did, + message: err instanceof Error ? err.message : String(err), + }); + }), + ); +} + async function requestCrawl(env: Env): Promise { const hosts = env.ROOKERY_RELAY_HOSTS?.split(",") .map((host) => host.trim()) @@ -1208,6 +1264,15 @@ app.delete("/admin/accounts/:did", async (c) => { collections: result.collections, }); + alertTakedown(c, { + actor, + did: account.did, + handle: account.handle, + recordsDeleted: result.recordsDeleted, + blobsDeleted: result.blobsDeleted, + collections: result.collections, + }); + return c.json({ did: account.did, handle: account.handle, diff --git a/test/commons.test.ts b/test/commons.test.ts index 8ae687d..f744701 100644 --- a/test/commons.test.ts +++ b/test/commons.test.ts @@ -2,6 +2,7 @@ // Copyright (c) 2026 sol pbc import { afterEach, beforeAll, describe, expect, it, vi } from "vitest"; +import { createExecutionContext, waitOnExecutionContext } from "cloudflare:test"; import { buildAccessToken, createDpopJwt, @@ -13,8 +14,10 @@ import { worker, } from "./helpers"; import { decodeFirst } from "@atcute/cbor"; +import app from "../src/worker"; import { __resetAccessJwksCache } from "../src/access"; import { AccountDurableObject } from "../src/account-do"; +import type { Env } from "../src/types"; import { deactivateAccount, initDirectory, @@ -53,6 +56,14 @@ type SignupResult = { accessToken: string; }; +type TakedownBody = { + did: string; + handle: string; + recordsDeleted: number; + blobsDeleted: number; + collections: string[]; +}; + type AccessFixture = { privateKey: CryptoKey; publicJwk: { @@ -94,12 +105,22 @@ function stubPlcDirectory(options: { body?: string; onPlc?: (url: string) => void | Promise; accessJwk?: AccessFixture["publicJwk"]; + webhook?: { + url: string; + handler: (request: Request) => Response | Promise; + }; } = {}) { const originalFetch = globalThis.fetch.bind(globalThis); return vi .spyOn(globalThis, "fetch") .mockImplementation(async (input, init) => { const url = fetchInputUrl(input); + if (options.webhook && url === options.webhook.url) { + const request = input instanceof Request && init === undefined + ? input + : new Request(input, init); + return options.webhook.handler(request); + } if (url.startsWith("https://plc.directory/")) { await options.onPlc?.(url); return new Response(options.body ?? null, { status: options.status ?? 200 }); @@ -384,6 +405,22 @@ async function adminFetch( return worker.fetch(new Request(`http://localhost${path}`, { ...init, headers })); } +async function directTakedownAccount( + account: SignupResult & { did: string }, + access: AccessFixture, + testEnv: Env, + ctx: ExecutionContext, +): Promise { + const headers = new Headers({ "content-type": "application/json" }); + headers.set("Cf-Access-Jwt-Assertion", await buildAccessAssertion(access)); + const request = new Request(`http://localhost/admin/accounts/${account.did}`, { + method: "DELETE", + headers, + body: JSON.stringify({ confirm: account.body.handle }), + }); + return app.fetch(request, testEnv, ctx); +} + afterEach(() => { __resetAccessJwksCache(); }); @@ -1584,4 +1621,195 @@ describe("commons operator account takedown", () => { fetchSpy.mockRestore(); } }); + + describe("security alert webhook", () => { + it("posts the completed takedown to the operator's alert endpoint", async () => { + const access = await generateAccessFixture(); + const webhookRequests: Request[] = []; + const fetchSpy = stubPlcDirectory({ + accessJwk: access.publicJwk, + webhook: { + url: "https://hub.test/alerts", + handler: (request) => { + webhookRequests.push(request); + return new Response(null, { status: 200 }); + }, + }, + }); + try { + const account = await enrollRookWithInvite("takedown-alert"); + await publishRecord(account, "app.bsky.feed.post", "alert-post"); + await uploadBlob(account, new TextEncoder().encode("alert blob")); + const testEnv: Env = { + ...env, + HUB_WEBHOOK_URL: "https://hub.test/alerts", + HUB_WEBHOOK_SECRET: "test-hub-secret", + }; + const ctx = createExecutionContext(); + + const response = await directTakedownAccount(account, access, testEnv, ctx); + expect(response.status).toBe(200); + const responseBody = await response.json() as TakedownBody; + await waitOnExecutionContext(ctx); + + expect(webhookRequests).toHaveLength(1); + const [request] = webhookRequests; + expect(request.url).toBe("https://hub.test/alerts"); + expect(request.method).toBe("POST"); + expect(request.headers.get("X-Hub-Secret")).toBe("test-hub-secret"); + expect(request.headers.get("Content-Type")).toBe("application/json"); + const body = await request.json() as { + office: string; + ts: string; + type: string; + tier: string; + actor: string; + did: string; + handle: string; + records_deleted: number; + blobs_deleted: number; + collections: string[]; + }; + expect(body).toEqual({ + office: "cso", + ts: body.ts, + type: "account_takedown", + tier: "T4", + actor: "operator@example.com", + did: account.did, + handle: account.body.handle, + records_deleted: responseBody.recordsDeleted, + blobs_deleted: responseBody.blobsDeleted, + collections: responseBody.collections, + }); + expect(new Date(body.ts).toISOString()).toBe(body.ts); + + const audit = await env.DIRECTORY.prepare( + "SELECT actor FROM takedowns WHERE did = ?", + ).bind(account.did).first<{ actor: string }>(); + expect(audit?.actor).toBe(body.actor); + } finally { + fetchSpy.mockRestore(); + } + }); + + it("fails open and logs when the alert fetch rejects", async () => { + const access = await generateAccessFixture(); + const webhookRequests: Request[] = []; + const fetchSpy = stubPlcDirectory({ + accessJwk: access.publicJwk, + webhook: { + url: "https://hub.test/alerts", + handler: (request) => { + webhookRequests.push(request); + return Promise.reject(new Error("test webhook rejected")); + }, + }, + }); + const consoleErrorSpy = vi.spyOn(console, "error").mockImplementation(() => {}); + try { + const account = await enrollRookWithInvite("takedown-alert-reject"); + await publishRecord(account, "app.bsky.feed.post", "reject-post"); + const testEnv: Env = { + ...env, + HUB_WEBHOOK_URL: "https://hub.test/alerts", + HUB_WEBHOOK_SECRET: "test-hub-secret", + }; + const ctx = createExecutionContext(); + + const response = await directTakedownAccount(account, access, testEnv, ctx); + expect(response.status).toBe(200); + const responseBody = await response.json() as TakedownBody; + expect(responseBody).toMatchObject({ recordsDeleted: 1, blobsDeleted: 0 }); + await waitOnExecutionContext(ctx); + + expect(webhookRequests).toHaveLength(1); + expect(consoleErrorSpy).toHaveBeenCalledWith( + "takedown security alert delivery failed", + { did: account.did, message: "test webhook rejected" }, + ); + expect(await accountExists(account.did)).toBe(false); + expect(await env.DIRECTORY.prepare( + "SELECT actor FROM takedowns WHERE did = ?", + ).bind(account.did).first()).toMatchObject({ actor: "operator@example.com" }); + } finally { + consoleErrorSpy.mockRestore(); + fetchSpy.mockRestore(); + } + }); + + it("fails open and logs the status when the alert endpoint returns non-2xx", async () => { + const access = await generateAccessFixture(); + const webhookRequests: Request[] = []; + const fetchSpy = stubPlcDirectory({ + accessJwk: access.publicJwk, + webhook: { + url: "https://hub.test/alerts", + handler: (request) => { + webhookRequests.push(request); + return new Response("hub unavailable", { status: 500 }); + }, + }, + }); + const consoleErrorSpy = vi.spyOn(console, "error").mockImplementation(() => {}); + try { + const account = await enrollRookWithInvite("takedown-alert-500"); + await publishRecord(account, "app.bsky.feed.post", "error-post"); + const testEnv: Env = { + ...env, + HUB_WEBHOOK_URL: "https://hub.test/alerts", + HUB_WEBHOOK_SECRET: "test-hub-secret", + }; + const ctx = createExecutionContext(); + + const response = await directTakedownAccount(account, access, testEnv, ctx); + expect(response.status).toBe(200); + const responseBody = await response.json() as TakedownBody; + expect(responseBody).toMatchObject({ recordsDeleted: 1, blobsDeleted: 0 }); + await waitOnExecutionContext(ctx); + + expect(webhookRequests).toHaveLength(1); + expect(consoleErrorSpy).toHaveBeenCalledWith( + "takedown security alert delivery failed", + { + did: account.did, + message: "takedown security alert delivery failed: 500 hub unavailable", + }, + ); + expect(await accountExists(account.did)).toBe(false); + expect(await env.DIRECTORY.prepare( + "SELECT actor FROM takedowns WHERE did = ?", + ).bind(account.did).first()).toMatchObject({ actor: "operator@example.com" }); + } finally { + consoleErrorSpy.mockRestore(); + fetchSpy.mockRestore(); + } + }); + + it("does not fetch the operator's alert endpoint when the URL is unconfigured", async () => { + const access = await generateAccessFixture(); + const fetchSpy = stubPlcDirectory({ accessJwk: access.publicJwk }); + try { + const account = await enrollRookWithInvite("takedown-alert-disabled"); + const testEnv: Env = { + ...env, + HUB_WEBHOOK_URL: undefined, + HUB_WEBHOOK_SECRET: "test-hub-secret", + }; + const ctx = createExecutionContext(); + + const response = await directTakedownAccount(account, access, testEnv, ctx); + expect(response.status).toBe(200); + await response.json(); + await waitOnExecutionContext(ctx); + + const hubRequests = fetchSpy.mock.calls.filter(([input]) => + fetchInputUrl(input).startsWith("https://hub.test/") + ); + expect(hubRequests).toHaveLength(0); + } finally { + fetchSpy.mockRestore(); + } + }); + }); });