Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions changelog.d/fix-kv-unified-open-flight-on-retry.fixed.md
Original file line number Diff line number Diff line change
@@ -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.
41 changes: 41 additions & 0 deletions src/agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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;
Expand All @@ -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) => {
Expand Down Expand Up @@ -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.
Expand Down
9 changes: 9 additions & 0 deletions src/framework.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}

Expand All @@ -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);
}
}

Expand Down
211 changes: 211 additions & 0 deletions test/kv-unified-open-flight-retry.test.ts
Original file line number Diff line number Diff line change
@@ -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<StreamEvent> {
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<void> {}
async stop(): Promise<void> {}
getTools(): ToolDefinition[] { return []; }
async handleToolCall(_call: ToolCall): Promise<ToolResult> {
return { success: false, error: 'no tools', isError: true };
}
async onProcess(event: ProcessEvent, _state: ProcessState): Promise<EventResponse> {
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<typeof AutobiographicalStrategy>[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 });
}
});
42 changes: 42 additions & 0 deletions test/kv-unified-wiring.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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');
});