/** * Soft client-side pacing gate: honor plan maxRpm between stream starts. */ import { createAssistantMessageEventStream, type AssistantMessageEventStream, } from "@earendil-works/pi-ai"; import type { BaseStreamFn } from "./retry-stream.ts"; import { minIntervalMsForRpm } from "./limits.ts"; export interface PacingState { maxRpm: number; lastStartMs: number; inFlight: number; } export function createPacingState(maxRpm = 60): PacingState { return { maxRpm, lastStartMs: 0, inFlight: 0 }; } export function setPacingRpm(state: PacingState, maxRpm: number): void { state.maxRpm = Math.max(1, maxRpm); } function defaultSleep(ms: number, signal?: AbortSignal): Promise { return new Promise((resolve, reject) => { if (signal?.aborted) return reject(new Error("aborted")); const onAbort = () => { clearTimeout(timer); reject(new Error("aborted")); }; const timer = setTimeout(() => { signal?.removeEventListener("abort", onAbort); resolve(); }, ms); signal?.addEventListener("abort", onAbort, { once: true }); }); } let gateTail: Promise = Promise.resolve(); /** * Acquire a pacing slot (serialized). Returns a release function. */ export async function acquirePacingSlot( state: PacingState, signal?: AbortSignal, sleep: (ms: number, signal?: AbortSignal) => Promise = defaultSleep, now: () => number = Date.now, ): Promise<() => void> { const singleFlight = state.maxRpm <= 2; const run = async () => { if (singleFlight) { while (state.inFlight > 0) { if (signal?.aborted) throw new Error("aborted"); await sleep(50, signal); } } const minInterval = minIntervalMsForRpm(state.maxRpm); const wait = Math.max(0, state.lastStartMs + minInterval - now()); if (wait > 0) await sleep(wait, signal); state.lastStartMs = now(); state.inFlight++; }; const myTurn = gateTail.then(run, run); gateTail = myTurn.then( () => undefined, () => undefined, ); await myTurn; return () => { state.inFlight = Math.max(0, state.inFlight - 1); }; } /** Pure helper for tests: compute wait ms before next start. */ export function waitMsBeforeNextStart(state: PacingState, nowMs: number): number { const minInterval = minIntervalMsForRpm(state.maxRpm); return Math.max(0, state.lastStartMs + minInterval - nowMs); } /** * Wrap streamSimple so starts are spaced at ~60/maxRpm seconds. * When maxRpm <= 2, also single-flight concurrent starts. */ export function createPacedStream( baseStream: BaseStreamFn, state: PacingState, sleep: (ms: number, signal?: AbortSignal) => Promise = defaultSleep, ): BaseStreamFn { return function pacedStream(model, context, options) { const out = createAssistantMessageEventStream(); const signal = options?.signal; (async () => { let release: (() => void) | undefined; try { release = await acquirePacingSlot(state, signal, sleep); const inner: AssistantMessageEventStream = baseStream(model, context, options); for await (const ev of inner) { out.push(ev); if (ev.type === "done" || ev.type === "error") return; } } catch (err) { if (signal?.aborted) { out.push({ type: "error", reason: "aborted", error: { role: "assistant", content: [], api: model.api, provider: model.provider, model: model.id, usage: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, totalTokens: 0, cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, }, stopReason: "aborted", errorMessage: err instanceof Error ? err.message : "aborted", timestamp: Date.now(), }, }); return; } throw err; } finally { release?.(); } })().catch((err) => { out.push({ type: "error", reason: "error", error: { role: "assistant", content: [], api: model.api, provider: model.provider, model: model.id, usage: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, totalTokens: 0, cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, }, stopReason: "error", errorMessage: err instanceof Error ? err.message : String(err), timestamp: Date.now(), }, }); }); return out; }; }