From 10be9a411adeb8361a896eee03c35f4e94f5dc70 Mon Sep 17 00:00:00 2001 From: Valera Brizhatiuk <19537764+plombeer31@users.noreply.github.com> Date: Sun, 13 Sep 2026 23:23:00 +0300 Subject: [PATCH] fix(memory): cancel a sub-call's request when its timeout fires Every memory sub-call wrapper in bootstrap (reflection, link generator, vote, query rewriter, distill) enforced its runner's timeout by racing the completion against the abort signal without handing that signal to llmComplete. The runner gave up while the HTTP request kept running: still billing on a cloud provider, still holding a llama-server slot. The five copies are replaced by abortableSubcall, which forwards the signal into the request and keeps the race as a backstop for a provider that ignores it. Each wrapper still sends exactly the fields it did. Forwarding alone was not enough, for two reasons found on the way: - OpenAI-compatible unary requests stopped listening to the signal once headers arrived (openAiFetch unlinks it because it also opens streams), so a provider that sends headers before the body could not be cancelled. openAiPostJson now reads the body through a reader the signal cancels. - An aborted llama-server request surfaces as a status-null LlamaServerError, which classifies as transport and makes shouldAdvance switch providers immediately. The unary fallback seam now rethrows an aborted request as the signal's reason, so it classifies as cancelled and never trips a breaker or flips the override, the rule completeStream already applied. --- AGENTS.md | 3 +- src/llm/provider/openai/openai-http.ts | 47 +++- src/runtime/abortable-subcall.network.test.ts | 215 ++++++++++++++++++ src/runtime/abortable-subcall.test.ts | 123 ++++++++++ src/runtime/abortable-subcall.ts | 71 ++++++ src/runtime/bootstrap.ts | 101 +++----- src/runtime/llm-fallback-seam.ts | 32 ++- 7 files changed, 510 insertions(+), 82 deletions(-) create mode 100644 src/runtime/abortable-subcall.network.test.ts create mode 100644 src/runtime/abortable-subcall.test.ts create mode 100644 src/runtime/abortable-subcall.ts diff --git a/AGENTS.md b/AGENTS.md index 21fa60c0..1ea7eb59 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -1163,6 +1163,7 @@ The three channels share one SQLite file `/memory.sqlite` (separate fr ### Reflection (async end-of-turn memory formation) - **When.** Fired at the end of every `AgentLoop.runTurn` after `assistant_reply` is emitted. **Fire-and-forget**, never awaited. `abortPending({ sessionId: state.id })` runs at the start of the next `runTurn` so at most one reflection is in flight **per session**; reflections on other sessions are never aborted as a side effect (load-bearing for cross-session parallelism — see §"Concurrency contract"). +- **A timeout cancels the request, not just the wait.** Every memory sub-call wrapper in bootstrap (reflection, link generator, vote, query rewriter, distill) goes through `abortableSubcall` ([src/runtime/abortable-subcall.ts](src/runtime/abortable-subcall.ts)), which forwards the runner's abort signal into `llmComplete` — so a fired timeout or `abortPending` closes the HTTP request and frees the slot — and still rejects the moment the signal aborts in case a provider ignores it. - **What.** A micro-prompt with its own small stable prefix asks the model to extract durable facts from the last `USER`/`ASSISTANT` exchange. Output is GBNF-constrained to either `NONE` or a bounded list of two flavours: - `SET key=value` (pinned fact) or `SET key=value [pinned=false; keywords=a,b,c]` (contextual fact). Caps at `memory.reflection.maxFactsPerCall` (default `3`). - `NOTE freeform observation [tag1, tag2]` → into `MemoryStore` with implicit `reflection` tag. Master switch `memory.reflection.autoStoreNotes` (default `true`); cap at `memory.reflection.maxNotesPerCall` (default `2`, set to `0` to disable). @@ -2314,7 +2315,7 @@ Every terminal failure the agent loop surfaces is normalised into a canonical `L ## Provider fallback chain -A cross-provider circuit breaker layered **above** the single-provider reliability policy. Where the two retry layers above recover a request on the *same* provider, the fallback chain switches to a *different* configured provider when the active one is unavailable. It lives in [src/llm/fallback/](src/llm/fallback/). The `llmComplete` / `llmCompleteStream` seams that wrap it are built by [src/runtime/llm-fallback-seam.ts](src/runtime/llm-fallback-seam.ts) (`createFallbackCompleter` / `createFallbackStreamer`, injected with `{ fallbackChain, resolveSlice, recordUnaryUsage, recordStreamUsage }`) and wired in [src/runtime/bootstrap.ts](src/runtime/bootstrap.ts) — strictly **after** the per-provider retry budget (PR #90 `runOpenAiWithRetry`, `LlamaServerClient.completionRetries`) is spent, never inside it. +A cross-provider circuit breaker layered **above** the single-provider reliability policy. Where the two retry layers above recover a request on the *same* provider, the fallback chain switches to a *different* configured provider when the active one is unavailable. It lives in [src/llm/fallback/](src/llm/fallback/). The `llmComplete` / `llmCompleteStream` seams that wrap it are built by [src/runtime/llm-fallback-seam.ts](src/runtime/llm-fallback-seam.ts) (`createFallbackCompleter` / `createFallbackStreamer`, injected with `{ fallbackChain, resolveSlice, recordUnaryUsage, recordStreamUsage }`) and wired in [src/runtime/bootstrap.ts](src/runtime/bootstrap.ts) — strictly **after** the per-provider retry budget (PR #90 `runOpenAiWithRetry`, `LlamaServerClient.completionRetries`) is spent, never inside it. The unary seam rethrows a request whose caller aborted as the signal's reason, so it classifies `cancelled` and never advances the chain — `LlamaServerClient` and `runOpenAiWithRetry` otherwise surface an abort as a `status: null` error that files as `transport`, which is how an abandoned memory sub-call could trip a breaker and flip the override. ### Chain unit and config diff --git a/src/llm/provider/openai/openai-http.ts b/src/llm/provider/openai/openai-http.ts index 4e4aeeca..cb964874 100644 --- a/src/llm/provider/openai/openai-http.ts +++ b/src/llm/provider/openai/openai-http.ts @@ -224,11 +224,56 @@ export async function openAiPostJson( if (!res.ok) { throw await httpErrorFromResponse(deps, path, res); } - return (await res.json()) as Record; + return readJsonBody(res, request.signal); }), ); } +/** + * Read a unary JSON body without going deaf to the caller's abort. + * + * `openAiFetch` unlinks the caller's signal the moment `fetch` resolves — + * it must, because it also opens streams, whose consumer owns the signal + * from then on. For a unary request that left the body read unabortable: + * a provider that sends headers first and the completion later kept the + * socket, the slot and the bill running after the caller had given up + * (a memory sub-call's timeout, a cancelled turn). Cancelling the reader + * tears the connection down; the caller gets `signal.reason`, which + * classifies `cancelled`. + * + * `res.json()` cannot be used for this: it locks the body, and cancelling + * a locked stream from outside is refused. Without a signal nothing can + * cancel the read, so that path keeps `res.json()` exactly as before. + */ +async function readJsonBody( + res: Response, + signal: AbortSignal | undefined, +): Promise> { + if (!signal || !res.body) return (await res.json()) as Record; + const reader = res.body.getReader(); + const onAbort = (): void => { + reader.cancel(signal.reason).catch(() => undefined); + }; + if (signal.aborted) onAbort(); + else signal.addEventListener("abort", onAbort, { once: true }); + try { + const chunks: Uint8Array[] = []; + for (;;) { + const { done, value } = await reader.read(); + if (done) break; + chunks.push(value); + } + if (signal.aborted) { + throw signal.reason ?? new DOMException("aborted", "AbortError"); + } + return JSON.parse( + new TextDecoder().decode(Buffer.concat(chunks)), + ) as Record; + } finally { + signal.removeEventListener("abort", onAbort); + } +} + /** * Recover the one HTTP failure that is fixable by changing the request * rather than by waiting or switching providers: a 402 that names how diff --git a/src/runtime/abortable-subcall.network.test.ts b/src/runtime/abortable-subcall.network.test.ts new file mode 100644 index 00000000..6e28c8f7 --- /dev/null +++ b/src/runtime/abortable-subcall.network.test.ts @@ -0,0 +1,215 @@ +import { createServer, type Server } from "node:http"; +import type { AddressInfo } from "node:net"; +import { afterEach, describe, expect, it } from "vitest"; + +import type { LlmStreamParams } from "../agent/step-executor.js"; +import { ProviderFallbackChain } from "../llm/fallback/index.js"; +import { DEFAULT_FALLBACK_TIMING } from "../llm/fallback/fallback-config.js"; +import { LlamaServerClient } from "../llm/llama-server-client.js"; +import { QWEN_THINK_PROFILE } from "../llm/model-profile.js"; +import { + fakeAnswer, + fakeProvider, +} from "../llm/provider/fake-provider.fixture.js"; +import { LlamaServerProvider } from "../llm/provider/llama-server/llama-server-provider.js"; +import type { LlmProvider } from "../llm/provider/llm-provider.js"; +import { OpenAiProvider } from "../llm/provider/openai/openai-provider.js"; +import { + createLinkGeneratorRunner, + type LinkGeneratorInput, + type LinkGeneratorLlmComplete, + type LinkGeneratorRunnerDeps, +} from "../memory/links/link-generator-runner.js"; +import { abortableSubcall } from "./abortable-subcall.js"; +import { createFallbackCompleter } from "./llm-fallback-seam.js"; + +/** + * End to end over a real socket: a memory sub-call runner whose timeout + * fires must close the HTTP request it started — not merely stop waiting + * for it — and the abandoned request must not advance the fallback chain. + * + * Wiring mirrors bootstrap: link-generator runner → `abortableSubcall` + * (bootstrap's link-gen request shape) → `createFallbackCompleter` → a + * real provider pointed at a local server that never finishes answering. + */ + +type Mode = "silent" | "headers-then-stall"; +type Kind = "openai" | "llama-server"; + +const servers: Server[] = []; + +afterEach(async () => { + await Promise.all( + servers.splice(0).map( + (server) => + new Promise((resolve) => { + server.closeAllConnections(); + server.close(() => resolve()); + }), + ), + ); +}); + +/** A server that accepts the request and never completes the response. */ +async function startStallingServer(mode: Mode) { + let requests = 0; + let closedSockets = 0; + const server = createServer((req, res) => { + requests += 1; + req.socket.once("close", () => { + closedSockets += 1; + }); + req.resume(); + if (mode === "headers-then-stall") { + // Headers plus a byte of body, then nothing: `fetch` has resolved, + // so only a cancelled body read can let go of this socket. + res.writeHead(200, { "content-type": "application/json" }); + res.write(" "); + } + }); + servers.push(server); + await new Promise((resolve) => + server.listen(0, "127.0.0.1", () => resolve()), + ); + const { port } = server.address() as AddressInfo; + return { + url: `http://127.0.0.1:${port}`, + requests: () => requests, + closedSockets: () => closedSockets, + }; +} + +function providerFor(kind: Kind, url: string): LlmProvider { + if (kind === "openai") { + return new OpenAiProvider({ + id: "primary", + baseUrl: url, + apiKey: "test-key", + defaultChatModel: "test-model", + }); + } + return new LlamaServerProvider( + new LlamaServerClient({ baseUrl: url, completionRetries: 1 }), + { + id: "primary", + getProfile: () => QWEN_THINK_PROFILE, + visionEnabledByConfig: false, + visionAutoDetect: false, + maxImageBytes: 1, + maxImagesPerCall: 1, + baseUrlOverride: url, + }, + ); +} + +/** Polls `predicate` until it holds or `ms` elapses. */ +async function eventually(predicate: () => boolean, ms: number) { + const deadline = Date.now() + ms; + while (!predicate()) { + if (Date.now() > deadline) return false; + await new Promise((resolve) => setTimeout(resolve, 10)); + } + return true; +} + +/** Settles with "hung" when `promise` has not settled within `ms`. */ +function within(promise: Promise, ms: number) { + return Promise.race([ + promise.then(() => "settled" as const), + new Promise<"hung">((resolve) => setTimeout(() => resolve("hung"), ms)), + ]); +} + +const INPUT: LinkGeneratorInput = { + sessionId: "s1", + userMessage: "where did we put the deploy notes?", + assistantReply: "In the ops wiki, next to the runbook.", + candidates: [ + { id: "1", body: "deploy notes live in the ops wiki" }, + { id: "2", body: "the runbook covers rollbacks" }, + ] as unknown as LinkGeneratorInput["candidates"], +}; + +const CASES: ReadonlyArray<{ kind: Kind; mode: Mode }> = [ + { kind: "openai", mode: "silent" }, + { kind: "openai", mode: "headers-then-stall" }, + { kind: "llama-server", mode: "silent" }, + { kind: "llama-server", mode: "headers-then-stall" }, +]; + +describe("a timed-out memory sub-call", () => { + for (const { kind, mode } of CASES) { + it(`closes its ${kind} request (${mode}) and leaves the fallback chain alone`, async () => { + const server = await startStallingServer(mode); + const chain = new ProviderFallbackChain({ + resolve: () => ({ + chain: ["primary", "backup"], + timing: DEFAULT_FALLBACK_TIMING, + }), + }); + let backupCalls = 0; + const providers = new Map([ + ["primary", providerFor(kind, server.url)], + [ + "backup", + fakeProvider("backup", "grammar", async () => { + backupCalls += 1; + return fakeAnswer("backup"); + }), + ], + ]); + const complete = createFallbackCompleter({ + fallbackChain: chain, + resolveSlice: (providerId) => { + const provider = providers.get(providerId)!; + return { provider, transport: provider.capabilities.toolTransport }; + }, + recordUnaryUsage: () => {}, + recordStreamUsage: () => {}, + }); + // The whole chain run — including any fallover an orphaned failure + // would trigger — is this promise; the runner stops watching it. + let chainRun: Promise | undefined; + const tracked = (params: LlmStreamParams) => { + const run = complete(params); + chainRun = run.catch(() => undefined); + return run; + }; + const llmComplete: LinkGeneratorLlmComplete = abortableSubcall( + tracked, + (params: Parameters[0]) => ({ + prompt: params.prompt, + grammar: params.grammar, + slotId: params.slotId, + sessionId: params.sessionId, + ...(params.responseFormat + ? { responseFormat: params.responseFormat } + : {}), + }), + ); + const outcomes: string[] = []; + const runner = createLinkGeneratorRunner({ + llmComplete, + linkStore: {} as LinkGeneratorRunnerDeps["linkStore"], + reflectionSlotId: -1, + timeoutMs: 150, + emitTrace: (event) => outcomes.push(event.outcome), + }); + + await expect(runner.generate(INPUT)).resolves.toBe(0); + + expect(outcomes).toEqual(["timeout"]); + expect(server.requests()).toBe(1); + // The load-bearing assertion: the server sees the socket go away. + // Without the signal reaching the request it stays open until the + // provider's own request timeout, minutes later. + expect(await eventually(() => server.closedSockets() > 0, 3_000)).toBe( + true, + ); + expect(chainRun).toBeDefined(); + expect(await within(chainRun!, 3_000)).toBe("settled"); + expect(backupCalls).toBe(0); + expect(chain.activeOverrideFor("link-gen:s1")).toBeNull(); + }); + } +}); diff --git a/src/runtime/abortable-subcall.test.ts b/src/runtime/abortable-subcall.test.ts new file mode 100644 index 00000000..4a7f7700 --- /dev/null +++ b/src/runtime/abortable-subcall.test.ts @@ -0,0 +1,123 @@ +import { describe, expect, it, vi } from "vitest"; + +import type { LlmStreamParams } from "../agent/step-executor.js"; +import { fakeAnswer } from "../llm/provider/fake-provider.fixture.js"; +import { abortableSubcall } from "./abortable-subcall.js"; + +interface RunnerParams { + prompt: string; + grammar: string; + slotId: number; + sessionId: string; + signal: AbortSignal; +} + +const shape = (params: RunnerParams) => ({ + prompt: params.prompt, + grammar: params.grammar, + slotId: params.slotId, + sessionId: params.sessionId, +}); + +function runnerParams(signal: AbortSignal): RunnerParams { + return { + prompt: "p", + grammar: 'root ::= "x"', + slotId: 3, + sessionId: "link-gen:s1", + signal, + }; +} + +/** Settles with "hung" when `promise` has not settled within `ms`. */ +function within(promise: Promise, ms: number) { + return Promise.race([ + promise.then( + (value) => ({ kind: "resolved" as const, value }), + (error: unknown) => ({ kind: "rejected" as const, error }), + ), + new Promise<{ kind: "hung" }>((resolve) => + setTimeout(() => resolve({ kind: "hung" }), ms), + ), + ]); +} + +describe("abortableSubcall", () => { + it("forwards the caller's own signal with the shaped fields untouched", async () => { + const seen: LlmStreamParams[] = []; + const complete = vi.fn(async (request: LlmStreamParams) => { + seen.push(request); + return fakeAnswer("cloud"); + }); + const controller = new AbortController(); + const params = runnerParams(controller.signal); + + const result = await abortableSubcall(complete, shape)(params); + + expect(result.modelId).toBe("cloud-model"); + expect(seen).toHaveLength(1); + expect(seen[0]!.signal).toBe(controller.signal); + expect(seen[0]).toEqual({ ...shape(params), signal: controller.signal }); + }); + + it("forwards the signal even when the shape strips it (reflection spreads the rest)", async () => { + const seen: LlmStreamParams[] = []; + const controller = new AbortController(); + const call = abortableSubcall( + async (request: LlmStreamParams) => { + seen.push(request); + return fakeAnswer("cloud"); + }, + ({ signal: _signal, ...rest }: RunnerParams) => rest, + ); + await call(runnerParams(controller.signal)); + expect(seen[0]!.signal).toBe(controller.signal); + }); + + it("rejects promptly on abort even if the completion never settles, and the inner signal is aborted", async () => { + let inner: AbortSignal | undefined; + const call = abortableSubcall((request: LlmStreamParams) => { + inner = request.signal; + return new Promise(() => {}); + }, shape); + const controller = new AbortController(); + const pending = call(runnerParams(controller.signal)); + controller.abort(); + + const outcome = await within(pending, 200); + expect(outcome.kind).toBe("rejected"); + const error = (outcome as { error: unknown }).error; + expect((error as Error).name).toBe("AbortError"); + expect(inner?.aborted).toBe(true); + }); + + it("never sends a request for a signal that is already aborted", async () => { + const complete = vi.fn(async () => fakeAnswer("cloud")); + const controller = new AbortController(); + controller.abort(); + await expect( + abortableSubcall(complete, shape)(runnerParams(controller.signal)), + ).rejects.toMatchObject({ name: "AbortError" }); + expect(complete).not.toHaveBeenCalled(); + }); + + it("passes a completion failure through unchanged when nothing aborted", async () => { + const boom = new Error("provider exploded"); + const call = abortableSubcall(async () => { + throw boom; + }, shape); + await expect(call(runnerParams(new AbortController().signal))).rejects.toBe( + boom, + ); + }); + + it("detaches its abort listener once the completion settles", async () => { + const controller = new AbortController(); + const remove = vi.spyOn(controller.signal, "removeEventListener"); + await abortableSubcall( + async () => fakeAnswer("cloud"), + shape, + )(runnerParams(controller.signal)); + expect(remove).toHaveBeenCalledWith("abort", expect.any(Function)); + }); +}); diff --git a/src/runtime/abortable-subcall.ts b/src/runtime/abortable-subcall.ts new file mode 100644 index 00000000..83260cbc --- /dev/null +++ b/src/runtime/abortable-subcall.ts @@ -0,0 +1,71 @@ +import type { LlmStreamParams } from "../agent/step-executor.js"; +import type { CompletionResult } from "../llm/provider/completion-types.js"; + +/** What a memory sub-call sends, minus the signal — the helper owns that. */ +export type SubcallRequest = Omit; + +export type SubcallComplete = ( + params: LlmStreamParams, +) => Promise; + +/** + * Adapt the runtime's `llmComplete` for a memory sub-call runner + * (reflection, link generator, vote, query rewriter, distill). + * + * Every one of those runners enforces its timeout by aborting the signal + * it passes in. The wrappers used to race the completion against that + * abort without handing the signal to `llmComplete`, so the runner gave + * up while the HTTP request kept going: still billing on a cloud + * provider, still holding a llama-server slot, and — because an orphan + * that later fails runs through the fallback chain like any request — + * able to trip a breaker or flip the sticky override minutes after + * anyone stopped caring about it. + * + * So the signal is forwarded into the request, which is what actually + * cancels it (the fallback seam turns that abort into a `cancelled` + * failure the chain never advances on). The race stays as a backstop: + * the returned promise rejects the moment the signal aborts even if a + * provider ignores the signal, so a runner's timeout is never hostage + * to one. + * + * `shape` builds the request fields each runner sends today; it is kept + * per call site so no wrapper's payload changes shape. + */ +export function abortableSubcall

( + complete: SubcallComplete, + shape: (params: P) => SubcallRequest, +): (params: P) => Promise { + return async (params) => { + const { signal } = params; + if (signal.aborted) throw subcallAbortError(); + const request: LlmStreamParams = { ...shape(params), signal }; + return new Promise((resolve, reject) => { + const onAbort = (): void => reject(subcallAbortError()); + signal.addEventListener("abort", onAbort, { once: true }); + const detach = (): void => signal.removeEventListener("abort", onAbort); + let pending: Promise; + try { + pending = complete(request); + } catch (err) { + detach(); + reject(err); + return; + } + pending.then( + (result) => { + detach(); + resolve(result); + }, + (err: unknown) => { + detach(); + reject(err); + }, + ); + }); + }; +} + +/** The same rejection the inline wrappers threw, so runners see no change. */ +function subcallAbortError(): DOMException { + return new DOMException("aborted", "AbortError"); +} diff --git a/src/runtime/bootstrap.ts b/src/runtime/bootstrap.ts index 263c1cee..82fc3fe1 100644 --- a/src/runtime/bootstrap.ts +++ b/src/runtime/bootstrap.ts @@ -99,6 +99,7 @@ import { createFallbackStreamer, type FallbackSeamDeps, } from "./llm-fallback-seam.js"; +import { abortableSubcall } from "./abortable-subcall.js"; import { MemoryStore } from "../memory/memory-store.js"; import { ProfileStore } from "../memory/profile-store.js"; @@ -2016,18 +2017,9 @@ export async function createAgentRuntime( ) { const reservedSlot = slotManager.reserveReflectionSlot(); const reflectionSlotId = reservedSlot ?? -1; - const linkGenLlmComplete: LinkGeneratorLlmComplete = async (params) => { - if (params.signal.aborted) { - throw new DOMException("aborted", "AbortError"); - } - const abortPromise = new Promise((_, reject) => { - params.signal.addEventListener( - "abort", - () => reject(new DOMException("aborted", "AbortError")), - { once: true }, - ); - }); - const completionPromise = llmComplete({ + const linkGenLlmComplete: LinkGeneratorLlmComplete = abortableSubcall( + llmComplete, + (params: Parameters[0]) => ({ prompt: params.prompt, grammar: params.grammar, slotId: params.slotId, @@ -2035,9 +2027,8 @@ export async function createAgentRuntime( ...(params.responseFormat ? { responseFormat: params.responseFormat } : {}), - }); - return Promise.race([completionPromise, abortPromise]); - }; + }), + ); const linkGenerator = createLinkGeneratorRunner({ llmComplete: linkGenLlmComplete, linkStore, @@ -2086,18 +2077,9 @@ export async function createAgentRuntime( if (reflectionRunner && voteStore) { const reservedSlot = slotManager.reserveReflectionSlot(); const voteSlotId = reservedSlot ?? -1; - const voteLlmComplete: VoteRunnerLlmComplete = async (params) => { - if (params.signal.aborted) { - throw new DOMException("aborted", "AbortError"); - } - const abortPromise = new Promise((_, reject) => { - params.signal.addEventListener( - "abort", - () => reject(new DOMException("aborted", "AbortError")), - { once: true }, - ); - }); - const completionPromise = llmComplete({ + const voteLlmComplete: VoteRunnerLlmComplete = abortableSubcall( + llmComplete, + (params: Parameters[0]) => ({ prompt: params.prompt, grammar: params.grammar, slotId: params.slotId, @@ -2105,9 +2087,8 @@ export async function createAgentRuntime( ...(params.responseFormat ? { responseFormat: params.responseFormat } : {}), - }); - return Promise.race([completionPromise, abortPromise]); - }; + }), + ); const voteRunner = createVoteRunner({ llmComplete: voteLlmComplete, voteStore, @@ -2223,18 +2204,9 @@ export async function createAgentRuntime( // pre-v18 chain. let memoryContextProvider = baseMemoryContextProvider; if (baseMemoryContextProvider && config.memory.retrieve.rewriter.enabled) { - const rewriterLlmComplete: RewriterLlmComplete = async (params) => { - if (params.signal.aborted) { - throw new DOMException("aborted", "AbortError"); - } - const abortPromise = new Promise((_, reject) => { - params.signal.addEventListener( - "abort", - () => reject(new DOMException("aborted", "AbortError")), - { once: true }, - ); - }); - const completionPromise = llmComplete({ + const rewriterLlmComplete: RewriterLlmComplete = abortableSubcall( + llmComplete, + (params: Parameters[0]) => ({ prompt: params.prompt, grammar: params.grammar, slotId: params.slotId, @@ -2242,9 +2214,8 @@ export async function createAgentRuntime( ...(params.responseFormat ? { responseFormat: params.responseFormat } : {}), - }); - return Promise.race([completionPromise, abortPromise]); - }; + }), + ); const rewriterCfg = config.memory.retrieve.rewriter; let gate: RewriterGate; if (rewriterCfg.gateMode === "embedding") { @@ -3044,18 +3015,9 @@ export async function createAgentRuntime( // `reserveReflectionSlot` again here is idempotent — the slot // manager returns the same id. const distillSlot = slotManager.reserveReflectionSlot() ?? -1; - const distillLlmComplete: ReflectionLlmComplete = async (params) => { - if (params.signal.aborted) { - throw new DOMException("aborted", "AbortError"); - } - const abortPromise = new Promise((_, reject) => { - params.signal.addEventListener( - "abort", - () => reject(new DOMException("aborted", "AbortError")), - { once: true }, - ); - }); - const completionPromise = llmComplete({ + const distillLlmComplete: ReflectionLlmComplete = abortableSubcall( + llmComplete, + (params: Parameters[0]) => ({ prompt: params.prompt, grammar: params.grammar, slotId: params.slotId, @@ -3063,9 +3025,8 @@ export async function createAgentRuntime( ...(params.responseFormat ? { responseFormat: params.responseFormat } : {}), - }); - return Promise.race([completionPromise, abortPromise]); - }; + }), + ); const distillRunner = new DistillRunner({ llmComplete: distillLlmComplete, slotId: distillSlot, @@ -3552,21 +3513,11 @@ function buildReflectionRunner(args: { { fallbackSlotId: reflectionSlotId }, ); } - const reflectionLlmComplete: ReflectionLlmComplete = async (params) => { - if (params.signal.aborted) { - throw new DOMException("aborted", "AbortError"); - } - const abortPromise = new Promise((_, reject) => { - params.signal.addEventListener( - "abort", - () => reject(new DOMException("aborted", "AbortError")), - { once: true }, - ); - }); - const { signal: _signal, ...rest } = params; - const completionPromise = args.llmComplete(rest); - return Promise.race([completionPromise, abortPromise]); - }; + const reflectionLlmComplete: ReflectionLlmComplete = abortableSubcall( + args.llmComplete, + ({ signal: _signal, ...rest }: Parameters[0]) => + rest, + ); const notesWriteEnabled = memory.notes.enabled && memory.reflection.autoStoreNotes && diff --git a/src/runtime/llm-fallback-seam.ts b/src/runtime/llm-fallback-seam.ts index cca29efc..5cfb1234 100644 --- a/src/runtime/llm-fallback-seam.ts +++ b/src/runtime/llm-fallback-seam.ts @@ -84,6 +84,17 @@ export interface FallbackSeamDeps extends LinkAttemptDeps { * chosen for. A pinned failure is rethrown as-is — the orchestrator is * the retry authority, and the chain's breaker state stays untouched by * a link it did not pick. + * + * **A request its caller aborted fails as a cancellation.** The unary + * clients do not surface an abort in one shape: `LlamaServerClient` wraps + * it as a `status: null` `LlamaServerError`, and `runOpenAiWithRetry` + * throws an `OpenAiHttpError` when it sees the signal already aborted — + * both classify `transport`, which `shouldAdvance` treats as an immediate + * provider-down signal. A memory sub-call whose timeout fired would then + * trip the breaker and flip the sticky override for a link that was fine. + * Rethrowing the signal's reason (abort-shaped by construction) makes it + * `cancelled`, which never advances — the same rule + * `OpenAiProvider.completeStream` applies on the streaming path. */ export function createFallbackCompleter( deps: FallbackSeamDeps, @@ -92,11 +103,13 @@ export function createFallbackCompleter( providerId: string, params: LlmStreamParams, ): Promise => { - const { result, transport } = await completeOnLink( - deps, - params, - providerId, - ); + let served: Awaited>; + try { + served = await completeOnLink(deps, params, providerId); + } catch (err) { + throw params.signal?.aborted ? cancellationOf(params.signal, err) : err; + } + const { result, transport } = served; deps.recordUnaryUsage(params, result, providerId); return { ...result, servedTransport: transport }; }; @@ -158,6 +171,15 @@ export function createFallbackStreamer( }; } +/** + * The error an aborted request fails with: the signal's own reason, or + * the original error for signal doubles that never populate `reason`. + */ +function cancellationOf(signal: AbortSignal, fallback: unknown): unknown { + const reason: unknown = signal.reason; + return reason ?? fallback; +} + /** * Stamp the serving link's transport on every chunk (see the * `StreamChunk.servedTransport` contract — the final result's stamp