{"version":3,"sources":["../../src/webhook/index.ts","../../src/send-result.ts"],"sourcesContent":["import { CALLBACK_MAP } from \"../send-result\";\nimport type { StreamCallbacks, StreamPart } from \"../types\";\n\nexport type SessionEventsFactory = (\n  sessionId: string,\n  runId: string,\n  metadata: Record<string, string>,\n) => StreamCallbacks;\n\nexport interface WebhookHandlerOptions {\n  secret: string;\n  /** Signature timestamp tolerance in seconds (default: 300). */\n  tolerance?: number;\n  onSessionEvents: SessionEventsFactory;\n}\n\nexport interface WebhookHandlerResult {\n  status: number;\n  body: string | null;\n}\n\nexport interface WebhookHandler {\n  /** Web standard Request/Response (Cloudflare Workers, Bun, Deno, Next.js). */\n  handle(req: Request): Promise<Response>;\n  /** Node.js raw http adapter (Express, Koa, plain http). */\n  express(\n    req: import(\"node:http\").IncomingMessage,\n    res: import(\"node:http\").ServerResponse,\n  ): Promise<void>;\n  /** Framework-agnostic: pass raw body + signature, get back status + body. */\n  handleRaw(\n    rawBody: string,\n    signatureHeader: string | null,\n  ): Promise<WebhookHandlerResult>;\n}\n\ninterface WebhookPayload {\n  sessionId: string;\n  /** Unique identifier for the originating `send()` invocation. */\n  runId: string;\n  sequence: number;\n  timestamp: number;\n  provider: string;\n  metadata: Record<string, string>;\n  event: StreamPart;\n}\n\nexport function createWebhookHandler(\n  options: WebhookHandlerOptions,\n): WebhookHandler {\n  const { secret, tolerance = 300, onSessionEvents } = options;\n\n  async function verifySignature(\n    rawBody: string,\n    signatureHeader: string,\n  ): Promise<boolean> {\n    const parts = signatureHeader.split(\",\");\n    const tPart = parts.find((p) => p.startsWith(\"t=\"));\n    const v1Part = parts.find((p) => p.startsWith(\"v1=\"));\n    if (!tPart || !v1Part) return false;\n\n    const timestamp = Number(tPart.slice(2));\n    const signature = v1Part.slice(3);\n\n    if (Number.isNaN(timestamp)) return false;\n\n    const now = Math.floor(Date.now() / 1000);\n    if (Math.abs(now - timestamp) > tolerance) return false;\n\n    const encoder = new TextEncoder();\n    const key = await crypto.subtle.importKey(\n      \"raw\",\n      encoder.encode(secret),\n      { name: \"HMAC\", hash: \"SHA-256\" },\n      false,\n      [\"sign\"],\n    );\n    const payload = `${timestamp}.${rawBody}`;\n    const expected = await crypto.subtle.sign(\n      \"HMAC\",\n      key,\n      encoder.encode(payload),\n    );\n    const expectedHex = Array.from(new Uint8Array(expected))\n      .map((b) => b.toString(16).padStart(2, \"0\"))\n      .join(\"\");\n\n    return timingSafeEqual(signature, expectedHex);\n  }\n\n  async function dispatch(\n    callbacks: StreamCallbacks,\n    event: StreamPart,\n  ): Promise<void> {\n    await callbacks.onPart?.(event);\n    const key = CALLBACK_MAP[event.type];\n    const cb = callbacks[key] as\n      | ((part: unknown) => void | Promise<void>)\n      | undefined;\n    if (cb) await cb(event);\n  }\n\n  async function processRequest(\n    rawBody: string,\n    signatureHeader: string | null,\n  ): Promise<WebhookHandlerResult> {\n    if (!signatureHeader) {\n      return {\n        status: 401,\n        body: JSON.stringify({ error: \"Missing X-Thalamus-Signature header\" }),\n      };\n    }\n\n    const valid = await verifySignature(rawBody, signatureHeader);\n    if (!valid) {\n      return {\n        status: 401,\n        body: JSON.stringify({ error: \"Invalid signature\" }),\n      };\n    }\n\n    let payload: WebhookPayload;\n    try {\n      payload = JSON.parse(rawBody) as WebhookPayload;\n    } catch {\n      return {\n        status: 400,\n        body: JSON.stringify({ error: \"Malformed JSON body\" }),\n      };\n    }\n\n    const { sessionId, runId, metadata, event } = payload;\n\n    if (!sessionId || !runId || !event?.type) {\n      return {\n        status: 400,\n        body: JSON.stringify({\n          error: \"Missing sessionId, runId, or event\",\n        }),\n      };\n    }\n\n    const callbacks = onSessionEvents(sessionId, runId, metadata ?? {});\n\n    try {\n      await dispatch(callbacks, event);\n    } catch {\n      return { status: 500, body: JSON.stringify({ error: \"Callback error\" }) };\n    }\n\n    return { status: 200, body: null };\n  }\n\n  return {\n    async handleRaw(rawBody, signatureHeader) {\n      return processRequest(rawBody, signatureHeader);\n    },\n\n    async handle(req: Request): Promise<Response> {\n      if (req.method !== \"POST\") {\n        return new Response(JSON.stringify({ error: \"Method not allowed\" }), {\n          status: 405,\n          headers: { \"Content-Type\": \"application/json\" },\n        });\n      }\n\n      const rawBody = await req.text();\n      const signatureHeader = req.headers.get(\"X-Thalamus-Signature\");\n      const result = await processRequest(rawBody, signatureHeader);\n      return new Response(result.body, {\n        status: result.status,\n        headers: result.body ? { \"Content-Type\": \"application/json\" } : {},\n      });\n    },\n\n    async express(req, res): Promise<void> {\n      if (req.method !== \"POST\") {\n        res.writeHead(405, { \"Content-Type\": \"application/json\" });\n        res.end(JSON.stringify({ error: \"Method not allowed\" }));\n        return;\n      }\n\n      const rawBody = await readNodeBody(req);\n      const signatureHeader =\n        (req.headers[\"x-thalamus-signature\"] as string) ?? null;\n      const result = await processRequest(rawBody, signatureHeader);\n      res.writeHead(\n        result.status,\n        result.body ? { \"Content-Type\": \"application/json\" } : {},\n      );\n      res.end(result.body);\n    },\n  };\n}\n\nfunction timingSafeEqual(a: string, b: string): boolean {\n  const encoder = new TextEncoder();\n  const bufA = encoder.encode(a);\n  const bufB = encoder.encode(b);\n  if (bufA.byteLength !== bufB.byteLength) return false;\n\n  let result = 0;\n  for (let i = 0; i < bufA.byteLength; i++) {\n    result |= bufA[i] ^ bufB[i];\n  }\n  return result === 0;\n}\n\nfunction readNodeBody(\n  req: import(\"node:http\").IncomingMessage,\n): Promise<string> {\n  return new Promise((resolve, reject) => {\n    const chunks: Buffer[] = [];\n    req.on(\"data\", (chunk: Buffer) => chunks.push(chunk));\n    req.on(\"end\", () => resolve(Buffer.concat(chunks).toString(\"utf-8\")));\n    req.on(\"error\", reject);\n  });\n}\n","import type {\n  Response,\n  SendResult,\n  StreamCallbacks,\n  StreamPart,\n} from \"./types\";\n\nexport const CALLBACK_MAP: Record<StreamPart[\"type\"], keyof StreamCallbacks> = {\n  \"text-delta\": \"onTextDelta\",\n  thinking: \"onThinking\",\n  refusal: \"onRefusal\",\n  \"tool-use-start\": \"onToolUseStart\",\n  \"tool-use-delta\": \"onToolUseDelta\",\n  \"tool-use-done\": \"onToolUseDone\",\n  \"tool-use-result\": \"onToolUseResult\",\n  \"mcp-tools-discovered\": \"onMcpToolsDiscovered\",\n  \"step-start\": \"onStepStart\",\n  \"step-done\": \"onStepDone\",\n  \"status-change\": \"onStatusChange\",\n  \"stream-start\": \"onStreamStart\",\n  finish: \"onFinish\",\n  error: \"onError\",\n  \"provider-event\": \"onProviderEvent\",\n};\n\nexport interface SendResultOptions {\n  autoStart?: boolean;\n}\n\nclass SendResultImpl implements SendResult {\n  private _promise: Promise<Response> | null = null;\n  private _sessionIdResolve!: (id: string) => void;\n  private readonly _sessionId: Promise<string>;\n\n  constructor(\n    private readonly source: AsyncIterable<StreamPart>,\n    readonly runId: string,\n    private readonly callbacks?: StreamCallbacks,\n    options?: SendResultOptions,\n  ) {\n    this._sessionId = new Promise<string>((resolve) => {\n      this._sessionIdResolve = resolve;\n    });\n\n    if (options?.autoStart) {\n      this._promise = this.run();\n    }\n  }\n\n  get sessionId(): Promise<string> {\n    this._promise ??= this.run();\n    return this._sessionId;\n  }\n\n  get response(): Promise<Response> {\n    this._promise ??= this.run();\n    return this._promise;\n  }\n\n  // biome-ignore lint/suspicious/noThenProperty: intentional PromiseLike implementation\n  then<TResult1 = Response, TResult2 = never>(\n    onfulfilled?:\n      | ((value: Response) => TResult1 | PromiseLike<TResult1>)\n      | null,\n    onrejected?: ((reason: unknown) => TResult2 | PromiseLike<TResult2>) | null,\n  ): Promise<TResult1 | TResult2> {\n    return this.response.then(onfulfilled, onrejected);\n  }\n\n  async text(): Promise<string> {\n    return (await this.response).content;\n  }\n\n  private async run(): Promise<Response> {\n    for await (const part of this.source) {\n      if (part.type === \"stream-start\" && part.sessionId) {\n        this._sessionIdResolve(part.sessionId);\n      }\n      await this.dispatch(part);\n      if (part.type === \"finish\") return part.response;\n      if (part.type === \"error\") throw part.error;\n    }\n    throw new Error(\"Stream ended without a finish event\");\n  }\n\n  private async dispatch(part: StreamPart): Promise<void> {\n    if (!this.callbacks) return;\n    await this.callbacks.onPart?.(part);\n    const key = CALLBACK_MAP[part.type];\n    const cb = this.callbacks[key] as\n      | ((part: StreamPart) => void | Promise<void>)\n      | undefined;\n    if (cb) await cb(part);\n  }\n}\n\nexport function createSendResult(\n  source: AsyncIterable<StreamPart>,\n  runId: string,\n  callbacks?: StreamCallbacks,\n  options?: SendResultOptions,\n): SendResult {\n  return new SendResultImpl(source, runId, callbacks, options);\n}\n"],"mappings":";;;;;;;;;;;;;;;;;;;;AAAA;AAAA;AAAA;AAAA;AAAA;;;ACOO,IAAM,eAAkE;AAAA,EAC7E,cAAc;AAAA,EACd,UAAU;AAAA,EACV,SAAS;AAAA,EACT,kBAAkB;AAAA,EAClB,kBAAkB;AAAA,EAClB,iBAAiB;AAAA,EACjB,mBAAmB;AAAA,EACnB,wBAAwB;AAAA,EACxB,cAAc;AAAA,EACd,aAAa;AAAA,EACb,iBAAiB;AAAA,EACjB,gBAAgB;AAAA,EAChB,QAAQ;AAAA,EACR,OAAO;AAAA,EACP,kBAAkB;AACpB;;;ADwBO,SAAS,qBACd,SACgB;AAChB,QAAM,EAAE,QAAQ,YAAY,KAAK,gBAAgB,IAAI;AAErD,iBAAe,gBACb,SACA,iBACkB;AAClB,UAAM,QAAQ,gBAAgB,MAAM,GAAG;AACvC,UAAM,QAAQ,MAAM,KAAK,CAAC,MAAM,EAAE,WAAW,IAAI,CAAC;AAClD,UAAM,SAAS,MAAM,KAAK,CAAC,MAAM,EAAE,WAAW,KAAK,CAAC;AACpD,QAAI,CAAC,SAAS,CAAC,OAAQ,QAAO;AAE9B,UAAM,YAAY,OAAO,MAAM,MAAM,CAAC,CAAC;AACvC,UAAM,YAAY,OAAO,MAAM,CAAC;AAEhC,QAAI,OAAO,MAAM,SAAS,EAAG,QAAO;AAEpC,UAAM,MAAM,KAAK,MAAM,KAAK,IAAI,IAAI,GAAI;AACxC,QAAI,KAAK,IAAI,MAAM,SAAS,IAAI,UAAW,QAAO;AAElD,UAAM,UAAU,IAAI,YAAY;AAChC,UAAM,MAAM,MAAM,OAAO,OAAO;AAAA,MAC9B;AAAA,MACA,QAAQ,OAAO,MAAM;AAAA,MACrB,EAAE,MAAM,QAAQ,MAAM,UAAU;AAAA,MAChC;AAAA,MACA,CAAC,MAAM;AAAA,IACT;AACA,UAAM,UAAU,GAAG,SAAS,IAAI,OAAO;AACvC,UAAM,WAAW,MAAM,OAAO,OAAO;AAAA,MACnC;AAAA,MACA;AAAA,MACA,QAAQ,OAAO,OAAO;AAAA,IACxB;AACA,UAAM,cAAc,MAAM,KAAK,IAAI,WAAW,QAAQ,CAAC,EACpD,IAAI,CAAC,MAAM,EAAE,SAAS,EAAE,EAAE,SAAS,GAAG,GAAG,CAAC,EAC1C,KAAK,EAAE;AAEV,WAAO,gBAAgB,WAAW,WAAW;AAAA,EAC/C;AAEA,iBAAe,SACb,WACA,OACe;AACf,UAAM,UAAU,SAAS,KAAK;AAC9B,UAAM,MAAM,aAAa,MAAM,IAAI;AACnC,UAAM,KAAK,UAAU,GAAG;AAGxB,QAAI,GAAI,OAAM,GAAG,KAAK;AAAA,EACxB;AAEA,iBAAe,eACb,SACA,iBAC+B;AAC/B,QAAI,CAAC,iBAAiB;AACpB,aAAO;AAAA,QACL,QAAQ;AAAA,QACR,MAAM,KAAK,UAAU,EAAE,OAAO,sCAAsC,CAAC;AAAA,MACvE;AAAA,IACF;AAEA,UAAM,QAAQ,MAAM,gBAAgB,SAAS,eAAe;AAC5D,QAAI,CAAC,OAAO;AACV,aAAO;AAAA,QACL,QAAQ;AAAA,QACR,MAAM,KAAK,UAAU,EAAE,OAAO,oBAAoB,CAAC;AAAA,MACrD;AAAA,IACF;AAEA,QAAI;AACJ,QAAI;AACF,gBAAU,KAAK,MAAM,OAAO;AAAA,IAC9B,QAAQ;AACN,aAAO;AAAA,QACL,QAAQ;AAAA,QACR,MAAM,KAAK,UAAU,EAAE,OAAO,sBAAsB,CAAC;AAAA,MACvD;AAAA,IACF;AAEA,UAAM,EAAE,WAAW,OAAO,UAAU,MAAM,IAAI;AAE9C,QAAI,CAAC,aAAa,CAAC,SAAS,CAAC,OAAO,MAAM;AACxC,aAAO;AAAA,QACL,QAAQ;AAAA,QACR,MAAM,KAAK,UAAU;AAAA,UACnB,OAAO;AAAA,QACT,CAAC;AAAA,MACH;AAAA,IACF;AAEA,UAAM,YAAY,gBAAgB,WAAW,OAAO,YAAY,CAAC,CAAC;AAElE,QAAI;AACF,YAAM,SAAS,WAAW,KAAK;AAAA,IACjC,QAAQ;AACN,aAAO,EAAE,QAAQ,KAAK,MAAM,KAAK,UAAU,EAAE,OAAO,iBAAiB,CAAC,EAAE;AAAA,IAC1E;AAEA,WAAO,EAAE,QAAQ,KAAK,MAAM,KAAK;AAAA,EACnC;AAEA,SAAO;AAAA,IACL,MAAM,UAAU,SAAS,iBAAiB;AACxC,aAAO,eAAe,SAAS,eAAe;AAAA,IAChD;AAAA,IAEA,MAAM,OAAO,KAAiC;AAC5C,UAAI,IAAI,WAAW,QAAQ;AACzB,eAAO,IAAI,SAAS,KAAK,UAAU,EAAE,OAAO,qBAAqB,CAAC,GAAG;AAAA,UACnE,QAAQ;AAAA,UACR,SAAS,EAAE,gBAAgB,mBAAmB;AAAA,QAChD,CAAC;AAAA,MACH;AAEA,YAAM,UAAU,MAAM,IAAI,KAAK;AAC/B,YAAM,kBAAkB,IAAI,QAAQ,IAAI,sBAAsB;AAC9D,YAAM,SAAS,MAAM,eAAe,SAAS,eAAe;AAC5D,aAAO,IAAI,SAAS,OAAO,MAAM;AAAA,QAC/B,QAAQ,OAAO;AAAA,QACf,SAAS,OAAO,OAAO,EAAE,gBAAgB,mBAAmB,IAAI,CAAC;AAAA,MACnE,CAAC;AAAA,IACH;AAAA,IAEA,MAAM,QAAQ,KAAK,KAAoB;AACrC,UAAI,IAAI,WAAW,QAAQ;AACzB,YAAI,UAAU,KAAK,EAAE,gBAAgB,mBAAmB,CAAC;AACzD,YAAI,IAAI,KAAK,UAAU,EAAE,OAAO,qBAAqB,CAAC,CAAC;AACvD;AAAA,MACF;AAEA,YAAM,UAAU,MAAM,aAAa,GAAG;AACtC,YAAM,kBACH,IAAI,QAAQ,sBAAsB,KAAgB;AACrD,YAAM,SAAS,MAAM,eAAe,SAAS,eAAe;AAC5D,UAAI;AAAA,QACF,OAAO;AAAA,QACP,OAAO,OAAO,EAAE,gBAAgB,mBAAmB,IAAI,CAAC;AAAA,MAC1D;AACA,UAAI,IAAI,OAAO,IAAI;AAAA,IACrB;AAAA,EACF;AACF;AAEA,SAAS,gBAAgB,GAAW,GAAoB;AACtD,QAAM,UAAU,IAAI,YAAY;AAChC,QAAM,OAAO,QAAQ,OAAO,CAAC;AAC7B,QAAM,OAAO,QAAQ,OAAO,CAAC;AAC7B,MAAI,KAAK,eAAe,KAAK,WAAY,QAAO;AAEhD,MAAI,SAAS;AACb,WAAS,IAAI,GAAG,IAAI,KAAK,YAAY,KAAK;AACxC,cAAU,KAAK,CAAC,IAAI,KAAK,CAAC;AAAA,EAC5B;AACA,SAAO,WAAW;AACpB;AAEA,SAAS,aACP,KACiB;AACjB,SAAO,IAAI,QAAQ,CAAC,SAAS,WAAW;AACtC,UAAM,SAAmB,CAAC;AAC1B,QAAI,GAAG,QAAQ,CAAC,UAAkB,OAAO,KAAK,KAAK,CAAC;AACpD,QAAI,GAAG,OAAO,MAAM,QAAQ,OAAO,OAAO,MAAM,EAAE,SAAS,OAAO,CAAC,CAAC;AACpE,QAAI,GAAG,SAAS,MAAM;AAAA,EACxB,CAAC;AACH;","names":[]}