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
8 changes: 8 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,14 @@

## 0.237.1

File-backed journal reads stream a fixed file prefix instead of allocating the full history as one string.
Observer restart verifies its complete digest chain while retaining only the last record.
Root-stream receipts hash and count the same raw-byte prefix, including an uncommitted tail in the hash.
Spawn, coordination, conversation, corpus and discovery readers share the same committed-record parser.
Corrupt committed lines, truncated snapshots and non-parse failures remain explicit errors; torn final writes retain their existing semantics.
Public array-returning methods still retain their selected history, and individual records remain subject to Node limits.
No execution policy, limits, credential behavior or durable wire format changes.

**A refused environment create fails closed instead of retrying forever.** A provider SDK throws its own error classes, so a `create` the platform refused arrived at the driver's classifier as neither a `BackendTransportError` nor an `AgentEvalError` and took the foreign-accident default: retry. `classifyDriverFailure` now reads a plain HTTP status off any thrown `Error` and applies the split the transport branch already promises — 408, 429 and 5xx are the upstream having a bad moment, any other 4xx is a request that will fail identically forever.

Measured 2026-09-16 on discovery-lab: a Tangle Sandbox create refused `HTTP 400 {"code":"CONFIG_ERROR"}`, for a key whose budget was fully reserved by an existing box, retried 14-22 times per node. Eleven of twelve roots showed a durable admission intent and no other event for 25 minutes, with healthy coordination servers and no error anywhere for an operator to read.
Expand Down
6 changes: 6 additions & 0 deletions docs/agent-managed-compute/reliability.md
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,12 @@ 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.
File-backed journal readers process a fixed file prefix one JSONL record at a time, rather than allocating the entire history as a string.
Observer restart verifies the complete digest chain while retaining only its last record.
Root-stream receipt hashing and event counting inspect the same raw-byte prefix, including any uncommitted tail in the hash.
A malformed committed record or a short snapshot read remains an error; only malformed unterminated tail records are ignored.
Public methods that return history arrays still retain their selected records in memory, and individual records remain subject to Node limits.
This removes the whole-file string limit; it is not constant-memory arbitrary replay or distributed snapshot isolation.
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 failure-attempt, deadline, cancellation, and resource bounds still apply.
Expand Down
19 changes: 4 additions & 15 deletions src/conversation/journal.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,12 +15,7 @@
* @stable
*/

import {
isNoEntError,
parseCommittedJsonLines,
prepareJsonlAppend,
writeAllBytes,
} from '../durable/jsonl-file'
import { prepareJsonlAppend, readCommittedJsonLines, writeAllBytes } from '../durable/jsonl-file'
import type { ConversationTurn, HaltReason } from './types'

export interface ConversationJournalEntry {
Expand Down Expand Up @@ -133,16 +128,10 @@ export class FileConversationJournal implements ConversationJournal {
constructor(private readonly path: string) {}

async loadRun(runId: string): Promise<ConversationJournalEntry | undefined> {
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
}
let entry: ConversationJournalEntry | undefined
for (const record of parseCommittedJsonLines<JournalRecord>(text, this.path)) {
for await (const record of readCommittedJsonLines<JournalRecord>(this.path, {
allowMissing: true,
})) {
if (record.runId !== runId) continue
if (record.kind === 'begin') {
entry = { runId, startedAt: record.startedAt, turns: [] }
Expand Down
82 changes: 71 additions & 11 deletions src/durable/jsonl-file.ts
Original file line number Diff line number Diff line change
@@ -1,24 +1,84 @@
import type { FileHandle } from 'node:fs/promises'

/** Parse an append-only JSONL file without treating a torn final write as committed data.
* A malformed newline-terminated or non-final record is corruption and fails loud. */
/** Parse in-memory evidence with the same commit boundary as streamed journal reads. */
export function parseCommittedJsonLines<T>(text: string, source: string): T[] {
const lines = text.split('\n')
const finalIndex = lines.length - 1
const records: T[] = []

for (const [index, line] of lines.entries()) {
if (line.length === 0) continue
const record = parseRecord<T>(line, source, index + 1, index < lines.length - 1)
if (record !== undefined) records.push(record.value)
}
return records
}

/** Read a fixed file prefix one record at a time, never allocating the whole journal as a string.
* A reader cannot chase concurrent appends forever. Missing files are empty only when requested;
* corruption, short reads, and other I/O errors remain failures. Closing the iterator closes the fd. */
export async function* readCommittedJsonLines<T>(
path: string,
options: { allowMissing?: boolean; onBytes?: (bytes: Uint8Array) => void } = {},
): AsyncGenerator<T> {
const fs = await import('node:fs/promises')
let handle: FileHandle
try {
handle = await fs.open(path, 'r')
} catch (error) {
if (options.allowMissing && isNoEntError(error)) return
throw error
}
try {
const { size } = await handle.stat()
if (size === 0) return
const { StringDecoder } = await import('node:string_decoder')
const decoder = new StringDecoder('utf8')
const stream = handle.createReadStream({ autoClose: false, end: size - 1 })
const fragments: string[] = []
let lineNumber = 1
try {
records.push(JSON.parse(line) as T)
} catch (cause) {
const isInvalidUnterminatedTail = index === finalIndex && !text.endsWith('\n')
if (isInvalidUnterminatedTail) break
throw new Error(`${source}: malformed JSONL record at line ${index + 1}`, { cause })
for await (const chunk of stream) {
options.onBytes?.(chunk as Buffer)
const text = decoder.write(chunk as Buffer)
let start = 0
let end = text.indexOf('\n', start)
while (end !== -1) {
fragments.push(text.slice(start, end))
const record = parseRecord<T>(fragments.join(''), path, lineNumber++, true)
fragments.length = 0
if (record !== undefined) yield record.value
start = end + 1
end = text.indexOf('\n', start)
}
if (start < text.length) fragments.push(text.slice(start))
}
if (stream.bytesRead !== size) {
throw new Error(`${path}: journal changed while reading its committed prefix`)
}
fragments.push(decoder.end())
const record = parseRecord<T>(fragments.join(''), path, lineNumber, false)
if (record !== undefined) yield record.value
} finally {
stream.destroy()
}
} finally {
await handle.close()
}
}

return records
/** An invalid unterminated tail may be a torn write; other parse failures are not evidence of one. */
function parseRecord<T>(
line: string,
source: string,
lineNumber: number,
terminated: boolean,
): { value: T } | undefined {
if (line.length === 0) return undefined
try {
return { value: JSON.parse(line) as T }
} catch (cause) {
if (!(cause instanceof SyntaxError)) throw cause
if (!terminated) return undefined
throw new Error(`${source}: malformed JSONL record at line ${lineNumber}`, { cause })
}
}

/** FileHandle.write may legally make a short write. Loop until every byte is appended. */
Expand Down
58 changes: 23 additions & 35 deletions src/durable/observer-journal.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
import { createHash } from 'node:crypto'
import { readFile } from 'node:fs/promises'
import { resolve } from 'node:path'
import { detachedSnapshot } from '../runtime/supervise/snapshot'
import {
Expand All @@ -8,7 +7,7 @@ import {
type RuntimeHooks,
withPursuitContext,
} from '../runtime-hooks'
import { parseCommittedJsonLines, prepareJsonlAppend, writeAllBytes } from './jsonl-file'
import { prepareJsonlAppend, readCommittedJsonLines, writeAllBytes } from './jsonl-file'

export type ObserverRecordKind = 'event' | 'decision'

Expand Down Expand Up @@ -84,17 +83,13 @@ export class FileObserverJournal implements ObserverJournal {
async read(): Promise<readonly ObserverRecord[]> {
await this.tail
this.assertComplete()
let text: string
try {
text = await readFile(this.path, 'utf8')
} catch (error) {
if (isNoEnt(error)) return []
throw error
const records: ObserverRecord[] = []
const verify = observerVerifier(this.pursuitId)
for await (const record of this.readExistingUnsafe()) {
verify(record)
records.push(record)
}
return verifyObserverRecords(
parseCommittedJsonLines<ObserverRecord>(text, this.path),
this.pursuitId,
)
return Object.freeze(records)
}

private enqueue(
Expand Down Expand Up @@ -165,25 +160,21 @@ export class FileObserverJournal implements ObserverJournal {

private async initialize(): Promise<void> {
if (this.initialized) return
const records = await this.readExistingUnsafe()
const verified = verifyObserverRecords(records, this.pursuitId)
const tail = verified.at(-1)
const verify = observerVerifier(this.pursuitId)
let tail: ObserverRecord | undefined
for await (const record of this.readExistingUnsafe()) {
verify(record)
tail = record
}
this.sequence = tail?.sequence ?? 0
this.previousDigest = tail?.digest
// Only latch after recovery + verification succeed. A transient read error or
// corruption must never leave an instance pretending it initialized cleanly.
this.initialized = true
}

private async readExistingUnsafe(): Promise<ObserverRecord[]> {
let text: string
try {
text = await readFile(this.path, 'utf8')
} catch (error) {
if (isNoEnt(error)) return []
throw error
}
return parseCommittedJsonLines<ObserverRecord>(text, this.path)
private readExistingUnsafe(): AsyncGenerator<ObserverRecord> {
return readCommittedJsonLines<ObserverRecord>(this.path, { allowMissing: true })
}

private async writeRecord(record: ObserverRecord): Promise<void> {
Expand Down Expand Up @@ -214,9 +205,16 @@ export function verifyObserverRecords(
records: readonly ObserverRecord[],
pursuitId?: string,
): readonly ObserverRecord[] {
const verify = observerVerifier(pursuitId)
for (const record of records) verify(record)
return Object.freeze([...records])
}

/** The same full chain validation is used for replay, startup and in-memory evidence. */
function observerVerifier(pursuitId?: string): (record: ObserverRecord) => void {
let previousDigest: string | undefined
let expectedSequence = 1
for (const record of records) {
return (record) => {
if (record.schemaVersion !== 1) throw new Error('observer journal: unsupported schemaVersion')
if (pursuitId !== undefined && record.pursuitId !== pursuitId) {
throw new Error(`observer journal: pursuit identity mismatch at sequence ${record.sequence}`)
Expand Down Expand Up @@ -253,7 +251,6 @@ export function verifyObserverRecords(
previousDigest = digest
expectedSequence += 1
}
return Object.freeze([...records])
}

/** Compute the canonical SHA-256 digest for an unsigned observer record. */
Expand All @@ -273,15 +270,6 @@ export function createFileObserverHooks(
return { journal, hooks: journal.hooks() }
}

function isNoEnt(error: unknown): boolean {
return (
typeof error === 'object' &&
error !== null &&
'code' in error &&
(error as { code?: unknown }).code === 'ENOENT'
)
}

function toError(error: unknown): Error {
return error instanceof Error ? error : new Error(String(error))
}
24 changes: 8 additions & 16 deletions src/durable/spawn-journal.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,8 +53,8 @@ import { addSpend, cloneSpend, zeroSpend } from '../runtime/util'
import { contentAddress } from './content-address'
import {
isNoEntError,
parseCommittedJsonLines,
prepareJsonlAppend,
readCommittedJsonLines,
writeAllBytes,
} from './jsonl-file'

Expand Down Expand Up @@ -282,18 +282,12 @@ export class FileSpawnJournal implements SpawnJournal {
constructor(private readonly path: string) {}

async loadTree(root: NodeId): Promise<SpawnEvent[] | undefined> {
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
}
let begun = false
const events: SpawnEvent[] = []
const index = new SpawnEventIndex(root)
for (const record of parseCommittedJsonLines<SpawnJournalRecord>(text, this.path)) {
for await (const record of readCommittedJsonLines<SpawnJournalRecord>(this.path, {
allowMissing: true,
})) {
if (record.root !== root) continue
if (record.kind === 'begin') {
begun = true
Expand Down Expand Up @@ -345,17 +339,12 @@ export class FileSpawnJournal implements SpawnJournal {
}

private async validationIndex() {
const fs = await import('node:fs/promises')
const stamp = await this.fileStamp()
if (this.appendIndex && this.appendIndex.stamp === stamp) return this.appendIndex
this.appendIndex = undefined
const trees = new Map<NodeId, { begunAt: string; index: SpawnEventIndex }>()
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<SpawnJournalRecord>(text, this.path)) {
for await (const record of readCommittedJsonLines<SpawnJournalRecord>(this.path)) {
if (record.kind === 'begin') {
if (!trees.has(record.root)) {
trees.set(record.root, { begunAt: record.at, index: new SpawnEventIndex(record.root) })
Expand All @@ -372,6 +361,9 @@ export class FileSpawnJournal implements SpawnJournal {
}
}
}
if ((await this.fileStamp()) !== stamp) {
throw new Error('spawn journal changed while rebuilding its append index')
}
this.appendIndex = { stamp, trees }
return this.appendIndex
}
Expand Down
Loading
Loading