import type { IncomingMessage, ServerResponse } from "node:http"; import type { AssistantMessage } from "./anthropic.js"; import type { CrawlWatchdogConfig } from "./config-types.js"; /** * The code `readBody` sets when it refused a body for exceeding the cap. Owned HERE, by the * thrower, since 2026-09-04 (contract review DR-005): `dashboard-routes.ts` re-exports it for its * own classifier, and `server.ts` classifies 413-vs-400 through `bodyReadStatus` below. Both * consumers read this TAG. Until this date the data plane still regex-matched the MESSAGE * ("too large") — the relay inferring its own intent from prose it had written itself, the * inference the dashboard route had already stopped making. */ export declare const BODY_TOO_LARGE_CODE = "ERR_DASHBOARD_BODY_TOO_LARGE"; /** 413 when `readBody` refused the body for size; 400 for any other body-read failure. */ export declare function bodyReadStatus(error: unknown): 413 | 400; export declare const MAX_VALIDATE_BYTES: number; /** 25 MiB decoded document × base64 expansion, plus JSON-envelope headroom. */ export declare const DEFAULT_MAX_BODY_BYTES: number; export declare const HOP_BY_HOP: Set; /** * Write a chunk to a ServerResponse while properly tracking backpressure. * Invokes onCommit on the first successful byte write. */ export declare function writeChunk(res: ServerResponse, bytes: Buffer | Uint8Array, onCommit?: () => void): Promise; /** * Read the full body of an IncomingMessage up to maxBytes. * Throws 413 error if body exceeds maxBytes. */ export declare function readBody(req: IncomingMessage, maxBytes?: number): Promise; /** * Wrap a streaming Response with an inter-chunk stall watchdog. */ export declare function withStallWatchdog(response: Response, controller: AbortController, stallTimeoutMs: number): Response; /** * Default crawl-abort threshold: ms per output token, sustained over a full trailing window, above * which a COMMITTED stream is judged CRAWLING rather than merely producing a long answer. * * Calibrated 2026-09-09, on the `latency-demotion.ts` precedent: 4 x `DEFAULT_LATENCY_MS_PER_TOKEN` * (250 ms/token, itself measured over 68 real requests on 2026-08-30 — see that file's own * comment). A crawl abort hands the client a failure it must retry — measured in * `docs/history/post-commit-stall-measurement-2026-09-09.md`: Claude Code retries once, downgraded to a * NON-STREAMING request; Codex retries up to five times, staying streaming, in all four measured * cells — so the bar must sit far above the demotion threshold, or a deployment merely slow enough * to be latency-demoted would also be aborted mid-response. 250 ms/token is itself ~3.5x the * healthy band measured that day, so 1000 ms/token sits well clear of ordinary slow-but-working * traffic. A tunable default, not a provider fact — the provenance invariant permits this. * Re-calibrate by reading `~/.llm-relay/usage/recent.json` the same way `latency-demotion.ts` * describes; there is no dedicated crawl-abort log field to read back yet (see the accepted gaps * in `docs/backlog.md`'s "Build the post-commit CRAWL abort" entry once it is amended). */ export declare const DEFAULT_CRAWL_MS_PER_TOKEN = 1000; /** * Width of the trailing window the crawl rate is measured over, in ms — and, since the rate is * `windowMs / tokensInWindow` (see `withCrawlWatchdog`), also the numerator of every rate this * watchdog ever computes. * * Corrected 2026-09-09 (fixing a same-day defect the packet that introduced this watchdog shipped * uncaught): the ORIGINAL pairing of `windowMs: 20_000` with `minTokens: 50` could never fire. * That version measured `spanMs = min(elapsed, windowMs)` — so `spanMs` never exceeded 20 000 — * and required `tokensInWindow >= minTokens` (50) before judging `rate = spanMs / tokensInWindow`. * The worst case allowed was 20 000 ms / 50 tokens = 400 ms/token, which can never clear a 1000 * ms/token threshold: the three defaults were mutually inconsistent by construction, and the * measured scenario this watchdog exists for (60 tokens over 90 s, one every 1.5 s) could not trip * it either. The rule is now: judge only once a FULL window has elapsed since commit * (`elapsed >= windowMs`), then take `tokensInWindow` as the tokens sampled in the trailing * `windowMs` and compute `rate = windowMs / tokensInWindow` (a fixed numerator, not a growing * `spanMs`) — so with `windowMs: 30_000` and `msPerToken: 1000`, fewer than 30 tokens landing in * any trailing 30 s window trips the abort, which the 90 s/60-token scenario clears easily (one * token per 1.5 s is roughly 20 tokens per 30 s window). */ export declare const DEFAULT_CRAWL_WINDOW_MS = 30000; /** * Minimum output tokens that must be observed SINCE COMMIT — across the whole stream, not just the * trailing window — before ANY judgement runs at all: evidence the stream is producing an answer * in the first place, distinct from `tokensInWindow` below. A window judged before this gate clears * would be able to abort a stream that has barely started, on the strength of a single early burst * falling silent — exactly the case the SEPARATE "full window holding zero tokens" rule below * already declines to judge, stated as its own gate so the two can be tested apart. */ export declare const DEFAULT_CRAWL_MIN_TOKENS = 20; export interface CrawlWatchdogSettings { enabled: boolean; msPerToken: number; windowMs: number; minTokens: number; } /** * Resolve `routing.crawl` into a total settings object. Absent, `{}`, or any missing key means * the tunable default for that key — the `resolveHedgeSettings`/`resolveLatencyDemotion` * precedent. * * The rule `withCrawlWatchdog` enforces, in words: `minTokens` (default 20) is the minimum number * of output tokens observed since commit — over the WHOLE stream — before any judgement runs at * all, evidence the stream is producing an answer. A window is judged only once it is FULL — * elapsed time since commit at least `windowMs` (default 30 000). `tokensInWindow` counts only the * tokens whose sample time falls inside the trailing `windowMs`; a full window holding zero tokens * yields no opinion at all (silence is `withStallWatchdog`'s job, and this watchdog must never * pre-empt it). Otherwise `rate = windowMs / tokensInWindow`, and the stream is CRAWLING — the * fetch is aborted — when `rate > msPerToken` (default 1000). */ export declare function resolveCrawlSettings(raw: CrawlWatchdogConfig | undefined): CrawlWatchdogSettings; /** * The CLIENT-facing wire shape the crawl watchdog reads deltas from. Deliberately the same three * members as `stream-commit.ts`'s `StreamCommitProtocol` (not imported from there — this module * stays a leaf the way `sse-frames.ts` does, and the membership is copied rather than re-exported * so a change to one is never mistaken for a change to the other). */ export type CrawlProtocol = "anthropic-messages" | "openai-chat" | "openai-responses"; /** * A committed stream the relay itself terminated because its measured per-token rate, over the * sliding window, stayed worse than the configured threshold. * * `message` is EXACTLY the text the client-visible SSE `error` frame carries — no * "backend stream failed mid-response" wrapping prefix — because the backlog entry ("Build the * post-commit CRAWL abort") pins the wording. `candidate-runner.ts` `handleMidStreamError` detects * this type via `AbortSignal.reason` (never by inspecting the caught exception, which may be an * opaque `AbortError` from the fetch machinery rather than this object) and uses `.message` * unwrapped, plus a distinct `errorKinds` member. */ export declare class CrawlAbortedError extends Error { readonly rateMsPerToken: number; readonly windowMs: number; readonly thresholdMsPerToken: number; constructor(rateMsPerToken: number, windowMs: number, thresholdMsPerToken: number); } /** * Wrap a COMMITTED streaming Response with a crawl watchdog: abort the backend fetch once a FULL * trailing window has elapsed since commit and that window's per-token rate — `settings.windowMs * / tokensInWindow` — exceeds `settings.msPerToken`. Two independent gates guard against judging * too early: `settings.minTokens` tokens must have been observed SINCE COMMIT (the whole stream, * not just the trailing window) before any judgement runs at all, and a full window holding ZERO * tokens yields no opinion rather than an abort (unmeasured is no opinion, never slow — the * `latency-demotion.ts` precedent; silence is `withStallWatchdog`'s job and this watchdog must * never pre-empt it). See `resolveCrawlSettings` for the rule spelled out in full. * * Installed at the SAME call site as `withStallWatchdog`, after the commit probe has already * rebuilt the response from its replayed prefix — so `now() - commitTime` (captured at * installation) approximates the true post-commit elapsed time. * * ⚠ The abort carries the `CrawlAbortedError` as its `AbortSignal.reason` — the same mechanism * `controller.abort(reason)` offers natively — rather than erroring the transform's own * `ReadableStreamDefaultController` directly. Enqueuing the triggering chunk and then erroring * the SAME controller in one microtask risks losing that already-enqueued-but-unread chunk (the * ReadableStream spec does not guarantee it survives an immediately following `error()`); routing * the abort through the existing `AbortController` — exactly how `withStallWatchdog` already ends * a stream — sidesteps that risk entirely and reuses a path already proven correct. */ export declare function withCrawlWatchdog(response: Response, controller: AbortController, protocol: CrawlProtocol, settings: CrawlWatchdogSettings, now?: () => number): Response; /** Fail-closed terminal error response for standard endpoints. */ export declare function failClosed(res: ServerResponse, status: number, message: string, extraHeaders?: Record, errorType?: string): void; /** Forward a local Response descriptor directly to ServerResponse. */ /** * Forward a response the RELAY authored — a `RequestMappingError` 400, a `DocumentError` 400, a * dialect destructive refusal — to the client. * * ⚠ The body is buffered BEFORE the head is committed, and that order is the point. This is the * local-failure exit of both candidate loops, where the contract is to fail CLEAN: while the head * is unsent the caller can still answer with a proper status, so a body read that rejects must * throw before `writeHead`, never after it. Committing the head first and then streaming leaves a * truncated body under an already-sent status, which is the one outcome this path exists to avoid. * The bodies are small and relay-authored, so buffering costs nothing. */ export declare function forwardLocalResponse(res: ServerResponse, local: Response): Promise; /** Parse an AssistantMessage from raw JSON string if valid. */ /** * Parse a buffered backend body into an `AssistantMessage`. * * ⚠ The shape test is `content` being an array — NOT `role === "assistant"`. `AssistantMessage` * declares no `role` field at all, so gating on one makes the parser demand something outside its * own contract and silently answer null (validation and repair are then skipped) for any body * that omits it. * * ⚠ The fields are copied deliberately rather than cast wholesale: `stop_reason` normalizes an * absent value to `null`, because `emitSse` writes the field back onto the wire and `undefined` * omits the key where `null` states it. */ export declare function parseAssistant(text: string): AssistantMessage | null; /** Find end of double newline in SSE stream. */ export declare function frameEnd(buf: Buffer): number; /** Check if an SSE frame opens a tool_use block. */ /** * If `frame` is a `content_block_start` event opening a `tool_use` block, return its block index; * otherwise null. This is the trigger to start withholding. * * ⚠ The `data:` lines are COLLECTED and JOINED before parsing, and the split tolerates CRLF — the * `sse-frames.ts` convention. SSE permits a payload spread over several `data:` lines, so parsing * each line on its own answers null for exactly the multi-line frame this exists to catch. */ export declare function frameOpensToolUse(frame: Buffer): number | null; /** Build Anthropic SSE error event frame. */ export declare function sseError(message: string): string; /** Build OpenAI SSE error event frame. */ export declare function openAiSseError(message: string): string;