diff --git a/CHANGELOG.md b/CHANGELOG.md index 7f2e5b69..fca79b1e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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. diff --git a/docs/agent-managed-compute/reliability.md b/docs/agent-managed-compute/reliability.md index deb36839..2b488f8b 100644 --- a/docs/agent-managed-compute/reliability.md +++ b/docs/agent-managed-compute/reliability.md @@ -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. diff --git a/src/conversation/journal.ts b/src/conversation/journal.ts index 2fc71111..d55a4047 100644 --- a/src/conversation/journal.ts +++ b/src/conversation/journal.ts @@ -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 { @@ -133,16 +128,10 @@ export class FileConversationJournal implements ConversationJournal { constructor(private readonly path: string) {} async loadRun(runId: string): Promise { - 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(text, this.path)) { + for await (const record of readCommittedJsonLines(this.path, { + allowMissing: true, + })) { if (record.runId !== runId) continue if (record.kind === 'begin') { entry = { runId, startedAt: record.startedAt, turns: [] } diff --git a/src/durable/jsonl-file.ts b/src/durable/jsonl-file.ts index f90ba114..28814cb9 100644 --- a/src/durable/jsonl-file.ts +++ b/src/durable/jsonl-file.ts @@ -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(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(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( + path: string, + options: { allowMissing?: boolean; onBytes?: (bytes: Uint8Array) => void } = {}, +): AsyncGenerator { + 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(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(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( + 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. */ diff --git a/src/durable/observer-journal.ts b/src/durable/observer-journal.ts index 85d8ee67..53793b3b 100644 --- a/src/durable/observer-journal.ts +++ b/src/durable/observer-journal.ts @@ -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 { @@ -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' @@ -84,17 +83,13 @@ export class FileObserverJournal implements ObserverJournal { async read(): Promise { 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(text, this.path), - this.pursuitId, - ) + return Object.freeze(records) } private enqueue( @@ -165,9 +160,12 @@ export class FileObserverJournal implements ObserverJournal { private async initialize(): Promise { 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 @@ -175,15 +173,8 @@ export class FileObserverJournal implements ObserverJournal { this.initialized = true } - private async readExistingUnsafe(): Promise { - let text: string - try { - text = await readFile(this.path, 'utf8') - } catch (error) { - if (isNoEnt(error)) return [] - throw error - } - return parseCommittedJsonLines(text, this.path) + private readExistingUnsafe(): AsyncGenerator { + return readCommittedJsonLines(this.path, { allowMissing: true }) } private async writeRecord(record: ObserverRecord): Promise { @@ -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}`) @@ -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. */ @@ -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)) } diff --git a/src/durable/spawn-journal.ts b/src/durable/spawn-journal.ts index 6c664fff..ddb8091e 100644 --- a/src/durable/spawn-journal.ts +++ b/src/durable/spawn-journal.ts @@ -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' @@ -282,18 +282,12 @@ export class FileSpawnJournal implements SpawnJournal { constructor(private readonly path: string) {} async loadTree(root: NodeId): Promise { - 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(text, this.path)) { + for await (const record of readCommittedJsonLines(this.path, { + allowMissing: true, + })) { if (record.root !== root) continue if (record.kind === 'begin') { begun = true @@ -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() 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)) { + for await (const record of readCommittedJsonLines(this.path)) { if (record.kind === 'begin') { if (!trees.has(record.root)) { trees.set(record.root, { begunAt: record.at, index: new SpawnEventIndex(record.root) }) @@ -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 } diff --git a/src/durable/supervision-discovery.ts b/src/durable/supervision-discovery.ts index 270dc9ec..044f10a6 100644 --- a/src/durable/supervision-discovery.ts +++ b/src/durable/supervision-discovery.ts @@ -1,6 +1,6 @@ import { resolve } from 'node:path' import type { NodeId } from '../runtime/supervise/types' -import { isNoEntError, parseCommittedJsonLines } from './jsonl-file' +import { readCommittedJsonLines } from './jsonl-file' export interface DurableCoordinationStreamIdentity { readonly runId: string @@ -57,67 +57,57 @@ export async function discoverDurableSupervisionRun( const canonicalRunDir = resolve(runDir) const spawnJournalPath = `${canonicalRunDir}/spawn-journal.jsonl` const coordinationLogPath = `${canonicalRunDir}/coordination-log.jsonl` - const [spawnText, coordinationText] = await Promise.all([ - readOptionalText(spawnJournalPath), - readOptionalText(coordinationLogPath), - ]) - const allRoots = new Set() const nestedRoots = new Set() const rootsBegunAt = new Map() - if (spawnText !== undefined) { - for (const record of parseCommittedJsonLines( - spawnText, - spawnJournalPath, - )) { - if (record.kind === 'begin') { - if (typeof record.root !== 'string' || record.root.length === 0) { - throw new Error(`${spawnJournalPath}: begin record has no non-empty string root identity`) - } - allRoots.add(record.root as NodeId) - rootsBegunAt.set(record.root as NodeId, typeof record.at === 'string' ? record.at : null) - continue - } - if (record.kind !== 'event') continue - const event = record.event - if (!isRecord(event) || event.kind !== 'spawned') continue - if (!Object.hasOwn(event, 'ownedTreeRoot')) continue - if (typeof event.ownedTreeRoot !== 'string' || event.ownedTreeRoot.length === 0) { - throw new Error( - `${spawnJournalPath}: spawned event ownedTreeRoot must be a non-empty string when present`, - ) + for await (const record of readCommittedJsonLines(spawnJournalPath, { + allowMissing: true, + })) { + if (record.kind === 'begin') { + if (typeof record.root !== 'string' || record.root.length === 0) { + throw new Error(`${spawnJournalPath}: begin record has no non-empty string root identity`) } - nestedRoots.add(event.ownedTreeRoot as NodeId) + allRoots.add(record.root as NodeId) + rootsBegunAt.set(record.root as NodeId, typeof record.at === 'string' ? record.at : null) + continue + } + if (record.kind !== 'event') continue + const event = record.event + if (!isRecord(event) || event.kind !== 'spawned') continue + if (!Object.hasOwn(event, 'ownedTreeRoot')) continue + if (typeof event.ownedTreeRoot !== 'string' || event.ownedTreeRoot.length === 0) { + throw new Error( + `${spawnJournalPath}: spawned event ownedTreeRoot must be a non-empty string when present`, + ) } + nestedRoots.add(event.ownedTreeRoot as NodeId) } const streams = new Map< string, { owners: Set; unscopedRecords: number; recordCount: number } >() - if (coordinationText !== undefined) { - for (const record of parseCommittedJsonLines( - coordinationText, - coordinationLogPath, - )) { - if (typeof record.runId !== 'string' || record.runId.length === 0) { - throw new Error(`${coordinationLogPath}: record has no non-empty string runId identity`) - } - const stream = streams.get(record.runId) ?? { - owners: new Set(), - unscopedRecords: 0, - recordCount: 0, - } - stream.recordCount += 1 - if (record.ownerId === undefined) { - stream.unscopedRecords += 1 - } else if (typeof record.ownerId === 'string' && record.ownerId.length > 0) { - stream.owners.add(record.ownerId) - } else { - throw new Error(`${coordinationLogPath}: ownerId must be a non-empty string when present`) - } - streams.set(record.runId, stream) + for await (const record of readCommittedJsonLines( + coordinationLogPath, + { allowMissing: true }, + )) { + if (typeof record.runId !== 'string' || record.runId.length === 0) { + throw new Error(`${coordinationLogPath}: record has no non-empty string runId identity`) + } + const stream = streams.get(record.runId) ?? { + owners: new Set(), + unscopedRecords: 0, + recordCount: 0, + } + stream.recordCount += 1 + if (record.ownerId === undefined) { + stream.unscopedRecords += 1 + } else if (typeof record.ownerId === 'string' && record.ownerId.length > 0) { + stream.owners.add(record.ownerId) + } else { + throw new Error(`${coordinationLogPath}: ownerId must be a non-empty string when present`) } + streams.set(record.runId, stream) } const coordinationStreams = [...streams.entries()] @@ -142,16 +132,6 @@ export async function discoverDurableSupervisionRun( }) } -async function readOptionalText(path: string): Promise { - const fs = await import('node:fs/promises') - try { - return await fs.readFile(path, 'utf8') - } catch (error) { - if (isNoEntError(error)) return undefined - throw error - } -} - function isRecord(value: unknown): value is Record { return typeof value === 'object' && value !== null && !Array.isArray(value) } diff --git a/src/durable/tests/streaming-journal.test.ts b/src/durable/tests/streaming-journal.test.ts new file mode 100644 index 00000000..0631b202 --- /dev/null +++ b/src/durable/tests/streaming-journal.test.ts @@ -0,0 +1,248 @@ +import { createHash } from 'node:crypto' +import * as fs from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { FileConversationJournal } from '../../conversation/journal' +import { FileCoordinationLog } from '../../runtime/supervise/coordination-log' +import { + createRootStreamSink, + readRootStream, + readRootStreamReceipt, +} from '../../runtime/supervise/root-stream' +import { parseCommittedJsonLines, readCommittedJsonLines } from '../jsonl-file' +import { FileObserverJournal } from '../observer-journal' +import { FileSpawnJournal } from '../spawn-journal' +import { discoverDurableSupervisionRun } from '../supervision-discovery' + +vi.mock('node:fs/promises', async (load) => { + const actual = await load() + return { ...actual, readFile: vi.fn(actual.readFile), open: vi.fn(actual.open) } +}) + +let root: string +beforeEach(async () => { + const actual = await vi.importActual('node:fs/promises') + vi.mocked(fs.open).mockImplementation(actual.open) + vi.mocked(fs.readFile).mockImplementation(actual.readFile) + vi.clearAllMocks() + root = await fs.mkdtemp(join(tmpdir(), 'streaming-journal-')) +}) +afterEach(async () => { + vi.restoreAllMocks() + await fs.rm(root, { recursive: true, force: true }) +}) +async function collect(records: AsyncIterable): Promise { + const out: T[] = [] + for await (const record of records) out.push(record) + return out +} + +const event = (id: string) => ({ + id, + pursuitId: 'pursuit', + runId: 'run', + target: 'agent.child' as const, + phase: 'event' as const, + timestamp: 1, +}) + +describe('streaming committed JSONL', () => { + it.each([ + '', + '\n', + 'null\nfalse\n0\n"text"\n', + '{"a":1}\n\n{"b":2}', + '{"a":1}\n{"torn":', + '{"first":', + '{"a":1}\r\n{"b":2}\r\n', + ])('matches in-memory parsing for %j', async (bytes) => { + const path = join(root, 'records.jsonl') + await fs.writeFile(path, bytes) + expect(await collect(readCommittedJsonLines(path))).toEqual( + parseCommittedJsonLines(bytes, path), + ) + }) + + it('preserves UTF-8 and records across multiple stream chunks', async () => { + const path = join(root, 'records.jsonl') + const records = [{ text: '🙂é'.repeat(50_001) }, { text: 'after boundary' }] + await fs.writeFile(path, records.map((record) => JSON.stringify(record)).join('\n')) + expect(await collect(readCommittedJsonLines(path))).toEqual(records) + }) + + it.each(['{}\n\nBAD\n', '{}\nBAD\n{"torn":', '{}\r{}\n', ' \n'])( + 'refuses committed corruption in %j', + async (bytes) => { + const path = join(root, 'records.jsonl') + await fs.writeFile(path, bytes) + const expected = (() => { + try { + parseCommittedJsonLines(bytes, path) + } catch (error) { + return (error as Error).message + } + })() + await expect(collect(readCommittedJsonLines(path))).rejects.toThrow(expected) + expect(await fs.readFile(path, 'utf8')).toBe(bytes) + }, + ) + + it('only treats a missing file as empty when explicitly requested', async () => { + const path = join(root, 'absent.jsonl') + await expect(collect(readCommittedJsonLines(path))).rejects.toMatchObject({ code: 'ENOENT' }) + expect(await collect(readCommittedJsonLines(path, { allowMissing: true }))).toEqual([]) + await expect(collect(readCommittedJsonLines(root, { allowMissing: true }))).rejects.toThrow() + }) + + it('never follows appends beyond the prefix opened by this read', async () => { + const path = join(root, 'records.jsonl') + await fs.writeFile(path, '1\n2\n') + const reader = readCommittedJsonLines(path) + expect(await reader.next()).toMatchObject({ value: 1, done: false }) + await fs.appendFile(path, '3\n') + expect(await collect(reader)).toEqual([2]) + expect(await collect(readCommittedJsonLines(path))).toEqual([1, 2, 3]) + }) + + it('does not mistake a concurrently truncated snapshot for a successful read', async () => { + const path = join(root, 'records.jsonl') + await fs.writeFile(path, `1\n${JSON.stringify('x'.repeat(1_000_000))}\n`) + const reader = readCommittedJsonLines(path) + expect(await reader.next()).toMatchObject({ value: 1 }) + await fs.truncate(path, 2) + await expect(collect(reader)).rejects.toThrow(/journal changed/) + }) + + it('closes its file descriptor when a consumer stops early', async () => { + const path = join(root, 'records.jsonl') + await fs.writeFile(path, '1\n2\n') + const { open } = await vi.importActual('node:fs/promises') + const handle = await open(path, 'r') + const close = vi.spyOn(handle, 'close') + vi.mocked(fs.open).mockResolvedValueOnce(handle) + const reader = readCommittedJsonLines(path) + await reader.next() + await reader.return(undefined) + expect(close).toHaveBeenCalled() + await expect(handle.stat()).rejects.toMatchObject({ code: 'EBADF' }) + }) + + it('does not hide an allocation failure as a torn final write', async () => { + const path = join(root, 'records.jsonl') + const bytes = '{"keep":true}' + await fs.writeFile(path, bytes) + const parse = JSON.parse + const failure = new RangeError('fixture allocation failure') + vi.spyOn(JSON, 'parse').mockImplementation((text, reviver) => { + if (text === bytes) throw failure + return parse(text, reviver) + }) + expect(() => parseCommittedJsonLines(bytes, path)).toThrow(failure) + await expect(collect(readCommittedJsonLines(path))).rejects.toBe(failure) + expect(await fs.readFile(path, 'utf8')).toBe(bytes) + }) +}) + +describe('durable consumers share streaming reads', () => { + it('resumes and verifies the observer without a whole-file read', async () => { + const path = join(root, 'observer.jsonl') + await new FileObserverJournal(path, 'pursuit').appendEvent(event('one')) + vi.mocked(fs.readFile).mockImplementation(async () => { + throw new Error('whole-file read') + }) + const resumed = new FileObserverJournal(path, 'pursuit') + const second = await resumed.appendEvent(event('two')) + const records = await resumed.read() + expect(records.map((record) => record.sequence)).toEqual([1, 2]) + expect(second.previousDigest).toBe(records[0]?.digest) + }) + + it('checks the entire observer chain before allowing a resumed append', async () => { + const path = join(root, 'observer.jsonl') + const original = new FileObserverJournal(path, 'pursuit') + await original.appendEvent(event('one')) + await original.appendEvent(event('two')) + const bytes = (await fs.readFile(path, 'utf8')).replace('"id":"one"', '"id":"tampered"') + await fs.writeFile(path, bytes) + const resumed = new FileObserverJournal(path, 'pursuit') + await expect(resumed.appendEvent(event('three'))).rejects.toThrow(/digest mismatch/) + expect(await fs.readFile(path, 'utf8')).toBe(bytes) + }) + + it('loads trees, rebuilds append indexes and discovers identities without whole-file reads', async () => { + const path = join(root, 'spawn-journal.jsonl') + await new FileSpawnJournal(path).beginTree('root', '2026-01-01T00:00:00Z') + vi.mocked(fs.readFile).mockImplementation(async () => { + throw new Error('whole-file read') + }) + const resumed = new FileSpawnJournal(path) + expect(await resumed.loadTree('root')).toEqual([]) + expect(await resumed.loadTree('missing')).toBeUndefined() + await resumed.beginTree('other', '2026-01-01T00:00:01Z') + expect((await discoverDurableSupervisionRun(root)).roots).toEqual(['other', 'root']) + }) + + it('keeps coordination and conversation filtering on the streamed path', async () => { + const path = join(root, 'coordination-log.jsonl') + await fs.writeFile( + path, + [ + { runId: 'run', ownerId: 'A' }, + { runId: 'run', ownerId: 'B' }, + ] + .map((identity) => + JSON.stringify({ + ...identity, + seq: 1, + at: 1, + priority: 0, + event: { type: 'finding', finding: { id: identity.ownerId } }, + }), + ) + .join('\n'), + ) + const conversation = new FileConversationJournal(join(root, 'conversation.jsonl')) + await conversation.beginRun('conversation', '2026-01-01T00:00:00Z') + vi.mocked(fs.readFile).mockImplementation(async () => { + throw new Error('whole-file read') + }) + expect((await new FileCoordinationLog(path).load('run', 'A')).findings).toEqual([{ id: 'A' }]) + expect(await conversation.loadRun('conversation')).toMatchObject({ + runId: 'conversation', + turns: [], + }) + expect(await conversation.loadRun('missing')).toBeUndefined() + }) + + it('hashes exact root-stream bytes, including an uncommitted non-UTF8 tail', async () => { + const path = join(root, 'root-stream.jsonl') + const bytes = Buffer.concat([Buffer.from('{"seq":1}\n{"torn":"'), Buffer.from([0xff])]) + await fs.writeFile(path, bytes) + vi.mocked(fs.readFile).mockImplementation(async () => { + throw new Error('whole-file read') + }) + expect(await readRootStreamReceipt(root)).toEqual({ + ref: `sha256:${createHash('sha256').update(bytes).digest('hex')}`, + events: 1, + }) + expect(await readRootStream(root)).toEqual([{ seq: 1 }]) + }) + + it('resumes the root stream and preserves its sequence and raw-byte receipt', async () => { + const sink = createRootStreamSink(root, () => 1) + await sink.beginAttempt() + sink.append({ kind: 'text_delta', text: 'first' }) + await sink.close() + vi.mocked(fs.readFile).mockImplementation(async () => { + throw new Error('whole-file read') + }) + const resumed = createRootStreamSink(root, () => 2) + await resumed.beginAttempt() + resumed.append({ kind: 'text_delta', text: 'second' }) + expect((await resumed.close())?.events).toBe(2) + expect((await readRootStream(root))?.map((record) => record.seq)).toEqual([1, 2]) + expect(await readRootStream(join(root, 'missing'))).toBeUndefined() + expect(await readRootStreamReceipt(join(root, 'missing'))).toBeUndefined() + }) +}) diff --git a/src/runtime/personify/corpus.ts b/src/runtime/personify/corpus.ts index daacdea9..9f7a1a37 100644 --- a/src/runtime/personify/corpus.ts +++ b/src/runtime/personify/corpus.ts @@ -18,12 +18,7 @@ */ import type { AgentProfile } from '@tangle-network/agent-interface' -import { - isNoEntError, - parseCommittedJsonLines, - prepareJsonlAppend, - writeAllBytes, -} from '../../durable/jsonl-file' +import { prepareJsonlAppend, readCommittedJsonLines, writeAllBytes } from '../../durable/jsonl-file' import { detachedFrozen } from '../supervise/snapshot' import type { Corpus, @@ -266,19 +261,10 @@ export class FileCorpus implements Corpus { } private async load(path = this.path): Promise> { - const fs = await import('node:fs/promises') - let text: string - try { - text = await fs.readFile(path, 'utf8') - } catch (err) { - if (isNoEntError(err)) return new Map() - throw err - } const byId = new Map() - const records = parseCommittedJsonLines(text, this.path) - for (let i = 0; i < records.length; i++) { - const parsed = records[i] - assertCorpusRecord(parsed, `corpus ${this.path} line ${i + 1}`) + let recordNumber = 0 + for await (const parsed of readCommittedJsonLines(path, { allowMissing: true })) { + assertCorpusRecord(parsed, `corpus ${this.path} line ${++recordNumber}`) const existing = byId.get(parsed.id) if (existing && !recordsEqual(existing, parsed)) { throw new Error( diff --git a/src/runtime/supervise/coordination-log.ts b/src/runtime/supervise/coordination-log.ts index c432fdf1..d932fc72 100644 --- a/src/runtime/supervise/coordination-log.ts +++ b/src/runtime/supervise/coordination-log.ts @@ -16,12 +16,7 @@ * @experimental */ -import { - isNoEntError, - parseCommittedJsonLines, - prepareJsonlAppend, - writeAllBytes, -} from '../../durable/jsonl-file' +import { prepareJsonlAppend, readCommittedJsonLines, writeAllBytes } from '../../durable/jsonl-file' import type { AnalystFindingEvent, ContinuationInstruction, @@ -148,14 +143,6 @@ export class FileCoordinationLog implements CoordinationLog { } async load(runId: string, ownerId?: CoordinationOwnerId): Promise { - const fs = await import('node:fs/promises') - let text: string - try { - text = await fs.readFile(this.path, 'utf8') - } catch (err) { - if (isNoEntError(err)) return emptyPriorCoordination(ownerId) - throw err - } const byId = new Map() const findings: AnalystFindingEvent[] = [] const escalations: QuestionEscalationRecord[] = [] @@ -165,9 +152,9 @@ export class FileCoordinationLog implements CoordinationLog { const mail: PeerMailEvent[] = [] const records: BusRecord[] = [] let legacySeq = 0 - for (const stored of parseCommittedJsonLines< + for await (const stored of readCommittedJsonLines< CoordinationLogRecord | LegacyCoordinationLogRecord - >(text, this.path)) { + >(this.path, { allowMissing: true })) { if (stored.runId !== runId) continue // Omitting ownerId preserves the historical all-run read for direct consumers. Runtime always // supplies one, so root and nested supervisors never receive one another's evidence. @@ -230,17 +217,3 @@ export class FileCoordinationLog implements CoordinationLog { } } } - -function emptyPriorCoordination(ownerId?: CoordinationOwnerId): PriorCoordination { - return { - ...(ownerId !== undefined ? { ownerId } : {}), - questions: [], - findings: [], - escalations: [], - analystDefinitions: [], - continuations: [], - deliveryEvidence: [], - mail: [], - records: [], - } -} diff --git a/src/runtime/supervise/root-stream.ts b/src/runtime/supervise/root-stream.ts index 63702db1..76632fb9 100644 --- a/src/runtime/supervise/root-stream.ts +++ b/src/runtime/supervise/root-stream.ts @@ -34,9 +34,9 @@ import { createHash } from 'node:crypto' import { closeSync, constants as fsConstants, fsyncSync, openSync, writeSync } from 'node:fs' -import { mkdir, readFile } from 'node:fs/promises' +import { mkdir } from 'node:fs/promises' import { resolve } from 'node:path' -import { isNoEntError, parseCommittedJsonLines, prepareJsonlAppend } from '../../durable/jsonl-file' +import { isNoEntError, prepareJsonlAppend, readCommittedJsonLines } from '../../durable/jsonl-file' import { assertNoSymlinkDescendant } from './durable-file' import type { ExecutorProgressEvent, RootStreamReceipt } from './types' @@ -90,7 +90,12 @@ export function createRootStreamSink(runDir: string, now: () => number): RootStr // A resumed run continues the same file. Recovery scans the tail once, so a torn line from a // process that died mid-write is truncated rather than followed by a record it would corrupt. const needsSeparator = await prepareJsonlAppend(path) - seq = ((await readRootStream(dir)) ?? []).length + seq = 0 + for await (const _record of readCommittedJsonLines(path, { + allowMissing: true, + })) { + seq += 1 + } fd = openSync( path, fsConstants.O_APPEND | @@ -177,31 +182,35 @@ export async function readRootStreamReceipt( runDir: string, ): Promise { const path = resolve(runDir, ROOT_STREAM_FILE) - let bytes: Buffer + const hash = createHash('sha256') + let events = 0 try { - bytes = await readFile(path) + for await (const _record of readCommittedJsonLines(path, { + onBytes: (bytes) => { + hash.update(bytes) + }, + })) { + events += 1 + } } catch (error) { if (isNoEntError(error)) return undefined throw error } - return Object.freeze({ - ref: `sha256:${createHash('sha256').update(bytes).digest('hex')}`, - events: parseCommittedJsonLines(bytes.toString('utf8'), path).length, - }) + return Object.freeze({ ref: `sha256:${hash.digest('hex')}`, events }) } /** Every committed line of the root stream, in order, or `undefined` when there is no file. A * torn final line from a process that died mid-write is not a record and is left out. */ export async function readRootStream(runDir: string): Promise { const path = resolve(runDir, ROOT_STREAM_FILE) - let text: string + const records: RootStreamRecord[] = [] try { - text = await readFile(path, 'utf8') + for await (const record of readCommittedJsonLines(path)) records.push(record) } catch (error) { if (isNoEntError(error)) return undefined throw error } - return parseCommittedJsonLines(text, path) + return records } function errorMessage(error: unknown): string {