import { createHmac, timingSafeEqual } from "node:crypto"; import http, { type IncomingHttpHeaders, type IncomingMessage, type ServerResponse } from "node:http"; import { sha256 } from "../core/json.js"; import type { JazzThoughtStore } from "../jazz/store.js"; import { isPermanentXActivityError, type XActivityConnector } from "./x-activity.js"; export interface XWebhookServerOptions { connector: XActivityConnector; store: JazzThoughtStore; consumerSecret: string; path: string; host?: "127.0.0.1" | "::1" | "localhost"; port?: number; maxBodyBytes?: number; requestTimeoutMs?: number; pendingRequestLimit?: number; onCrcDiagnostic?: (diagnostic: XCrcDiagnostic) => void; } export interface XCrcDiagnostic { accepted: boolean; queryKeys: string[]; tokenCount: number; tokenLength: number; nonceCount: number; nonceLength: number; } export interface XWebhookServerHandle { server: http.Server; host: string; port: number; path: string; drain(): Promise; close(): Promise; } export async function startXWebhookServer(options: XWebhookServerOptions): Promise { const consumerSecret = required(options.consumerSecret, "X consumer secret"); const webhookPath = normalizedPath(options.path); const host = options.host ?? "127.0.0.1"; if (!["127.0.0.1", "::1", "localhost"].includes(host)) { throw new Error("X webhook receiver may bind only to loopback"); } const requestedPort = boundedInteger(options.port ?? 4_319, "X webhook port", 0, 65_535); const maxBodyBytes = boundedInteger(options.maxBodyBytes ?? 2 * 1024 * 1024, "X webhook maxBodyBytes", 1_024, 8 * 1024 * 1024); const requestTimeoutMs = boundedInteger(options.requestTimeoutMs ?? 8_000, "X webhook requestTimeoutMs", 1_000, 9_000); const pendingRequestLimit = boundedInteger(options.pendingRequestLimit ?? 100, "X webhook pendingRequestLimit", 1, 10_000); let accepting = true; let ingestTail: Promise = Promise.resolve(); let pendingIngests = 0; const server = http.createServer((request, response) => { void handleRequest(request, response).catch(() => { if (!response.headersSent) sendText(response, 500, "webhook request failed"); else response.destroy(); }); }); server.headersTimeout = 10_000; server.requestTimeout = requestTimeoutMs + 1_000; server.keepAliveTimeout = 5_000; async function handleRequest(request: IncomingMessage, response: ServerResponse): Promise { if (!accepting) { sendText(response, 503, "webhook receiver is stopping"); return; } const url = new URL(request.url ?? "/", "http://127.0.0.1"); if (url.pathname !== webhookPath) { sendText(response, 404, "not found"); return; } if (request.method === "GET") { handleCrc(url, response, consumerSecret, options.onCrcDiagnostic); return; } if (request.method !== "POST") { response.setHeader("allow", "GET, POST"); sendText(response, 405, "method not allowed"); return; } if (url.search) { sendText(response, 404, "not found"); return; } const contentType = request.headers["content-type"]?.split(";", 1)[0]?.trim().toLowerCase(); if (contentType !== "application/json") { sendText(response, 415, "application/json required"); return; } const suppliedSignature = singleHeader(request.headers, "x-twitter-webhooks-signature"); if (!suppliedSignature) { sendText(response, 401, "signature required"); return; } let rawBody: Buffer; try { rawBody = await readBoundedBody(request, maxBodyBytes); } catch (error) { if (error instanceof XWebhookRequestError) { sendText(response, error.status, error.message); return; } sendText(response, 400, "invalid request body"); return; } const expectedSignature = xWebhookSignature(rawBody, consumerSecret); if (!constantTimeEqual(suppliedSignature, expectedSignature)) { sendText(response, 401, "invalid signature"); return; } let envelope: unknown; try { envelope = JSON.parse(rawBody.toString("utf8")) as unknown; } catch { sendText(response, 400, "invalid JSON"); return; } const bodySha256 = sha256(rawBody); if (pendingIngests >= pendingRequestLimit) { sendText(response, 503, "webhook ingestion queue is full"); return; } pendingIngests += 1; const ingest = ingestTail.then(async () => { await options.connector.ingest(options.store, envelope, bodySha256); }); const trackedIngest = ingest.finally(() => { pendingIngests -= 1; }); ingestTail = trackedIngest.then(() => undefined, () => undefined); try { await settleWithin(trackedIngest, requestTimeoutMs); response.writeHead(200, { ...securityHeaders(), "content-length": "0" }); response.end(); } catch (error) { if (isPermanentXActivityError(error)) { sendText(response, 400, "invalid X activity payload"); return; } const timedOut = error instanceof XWebhookRequestError && error.status === 503; sendText(response, 503, timedOut ? "durable ingestion timed out" : "durable ingestion failed"); } } await new Promise((resolve, reject) => { server.once("error", reject); server.listen(requestedPort, host, () => { server.off("error", reject); resolve(); }); }); const address = server.address(); if (!address || typeof address === "string") { server.closeAllConnections(); throw new Error("X webhook receiver did not obtain a TCP address"); } return { server, host, port: address.port, path: webhookPath, drain: async () => { await ingestTail; }, close: async () => { accepting = false; const closed = new Promise((resolve) => server.close(() => resolve())); await ingestTail; server.closeIdleConnections(); await Promise.race([closed, new Promise((resolve) => setTimeout(resolve, 2_000))]); server.closeAllConnections(); }, }; } export function xWebhookSignature(body: Buffer | string, consumerSecret: string): string { return `sha256=${createHmac("sha256", required(consumerSecret, "X consumer secret")) .update(body) .digest("base64")}`; } export function xCrcResponseToken(crcToken: string, consumerSecret: string): string { return xWebhookSignature(required(crcToken, "X CRC token"), consumerSecret); } 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"); const nonces = url.searchParams.getAll("nonce"); const token = tokens[0] ?? ""; const nonce = nonces[0] ?? ""; const accepted = keys.length === 1 + nonces.length && keys.every((key) => key === "crc_token" || key === "nonce") && tokens.length === 1 && Boolean(token) && token.length <= 1_024 && nonces.length <= 1 && (nonces.length === 0 || (Boolean(nonce) && nonce.length <= 1_024)); try { onDiagnostic?.({ accepted, queryKeys: keys.slice(0, 10).map((key) => key.slice(0, 100)), tokenCount: tokens.length, tokenLength: token.length, nonceCount: nonces.length, nonceLength: nonce.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(token, consumerSecret) }); response.writeHead(200, { ...securityHeaders(), "content-type": "application/json; charset=utf-8", "content-length": Buffer.byteLength(body), }); response.end(body); } class XWebhookRequestError extends Error { constructor(readonly status: number, message: string) { super(message); this.name = "XWebhookRequestError"; } } async function readBoundedBody(request: IncomingMessage, maxBodyBytes: number): Promise { const declaredLength = request.headers["content-length"]; if (declaredLength !== undefined) { const parsed = Number(declaredLength); if (!Number.isSafeInteger(parsed) || parsed < 0) throw new XWebhookRequestError(400, "invalid content length"); if (parsed > maxBodyBytes) throw new XWebhookRequestError(413, "request body too large"); } const chunks: Buffer[] = []; let received = 0; for await (const chunk of request) { const buffer = Buffer.from(chunk); received += buffer.length; if (received > maxBodyBytes) { request.resume(); throw new XWebhookRequestError(413, "request body too large"); } chunks.push(buffer); } if (received === 0) throw new XWebhookRequestError(400, "request body required"); return Buffer.concat(chunks, received); } async function settleWithin(promise: Promise, timeoutMs: number): Promise { let timer: NodeJS.Timeout | undefined; try { return await Promise.race([ promise, new Promise((_, reject) => { timer = setTimeout(() => reject(new XWebhookRequestError(503, "durable ingestion timed out")), timeoutMs); }), ]); } finally { if (timer) clearTimeout(timer); } } function normalizedPath(value: string): string { const path = required(value, "X webhook path"); if (!/^\/[A-Za-z0-9/_-]+$/.test(path) || path.includes("//") || path.split("/").includes("..")) { throw new Error("X webhook path must be a normalized absolute path"); } return path; } function singleHeader(headers: IncomingHttpHeaders, name: string): string | undefined { const value = headers[name]; return typeof value === "string" ? value : undefined; } function constantTimeEqual(left: string, right: string): boolean { return timingSafeEqual(Buffer.from(sha256(left), "hex"), Buffer.from(sha256(right), "hex")); } function sendText(response: ServerResponse, status: number, body: string): void { response.writeHead(status, { ...securityHeaders(), "content-type": "text/plain; charset=utf-8", "content-length": Buffer.byteLength(body), }); response.end(body); } function securityHeaders(): Record { return { "cache-control": "no-store", "content-security-policy": "default-src 'none'; frame-ancestors 'none'", "referrer-policy": "no-referrer", "x-content-type-options": "nosniff", }; } function boundedInteger(value: number, label: string, minimum: number, maximum: number): number { if (!Number.isSafeInteger(value) || value < minimum || value > maximum) { throw new Error(`${label} must be an integer between ${minimum} and ${maximum}`); } return value; } function required(value: string, label: string): string { const normalized = value.trim(); if (!normalized) throw new Error(`${label} is required`); return normalized; }