From 0cb003b8e270d04ee67d67c8615deb27b16949bd Mon Sep 17 00:00:00 2001 From: drewstone Date: Wed, 16 Sep 2026 05:33:37 +0000 Subject: [PATCH] fix(runtime): keep long-running journals and retries reliable --- api-surface.json | 2 +- docs/agent-managed-compute/reliability.md | 9 + docs/api/primitive-catalog.md | 2 +- docs/canonical-api.md | 2 +- package.json | 2 +- src/durable/observer-journal.ts | 14 +- src/durable/spawn-journal.ts | 381 +++++++++++------- src/durable/tests/long-journal.test.ts | 220 ++++++++++ src/runtime/supervise/driver-retry.test.ts | 68 ++++ src/runtime/supervise/driver-retry.ts | 3 + .../fixtures/agent-improvement-proposal.json | 10 +- .../agent-profile-improvement-proposal.json | 6 +- 12 files changed, 561 insertions(+), 158 deletions(-) create mode 100644 src/durable/tests/long-journal.test.ts diff --git a/api-surface.json b/api-surface.json index 871e97b3..8e9f205c 100644 --- a/api-surface.json +++ b/api-surface.json @@ -940,7 +940,7 @@ "FileCoordinationLog": "value bd8f02ee8c2c", "FileCorpus": "value ac872ff9718b", "FileResultBlobStore": "value a464e1a75d77", - "FileSpawnJournal": "value aead755d0085", + "FileSpawnJournal": "value e9bda7d80f7c", "FinalizeContext": "type e45cea6eacef", "FinalizerSettled": "type c525a06427fc", "FlatWidenGate": "type 825b8b585b72", diff --git a/docs/agent-managed-compute/reliability.md b/docs/agent-managed-compute/reliability.md index a4ff09a9..38b5647c 100644 --- a/docs/agent-managed-compute/reliability.md +++ b/docs/agent-managed-compute/reliability.md @@ -28,6 +28,15 @@ Cancellation and cleanup failures preserve accepted output; budget and execution The built-in provider recovery path supports one-shot, nonsteering execution. Retained execution does not create a steering capability the provider lacks. +Spawn-journal append validation keeps only rebuildable node, admission, and cursor indexes. +Healthy appends do not reload prior events; cold reads validate each record once. +File replacement, truncation, changed metadata, or a failed write invalidates the append index. +Recovery validates the actual retained bytes before accepting another append. +The index is not a second durable record or a cross-process ownership fence. +Observer records detach their inputs before queued I/O, so later hook mutations cannot rewrite the evidence being saved. +Completed director invocations reset the consecutive transport-failure counter even when the pursuit remains incomplete. +The explicit total-attempt, deadline, cancellation, and resource bounds still apply. + The file run lock protects one local coordinator. It does not fence provider mutations from a partitioned coordinator on another machine. No deployed or live multi-provider recovery proof is claimed here. diff --git a/docs/api/primitive-catalog.md b/docs/api/primitive-catalog.md index 43b00003..87cb5019 100644 --- a/docs/api/primitive-catalog.md +++ b/docs/api/primitive-catalog.md @@ -7,7 +7,7 @@ # Primitive catalog — the never-stale anti-reinvention inventory -> **GENERATED** from `@tangle-network/agent-runtime@0.231.1` and `@tangle-network/agent-eval@0.182.0` by `scripts/gen-primitive-catalog.mjs`. Do NOT hand-edit — run `pnpm run docs:api`. This is the mechanical companion to the JUDGMENT in `canonical-api.md` (§2 decision table + §1.5 AgentProfile law): that doc says WHICH primitive to reach for and what NOT to build; this catalog proves WHAT exists. Per-symbol signatures + `file:line` live in the per-module pages under `docs/api/`. +> **GENERATED** from `@tangle-network/agent-runtime@0.232.0` and `@tangle-network/agent-eval@0.182.0` by `scripts/gen-primitive-catalog.mjs`. Do NOT hand-edit — run `pnpm run docs:api`. This is the mechanical companion to the JUDGMENT in `canonical-api.md` (§2 decision table + §1.5 AgentProfile law): that doc says WHICH primitive to reach for and what NOT to build; this catalog proves WHAT exists. Per-symbol signatures + `file:line` live in the per-module pages under `docs/api/`. ## 1. agent-runtime — own public surface diff --git a/docs/canonical-api.md b/docs/canonical-api.md index 3570a6fc..2ecb99b6 100644 --- a/docs/canonical-api.md +++ b/docs/canonical-api.md @@ -4,7 +4,7 @@ Generated signatures and the complete export list live in docs/api/. Run pnpm docs:freshness after editing this file. --> -> **Version 0.231.1.** +> **Version 0.232.0.** > [`docs/api/primitive-catalog.md`](./api/primitive-catalog.md) lists every export and import path. > `agent-eval` must satisfy `>=0.182.0 <0.183.0`. > `sandbox` must satisfy `>=0.36.4 <0.41.0`. diff --git a/package.json b/package.json index f2ebf77a..a266f53f 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@tangle-network/agent-runtime", - "version": "0.231.1", + "version": "0.232.0", "description": "Shared task-lifecycle skeleton for agents: a recursive loop kernel for chat turns, one-shot tasks, and multi-attempt loops, with trace capture and eval-gated self-improvement. Domain behavior lives in adapters; scoring and ship-gates in @tangle-network/agent-eval.", "homepage": "https://github.com/tangle-network/agent-runtime#readme", "repository": { diff --git a/src/durable/observer-journal.ts b/src/durable/observer-journal.ts index 6f9d9924..85d8ee67 100644 --- a/src/durable/observer-journal.ts +++ b/src/durable/observer-journal.ts @@ -1,6 +1,7 @@ import { createHash } from 'node:crypto' import { readFile } from 'node:fs/promises' import { resolve } from 'node:path' +import { detachedSnapshot } from '../runtime/supervise/snapshot' import { type RuntimeDecisionPoint, type RuntimeHookEvent, @@ -108,6 +109,15 @@ export class FileObserverJournal implements ObserverJournal { ) } + // Other hooks receive this same input and may mutate it before queued I/O runs. + // Detach now, not after awaiting the previous append or computing the digest. + let snapshot: typeof value + try { + snapshot = detachedSnapshot(value, 'observer record') + } catch (error) { + return Promise.reject(error) + } + let result: ObserverRecord | undefined const operation = this.tail.then(async () => { this.assertComplete() @@ -126,8 +136,8 @@ export class FileObserverJournal implements ObserverJournal { observedAt: Date.now(), ...(this.previousDigest ? { previousDigest: this.previousDigest } : {}), ...(kind === 'event' - ? { event: value as RuntimeHookEvent } - : { decision: value as RuntimeDecisionPoint }), + ? { event: snapshot as RuntimeHookEvent } + : { decision: snapshot as RuntimeDecisionPoint }), } const record: ObserverRecord = Object.freeze({ ...unsigned, diff --git a/src/durable/spawn-journal.ts b/src/durable/spawn-journal.ts index f156953c..fd772eef 100644 --- a/src/durable/spawn-journal.ts +++ b/src/durable/spawn-journal.ts @@ -20,6 +20,7 @@ * @stable */ +import type { BigIntStats } from 'node:fs' import { assertValidSpend } from '../runtime/supervise/budget' import { assertNoSymlinkDescendant, @@ -224,7 +225,10 @@ function assertContentAddress(outRef: string, artifact: unknown): void { * @stable */ export class InMemorySpawnJournal implements SpawnJournal { - private readonly trees = new Map() + private readonly trees = new Map< + NodeId, + { begunAt: string; events: SpawnEvent[]; index: SpawnEventIndex } + >() async loadTree(root: NodeId): Promise { const tree = this.trees.get(root) @@ -242,7 +246,7 @@ export class InMemorySpawnJournal implements SpawnJournal { } return } - this.trees.set(root, { begunAt: at, events: [] }) + this.trees.set(root, { begunAt: at, events: [], index: new SpawnEventIndex(root) }) } async appendEvent(root: NodeId, ev: SpawnEvent): Promise { @@ -250,8 +254,10 @@ export class InMemorySpawnJournal implements SpawnJournal { if (!tree) { throw new Error(`appendEvent called for unknown spawn tree '${root}'; call beginTree first`) } - assertSeqUnique(root, tree.events, ev) - tree.events.push(detachedSnapshot(ev, 'spawn event')) + const event = detachedSnapshot(ev, 'spawn event') + tree.index.assert(event) + tree.events.push(event) + tree.index.add(event) } } @@ -266,6 +272,10 @@ export class InMemorySpawnJournal implements SpawnJournal { */ export class FileSpawnJournal implements SpawnJournal { private appendTail: Promise = Promise.resolve() + // Rebuildable validation state, not another journal or a copy of event payloads. + private appendIndex: + | { stamp: string | undefined; trees: Map } + | undefined constructor(private readonly path: string) {} @@ -280,6 +290,7 @@ export class FileSpawnJournal implements SpawnJournal { } let begun = false const events: SpawnEvent[] = [] + const index = new SpawnEventIndex(root) for (const record of parseCommittedJsonLines(text, this.path)) { if (record.root !== root) continue if (record.kind === 'begin') { @@ -290,7 +301,8 @@ export class FileSpawnJournal implements SpawnJournal { `spawn journal corrupted: event for tree '${root}' precedes its begin record`, ) } - assertSeqUnique(root, events, record.event) + index.assert(record.event) + index.add(record.event) events.push(record.event) } } @@ -299,44 +311,77 @@ export class FileSpawnJournal implements SpawnJournal { async beginTree(root: NodeId, at: string): Promise { return this.serializeAppend(async () => { - const existing = await this.loadTreeBegin(root) + const state = await this.validationIndex() + const existing = state.trees.get(root) if (existing) { - if (existing !== at) { + if (existing.begunAt !== at) { throw new Error( - `spawn tree '${root}' already begun in ${this.path} at ${existing}; refusing to overwrite with ${at}`, + `spawn tree '${root}' already begun in ${this.path} at ${existing.begunAt}; refusing to overwrite with ${at}`, ) } return } - await this.writeRecord({ kind: 'begin', root, at }) + state.stamp = await this.writeRecord({ kind: 'begin', root, at }) + state.trees.set(root, { begunAt: at, index: new SpawnEventIndex(root) }) }) } async appendEvent(root: NodeId, ev: SpawnEvent): Promise { const event = detachedSnapshot(ev, 'spawn event') return this.serializeAppend(async () => { - const events = await this.loadTree(root) - if (events === undefined) { + const state = await this.validationIndex() + const tree = state.trees.get(root) + if (tree === undefined) { throw new Error(`appendEvent called for unknown spawn tree '${root}'; call beginTree first`) } - assertSeqUnique(root, events, event) - await this.writeRecord({ kind: 'event', root, event }) + tree.index.assert(event) + state.stamp = await this.writeRecord({ kind: 'event', root, event }) + // An unacknowledged append cannot advance validation state. A failed write invalidates + // the index, so recovery reads the actual committed bytes before admitting another event. + tree.index.add(event) }) } - private async loadTreeBegin(root: NodeId): Promise { + private async validationIndex() { const fs = await import('node:fs/promises') - let text: string - try { - text = await fs.readFile(this.path, 'utf8') - } catch (err) { - if (isNoEntError(err)) return undefined - throw err + const stamp = await this.fileStamp() + if (this.appendIndex && this.appendIndex.stamp === stamp) return this.appendIndex + this.appendIndex = undefined + const trees = new Map() + if (stamp !== undefined) { + const text = await fs.readFile(this.path, 'utf8') + if ((await this.fileStamp()) !== stamp) { + throw new Error('spawn journal changed while rebuilding its append index') + } + for (const record of parseCommittedJsonLines(text, this.path)) { + if (record.kind === 'begin') { + if (!trees.has(record.root)) { + trees.set(record.root, { begunAt: record.at, index: new SpawnEventIndex(record.root) }) + } + } else { + const tree = trees.get(record.root) + if (!tree) { + throw new Error( + `spawn journal corrupted: event for tree '${record.root}' precedes its begin record`, + ) + } + tree.index.assert(record.event) + tree.index.add(record.event) + } + } } - for (const record of parseCommittedJsonLines(text, this.path)) { - if (record.root === root && record.kind === 'begin') return record.at + this.appendIndex = { stamp, trees } + return this.appendIndex + } + + private async fileStamp(): Promise { + const fs = await import('node:fs/promises') + try { + return journalFileStamp(await fs.stat(this.path, { bigint: true })) + } catch (error) { + if (isNoEntError(error)) return undefined + throw error } - return undefined } private async serializeAppend(operation: () => Promise): Promise { @@ -345,21 +390,36 @@ export class FileSpawnJournal implements SpawnJournal { return append } - private async writeRecord(record: SpawnJournalRecord): Promise { + private async writeRecord(record: SpawnJournalRecord): Promise { const fs = await import('node:fs/promises') const path = await import('node:path') - await fs.mkdir(path.dirname(this.path), { recursive: true }) - const needsSeparator = await prepareJsonlAppend(this.path) - const fh = await fs.open(this.path, 'a') try { - await writeAllBytes(fh, `${needsSeparator ? '\n' : ''}${JSON.stringify(record)}\n`) - await fh.sync() - } finally { - await fh.close() + await fs.mkdir(path.dirname(this.path), { recursive: true }) + const needsSeparator = await prepareJsonlAppend(this.path) + const fh = await fs.open(this.path, 'a') + try { + await writeAllBytes(fh, `${needsSeparator ? '\n' : ''}${JSON.stringify(record)}\n`) + await fh.sync() + const stamp = journalFileStamp(await fh.stat({ bigint: true })) + if ((await this.fileStamp()) !== stamp) { + throw new Error('spawn journal changed while acknowledging its append') + } + return stamp + } finally { + await fh.close() + } + } catch (error) { + this.appendIndex = undefined + throw error } } } +/** Detect replacement, truncation, and same-size edits; this is not a cross-process write lock. */ +function journalFileStamp(stat: BigIntStats): string { + return `${stat.dev}:${stat.ino}:${stat.size}:${stat.mtimeNs}:${stat.ctimeNs}` +} + /** * Load every journal tree owned by one recursive supervision run and flatten its nodes/events. * @@ -735,131 +795,164 @@ type SpawnJournalRecord = | { kind: 'begin'; root: NodeId; at: string } | { kind: 'event'; root: NodeId; event: SpawnEvent } -/** Retained records form one ordered admission chain; older journals need no such chain. */ -function assertRetainedExecutionOrder(events: SpawnEvent[], event: SpawnEvent): void { - if ( - event.kind !== 'execution-input' && - event.kind !== 'execution-admitted' && - event.kind !== 'execution-result' - ) - return - const nodeEvents = events.filter((item) => item.id === event.id) - let inputIndex = -1 - for (let index = 0; index < nodeEvents.length; index++) { - if (nodeEvents[index]!.kind === 'execution-input') inputIndex = index - } - const prior = inputIndex < 0 ? nodeEvents : nodeEvents.slice(inputIndex) - function fail(reason: string): never { - throw new Error(`spawn journal corrupted: retained execution '${event.id}' ${reason}`) - } - if (!nodeEvents.some((item) => item.kind === 'spawned')) fail('precedes its spawn') - if (nodeEvents.some(closesCursorSlot)) fail('follows its terminal settlement') - if (event.kind === 'execution-input') { - if (!/^sha256:[0-9a-f]{64}$/.test(event.taskRef)) fail('has an invalid task reference') - if (nodeEvents.some((item) => item.kind === 'execution-input' && item.seq === event.seq)) - fail('has duplicate input sequence') - if (inputIndex >= 0 && !prior.some((item) => item.kind === 'execution-result')) - fail('input replaces an unfinished invocation') - return - } - const admissions = prior.flatMap((item) => - item.kind === 'execution-admitted' ? [item.admission] : [], - ) - if (event.kind === 'execution-result') { - if (prior.some((item) => item.kind === 'execution-result')) fail('has duplicate result') - if (!admissions.some((item) => item.phase === 'dispatched')) fail('result precedes dispatch') - if (!/^sha256:[0-9a-f]{64}$/.test(event.outRef)) fail('has an invalid result reference') - assertValidSpend(event.spent, 'retained execution result') - executorFailureReason(event) - return - } - const admission = event.admission - if (!admission || !['intent', 'environment', 'dispatched'].includes(admission.phase)) - fail('has an invalid admission phase') - if (admissions.some((item) => item.phase === admission.phase)) - fail('has duplicate admission phase') - if (prior.some((item) => item.kind === 'execution-result')) fail('admission follows result') - if (!prior.some((item) => item.kind === 'execution-input')) fail('admission precedes input') - if (admission.phase === 'intent') { - if (admissions.length !== 0) fail('intent follows another admission') - return - } - const intent = admissions.find((item) => item.phase === 'intent') - if (intent?.phase !== 'intent') fail('admission precedes intent') - if (admission.idempotencyKey !== intent.idempotencyKey || admission.turnId !== intent.turnId) - fail('changes its admitted request identity') - if (admission.phase === 'environment') { - if ( - admission.provider !== intent.provider || - admission.sessionId !== intent.sessionId || - admission.executionId !== intent.executionId - ) - fail('changes its admitted execution identity') - return - } - const environment = admissions.find((item) => item.phase === 'environment') - if (environment?.phase !== 'environment') fail('dispatch precedes environment') - const control = admission.controlRef - if ( - !control || - control.provider !== intent.provider || - control.environmentId !== environment.environmentId || - control.sessionId !== intent.sessionId || - control.executionId !== intent.executionId - ) - fail('dispatch changes its admitted execution identity') +/** Minimal state needed to validate a node; progress and trace payloads are never retained here. */ +interface JournalNodeIndex { + spawned: boolean + closed: boolean + materialized: boolean + bindings: Set + inputs: Set + input: boolean + result: boolean + admissions: Map< + Extract['admission']['phase'], + Extract['admission'] + > } -/** - * Two `seq` namespaces share the journal: a `spawned` event's `seq` is the spawn ordinal - * (the order children were created), and a `settled`/`cancelled` event's `seq` is the - * monotonic CURSOR order `scope.next()` yielded that settlement (B2). The uniqueness - * replay rests on is the cursor namespace — two settlements cannot share the position - * replay orders by — so the guard checks only settled/cancelled events. A `spawned` - * ordinal legitimately equals a later `settled` cursor seq and is not a collision. - */ -function assertSeqUnique(root: NodeId, events: SpawnEvent[], ev: SpawnEvent): void { - assertRetainedExecutionOrder(events, ev) - if (ev.kind === 'materialized') { - if (events.some((event) => event.kind === 'materialized' && event.id === ev.id)) { - throw new Error( - `spawn journal corrupted: duplicate materialization receipt for node '${ev.id}' in tree '${root}'`, - ) +/** One validation implementation for memory, durable append, and cold replay. */ +class SpawnEventIndex { + private readonly nodes = new Map() + private readonly cursors = new Set() + + constructor(private readonly root: NodeId) {} + + assert(event: SpawnEvent): void { + const node = this.nodes.get(event.id) + this.assertRetained(node, event) + if (event.kind === 'materialized') { + if (node?.materialized) { + throw new Error( + `spawn journal corrupted: duplicate materialization receipt for node '${event.id}' in tree '${this.root}'`, + ) + } + if (!node?.spawned) { + throw new Error( + `spawn journal corrupted: materialization for node '${event.id}' precedes its spawn in tree '${this.root}'`, + ) + } + } + if (event.kind === 'execution-bound') { + if (node?.bindings.has(event.binding.attemptId)) { + throw new Error( + `spawn journal corrupted: duplicate execution binding for node '${event.id}' attempt '${event.binding.attemptId}' in tree '${this.root}'`, + ) + } + if (!node?.materialized) { + throw new Error( + `spawn journal corrupted: execution binding for node '${event.id}' precedes materialization in tree '${this.root}'`, + ) + } } - if (!events.some((event) => event.kind === 'spawned' && event.id === ev.id)) { + if (!outsideCursorNamespace(event) && this.cursors.has(event.seq)) { throw new Error( - `spawn journal corrupted: materialization for node '${ev.id}' precedes its spawn in tree '${root}'`, + `spawn journal corrupted: duplicate cursor seq ${event.seq} in tree '${this.root}'; ` + + 'the cursor order replay relies on is not unique', ) } } - if (ev.kind === 'execution-bound') { + + /** Called only after validation and, for the file writer, a successful durable append. */ + add(event: SpawnEvent): void { + if (!outsideCursorNamespace(event)) this.cursors.add(event.seq) if ( - events.some( - (event) => - event.kind === 'execution-bound' && - event.id === ev.id && - event.binding.attemptId === ev.binding.attemptId, - ) - ) { - throw new Error( - `spawn journal corrupted: duplicate execution binding for node '${ev.id}' attempt '${ev.binding.attemptId}' in tree '${root}'`, - ) + event.kind !== 'spawned' && + event.kind !== 'materialized' && + event.kind !== 'execution-bound' && + event.kind !== 'execution-input' && + event.kind !== 'execution-admitted' && + event.kind !== 'execution-result' && + !closesCursorSlot(event) + ) + return + let node = this.nodes.get(event.id) + if (!node) { + node = { + spawned: false, + closed: false, + materialized: false, + bindings: new Set(), + inputs: new Set(), + input: false, + result: false, + admissions: new Map(), + } + this.nodes.set(event.id, node) } - if (!events.some((event) => event.kind === 'materialized' && event.id === ev.id)) { - throw new Error( - `spawn journal corrupted: execution binding for node '${ev.id}' precedes materialization in tree '${root}'`, + if (event.kind === 'spawned') node.spawned = true + else if (event.kind === 'materialized') node.materialized = true + else if (event.kind === 'execution-bound') node.bindings.add(event.binding.attemptId) + else if (event.kind === 'execution-input') { + node.inputs.add(event.seq) + node.input = true + node.result = false + node.admissions.clear() + } else if (event.kind === 'execution-admitted') { + node.admissions.set(event.admission.phase, event.admission) + } else if (event.kind === 'execution-result') node.result = true + if (closesCursorSlot(event)) node.closed = true + } + + private assertRetained(node: JournalNodeIndex | undefined, event: SpawnEvent): void { + if ( + event.kind !== 'execution-input' && + event.kind !== 'execution-admitted' && + event.kind !== 'execution-result' + ) + return + function fail(reason: string): never { + throw new Error(`spawn journal corrupted: retained execution '${event.id}' ${reason}`) + } + if (!node?.spawned) fail('precedes its spawn') + if (node.closed) fail('follows its terminal settlement') + if (event.kind === 'execution-input') { + if (!/^sha256:[0-9a-f]{64}$/.test(event.taskRef)) fail('has an invalid task reference') + if (node.inputs.has(event.seq)) fail('has duplicate input sequence') + if (node.input && !node.result) fail('input replaces an unfinished invocation') + return + } + if (event.kind === 'execution-result') { + if (node.result) fail('has duplicate result') + if (!node.admissions.has('dispatched')) fail('result precedes dispatch') + if (!/^sha256:[0-9a-f]{64}$/.test(event.outRef)) fail('has an invalid result reference') + assertValidSpend(event.spent, 'retained execution result') + executorFailureReason(event) + return + } + const admission = event.admission + if (!admission || !['intent', 'environment', 'dispatched'].includes(admission.phase)) + fail('has an invalid admission phase') + if (node.admissions.has(admission.phase)) fail('has duplicate admission phase') + if (node.result) fail('admission follows result') + if (!node.input) fail('admission precedes input') + if (admission.phase === 'intent') { + if (node.admissions.size !== 0) fail('intent follows another admission') + return + } + const intent = node.admissions.get('intent') + if (intent?.phase !== 'intent') fail('admission precedes intent') + if (admission.idempotencyKey !== intent.idempotencyKey || admission.turnId !== intent.turnId) + fail('changes its admitted request identity') + if (admission.phase === 'environment') { + if ( + admission.provider !== intent.provider || + admission.sessionId !== intent.sessionId || + admission.executionId !== intent.executionId ) + fail('changes its admitted execution identity') + return } - } - // `spawned` (ordinal namespace), `waiting` (the wait-ordinal namespace — it CREATES a node, it - // does not settle one), and `metered` (informational spend, no settlement order) live outside - // the cursor-uniqueness namespace replay relies on. `woken` IS a settlement and does not. - if (outsideCursorNamespace(ev)) return - if (events.some((e) => !outsideCursorNamespace(e) && e.seq === ev.seq)) { - throw new Error( - `spawn journal corrupted: duplicate cursor seq ${ev.seq} in tree '${root}'; ` + - 'the cursor order replay relies on is not unique', + const environment = node.admissions.get('environment') + if (environment?.phase !== 'environment') fail('dispatch precedes environment') + const control = admission.controlRef + if ( + !control || + control.provider !== intent.provider || + control.environmentId !== environment.environmentId || + control.sessionId !== intent.sessionId || + control.executionId !== intent.executionId ) + fail('dispatch changes its admitted execution identity') } } diff --git a/src/durable/tests/long-journal.test.ts b/src/durable/tests/long-journal.test.ts new file mode 100644 index 00000000..4d3c4b22 --- /dev/null +++ b/src/durable/tests/long-journal.test.ts @@ -0,0 +1,220 @@ +import { appendFile, mkdtemp, readFile, rename, rm, stat, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, describe, expect, it, vi } from 'vitest' +import type { SpawnEvent } from '../../runtime/supervise/types' +import type { RuntimeHookEvent } from '../../runtime-hooks' +import { composeRuntimeHooks } from '../../runtime-hooks' +import { FileObserverJournal } from '../observer-journal' +import { FileSpawnJournal } from '../spawn-journal' + +const io = vi.hoisted(() => ({ reads: 0, failSync: false })) +vi.mock('node:fs/promises', async (original) => { + const fs = await original() + return { + ...fs, + open: async (...args: Parameters) => { + const handle = await fs.open(...args) + if (args[1] === 'a' && io.failSync) { + io.failSync = false + handle.sync = async () => { + throw new Error('injected fsync failure') + } + } + return handle + }, + readFile: (...args: Parameters) => { + io.reads++ + return fs.readFile(...args) + }, + } +}) +const dirs: string[] = [] +afterEach(async () => { + io.failSync = false + await Promise.all(dirs.splice(0).map((dir) => rm(dir, { recursive: true, force: true }))) +}) +async function file(name = 'journal.jsonl') { + const dir = await mkdtemp(join(tmpdir(), 'long-journal-')) + dirs.push(dir) + return join(dir, name) +} +const at = '2026-09-16T00:00:00.000Z' +const progress = (seq: number): SpawnEvent => ({ + kind: 'progress', + id: 'worker', + seq, + at, + spend: { iterations: seq, tokens: { input: seq, output: 0 }, usd: 0, ms: seq }, +}) +const cancelled = (seq: number): SpawnEvent => ({ + kind: 'cancelled', + id: `worker-${seq}`, + reason: 'test', + seq, + at, + spent: { iterations: 0, tokens: { input: 0, output: 0 }, usd: 0, ms: 0 }, +}) + +describe('long-lived spawn journal', () => { + it('does not reread the whole history for every owned append', async () => { + const journal = new FileSpawnJournal(await file()) + await journal.beginTree('root', at) + io.reads = 0 + for (let seq = 0; seq < 100; seq++) await journal.appendEvent('root', progress(seq)) + expect(io.reads).toBeLessThanOrEqual(1) + expect(await journal.loadTree('root')).toHaveLength(100) + }) + + it('rebuilds validation on restart and rejects already committed cursor positions', async () => { + const path = await file() + const journal = new FileSpawnJournal(path) + await journal.beginTree('root', at) + await journal.appendEvent('root', cancelled(1)) + const restarted = new FileSpawnJournal(path) + await expect(restarted.appendEvent('root', cancelled(1))).rejects.toThrow('duplicate cursor') + await restarted.appendEvent('root', cancelled(2)) + expect(await restarted.loadTree('root')).toHaveLength(2) + }) + + it('notices sequential writes by another journal instance', async () => { + const path = await file() + const first = new FileSpawnJournal(path) + await first.beginTree('root', at) + const second = new FileSpawnJournal(path) + await second.appendEvent('root', cancelled(1)) + await expect(first.appendEvent('root', cancelled(1))).rejects.toThrow('duplicate cursor') + await first.appendEvent('root', cancelled(2)) + expect(await second.loadTree('root')).toHaveLength(2) + }) + + it('invalidates a warm index when its file is replaced with equal-sized bytes', async () => { + const path = await file() + const journal = new FileSpawnJournal(path) + await journal.beginTree('aaaa', at) + const text = await readFile(path, 'utf8') + await writeFile(`${path}.new`, text.replace('aaaa', 'bbbb')) + expect((await stat(path)).size).toBe((await stat(`${path}.new`)).size) + await rename(`${path}.new`, path) + await expect(journal.appendEvent('aaaa', progress(0))).rejects.toThrow('unknown spawn tree') + await journal.appendEvent('bbbb', progress(0)) + expect(await journal.loadTree('bbbb')).toHaveLength(1) + }) + + it('detects an equal-sized in-place rewrite without relying on inode replacement', async () => { + const path = await file() + const journal = new FileSpawnJournal(path) + await journal.beginTree('aaaa', at) + const before = await stat(path, { bigint: true }) + await writeFile(path, (await readFile(path, 'utf8')).replace('aaaa', 'bbbb')) + const after = await stat(path, { bigint: true }) + expect(after.ino).toBe(before.ino) + expect(after.size).toBe(before.size) + await expect(journal.appendEvent('aaaa', progress(0))).rejects.toThrow('unknown spawn tree') + await journal.appendEvent('bbbb', progress(0)) + expect(await journal.loadTree('bbbb')).toHaveLength(1) + }) + + it('rebuilds after a failed fsync instead of forgetting bytes already written', async () => { + const path = await file() + const journal = new FileSpawnJournal(path) + await journal.beginTree('root', at) + io.failSync = true + await expect(journal.appendEvent('root', cancelled(1))).rejects.toThrow('injected fsync') + // The append was not acknowledged, but its bytes are present. Re-read instead of admitting + // a duplicate cursor based on the pre-write index or inventing that the write never happened. + await expect(journal.appendEvent('root', cancelled(1))).rejects.toThrow('duplicate cursor') + await journal.appendEvent('root', cancelled(2)) + expect(await journal.loadTree('root')).toHaveLength(2) + }) + + it('invalidates a warm index on truncation and does not resurrect its old tree', async () => { + const path = await file() + const journal = new FileSpawnJournal(path) + await journal.beginTree('root', at) + await writeFile(path, '') + await expect(journal.appendEvent('root', progress(0))).rejects.toThrow('unknown spawn tree') + await journal.beginTree('fresh', at) + expect(await journal.loadTree('root')).toBeUndefined() + }) + + it('recovers a torn tail without forgetting cursor uniqueness or another tree', async () => { + const path = await file() + const journal = new FileSpawnJournal(path) + await journal.beginTree('root', at) + await journal.beginTree('other', at) + await journal.appendEvent('root', cancelled(1)) + await appendFile(path, '{"kind":"event"') + await journal.appendEvent('other', cancelled(1)) + await expect(journal.appendEvent('root', cancelled(1))).rejects.toThrow('duplicate cursor') + expect(await journal.loadTree('root')).toHaveLength(1) + expect(await journal.loadTree('other')).toHaveLength(1) + }) + + it('still refuses committed corruption after its append index is warm', async () => { + const path = await file() + const journal = new FileSpawnJournal(path) + await journal.beginTree('root', at) + await appendFile(path, '{bad}\n') + await expect(journal.appendEvent('root', progress(0))).rejects.toThrow('malformed JSONL') + }) + + it('serializes validation with append and accepts only one duplicate cursor', async () => { + const journal = new FileSpawnJournal(await file()) + await journal.beginTree('root', at) + const results = await Promise.allSettled([ + journal.appendEvent('root', cancelled(1)), + journal.appendEvent('root', cancelled(1)), + ]) + expect(results.map((r) => r.status).sort()).toEqual(['fulfilled', 'rejected']) + expect(await journal.loadTree('root')).toHaveLength(1) + }) +}) + +describe('observer input ownership', () => { + const event = (): RuntimeHookEvent => ({ + id: 'event', + pursuitId: 'pursuit', + runId: 'run', + target: 'agent.child', + phase: 'event', + timestamp: 0, + payload: { outcome: { text: 'original' } }, + }) + it('snapshots an event before yielding to the append queue', async () => { + const journal = new FileObserverJournal(await file('observer.jsonl'), 'pursuit') + const input = event() + const pending = journal.appendEvent(input) + ;(input.payload as { outcome: { text: string } }).outcome.text = 'mutated' + await pending + expect((await journal.read())[0]?.event?.payload).toEqual({ outcome: { text: 'original' } }) + }) + + it('isolates the journal from another hook that mutates the shared event', async () => { + const journal = new FileObserverJournal(await file('observer.jsonl'), 'pursuit') + const hooks = composeRuntimeHooks(journal.hooks(), { + onEvent: (input) => { + ;(input.payload as { outcome: { text: string } }).outcome.text = 'rewritten' + }, + }) + await hooks.onEvent!(event(), {}) + expect((await journal.read())[0]?.event?.payload).toEqual({ outcome: { text: 'original' } }) + }) + + it('retains the original decision when a caller reuses its input objects', async () => { + const journal = new FileObserverJournal(await file('observer.jsonl'), 'pursuit') + const input = { + id: 'decision', + pursuitId: 'pursuit', + runId: 'run', + stepIndex: 0, + kind: 'continue' as const, + candidateActions: ['continue'], + evidence: [], + } + const pending = journal.appendDecision(input) + input.candidateActions[0] = 'stop' + await pending + expect((await journal.read())[0]?.decision?.candidateActions).toEqual(['continue']) + }) +}) diff --git a/src/runtime/supervise/driver-retry.test.ts b/src/runtime/supervise/driver-retry.test.ts index 242cea72..f775edf5 100644 --- a/src/runtime/supervise/driver-retry.test.ts +++ b/src/runtime/supervise/driver-retry.test.ts @@ -759,3 +759,71 @@ describe('driver admission after asynchronous callbacks', () => { expect(script.attempts).toEqual([1]) }) }) + +describe('long-run retry streaks', () => { + it('does not accumulate non-consecutive transport failures across completed continuations', async () => { + let calls = 0 + const records: DriverAttemptRecord[] = [] + await runDriverWithRetry({ + drive: async () => { + calls++ + if (calls % 2 === 1) throw new Error('temporary transport interruption') + }, + progress: () => mark({ contract: calls >= 12 ? 'met' : 'unmet' }), + budget: () => budget(), + signal: new AbortController().signal, + policy: { maxAttempts: 20, maxConsecutiveFailures: 3 }, + reprompt: { maxReprompts: 10 }, + onAttempt: (record) => { + records.push(record) + }, + sleep: instantSleep, + }) + expect(calls).toBe(12) + expect(records.filter((record) => record.error).map((record) => record.retryInMs)).toEqual([ + 2000, 2000, 2000, 2000, 2000, 2000, + ]) + expect(records.at(-1)?.contract).toBe('met') + }) + + it('still stops a truly consecutive failure streak after a successful continuation', async () => { + let calls = 0 + await expect( + runDriverWithRetry({ + drive: async () => { + calls++ + if (calls > 1) throw new Error('transport unavailable') + }, + progress: () => mark(), + budget: () => budget(), + signal: new AbortController().signal, + policy: { maxAttempts: 20, maxConsecutiveFailures: 3 }, + reprompt: { maxReprompts: 10 }, + sleep: instantSleep, + }), + ).rejects.toMatchObject({ stop: 'no-progress' }) + expect(calls).toBe(4) + }) + + it('does not weaken an explicit total attempt ceiling', async () => { + let calls = 0 + const records: DriverAttemptRecord[] = [] + await runDriverWithRetry({ + drive: async () => { + calls++ + if (calls % 2 === 1) throw new Error('temporary transport interruption') + }, + progress: () => mark(), + budget: () => budget(), + signal: new AbortController().signal, + policy: { maxAttempts: 4, maxConsecutiveFailures: 3 }, + reprompt: { maxReprompts: 10 }, + onAttempt: (record) => { + records.push(record) + }, + sleep: instantSleep, + }) + expect(calls).toBe(4) + expect(records.at(-1)?.repromptRefusedBy).toBe('max-attempts') + }) +}) diff --git a/src/runtime/supervise/driver-retry.ts b/src/runtime/supervise/driver-retry.ts index 193f128b..45f59427 100644 --- a/src/runtime/supervise/driver-retry.ts +++ b/src/runtime/supervise/driver-retry.ts @@ -513,6 +513,9 @@ export async function runDriverWithRetry(run: DriverRetryRun): Promise { continue } + // A completed invocation ends the transport-failure streak, even while the pursuit's + // independent completion check remains unmet. Successful continuation is not a failure. + consecutiveBarren = 0 const durationMs = now() - startedAt const after = run.progress() const progressed = madeProgress(before, after) diff --git a/src/testing/fixtures/agent-improvement-proposal.json b/src/testing/fixtures/agent-improvement-proposal.json index 7699974d..8ad18e32 100644 --- a/src/testing/fixtures/agent-improvement-proposal.json +++ b/src/testing/fixtures/agent-improvement-proposal.json @@ -1,6 +1,6 @@ { "changedSurfaces": ["prompt"], - "digest": "sha256:6e82ad15ab3d4ce4f73913de01f1e13ebc68bb4b4d9120b43382f01954dd956e", + "digest": "sha256:201a9a98fbdc078d03c35314e4e03f34abaa36d7db394a11c4b1110bd036521c", "evaluation": { "decision": { "contributingChecks": [ @@ -4882,7 +4882,7 @@ ], "metadata": { "fixture": "agent-improvement-proposal", - "runtimeVersion": "0.231.1" + "runtimeVersion": "0.232.0" }, "objectives": [ { @@ -4993,8 +4993,8 @@ "baselineContentHash": "sha256:5c21ee53e513fc604cb09754e21c392b24a424da0ef37dbf8f1ee4a8a0b08f09", "candidateContentHash": "sha256:60fcbb1c728194bd51d7d19cb732d1c3f1881dce7e0a6266b41c8b98cfd65693", "kind": "agent-eval-loop", - "recordDigest": "sha256:e26e7e4ff57e2509e042c553c638f579680f438553fbede8b8d51c35b2c6be68", - "runId": "agent-runtime-0.231.1-proposal-fixture", + "recordDigest": "sha256:d26f696c1a140606d4e75b1ade7a03f08ce13f4fc4fdfe0527d2fbb4717e438d", + "runId": "agent-runtime-0.232.0-proposal-fixture", "schema": "agent-candidate-experiment" } }, @@ -5021,5 +5021,5 @@ ], "kind": "agent-improvement-proposal", "proposedAt": "2026-07-10T01:00:00.000Z", - "runId": "agent-runtime-0.231.1-proposal-fixture" + "runId": "agent-runtime-0.232.0-proposal-fixture" } diff --git a/src/testing/fixtures/agent-profile-improvement-proposal.json b/src/testing/fixtures/agent-profile-improvement-proposal.json index 65432b0f..1d92c27a 100644 --- a/src/testing/fixtures/agent-profile-improvement-proposal.json +++ b/src/testing/fixtures/agent-profile-improvement-proposal.json @@ -1,6 +1,6 @@ { "changedSurfaces": ["prompt", "skills"], - "digest": "sha256:d02c9c48865e63a968095846d85416c925186c770caca400680b2448988414f8", + "digest": "sha256:702dc1a85ce46804d67de7fd179bd69f6af0027dca7b6a99c46ab7f159bf7f5f", "evaluation": { "decision": { "contributingChecks": [ @@ -1715,7 +1715,7 @@ ], "metadata": { "fixture": "agent-profile-improvement-proposal", - "runtimeVersion": "0.231.1" + "runtimeVersion": "0.232.0" }, "objectives": [ { @@ -1826,7 +1826,7 @@ "baselineContentHash": "sha256:21c495a37c418c10bde64fbaa188beddeed31f1f051ea60a6a6582a9ee0db704", "candidateContentHash": "sha256:103f77bc8481601eef1ad5fe6ba84a40dffabc3a44f421f8c8559121edab84e9", "kind": "agent-eval-loop", - "recordDigest": "sha256:5cad90b609c3f39e1fe23ea5a140b10fb58221ee0b7aabf12b62a8e5c0e3b58c", + "recordDigest": "sha256:6a55eeeef02f261ce41fa35063839ec37d0b64815dae2517e9d92c75dd2585ce", "runId": "profile-improvement-1", "schema": "agent-profile-improvement-experiment" }