/**
* HTTP helpers for LLM providers.
*
* Similar in spirit to `core/embedding/fetcher.ts`, but LLM calls differ:
* - Retries on 5xx / 429 / transient network errors with exponential backoff.
* - Timeouts are per-call, not per-request, so streaming can take minutes.
* - Errors are mapped to `llm_unavailable` / `llm_rate_limited` /
* `llm_timeout` — the client cares which one it is.
* - A small SSE decoder is provided for providers that return
* `text/event-stream` (openai_compatible, anthropic).
*/
import { ERROR_CODES, MemosError } from "../../agent-contract/errors.js";
import {
getRetryCooldown,
parseRetryAfterMs,
planRetry,
recordRetryCooldown,
retryCooldownKey,
type RetryPlan,
waitForRetry,
} from "../util/retry-after.js";
import type { LlmProviderLogger, LlmProviderName } from "./types.js";
export interface HttpPostOpts
{
url: string;
body: TBody;
headers?: Record;
timeoutMs: number;
maxRetries: number;
signal?: AbortSignal;
/** Absolute end-to-end deadline. Unlike timeoutMs, this is not renewed per attempt. */
deadlineAt?: number;
/** Model/deployment scope; prevents one model cooldown from blocking another. */
cooldownScope?: string;
provider: LlmProviderName;
log: LlmProviderLogger;
onRetry?: (attempt: number) => void;
}
/**
* Single JSON POST with retry + timeout. For streaming, see `httpPostStream`.
*/
export async function httpPostJson(opts: HttpPostOpts): Promise<{
json: TResp;
status: number;
durationMs: number;
}> {
let attempt = 0;
let lastErr: unknown = null;
const cooldownKey = retryCooldownKey("llm", opts.provider, opts.url, opts.cooldownScope);
while (attempt <= opts.maxRetries) {
attempt++;
const start = Date.now();
try {
const cooldown = getRetryCooldown(cooldownKey, start);
if (cooldown) {
const details = {
provider: opts.provider,
url: opts.url,
status: cooldown.status,
attempt,
maxRetries: opts.maxRetries,
retryAfterMs: cooldown.retryAfterMs,
retryAt: cooldown.retryAt,
retryDecision: "defer",
retryReason: "cooldown_active",
remainingDeadlineMs: remainingDeadlineMs(opts.deadlineAt, start),
};
opts.log.warn("http.retry_cooldown", details);
throw new MemosError(
errCodeForStatus(cooldown.status),
`${opts.provider} is cooling down until ${new Date(cooldown.retryAt).toISOString()}`,
details,
);
}
const signal = mergeSignals(opts.signal, AbortSignal.timeout(opts.timeoutMs));
const resp = await fetch(opts.url, {
method: "POST",
headers: {
"Content-Type": "application/json",
Accept: "application/json",
...opts.headers,
},
body: JSON.stringify(opts.body),
signal,
});
const ms = Date.now() - start;
if (!resp.ok) {
const text = await safeText(resp);
const transient = resp.status >= 500 || resp.status === 429;
const retryAfterMs = resp.status === 429 || resp.status === 503
? parseRetryAfterMs(resp.headers.get("Retry-After"))
: null;
if (retryAfterMs !== null) {
recordRetryCooldown(cooldownKey, {
retryAfterMs,
retryAt: Date.now() + retryAfterMs,
status: resp.status,
});
}
opts.log.warn("http.non_ok", {
status: resp.status,
attempt,
transient,
durationMs: ms,
retryAfterMs,
body: truncateLogBody(text),
});
if (transient && attempt <= opts.maxRetries) {
const plan = planRetry({
attempt,
baseMs: 250,
jitterMaxMs: 120,
retryAfterMs,
deadlineAt: opts.deadlineAt,
});
const retryDetails = retryPlanDetails(plan, opts, resp.status, attempt);
if (plan.action === "defer") {
opts.log.warn("http.retry_deferred", retryDetails);
throw new MemosError(
errCodeForStatus(resp.status),
`HTTP ${resp.status} from ${opts.provider}; retry deferred until ${new Date(plan.retryAt).toISOString()}`,
retryDetails,
);
}
opts.log.warn("http.retry_scheduled", retryDetails);
opts.onRetry?.(attempt);
await waitForRetry(plan.delayMs, opts.signal);
continue;
}
throw new MemosError(
errCodeForStatus(resp.status),
`HTTP ${resp.status} from ${opts.provider}`,
{
provider: opts.provider,
url: opts.url,
status: resp.status,
body: text,
...(retryAfterMs === null
? {}
: {
retryAfterMs,
retryAt: Date.now() + retryAfterMs,
retryDecision: "stop",
retryReason: "retries_exhausted",
}),
},
);
}
const json = (await resp.json()) as TResp;
opts.log.debug("http.ok", {
status: resp.status,
attempt,
durationMs: ms,
});
return { json, status: resp.status, durationMs: ms };
} catch (err) {
lastErr = err;
if (err instanceof MemosError) throw err;
if (opts.signal?.aborted) {
throw new MemosError(
ERROR_CODES.LLM_TIMEOUT,
`${opts.provider} request was cancelled`,
{ provider: opts.provider, url: opts.url, cancelled: true },
);
}
const transient = isTransientError(err);
const timedOut = isTimeout(err) || opts.signal?.aborted === true;
opts.log.warn("http.exception", {
attempt,
transient,
timedOut,
err: toErrDetail(err),
});
if ((transient || timedOut) && attempt <= opts.maxRetries) {
const plan = planRetry({
attempt,
baseMs: 250,
jitterMaxMs: 120,
deadlineAt: opts.deadlineAt,
});
const retryDetails = retryPlanDetails(plan, opts, null, attempt);
if (plan.action === "defer") {
opts.log.warn("http.retry_deferred", retryDetails);
throw new MemosError(
timedOut ? ERROR_CODES.LLM_TIMEOUT : ERROR_CODES.LLM_UNAVAILABLE,
`${opts.provider} retry cannot fit the request deadline`,
retryDetails,
);
}
opts.log.warn("http.retry_scheduled", retryDetails);
opts.onRetry?.(attempt);
await waitForRetry(plan.delayMs, opts.signal);
continue;
}
if (timedOut) {
throw new MemosError(
ERROR_CODES.LLM_TIMEOUT,
`${opts.provider} timed out after ${opts.timeoutMs} ms`,
{ provider: opts.provider, url: opts.url, timeoutMs: opts.timeoutMs },
);
}
throw new MemosError(
ERROR_CODES.LLM_UNAVAILABLE,
`${opts.provider} request failed: ${(err as Error).message ?? String(err)}`,
{ provider: opts.provider, url: opts.url },
);
}
}
throw new MemosError(
ERROR_CODES.LLM_UNAVAILABLE,
`Exhausted retries to ${opts.provider}`,
{
provider: opts.provider,
url: opts.url,
cause: lastErr instanceof Error ? lastErr.message : String(lastErr),
},
);
}
/**
* Open an HTTP POST and return the raw streaming body. The caller is
* responsible for parsing SSE. No retries here — streaming is "either works
* or you start over from scratch".
*/
export async function httpPostStream(opts: {
url: string;
body: unknown;
headers?: Record;
timeoutMs: number;
signal?: AbortSignal;
provider: LlmProviderName;
log: LlmProviderLogger;
}): Promise {
const start = Date.now();
const signal = mergeSignals(opts.signal, AbortSignal.timeout(opts.timeoutMs));
const resp = await fetch(opts.url, {
method: "POST",
headers: {
"Content-Type": "application/json",
Accept: "text/event-stream",
...opts.headers,
},
body: JSON.stringify(opts.body),
signal,
});
if (!resp.ok) {
const text = await safeText(resp);
opts.log.warn("http.non_ok", {
status: resp.status,
transient: resp.status >= 500 || resp.status === 429,
durationMs: Date.now() - start,
body: truncateLogBody(text),
});
throw new MemosError(
errCodeForStatus(resp.status),
`HTTP ${resp.status} from ${opts.provider} (stream)`,
{ provider: opts.provider, url: opts.url, status: resp.status, body: text },
);
}
if (!resp.body) {
throw new MemosError(
ERROR_CODES.LLM_UNAVAILABLE,
`${opts.provider} returned empty streaming body`,
{ provider: opts.provider, url: opts.url },
);
}
return resp;
}
/**
* Parse a `text/event-stream` body into its raw `data:` payloads.
* Yields each `data: …` payload as a string. Handles the "[DONE]" sentinel
* common to OpenAI-shape providers.
*/
export async function* decodeSse(body: ReadableStream): AsyncGenerator {
const reader = body.getReader();
const decoder = new TextDecoder("utf-8");
let buf = "";
while (true) {
const { value, done } = await reader.read();
if (done) break;
buf += decoder.decode(value, { stream: true });
// SSE events are separated by blank lines (\n\n).
let idx = buf.indexOf("\n\n");
while (idx !== -1) {
const event = buf.slice(0, idx);
buf = buf.slice(idx + 2);
for (const line of event.split("\n")) {
if (line.startsWith("data:")) {
const payload = line.slice(5).trim();
if (payload.length > 0) yield payload;
}
}
idx = buf.indexOf("\n\n");
}
}
// Flush whatever's left in buf.
for (const line of buf.split("\n")) {
if (line.startsWith("data:")) {
const payload = line.slice(5).trim();
if (payload.length > 0) yield payload;
}
}
}
// ─── Helpers ─────────────────────────────────────────────────────────────────
function errCodeForStatus(status: number): "llm_rate_limited" | "llm_unavailable" {
if (status === 429) return ERROR_CODES.LLM_RATE_LIMITED;
return ERROR_CODES.LLM_UNAVAILABLE;
}
async function safeText(resp: Response): Promise {
try {
return await resp.text();
} catch {
return undefined;
}
}
function truncateLogBody(text: string | undefined): string | undefined {
return text?.slice(0, 512);
}
function isTransientError(err: unknown): boolean {
if (!(err instanceof Error)) return false;
const msg = err.message ?? "";
if (/ECONNRESET|EAI_AGAIN|socket hang up/i.test(msg)) return true;
return false;
}
function isTimeout(err: unknown): boolean {
if (err instanceof Error) {
if (err.name === "TimeoutError") return true;
if ((err as { code?: string }).code === "ABORT_ERR") return true;
if (/timeout|ETIMEDOUT/i.test(err.message ?? "")) return true;
}
return false;
}
function retryPlanDetails(
plan: RetryPlan,
opts: HttpPostOpts,
status: number | null,
attempt: number,
): Record {
return {
provider: opts.provider,
url: opts.url,
status,
attempt,
maxRetries: opts.maxRetries,
backoffMs: plan.backoffMs,
plannedDelayMs: plan.delayMs,
retryAfterMs: plan.retryAfterMs,
retryAt: plan.retryAt,
retrySource: plan.source,
retryDecision: plan.action,
...(plan.action === "defer" ? { retryReason: plan.reason } : {}),
remainingDeadlineMs: remainingDeadlineMs(opts.deadlineAt),
};
}
function remainingDeadlineMs(deadlineAt?: number, nowMs: number = Date.now()): number | null {
return deadlineAt === undefined ? null : Math.max(0, deadlineAt - nowMs);
}
function mergeSignals(a: AbortSignal | undefined, b: AbortSignal): AbortSignal {
if (!a) return b;
const ctrl = new AbortController();
const forward = () => ctrl.abort();
if (a.aborted || b.aborted) ctrl.abort();
a.addEventListener("abort", forward, { once: true });
b.addEventListener("abort", forward, { once: true });
return ctrl.signal;
}
function toErrDetail(e: unknown): Record {
if (e instanceof Error) return { name: e.name, message: e.message };
return { value: String(e) };
}