Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
11 kB · 317 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318import { 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<void>; close(): Promise<void>;}
export async function startXWebhookServer(options: XWebhookServerOptions): Promise<XWebhookServerHandle> { 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<void> = 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<void> { 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<void>((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<void>((resolve) => server.close(() => resolve())); await ingestTail; server.closeIdleConnections(); await Promise.race([closed, new Promise<void>((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<Buffer> { 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<T>(promise: Promise<T>, timeoutMs: number): Promise<T> { let timer: NodeJS.Timeout | undefined; try { return await Promise.race([ promise, new Promise<never>((_, 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<string, string> { 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;}