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'); +});