File size: 20,538 Bytes
67d18ac | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 | /**
* Async route handlers — wrap the bridge with format-specific translation.
*
* Three handlers:
* - `handleAsyncMessages` (Anthropic client, POST /async/v1/messages) — passthrough
* - `handleAsyncChat` (OpenAI client, POST /async/v1/chat/completions) — request OAI→ANT, response ANT SSE→OAI SSE
* - `handleAsyncHealth` (GET /async/v1/health) — probe queue availability
*
* Common pre-flight (B1 fix: validate BEFORE takeTicket so we never leak a ticket
* on JSON parse / model-missing / translation failures):
* 1. Verify credential has `jwt` (login-captured JWT; absent on imported keys)
* 2. Read + parse client body (skip for health)
* 3. Validate required fields + build the Anthropic-format upstream body
* 4. ONLY THEN takeTicket (any failure above returns 4xx WITHOUT a ticket)
*
* For non-stream (B5+B10): internally force `stream:true` upstream; return a
* chunked `application/json` response that emits legal leading whitespace during
* wait (defeats client TCP idle) and writes the final aggregated JSON at the end.
*
* @see .omo/plans/async-off-peak-bridge.md §3 for full design.
*/
import type { ProxyConfig } from "../config/types.js";
import type { AuthManager } from "../auth/manager.js";
import type { Credential } from "../auth/types.js";
import { credentialString } from "../auth/types.js";
import { errorResponse } from "../proxy/handler.js";
import { transformRequestBody } from "../proxy/body-transformer.js";
import { inflateWithCap } from "../proxy/inflate.js";
import { translateRequestOpenAIToAnthropic, translateResponseAnthropicToOpenAI } from "../translator/openai-to-anthropic.js";
import { anthropicSseToOpenaiSseWithKeepalive } from "./openai-stream-adapter.js";
import type { AnthropicMessagesRequest, OpenAIChatRequest, AnthropicMessagesResponse } from "../translator/types.js";
import { createOffPeakClient, type OffPeakClient } from "./client.js";
import type { OffPeakCredentials, TakeTicketResult } from "./types.js";
import { runAsyncBridge } from "./bridge.js";
/** Cap request body size to prevent memory exhaustion (B16). */
const MAX_REQUEST_BODY_BYTES = 4 * 1024 * 1024;
/** Single space byte — used for non-stream chunked JSON whitespace keepalive. */
const SINGLE_SPACE = new Uint8Array([32]);
export interface AsyncHandlerOptions {
config: ProxyConfig;
auth: AuthManager;
fetchImpl?: (url: string | URL | Request, init?: RequestInit) => Promise<Response>;
debug?: boolean;
}
function buildCredentials(cred: Credential): OffPeakCredentials {
return {
jwt: cred.jwt ?? "",
codingPlanApiKey: credentialString(cred),
};
}
function generateTaskId(): string {
return `proxy-${Date.now().toString(36)}-${Math.random().toString(36).slice(2, 10)}`;
}
function resolveModel(req: { model?: string }, config: ProxyConfig): string {
const explicit = typeof req.model === "string" ? req.model.trim() : "";
if (explicit) return explicit;
if (config.async.defaultModel && config.async.defaultModel.trim()) return config.async.defaultModel.trim();
return config.defaultModel;
}
async function readBody(req: Request): Promise<{ ok: true; body: string } | { ok: false; response: Response }> {
// Reject oversized Content-Length up front; otherwise drain the stream incrementally
// and abort as soon as we exceed the cap. This prevents an attacker from exhausting
// memory by sending a huge chunked body with no Content-Length.
const contentLength = req.headers.get("content-length");
if (contentLength) {
const cl = parseInt(contentLength, 10);
if (Number.isFinite(cl) && cl > MAX_REQUEST_BODY_BYTES) {
// Cancel the request body stream so the underlying socket releases; otherwise
// the client can keep the connection alive despite the 413 response.
req.body?.cancel().catch(() => {});
return { ok: false, response: errorResponse(413, "request_too_large", `body exceeds ${MAX_REQUEST_BODY_BYTES} byte cap`) };
}
}
if (!req.body) {
return { ok: false, response: errorResponse(400, "invalid_request_error", "missing request body") };
}
const reader = req.body.getReader();
const chunks: Uint8Array[] = [];
let total = 0;
try {
while (true) {
const { done, value } = await reader.read();
if (done) break;
total += value.byteLength;
if (total > MAX_REQUEST_BODY_BYTES) {
await reader.cancel().catch(() => {});
return { ok: false, response: errorResponse(413, "request_too_large", `body exceeds ${MAX_REQUEST_BODY_BYTES} byte cap`) };
}
chunks.push(value);
}
} catch {
return { ok: false, response: errorResponse(400, "invalid_request_error", "could not read request body") };
} finally {
reader.releaseLock?.();
}
// Inflate `content-encoding: gzip` request bodies with the cap enforced on
// the DECOMPRESSED size — a small gzip bomb must not bypass the byte cap.
const encoding = req.headers.get("content-encoding")?.toLowerCase().trim() ?? "";
let bytes: Uint8Array = Buffer.concat(chunks);
if (encoding === "gzip" || encoding === "x-gzip") {
const inflated = await inflateWithCap(bytes, MAX_REQUEST_BODY_BYTES);
if (!inflated.ok) {
if (inflated.reason === "too_large") {
return { ok: false, response: errorResponse(413, "request_too_large", `decompressed body exceeds ${MAX_REQUEST_BODY_BYTES} byte cap`) };
}
return { ok: false, response: errorResponse(400, "invalid_request_error", "could not decompress gzip request body") };
}
bytes = inflated.bytes;
}
const body = new TextDecoder().decode(bytes);
if (!body || body.length === 0) {
return { ok: false, response: errorResponse(400, "invalid_request_error", "empty request body") };
}
return { ok: true, body };
}
async function resolveCredential(opts: AsyncHandlerOptions): Promise<{ ok: true; cred: Credential; credentials: OffPeakCredentials } | { ok: false; response: Response }> {
let cred: Credential;
try {
cred = await opts.auth.getCredential();
} catch (err) {
return { ok: false, response: errorResponse(401, "authentication_error", `credential resolution failed: ${(err as Error).message}`) };
}
if (!cred.jwt) {
return {
ok: false,
response: errorResponse(
400,
"async_credentials_unavailable",
"async endpoints require a logged-in oauth credential (JWT missing). Re-run `auth login` or use sync /v1/* endpoints.",
),
};
}
return { ok: true, cred, credentials: buildCredentials(cred) };
}
function buildClient(opts: AsyncHandlerOptions, credentials: OffPeakCredentials): OffPeakClient {
return createOffPeakClient({
origin: opts.config.async.origin,
credentials,
controlTimeoutMs: opts.config.async.controlTimeoutMs,
settleTimeoutMs: opts.config.async.settleTimeoutMs,
fetchImpl: opts.fetchImpl,
});
}
async function takeTicketOr502(client: OffPeakClient, taskId: string, opts: AsyncHandlerOptions, signal: AbortSignal | undefined): Promise<{ ok: true; ticket: TakeTicketResult } | { ok: false; response: Response }> {
try {
const ticket = await client.takeTicket(taskId, signal);
return { ok: true, ticket };
} catch (err) {
return { ok: false, response: errorResponse(502, "async_take_ticket_failed", `off-peak takeTicket failed: ${(err as Error).message}`) };
}
}
function buildBridge(opts: AsyncHandlerOptions, client: OffPeakClient, credentials: OffPeakCredentials, llmRequestBody: string, initialTicket: TakeTicketResult, taskId: string, req: Request) {
return runAsyncBridge({
client,
credentials,
origin: opts.config.async.origin,
identity: opts.config.identity,
llmRequestBody,
initialTicket,
taskId,
pollIntervalMs: opts.config.async.pollIntervalMs,
keepAliveIntervalMs: opts.config.async.keepAliveIntervalMs,
maxRetries: opts.config.async.maxRetries,
maxWaitMs: opts.config.async.maxWaitMs,
clientSignal: req.signal,
fetchImpl: opts.fetchImpl,
onTransition: opts.debug
? (info) => {
console.log(`[async] task=${taskId} ticket=${info.ticketId} phase=${info.phase} attempt=${info.attempt}${info.state ? ` state=${info.state}` : ""}${info.message ? ` msg=${info.message}` : ""}`);
}
: undefined,
});
}
function sseHeaders(): Record<string, string> {
return {
"content-type": "text/event-stream; charset=utf-8",
"cache-control": "no-cache",
connection: "keep-alive",
};
}
export async function handleAsyncMessages(req: Request, opts: AsyncHandlerOptions): Promise<Response> {
// B1: validate everything before ticket acquisition
const cred = await resolveCredential(opts);
if (!cred.ok) return cred.response;
const bodyResult = await readBody(req);
if (!bodyResult.ok) return bodyResult.response;
let parsedBody: Record<string, unknown>;
try {
const raw = JSON.parse(bodyResult.body);
if (raw === null || typeof raw !== "object" || Array.isArray(raw)) {
return errorResponse(400, "invalid_request_error", "request body must be a JSON object");
}
parsedBody = raw as Record<string, unknown>;
} catch {
return errorResponse(400, "invalid_request_error", "request body is not valid JSON");
}
if (!Array.isArray(parsedBody.messages) || parsedBody.messages.length === 0) {
return errorResponse(400, "invalid_request_error", "missing or invalid `messages` field");
}
// Anthropic spec: omitted `stream` defaults to non-streaming (false).
const clientWantsStream = parsedBody.stream === true;
const modelStr = typeof parsedBody.model === "string" ? parsedBody.model : undefined;
// Anthropic spec: omitted `stream` defaults to non-streaming (false).
// Validated: parsedBody is a plain object with messages[]. Remaining fields
// (max_tokens, tools, etc.) are forwarded as-is — upstream rejects invalid shapes.
const upstreamBody = {
...parsedBody,
model: resolveModel({ model: modelStr }, opts.config),
stream: true,
} as AnthropicMessagesRequest;
const upstreamBodyText = transformRequestBody(JSON.stringify(upstreamBody), { format: "anthropic", userId: cred.cred.userId }) ?? JSON.stringify(upstreamBody);
// Now we're safe to take a ticket
const client = buildClient(opts, cred.credentials);
const taskId = generateTaskId();
const ticket = await takeTicketOr502(client, taskId, opts, req.signal);
if (!ticket.ok) return ticket.response;
const { stream, outcome } = buildBridge(opts, client, cred.credentials, upstreamBodyText, ticket.ticket, taskId, req);
void outcome;
if (clientWantsStream) {
return new Response(stream, { status: 200, headers: sseHeaders() });
}
// Non-stream: chunked response with leading whitespace during wait + final JSON (B10).
// NOTE: no explicit `transfer-encoding` header — it is a forbidden Response
// constructor header (runtimes drop/override it) and Node http already sends
// chunked when no content-length is set.
return new Response(nonStreamChunkedJson(stream), {
status: 200,
headers: { "content-type": "application/json; charset=utf-8", "cache-control": "no-cache" },
});
}
export async function handleAsyncChat(req: Request, opts: AsyncHandlerOptions): Promise<Response> {
const cred = await resolveCredential(opts);
if (!cred.ok) return cred.response;
const bodyResult = await readBody(req);
if (!bodyResult.ok) return bodyResult.response;
let openaiReq: OpenAIChatRequest;
try {
openaiReq = JSON.parse(bodyResult.body) as OpenAIChatRequest;
} catch {
return errorResponse(400, "invalid_request_error", "request body is not valid JSON");
}
if (!Array.isArray(openaiReq.messages) || openaiReq.messages.length === 0) {
return errorResponse(400, "invalid_request_error", "missing or invalid `messages` field");
}
openaiReq.model = resolveModel(openaiReq, opts.config);
const clientWantsStream = openaiReq.stream === true;
let anthropicReq: AnthropicMessagesRequest;
try {
anthropicReq = translateRequestOpenAIToAnthropic(openaiReq);
} catch (err) {
return errorResponse(400, "invalid_request_error", `OpenAI→Anthropic translation failed: ${(err as Error).message}`);
}
anthropicReq.stream = true;
const upstreamBodyText = transformRequestBody(JSON.stringify(anthropicReq), { format: "anthropic", userId: cred.cred.userId }) ?? JSON.stringify(anthropicReq);
const client = buildClient(opts, cred.credentials);
const taskId = generateTaskId();
const ticket = await takeTicketOr502(client, taskId, opts, req.signal);
if (!ticket.ok) return ticket.response;
const { stream: rawStream, outcome } = buildBridge(opts, client, cred.credentials, upstreamBodyText, ticket.ticket, taskId, req);
void outcome;
if (clientWantsStream) {
// B4: custom translator that preserves `: keepalive` comments and converts Anthropic errors
const openaiStream = anthropicSseToOpenaiSseWithKeepalive(rawStream, openaiReq.model);
return new Response(openaiStream, { status: 200, headers: sseHeaders() });
}
return new Response(nonStreamChunkedJson(rawStream, { translate: "openai", model: openaiReq.model }), {
status: 200,
headers: { "content-type": "application/json; charset=utf-8", "cache-control": "no-cache" },
});
}
export async function handleAsyncHealth(_req: Request, opts: AsyncHandlerOptions): Promise<Response> {
const cred = await resolveCredential(opts);
if (!cred.ok) return cred.response;
const client = buildClient(opts, cred.credentials);
try {
const avail = await client.getAvailability();
return new Response(JSON.stringify(avail), { status: 200, headers: { "content-type": "application/json" } });
} catch (err) {
return errorResponse(502, "async_health_failed", (err as Error).message);
}
}
/**
* Wrap the SSE byte stream as a non-stream JSON response. Emits leading whitespace
* during ticket-queue wait (defeats client TCP idle), then a single JSON document.
*
* Two modes:
* - default: Anthropic batch JSON shape
* - {translate: "openai"}: OpenAI batch JSON shape (translated from Anthropic)
*/
function nonStreamChunkedJson(
bridgeStream: ReadableStream<Uint8Array>,
translateOpts?: { translate: "openai"; model: string },
): ReadableStream<Uint8Array> {
const encoder = new TextEncoder();
return new ReadableStream<Uint8Array>({
async start(controller) {
const reader = bridgeStream.getReader();
const decoder = new TextDecoder();
let sseBuffer = "";
try {
while (true) {
const { done, value } = await reader.read();
if (done) break;
// One space byte per received chunk (not per char). Resets TCP idle timer
// while keeping allocation count proportional to chunk count, not byte count.
try {
controller.enqueue(SINGLE_SPACE);
} catch {
return;
}
sseBuffer += decoder.decode(value, { stream: true });
}
sseBuffer += decoder.decode();
} finally {
reader.releaseLock?.();
}
// Reconstruct Anthropic batch JSON from accumulated SSE
const anthropicMsg = reconstructAnthropicBatch(sseBuffer);
if (!anthropicMsg) {
const errPayload = { error: { type: "async_aggregation_failed", message: "could not reconstruct response from bridge stream" } };
try {
controller.enqueue(encoder.encode(JSON.stringify(errPayload)));
} catch {
// closed
}
controller.close();
return;
}
const finalJson = translateOpts?.translate === "openai"
? JSON.stringify(translateResponseAnthropicToOpenAI(anthropicMsg, translateOpts.model))
: JSON.stringify(anthropicMsg);
try {
controller.enqueue(encoder.encode(finalJson));
} catch {
// closed
}
controller.close();
},
});
}
/**
* Reconstruct a synthetic `AnthropicMessagesResponse` from a stream of Anthropic SSE bytes.
* Handles message_start, content_block_start/delta/stop, message_delta, message_stop.
*
* Fail-closed: returns null if `message_stop` not seen, or on `event: error`.
* Preserves `signature_delta` for thinking blocks. No production `any`.
*/
function reconstructAnthropicBatch(sseText: string): AnthropicMessagesResponse | null {
const blocks = sseText.split("\n\n");
type ContentBlock =
| { type: "text"; text: string }
| { type: "thinking"; thinking: string; signature?: string }
| { type: "tool_use"; id: string; name: string; input: unknown };
let message: Partial<AnthropicMessagesResponse> | null = null;
const content: ContentBlock[] = [];
let currentBlock: ContentBlock | null = null;
let currentToolJson = "";
let sawMessageStop = false;
let sawError = false;
for (const block of blocks) {
const lines = block.split("\n");
let eventType: string | undefined;
let data: string | undefined;
for (const line of lines) {
if (line.startsWith("event:")) eventType = line.slice(6).trim();
else if (line.startsWith("data:")) data = line.slice(5).trim();
}
if (!data) continue;
let parsed: Record<string, unknown>;
try {
parsed = JSON.parse(data) as Record<string, unknown>;
} catch {
continue;
}
const type = (eventType ?? parsed.type) as string;
switch (type) {
case "message_start": {
const msg = parsed.message as Partial<AnthropicMessagesResponse> | undefined;
message = { ...(msg ?? {}) };
break;
}
case "content_block_start": {
const cb = parsed.content_block as Partial<ContentBlock> | undefined;
if (!cb || !cb.type) break;
if (cb.type === "text") currentBlock = { type: "text", text: "" };
else if (cb.type === "thinking") currentBlock = { type: "thinking", thinking: "" };
else if (cb.type === "tool_use" && typeof cb.id === "string" && typeof cb.name === "string") {
currentBlock = { type: "tool_use", id: cb.id, name: cb.name, input: {} };
currentToolJson = "";
}
break;
}
case "content_block_delta": {
const delta = parsed.delta as Record<string, unknown> | undefined;
if (!currentBlock || !delta) break;
if (delta.type === "text_delta" && currentBlock.type === "text" && typeof delta.text === "string") {
currentBlock.text += delta.text;
} else if (delta.type === "thinking_delta" && currentBlock.type === "thinking" && typeof delta.thinking === "string") {
currentBlock.thinking += delta.thinking;
} else if (delta.type === "signature_delta" && currentBlock.type === "thinking" && typeof delta.signature === "string") {
currentBlock.signature = (currentBlock.signature ?? "") + delta.signature;
} else if (delta.type === "input_json_delta" && currentBlock.type === "tool_use" && typeof delta.partial_json === "string") {
currentToolJson += delta.partial_json;
}
break;
}
case "content_block_stop": {
if (currentBlock) {
if (currentBlock.type === "tool_use") {
try {
currentBlock.input = JSON.parse(currentToolJson || "{}");
} catch {
currentBlock.input = {};
}
currentToolJson = "";
}
content.push(currentBlock);
currentBlock = null;
}
break;
}
case "message_delta": {
const delta = parsed.delta as Partial<AnthropicMessagesResponse> | undefined;
const usage = parsed.usage as Record<string, number> | undefined;
if (delta && message) Object.assign(message, delta);
if (usage && message) message.usage = { ...(message.usage ?? { input_tokens: 0, output_tokens: 0 }), ...usage } as AnthropicMessagesResponse["usage"];
break;
}
case "message_stop":
sawMessageStop = true;
break;
case "error":
sawError = true;
break;
default:
// ignore ping / unknown
break;
}
}
if (sawError || !sawMessageStop || !message) return null;
message.content = content as AnthropicMessagesResponse["content"];
if (!message.stop_reason) message.stop_reason = "end_turn";
if (!message.role) message.role = "assistant";
if (!message.usage) message.usage = { input_tokens: 0, output_tokens: 0 };
return message as AnthropicMessagesResponse;
}
|