diff --git a/src/cli.ts b/src/cli.ts index f6def0c..1d1b112 100644 --- a/src/cli.ts +++ b/src/cli.ts @@ -304,6 +304,12 @@ try { maxBodyBytes: sourceConfig.maxBodyBytes, requestTimeoutMs: sourceConfig.requestTimeoutMs, pendingRequestLimit: sourceConfig.pendingRequestLimit, + onCrcDiagnostic: (diagnostic) => print({ + xWebhookCrc: { + source: sourceConfig.id, + ...diagnostic, + }, + }), }); print({ xWebhook: { diff --git a/src/connectors/x-webhook.ts b/src/connectors/x-webhook.ts index f3d1e0b..037fabf 100644 --- a/src/connectors/x-webhook.ts +++ b/src/connectors/x-webhook.ts @@ -14,6 +14,14 @@ export interface XWebhookServerOptions { maxBodyBytes?: number; requestTimeoutMs?: number; pendingRequestLimit?: number; + onCrcDiagnostic?: (diagnostic: XCrcDiagnostic) => void; +} + +export interface XCrcDiagnostic { + accepted: boolean; + queryKeys: string[]; + tokenCount: number; + tokenLength: number; } export interface XWebhookServerHandle { @@ -61,7 +69,7 @@ export async function startXWebhookServer(options: XWebhookServerOptions): Promi return; } if (request.method === "GET") { - handleCrc(url, response, consumerSecret); + handleCrc(url, response, consumerSecret, options.onCrcDiagnostic); return; } if (request.method !== "POST") { @@ -172,14 +180,35 @@ export function xCrcResponseToken(crcToken: string, consumerSecret: string): str return xWebhookSignature(required(crcToken, "X CRC token"), consumerSecret); } -function handleCrc(url: URL, response: ServerResponse, consumerSecret: string): void { +function handleCrc( + url: URL, + response: ServerResponse, + consumerSecret: string, + onDiagnostic?: (diagnostic: XCrcDiagnostic) => void, +): void { const keys = [...url.searchParams.keys()]; const tokens = url.searchParams.getAll("crc_token"); - if (keys.length !== 1 || keys[0] !== "crc_token" || tokens.length !== 1 || tokens[0]!.length > 1_024 || !tokens[0]) { + const token = tokens[0] ?? ""; + const accepted = keys.length === 1 + && keys[0] === "crc_token" + && tokens.length === 1 + && Boolean(token) + && token.length <= 1_024; + try { + onDiagnostic?.({ + accepted, + queryKeys: keys.slice(0, 10).map((key) => key.slice(0, 100)), + tokenCount: tokens.length, + tokenLength: token.length, + }); + } catch { + // Observability must not alter X's challenge-response protocol. + } + if (!accepted) { sendText(response, 400, "one crc_token is required"); return; } - const body = JSON.stringify({ response_token: xCrcResponseToken(tokens[0], consumerSecret) }); + const body = JSON.stringify({ response_token: xCrcResponseToken(token, consumerSecret) }); response.writeHead(200, { ...securityHeaders(), "content-type": "application/json; charset=utf-8", diff --git a/test/x-webhook.test.ts b/test/x-webhook.test.ts index 790b3ac..3eb7d13 100644 --- a/test/x-webhook.test.ts +++ b/test/x-webhook.test.ts @@ -2,7 +2,12 @@ import fs from "node:fs/promises"; import path from "node:path"; import { afterEach, describe, expect, test, vi } from "vitest"; import { XActivityConnector } from "../src/connectors/x-activity.js"; -import { xCrcResponseToken, startXWebhookServer, xWebhookSignature } from "../src/connectors/x-webhook.js"; +import { + xCrcResponseToken, + startXWebhookServer, + xWebhookSignature, + type XCrcDiagnostic, +} from "../src/connectors/x-webhook.js"; import { sha256 } from "../src/core/json.js"; import type { JazzThoughtStore } from "../src/jazz/store.js"; import { temporaryProject, testStore } from "./helpers.js"; @@ -51,7 +56,8 @@ describe("X Activity webhook", () => { }); test("answers CRC without opening or mutating the event store", async () => { - const fixture = await setup(); + const diagnostics: XCrcDiagnostic[] = []; + const fixture = await setup({ onCrcDiagnostic: (diagnostic) => diagnostics.push(diagnostic) }); const append = vi.spyOn(fixture.store, "appendEvent"); const response = await fetch(`${fixture.endpoint}?crc_token=fixture-token`); expect(response.status).toBe(200); @@ -61,6 +67,12 @@ describe("X Activity webhook", () => { expect(append).not.toHaveBeenCalled(); expect((await fetch(`${fixture.endpoint}?crc_token=one&extra=two`)).status).toBe(400); expect((await fetch(`${fixture.endpoint}?crc_token=one&crc_token=two`)).status).toBe(400); + expect(diagnostics).toEqual([ + { accepted: true, queryKeys: ["crc_token"], tokenCount: 1, tokenLength: 13 }, + { accepted: false, queryKeys: ["crc_token", "extra"], tokenCount: 1, tokenLength: 3 }, + { accepted: false, queryKeys: ["crc_token", "crc_token"], tokenCount: 2, tokenLength: 3 }, + ]); + expect(JSON.stringify(diagnostics)).not.toContain("fixture-token"); await fixture.receiver.close(); }); @@ -328,6 +340,7 @@ async function setup(options: { maxBodyBytes?: number; requestTimeoutMs?: number; pendingRequestLimit?: number; + onCrcDiagnostic?: (diagnostic: XCrcDiagnostic) => void; } = {}) { const root = await temporaryProject("thoughtstream-x-webhook-"); roots.push(root); @@ -351,6 +364,7 @@ async function setup(options: { maxBodyBytes: options.maxBodyBytes ?? 2 * 1024 * 1024, requestTimeoutMs: options.requestTimeoutMs ?? 8_000, pendingRequestLimit: options.pendingRequestLimit ?? 100, + ...(options.onCrcDiagnostic ? { onCrcDiagnostic: options.onCrcDiagnostic } : {}), }); return { store, connector, receiver, endpoint: `http://${receiver.host}:${receiver.port}${receiver.path}` }; }