Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
7.5 kB · 218 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219import { timingSafeEqual } from "node:crypto";import http, { type IncomingMessage, type ServerResponse } from "node:http";import { sha256 } from "../core/json.js";import type { JazzThoughtStore } from "../jazz/store.js";import { parseTelegramBotUpdate, type TelegramBotConnector } from "./telegram-bot.js";
export interface TelegramWebhookServerOptions { connector: TelegramBotConnector; store: JazzThoughtStore; secretToken: string; path: string; host?: "127.0.0.1" | "::1" | "localhost"; port?: number; maxBodyBytes?: number;}
export interface TelegramWebhookServerHandle { server: http.Server; host: string; port: number; path: string; drain(): Promise<void>; close(): Promise<void>;}
export async function startTelegramWebhookServer( options: TelegramWebhookServerOptions,): Promise<TelegramWebhookServerHandle> { const secretToken = webhookSecret(options.secretToken); 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("Telegram webhook receiver may bind only to loopback"); } const requestedPort = boundedInteger(options.port ?? 4_318, "Telegram webhook port", 0, 65_535); const maxBodyBytes = boundedInteger(options.maxBodyBytes ?? 1_048_576, "Telegram webhook maxBodyBytes", 1_024, 8 * 1024 * 1024); let accepting = true; let ingestTail: Promise<void> = Promise.resolve();
const server = http.createServer((request, response) => { void handleRequest(request, response).catch(() => { if (!response.headersSent) send(response, 500, "webhook delivery failed"); else response.destroy(); }); }); server.headersTimeout = 10_000; server.requestTimeout = 15_000; server.keepAliveTimeout = 5_000;
async function handleRequest(request: IncomingMessage, response: ServerResponse): Promise<void> { if (!accepting) { send(response, 503, "webhook receiver is stopping"); return; } const url = new URL(request.url ?? "/", "http://127.0.0.1"); if (url.pathname !== webhookPath || url.search) { send(response, 404, "not found"); return; } if (request.method !== "POST") { response.setHeader("allow", "POST"); send(response, 405, "method not allowed"); return; } const contentType = request.headers["content-type"]?.split(";", 1)[0]?.trim().toLowerCase(); if (contentType !== "application/json") { send(response, 415, "application/json required"); return; } const suppliedSecret = singleHeader(request.headers["x-telegram-bot-api-secret-token"]); if (!suppliedSecret || !constantTimeEqual(suppliedSecret, secretToken)) { send(response, 401, "unauthorized"); return; }
let update: ReturnType<typeof parseTelegramBotUpdate>; try { const body = await readBoundedJson(request, maxBodyBytes); update = parseTelegramBotUpdate(body); } catch (error) { if (error instanceof WebhookRequestError) { send(response, error.status, error.message); return; } send(response, 400, "invalid Telegram update"); return; }
const ingest = ingestTail.then(async () => { await options.connector.ingest(options.store, update); }); ingestTail = ingest.then(() => undefined, () => undefined); try { await ingest; response.writeHead(204, securityHeaders()); response.end(); } catch { send(response, 503, "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("Telegram 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(); }, };}
class WebhookRequestError extends Error { constructor(readonly status: number, message: string) { super(message); this.name = "WebhookRequestError"; }}
async function readBoundedJson(request: IncomingMessage, maxBodyBytes: number): Promise<unknown> { const declaredLength = request.headers["content-length"]; if (declaredLength !== undefined) { const parsed = Number(declaredLength); if (!Number.isSafeInteger(parsed) || parsed < 0) throw new WebhookRequestError(400, "invalid content length"); if (parsed > maxBodyBytes) throw new WebhookRequestError(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 WebhookRequestError(413, "request body too large"); } chunks.push(buffer); } if (received === 0) throw new WebhookRequestError(400, "request body required"); try { return JSON.parse(Buffer.concat(chunks, received).toString("utf8")) as unknown; } catch { throw new WebhookRequestError(400, "invalid JSON"); }}
function normalizedPath(value: string): string { const path = required(value, "Telegram webhook path"); if (!/^\/[A-Za-z0-9/_-]+$/.test(path) || path.includes("//") || path.split("/").includes("..")) { throw new Error("Telegram webhook path must be a normalized absolute path"); } return path;}
function singleHeader(value: string | string[] | undefined): string | undefined { 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 send(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;}
function webhookSecret(value: string): string { const secret = required(value, "Telegram webhook secret"); if (!/^[A-Za-z0-9_-]{1,256}$/.test(secret)) { throw new Error("Telegram webhook secret must use 1-256 ASCII letters, digits, underscores, or hyphens"); } return secret;}