From caec770447c164f3873d303e548f43f8dcbccffe Mon Sep 17 00:00:00 2001 From: Anarchid Date: Wed, 16 Sep 2026 18:19:20 +0300 Subject: [PATCH] fix(kv-unified): close the predecessor's open receipt flight before a new activation submits MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The kv-unified wire receipt fires per provider attempt, before the adapter call, and only that attempt's usage event settles it (driveStream's usage handler). An attempt that dies without usage — transport error, idle timeout, a framework cancel for a budget restart or endTurn — left its flight open until the old driveStream's `finally` drained it, and every successor path (error-policy retry from inside the failed stream's event loop, budget restart, the next wake after an abort) starts the new stream first. The successor's first receipt then hit the ledger's single-flight guard and the recovery itself failed: [inference-failed] agent=devops consecutive=1: kv-unified submission devops:46:1789544019150:4 is still in flight (2026-09-16 07:33Z) Agent now keeps the live receipt queue and, at the start of every activation, fails whatever the previous one left open before any new receipt can begin. The per-stream `finally` drain stays as the backstop; whichever runs first empties the shared queue. Also: inference-log records below the blob threshold are persisted through their JSON view. The compiled request carries the receipt hook (a function), which Chronicle rejected with "JS functions cannot be represented as a serde_json::Value" — thrown from the failure path it was logging. Test: framework-level repro (real kv-unified strategy, mock membrane whose first attempt dies after the receipt) and an Agent-level ordering test. Co-Authored-By: Claude Fable 5.1 --- ...x-kv-unified-open-flight-on-retry.fixed.md | 2 + src/agent.ts | 41 ++++ src/framework.ts | 9 + test/kv-unified-open-flight-retry.test.ts | 211 ++++++++++++++++++ test/kv-unified-wiring.test.ts | 42 ++++ 5 files changed, 305 insertions(+) create mode 100644 changelog.d/fix-kv-unified-open-flight-on-retry.fixed.md create mode 100644 test/kv-unified-open-flight-retry.test.ts diff --git a/changelog.d/fix-kv-unified-open-flight-on-retry.fixed.md b/changelog.d/fix-kv-unified-open-flight-on-retry.fixed.md new file mode 100644 index 00000000..6e3879c6 --- /dev/null +++ b/changelog.d/fix-kv-unified-open-flight-on-retry.fixed.md @@ -0,0 +1,2 @@ +- kv-unified: a new activation now closes any receipt flight its predecessor left unsettled before it submits its own. A provider attempt that died without a usage event (transport error, idle timeout, budget-restart or endTurn cancel) used to keep its flight open until the old stream's teardown, and every successor path started first — so the retry itself failed with `kv-unified submission … is still in flight` (devops agent, 2026-09-16). +- Inference-log records below the blob threshold are persisted through their JSON view; the kv-unified receipt hook on a compiled request (a function) made Chronicle reject the record with `JS functions cannot be represented as a serde_json::Value`, throwing from the failure path it was logging. diff --git a/src/agent.ts b/src/agent.ts index c3d6d160..c8ea7b60 100644 --- a/src/agent.ts +++ b/src/agent.ts @@ -99,6 +99,11 @@ export class Agent { private _state: AgentState = { status: 'idle' }; private _inferenceStartedAt = 0; private _streamId = 0; + /** kv-unified receipt flights opened by the CURRENT activation and not yet + * settled by a usage event. Held on the agent (not only in the stream's + * closure) so the next activation can close whatever its predecessor left + * open — see failOpenKvSubmissions. */ + private kvOpenQueue: Array<{ submissionId: string; wireReceipt: CacheWireReceipt }> | null = null; lastStreamInputTokens = 0; /** Real prefix size of the last usage event: fresh + cache creation + * cache read. THE window-shaped number — `lastStreamInputTokens` alone @@ -730,6 +735,19 @@ export class Agent { throw new Error(`Agent ${this.name} cannot start stream in state ${this._state.status}`); } + // A new activation supersedes every receipt flight the previous one left + // open. Membrane fires the kv-unified wire receipt per provider attempt, + // BEFORE the adapter call; only that attempt's usage event settles it. + // When the attempt dies without usage (transport error, idle timeout, + // framework cancel for a budget restart or endTurn) the flight stays open + // until the old driveStream's `finally` runs — and every successor path + // (error-policy retry, budget restart, the next wake after an abort) + // starts the new stream BEFORE that `finally`. The successor's first + // receipt then met the ledger's single-flight guard and failed the + // recovery itself: "kv-unified submission devops:46:…:4 is still in + // flight" (devops agent, 2026-09-16 07:33Z, no provider error logged). + this.failOpenKvSubmissions(); + this._streamId++; this._inferenceStartedAt = Date.now(); this.lastStreamInputTokens = 0; @@ -753,6 +771,7 @@ export class Agent { const kvQueue: Array<{ submissionId: string; wireReceipt: CacheWireReceipt }> = []; let kvCall = 0; if (kvEnabled) { + this.kvOpenQueue = kvQueue; const layoutHash = stableHash(request.messages); const kvRequest = request as NormalizedRequest & KvUnifiedRequestHooks; kvRequest.onCacheWireReceipt = (receipt) => { @@ -791,6 +810,28 @@ export class Agent { }; } + /** Close every kv-unified receipt flight the previous activation left + * unsettled. Idempotent: the framework's driveStream `finally` drains the + * same per-stream queue, so whichever runs first empties it and the other + * finds nothing. Failing an already-settled id is a no-op in the ledger. */ + private failOpenKvSubmissions(): void { + const open = this.kvOpenQueue?.splice(0) ?? []; + this.kvOpenQueue = null; + if (open.length === 0) return; + const strategy = (this.contextManager as unknown as { getStrategy?: () => unknown }) + .getStrategy?.() as { reportKvUnifiedFailed?: (submissionId: string) => void } | undefined; + for (const { submissionId } of open) { + console.error( + `[kv-unified] ${this.name}: closing receipt flight ${submissionId} left open by the previous activation`, + ); + try { + strategy?.reportKvUnifiedFailed?.(submissionId); + } catch (err) { + console.error(`[kv-unified] ${this.name}: could not close flight ${submissionId}:`, err); + } + } + } + /** * Transition to waiting_for_tools when stream yields tool calls. * Called by framework's driveStream. diff --git a/src/framework.ts b/src/framework.ts index db55fc98..911e59fe 100644 --- a/src/framework.ts +++ b/src/framework.ts @@ -8853,11 +8853,18 @@ export class AgentFramework { // Blob threshold: 10KB - typical context-heavy requests exceed this const BLOB_THRESHOLD = 10000; + // Both branches persist the JSON view, never the live object: a compiled + // request carries non-JSON members (the kv-unified `onCacheWireReceipt` + // hook is a function) and Chronicle's JSON bridge rejects those with + // "JS functions cannot be represented as a serde_json::Value" — which, + // thrown from the failure path, masked the failure being logged. if (entry.request && typeof entry.request === 'object') { const requestJson = JSON.stringify(entry.request); if (requestJson.length > BLOB_THRESHOLD) { const blobId = this.store.storeBlob(Buffer.from(requestJson), 'application/json'); entryToStore.request = { blobId }; + } else { + entryToStore.request = JSON.parse(requestJson); } } @@ -8866,6 +8873,8 @@ export class AgentFramework { if (responseJson.length > BLOB_THRESHOLD) { const blobId = this.store.storeBlob(Buffer.from(responseJson), 'application/json'); entryToStore.response = { blobId }; + } else { + entryToStore.response = JSON.parse(responseJson); } } diff --git a/test/kv-unified-open-flight-retry.test.ts b/test/kv-unified-open-flight-retry.test.ts new file mode 100644 index 00000000..9f38473c --- /dev/null +++ b/test/kv-unified-open-flight-retry.test.ts @@ -0,0 +1,211 @@ +/** + * kv-unified receipt flight vs. the framework's error-policy retry. + * + * Field incident (devops agent on gpt-6-astra, 2026-09-16 10:42 local): + * [inference-failed] kv-unified submission devops:46:1789544019150:4 is still in flight + * + * Membrane fires the cache-wire receipt once per provider attempt, BEFORE the + * adapter call, and the framework only settles that flight on the attempt's + * usage event (accept) or in driveStream's `finally` (fail). A provider call + * that dies before its usage event therefore leaves the flight open until the + * stream is torn down — but the error-policy retry starts the successor stream + * from INSIDE the failed stream's event loop, before that `finally` runs. The + * successor's first receipt then hits the ledger's single-flight guard and the + * retry itself fails with "is still in flight", turning one transient provider + * error into a lost turn. + */ +import { test } from 'node:test'; +import assert from 'node:assert/strict'; +import { mkdtempSync, rmSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import type { NormalizedRequest, StreamEvent, YieldingStream } from '@animalabs/membrane'; +import { AutobiographicalStrategy } from '@animalabs/context-manager'; +import { AgentFramework } from '../src/index.js'; +import type { + ErrorAction, + ErrorPolicy, + EventResponse, + Module, + ModuleContext, + ProcessEvent, + ProcessState, + ToolCall, + ToolDefinition, + ToolResult, + TraceEvent, +} from '../src/index.js'; +import type { KvUnifiedRequestHooks } from '../src/kv-unified-wire.js'; +import { createMockResponse, MockMembrane } from './helpers/mock-membrane.js'; + +/** Mirrors membrane's yielding stream closely enough for the receipt seam: + * the wire receipt fires inside streamOnce, after `await applyBeforeRequestHook`, + * i.e. a microtask or two after the consumer's first pull, and before any + * usage event. A receipt hook that throws surfaces as an `error` event (the + * real stream's startInference catch), never as a thrown `next()`. */ +class ScriptedStream { + private done = false; + constructor( + private readonly onFirstPull: () => void, + private readonly script: StreamEvent[], + ) {} + cancel(): void { this.done = true; } + get isWaitingForTools(): boolean { return false; } + get pendingToolCallIds(): string[] { return []; } + get toolDepth(): number { return 0; } + provideToolResults(): void { throw new Error('ScriptedStream is not waiting for tools'); } + async *[Symbol.asyncIterator](): AsyncIterator { + await Promise.resolve(); + await Promise.resolve(); + try { + this.onFirstPull(); + } catch (error) { + yield { type: 'error', error: error as Error } as StreamEvent; + return; + } + for (const event of this.script) { + if (this.done) return; + yield event; + } + } +} + +/** First provider attempt dies after the receipt fired and before any usage + * event (a transport error); every later attempt completes normally. */ +class DeadThenAliveMembrane extends MockMembrane { + readonly order: string[] = []; + private attempts = 0; + + streamYielding(request: NormalizedRequest): YieldingStream { + this.calls.push(request); + const attempt = this.attempts++; + const hooks = request as NormalizedRequest & KvUnifiedRequestHooks; + const fireReceipt = () => { + this.order.push(`pull:${attempt}`); + hooks.onCacheWireReceipt?.({ requestHash: `wire-${attempt}`, markers: [] }); + this.order.push(`receipt:${attempt}`); + }; + const script: StreamEvent[] = attempt === 0 + ? [{ type: 'error', error: new Error('socket hang up') } as StreamEvent] + : [ + { type: 'usage', usage: { inputTokens: 10, outputTokens: 5 } } as StreamEvent, + { + type: 'complete', + response: createMockResponse([{ type: 'text', text: 'recovered' }]), + } as StreamEvent, + ]; + return new ScriptedStream(fireReceipt, script) as unknown as YieldingStream; + } +} + +class WakeModule implements Module { + readonly name = 'wake'; + async start(_ctx: ModuleContext): Promise {} + async stop(): Promise {} + getTools(): ToolDefinition[] { return []; } + async handleToolCall(_call: ToolCall): Promise { + return { success: false, error: 'no tools', isError: true }; + } + async onProcess(event: ProcessEvent, _state: ProcessState): Promise { + if (event.type !== 'external-message') return {}; + return { + addMessages: [{ participant: 'User', content: [{ type: 'text', text: String(event.content) }] }], + requestInference: true, + }; + } +} + +class RetryOncePolicy implements ErrorPolicy { + maxRetries = 1; + onInferenceError(_error: Error, _agentName: string, attempt: number): ErrorAction { + return attempt < this.maxRetries ? { retry: true, delayMs: 0 } : { retry: false }; + } +} + +function kvUnifiedStrategy(): AutobiographicalStrategy { + return new AutobiographicalStrategy({ + adaptiveResolution: true, + foldingStrategy: 'kv-unified', + headWindowTokens: 0, + recentWindowTokens: 100, + kvUnified: { + policy: { + alpha: 0.7, + budgetLowRatio: 0.5, + budgetHighRatio: 0.9, + budgetUnderLambda: 10, + budgetOverLambda: 10, + cacheLambda: 1, + cacheScale: 1000, + cacheReadPrice: 0.1, + cacheWritePrice: 1.25, + continuityLambda: 1, + continuityScale: 1000, + continuityRecencyHalfLifeTokens: 1000, + continuityRecencyFloor: 0.2, + continuityStableHalfLife: 10, + continuityStableFloor: 0.25, + }, + tokenBucketSize: 100, + continuityBucketSize: 100, + fidelityBucketSize: 100, + labelCeiling: 10_000, + adoptEpsilon: 0, + treeifyNonContiguousSummaries: false, + preserveGapBearingSummaries: true, + }, + } as ConstructorParameters[0]); +} + +test('error-policy retry settles the dead attempt\'s kv-unified flight before the successor submits', async () => { + const tempDir = mkdtempSync(join(tmpdir(), 'kv-open-flight-')); + const membrane = new DeadThenAliveMembrane(); + const strategy = kvUnifiedStrategy(); + // Observe the ledger from the outside: the framework reaches the strategy + // through getStrategy(), so instance-level wrapping sees every call. + const spied = strategy as unknown as { + reportKvUnifiedFailed: (id: string) => void; + reportKvUnifiedAccepted: (args: { submissionId: string }) => void; + }; + const origFail = spied.reportKvUnifiedFailed.bind(strategy); + spied.reportKvUnifiedFailed = (id: string) => { membrane.order.push(`fail:${id}`); origFail(id); }; + const origAccept = spied.reportKvUnifiedAccepted.bind(strategy); + spied.reportKvUnifiedAccepted = (args: { submissionId: string }) => { + membrane.order.push(`accept:${args.submissionId}`); + origAccept(args); + }; + + const framework = await AgentFramework.create({ + storePath: join(tempDir, 'store.chronicle'), + membrane: membrane.asMembrane(), + agents: [{ name: 'devops', model: 'test-model', systemPrompt: 'Test', strategy }], + modules: [new WakeModule()], + errorPolicy: new RetryOncePolicy(), + syncIntervalMs: 0, + }); + const traces: TraceEvent[] = []; + framework.onTrace((event) => traces.push(event)); + + try { + framework.pushEvent({ type: 'external-message', source: 'test', content: 'hello', metadata: {} }); + await framework.runUntilIdle(); + + const failures = traces + .filter((t): t is TraceEvent & { error: string } => t.type === 'inference:failed') + .map((t) => t.error); + assert.deepEqual(failures, ['socket hang up'], `only the provider error may fail: ${JSON.stringify(membrane.order)}`); + assert.ok(traces.some((t) => t.type === 'inference:completed'), `retry must complete the turn: ${JSON.stringify(membrane.order)}`); + assert.ok(!traces.some((t) => t.type === 'inference:exhausted'), 'a single transient provider error must not exhaust the turn'); + + // The dead attempt's flight is closed BEFORE the successor's first pull, + // so the successor's receipt never meets an open flight. + const failIndex = membrane.order.findIndex((entry) => entry.startsWith('fail:')); + const successorPull = membrane.order.indexOf('pull:1'); + assert.ok(failIndex >= 0, `dead flight was never failed: ${JSON.stringify(membrane.order)}`); + assert.ok(failIndex < successorPull, `dead flight failed after the successor started: ${JSON.stringify(membrane.order)}`); + assert.ok(membrane.order.some((entry) => entry.startsWith('accept:')), 'successor flight must be accepted on its usage event'); + } finally { + await framework.stop(); + rmSync(tempDir, { recursive: true, force: true }); + } +}); diff --git a/test/kv-unified-wiring.test.ts b/test/kv-unified-wiring.test.ts index 3c76e622..4503cd8e 100644 --- a/test/kv-unified-wiring.test.ts +++ b/test/kv-unified-wiring.test.ts @@ -66,3 +66,45 @@ test('agent leaves marker ownership unchanged for non-kv strategies', async () = assert.equal((request as typeof request & KvUnifiedRequestHooks).cacheMarkers, undefined); assert.equal(compileOptions, undefined); }); + +test('a new activation closes the receipt flight its predecessor left open before it submits', async () => { + const order: string[] = []; + const strategy = { + isKvUnifiedEnabled: () => true, + beginKvUnifiedSubmission: (args: { submissionId: string }) => { order.push(`begin:${args.submissionId}`); }, + reportKvUnifiedFailed: (submissionId: string) => { order.push(`fail:${submissionId}`); }, + }; + const cm = { + getStrategy: () => strategy, + setToolDefinitions: () => {}, + compile: async () => ({ + messages: [{ participant: 'user', content: [{ type: 'text', text: 'hello' }] }], + systemInjections: [], + }), + } as unknown as ContextManager; + const membrane = { streamYielding: () => ({}) as never } as unknown as Membrane; + const agent = new Agent({ name: 'devops', model: 'test', systemPrompt: 'system' }, cm, membrane); + + // Activation 1: the provider attempt fires its receipt and then dies with no + // usage event (transport error / idle timeout / framework cancel). Nothing + // takes the submission; the stream's own teardown has not run yet. + const first = await agent.startStreamWithInjections([], undefined); + (first.request as typeof first.request & KvUnifiedRequestHooks) + .onCacheWireReceipt?.({ requestHash: 'wire-1', markers: [] }); + const firstId = order[0]?.slice('begin:'.length); + assert.ok(firstId, 'first activation must have begun a submission'); + agent.reset(); + + // Activation 2 (retry, restart, or the next wake) must supersede that flight + // BEFORE its own receipt reaches the ledger. + const second = await agent.startStreamWithInjections([], undefined); + (second.request as typeof second.request & KvUnifiedRequestHooks) + .onCacheWireReceipt?.({ requestHash: 'wire-2', markers: [] }); + assert.equal(order.length, 3, JSON.stringify(order)); + assert.equal(order[1], `fail:${firstId}`); + assert.ok(order[2]!.startsWith('begin:') && order[2] !== order[0]); + // The predecessor's late teardown finds nothing left to fail. + assert.deepEqual(first.drainKvSubmissionIds?.(), []); + // The successor's own flight is still live for its usage event. + assert.equal(second.takeKvSubmission?.()?.wireReceipt.requestHash, 'wire-2'); +});