From 2d70e9c64fa6bf266013324c99f2dd741aea45a1 Mon Sep 17 00:00:00 2001 From: LunarFeller Date: Thu, 10 Sep 2026 14:23:17 +0800 Subject: [PATCH 1/3] fix(agent-core-v2): recover session persistence after disk-full instead of losing buffered records An ENOSPC append failure put AppendLogStore into a sticky failure state that never retried, so all subsequently appended records lived only in memory and were lost on process exit. The store now retries in the background with backoff after reconciling the on-disk tail against the failed batch (never duplicating or rewriting user data), and close/retirement make a final recovery attempt. SessionEventJournal no longer drops buffered events (or the pending header) on a failed write and retries with backoff. --- .changeset/rich-pandas-heal.md | 5 + .../backends/node-fs/appendLogStore.ts | 182 ++++++++- .../persistence/interface/appendLogStore.ts | 1 + .../agent-core-v2/src/wire/wireService.ts | 21 +- packages/agent-core-v2/test/harness/agent.ts | 1 + .../backends/node-fs/appendLogStore.test.ts | 359 +++++++++++++++++- .../agentLifecycle/agentLifecycle.test.ts | 1 + .../test/session/subagent/forkParity.test.ts | 1 + packages/agent-core-v2/test/wire/stubs.ts | 2 + .../transport/ws/v1/sessionEventJournal.ts | 53 ++- .../test/sessionEventJournal.test.ts | 77 +++- 11 files changed, 686 insertions(+), 17 deletions(-) create mode 100644 .changeset/rich-pandas-heal.md diff --git a/.changeset/rich-pandas-heal.md b/.changeset/rich-pandas-heal.md new file mode 100644 index 00000000000..69c94bb24f9 --- /dev/null +++ b/.changeset/rich-pandas-heal.md @@ -0,0 +1,5 @@ +--- +"@moonshot-ai/kimi-code": patch +--- + +Fix loss of session history when the disk runs out of space mid-session. diff --git a/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts b/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts index e38c7aeec58..88078bddd4f 100644 --- a/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts +++ b/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts @@ -3,7 +3,11 @@ import { LifecycleScope } from '#/app/scopes'; import { ScopeActivation, registerScopedService } from '#/_base/di/scope'; import { Emitter, type Event } from '#/_base/event'; -import { IFileSystemStorageService } from '#/persistence/interface/storage'; +import { + IFileSystemStorageService, + StorageError, + StorageErrors, +} from '#/persistence/interface/storage'; import { AppendLogCorruptedError, IAppendLogStore, @@ -14,6 +18,24 @@ import { const textEncoder = new TextEncoder(); +const RECOVERY_BASE_DELAY_MS = 1_000; +const RECOVERY_MAX_DELAY_MS = 30_000; +const NEWLINE = 0x0a; + +interface RecoveryState { + readonly failedBatch: readonly unknown[]; + attempts: number; + timer: ReturnType | undefined; +} + +function isRecoverableStorageError(error: unknown): boolean { + return ( + error instanceof StorageError && + (error.code === StorageErrors.codes.STORAGE_DISK_FULL || + (StorageErrors.retryable as readonly string[]).includes(error.code)) + ); +} + const pendingRetirements = new Set>(); export async function drainAppendLogRetirements(): Promise { @@ -27,6 +49,7 @@ interface LogState { flushPromise: Promise | undefined; flushScheduled: boolean; storageFailure: { readonly error: unknown } | undefined; + recovery: RecoveryState | undefined; cutoverEpoch: number; refCount: number; retired: boolean; @@ -41,9 +64,18 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { private readonly logs = new Map(); private readonly writeEmitter = this._register(new Emitter()); readonly onDidWrite: Event = this.writeEmitter.event; + private readonly recoverEmitter = this._register(new Emitter()); + readonly onDidRecover: Event = this.recoverEmitter.event; constructor(@IFileSystemStorageService private readonly storage: IFileSystemStorageService) { super(); + this._register( + toDisposable(() => { + for (const state of this.logs.values()) { + if (state.recovery?.timer !== undefined) clearTimeout(state.recovery.timer); + } + }), + ); } append(scope: string, key: string, record: R, options?: AppendLogOptions): void { @@ -114,6 +146,7 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { const encoded = encodeBatch(records); const state = this.state(scope, key); state.cutoverEpoch++; + this.cancelRecovery(state); const prior = state.flushPromise ?? state.ready; const priorSettled = prior.then( () => undefined, @@ -126,6 +159,10 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { return true; } catch (error) { state.storageFailure = { error }; + if (state.recovery === undefined && isRecoverableStorageError(error)) { + state.recovery = { failedBatch: state.pending.slice(), attempts: 0, timer: undefined }; + this.scheduleRecovery(scope, key, state); + } throw error; } }); @@ -144,6 +181,15 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { } async close(): Promise { + const recoveries: Promise[] = []; + for (const [id, state] of this.logs) { + if (state.recovery === undefined) continue; + if (state.recovery.timer !== undefined) clearTimeout(state.recovery.timer); + state.recovery.timer = undefined; + const { scope, key } = fromLogId(id); + recoveries.push(this.attemptRecovery(scope, key, state)); + } + await Promise.all(recoveries); await this.flush(); } @@ -169,6 +215,7 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { flushPromise: undefined, flushScheduled: false, storageFailure: undefined, + recovery: undefined, cutoverEpoch: 0, refCount: 0, retired: false, @@ -182,6 +229,7 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { private scheduleFlush(scope: string, key: string, state: LogState): void { if (state.flushScheduled || state.flushPromise !== undefined) return; + if (state.storageFailure !== undefined) return; state.flushScheduled = true; queueMicrotask(() => { state.flushScheduled = false; @@ -213,6 +261,11 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { private async settleRetiredState(scope: string, key: string, state: LogState): Promise { try { + if (state.recovery !== undefined) { + if (state.recovery.timer !== undefined) clearTimeout(state.recovery.timer); + state.recovery.timer = undefined; + await this.attemptRecovery(scope, key, state); + } await this.flushState(scope, key, state); } finally { const id = logId(scope, key); @@ -220,6 +273,117 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { } } + private cancelRecovery(state: LogState): void { + if (state.recovery?.timer !== undefined) clearTimeout(state.recovery.timer); + state.recovery = undefined; + } + + private scheduleRecovery(scope: string, key: string, state: LogState): void { + const recovery = state.recovery; + if (recovery === undefined || state.retired || recovery.timer !== undefined) return; + const delay = Math.min(RECOVERY_BASE_DELAY_MS * 2 ** recovery.attempts, RECOVERY_MAX_DELAY_MS); + recovery.timer = setTimeout(() => { + recovery.timer = undefined; + void this.attemptRecovery(scope, key, state); + }, delay); + recovery.timer.unref(); + } + + private async attemptRecovery(scope: string, key: string, state: LogState): Promise { + const recovery = state.recovery; + if (recovery === undefined) return; + state.recovery = undefined; + try { + if (state.flushPromise !== undefined) { + await state.flushPromise.catch(() => undefined); + } + if (state.storageFailure === undefined) return; + if (state.recovery !== undefined) return; + if (state.flushPromise !== undefined) { + state.recovery = recovery; + this.scheduleRecovery(scope, key, state); + return; + } + await this.ownFlush( + scope, + key, + state, + this.reconcileCommittedTail(scope, key, state, recovery.failedBatch), + { value: false }, + ); + state.storageFailure = undefined; + this.recoverEmitter.fire({ scope, key }); + } catch (error) { + if (state.recovery !== undefined) { + state.recovery.attempts = recovery.attempts + 1; + if (state.recovery.timer !== undefined) clearTimeout(state.recovery.timer); + state.recovery.timer = undefined; + this.scheduleRecovery(scope, key, state); + return; + } + if (isRecoverableStorageError(error)) { + state.recovery = recovery; + recovery.attempts++; + this.scheduleRecovery(scope, key, state); + } + } + } + + private async reconcileCommittedTail( + scope: string, + key: string, + state: LogState, + failedBatch: readonly unknown[], + ): Promise { + const encoded = encodeBatch(failedBatch); + if (encoded.byteLength === 0) return false; + const epoch = state.cutoverEpoch; + const size = await this.storage.size(scope, key); + if (size === undefined || size === 0) return false; + const window = Math.min(size, encoded.byteLength); + const tail = new Uint8Array(window); + let offset = 0; + for await (const chunk of this.storage.readStream(scope, key, { + start: size - window, + end: size - 1, + })) { + tail.set(chunk, offset); + offset += chunk.byteLength; + } + if (state.cutoverEpoch !== epoch) { + throw new StorageError( + StorageErrors.codes.STORAGE_CORRUPTED, + 'append-log recovery aborted: log rewritten concurrently', + { details: { scope, key } }, + ); + } + const overlap = committedTailOverlap(tail, encoded); + if (overlap === encoded.byteLength) { + state.pending.splice(0, failedBatch.length); + return false; + } + if (overlap > 0) { + const full = await this.storage.read(scope, key); + if (full === undefined || full.byteLength !== size) { + throw new StorageError( + StorageErrors.codes.STORAGE_CORRUPTED, + 'append-log recovery aborted: log changed while reconciling', + { details: { scope, key } }, + ); + } + await this.storage.write(scope, key, full.subarray(0, size - overlap), { atomic: true }); + return false; + } + if (tail[tail.byteLength - 1] !== NEWLINE) { + throw new StorageError( + StorageErrors.codes.STORAGE_CORRUPTED, + 'append-log tail does not match the failed batch; automatic recovery aborted', + { details: { scope, key } }, + ); + } + return false; + } + private ownFlush( scope: string, key: string, @@ -280,6 +444,10 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { if (wroteBox !== undefined) wroteBox.value = true; } catch (error) { const failure = (state.storageFailure ??= { error }); + if (state.recovery === undefined && isRecoverableStorageError(error)) { + state.recovery = { failedBatch: batch, attempts: 0, timer: undefined }; + this.scheduleRecovery(scope, key, state); + } throw failure.error; } if (state.cutoverEpoch !== cutoverEpoch) return wrote; @@ -304,6 +472,18 @@ function encodeBatch(records: readonly unknown[]): Uint8Array { return textEncoder.encode(content); } +function committedTailOverlap(existing: Uint8Array, batch: Uint8Array): number { + const max = Math.min(existing.byteLength, batch.byteLength); + for (let k = max; k > 0; k--) { + const before = existing.byteLength - k; + if (before !== 0 && existing[before - 1] !== NEWLINE) continue; + let i = 0; + while (i < k && existing[before + i] === batch[i]) i++; + if (i === k) return k; + } + return 0; +} + registerScopedService( LifecycleScope.App, IAppendLogStore, diff --git a/packages/agent-core-v2/src/persistence/interface/appendLogStore.ts b/packages/agent-core-v2/src/persistence/interface/appendLogStore.ts index 973c0afd6e4..ee189d15f78 100644 --- a/packages/agent-core-v2/src/persistence/interface/appendLogStore.ts +++ b/packages/agent-core-v2/src/persistence/interface/appendLogStore.ts @@ -41,6 +41,7 @@ export interface IAppendLogStore { readonly _serviceBrand: undefined; readonly onDidWrite: Event; + readonly onDidRecover: Event; append(scope: string, key: string, record: R, options?: AppendLogOptions): void; read(scope: string, key: string, options?: AppendLogReadOptions): AsyncIterable; diff --git a/packages/agent-core-v2/src/wire/wireService.ts b/packages/agent-core-v2/src/wire/wireService.ts index 92588253eb3..c0ff38ad0de 100644 --- a/packages/agent-core-v2/src/wire/wireService.ts +++ b/packages/agent-core-v2/src/wire/wireService.ts @@ -11,7 +11,7 @@ import { type AppendLogTruncation, IAppendLogStore, } from '#/persistence/interface/appendLogStore'; -import { IFileSystemStorageService, StorageError, StorageErrors } from '#/persistence/interface/storage'; +import { IFileSystemStorageService, isStorageError, StorageError, StorageErrors } from '#/persistence/interface/storage'; import { IWireService } from './wire'; import { WireError, WireErrors } from './errors'; @@ -59,6 +59,15 @@ export class WireService extends Service implements IWireService { this.wireScope = scopeContext.scope(); this.agentId = scopeContext.agentId; this._register(this.log.acquire(this.wireScope, AGENT_WIRE_RECORD_KEY)); + this._register( + this.log.onDidRecover((write) => { + if (write.scope === this.wireScope && write.key === AGENT_WIRE_RECORD_KEY) { + this.logger.info('session wire persistence recovered after storage failure', { + scope: this.wireScope, + }); + } + }), + ); } async seal(): Promise { @@ -343,7 +352,15 @@ export class WireService extends Service implements IWireService { private appendRecordLow(record: WireRecord): void { this.log.append(this.wireScope, AGENT_WIRE_RECORD_KEY, record, { - onError: onUnexpectedError, + onError: (error) => { + onUnexpectedError(error); + if (isStorageError(error, StorageErrors.codes.STORAGE_DISK_FULL)) { + this.logger.error( + 'session wire persistence degraded: no space left on device; new records stay buffered in memory and the store retries in the background — free disk space to resume durable writes', + { scope: this.wireScope }, + ); + } + }, }); this.lines += 1; if (record.type === 'context.clear') this.lastClearLine = this.lines; diff --git a/packages/agent-core-v2/test/harness/agent.ts b/packages/agent-core-v2/test/harness/agent.ts index 5ef5a12bb15..ba54db0df8a 100644 --- a/packages/agent-core-v2/test/harness/agent.ts +++ b/packages/agent-core-v2/test/harness/agent.ts @@ -963,6 +963,7 @@ function reassertServiceOverrides( class PersistenceAppendLogStore implements IAppendLogStore { declare readonly _serviceBrand: undefined; readonly onDidWrite: IAppendLogStore['onDidWrite'] = Event.None as IAppendLogStore['onDidWrite']; + readonly onDidRecover: IAppendLogStore['onDidRecover'] = Event.None as IAppendLogStore['onDidRecover']; private readonly history: WireRecord[] = []; private historySeeded = false; diff --git a/packages/agent-core-v2/test/persistence/backends/node-fs/appendLogStore.test.ts b/packages/agent-core-v2/test/persistence/backends/node-fs/appendLogStore.test.ts index 422345570c7..ba12f7c0773 100644 --- a/packages/agent-core-v2/test/persistence/backends/node-fs/appendLogStore.test.ts +++ b/packages/agent-core-v2/test/persistence/backends/node-fs/appendLogStore.test.ts @@ -1,10 +1,10 @@ -import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; import { SyncDescriptor } from '#/_base/di/descriptors'; import { DisposableStore } from '#/_base/di/lifecycle'; import { TestInstantiationService } from '#/_base/di/test'; import { AppendLogCorruptedError, IAppendLogStore, type AppendLogTruncation } from '#/persistence/interface/appendLogStore'; -import { IFileSystemStorageService } from '#/persistence/interface/storage'; +import { IFileSystemStorageService, StorageError, StorageErrors } from '#/persistence/interface/storage'; import { AppendLogStore } from '#/persistence/backends/node-fs/appendLogStore'; import { InMemoryStorageService } from '#/persistence/backends/memory/inMemoryStorageService'; @@ -791,4 +791,359 @@ describe('AppendLogStore', () => { await expect(record.flush()).rejects.toThrow('disk full'); expect(events).toEqual([`${SCOPE}/${KEY}`]); }); + + describe('disk-full recovery', () => { + const diskFull = (): StorageError => + new StorageError( + StorageErrors.codes.STORAGE_DISK_FULL, + 'storage append failed: no space left on device', + { details: { errno: 'ENOSPC' } }, + ); + + const settle = async (): Promise => { + for (let i = 0; i < 10; i++) await Promise.resolve(); + }; + + beforeEach(() => { + vi.useFakeTimers(); + }); + + afterEach(() => { + vi.useRealTimers(); + }); + + function watchRecovery(): Promise { + return new Promise((resolve) => { + const subscription = record.onDidRecover((write) => { + if (write.scope === SCOPE && write.key === KEY) { + subscription.dispose(); + resolve(); + } + }); + }); + } + + it('retries a disk-full append after a backoff and persists buffered records exactly once', async () => { + const failure = diskFull(); + let attempts = 0; + const originalAppend = storage.append.bind(storage); + storage.append = async (...args) => { + attempts++; + if (attempts === 1) throw failure; + return originalAppend(...args); + }; + let failures = 0; + let reportFailure!: () => void; + const reportedFailure = new Promise((resolve) => { + reportFailure = resolve; + }); + const recovered = watchRecovery(); + + record.append(SCOPE, KEY, { n: 1 }, { + onError: () => { + failures++; + reportFailure(); + }, + }); + await reportedFailure; + + await expect(record.flush()).rejects.toBe(failure); + expect(attempts).toBe(1); + + record.append(SCOPE, KEY, { n: 2 }); + await settle(); + expect(failures).toBe(1); + + await vi.advanceTimersByTimeAsync(1000); + await recovered; + await record.flush(); + + expect(await collect(SCOPE, KEY)).toEqual([{ n: 1 }, { n: 2 }]); + expect(attempts).toBe(2); + expect(failures).toBe(1); + expect(new TextDecoder().decode(await storage.read(SCOPE, KEY))).toBe('{"n":1}\n{"n":2}\n'); + }); + + it('repairs a torn tail from a partially committed batch before retrying', async () => { + const failure = diskFull(); + let attempts = 0; + const originalAppend = storage.append.bind(storage); + storage.append = async (scope, key, data, options) => { + attempts++; + if (attempts === 1) { + await originalAppend(scope, key, data.subarray(0, 9), options); + throw failure; + } + return originalAppend(scope, key, data, options); + }; + let reportFailure!: () => void; + const reportedFailure = new Promise((resolve) => { + reportFailure = resolve; + }); + const recovered = watchRecovery(); + + record.append(SCOPE, KEY, { n: 1 }, { onError: reportFailure }); + record.append(SCOPE, KEY, { n: 2 }); + await reportedFailure; + + expect(new TextDecoder().decode(await storage.read(SCOPE, KEY))).toBe('{"n":1}\n{'); + + await vi.advanceTimersByTimeAsync(1000); + await recovered; + await record.flush(); + + expect(await collect(SCOPE, KEY)).toEqual([{ n: 1 }, { n: 2 }]); + expect(new TextDecoder().decode(await storage.read(SCOPE, KEY))).toBe('{"n":1}\n{"n":2}\n'); + expect(attempts).toBe(2); + }); + + it('drops an ambiguously committed batch instead of duplicating it', async () => { + const failure = diskFull(); + let attempts = 0; + const originalAppend = storage.append.bind(storage); + storage.append = async (scope, key, data, options) => { + attempts++; + if (attempts === 1) { + await originalAppend(scope, key, data, options); + throw failure; + } + return originalAppend(scope, key, data, options); + }; + let reportFailure!: () => void; + const reportedFailure = new Promise((resolve) => { + reportFailure = resolve; + }); + const recovered = watchRecovery(); + + record.append(SCOPE, KEY, { n: 1 }, { onError: reportFailure }); + record.append(SCOPE, KEY, { n: 2 }); + await reportedFailure; + + expect(new TextDecoder().decode(await storage.read(SCOPE, KEY))).toBe('{"n":1}\n{"n":2}\n'); + + await vi.advanceTimersByTimeAsync(1000); + await recovered; + await record.flush(); + + expect(await collect(SCOPE, KEY)).toEqual([{ n: 1 }, { n: 2 }]); + expect(attempts).toBe(1); + }); + + it('stays sticky when the torn tail does not match the failed batch', async () => { + const failure = diskFull(); + await storage.append(SCOPE, KEY, enc.encode('{"n":0}\n{"n":99')); + let attempts = 0; + const originalAppend = storage.append.bind(storage); + storage.append = async (...args) => { + attempts++; + if (attempts === 1) throw failure; + return originalAppend(...args); + }; + let reportFailure!: () => void; + const reportedFailure = new Promise((resolve) => { + reportFailure = resolve; + }); + + record.append(SCOPE, KEY, { n: 1 }, { onError: reportFailure }); + await reportedFailure; + + await vi.advanceTimersByTimeAsync(5000); + await settle(); + expect(attempts).toBe(1); + await expect(record.flush()).rejects.toBe(failure); + await expect(collect(SCOPE, KEY)).rejects.toBe(failure); + + await record.rewrite(SCOPE, KEY, [{ n: 9 }]); + expect(await collect(SCOPE, KEY)).toEqual([{ n: 9 }, { n: 1 }]); + }); + + it('re-enters backoff with a growing delay when the recovery attempt itself hits disk-full', async () => { + const failure = diskFull(); + let attempts = 0; + const originalAppend = storage.append.bind(storage); + storage.append = async (...args) => { + attempts++; + if (attempts <= 2) throw failure; + return originalAppend(...args); + }; + let reportFailure!: () => void; + const reportedFailure = new Promise((resolve) => { + reportFailure = resolve; + }); + const recovered = watchRecovery(); + + record.append(SCOPE, KEY, { n: 1 }, { onError: reportFailure }); + await reportedFailure; + + await vi.advanceTimersByTimeAsync(1000); + await settle(); + expect(attempts).toBe(2); + + await vi.advanceTimersByTimeAsync(1000); + await settle(); + expect(attempts).toBe(2); + + await vi.advanceTimersByTimeAsync(1000); + await recovered; + await record.flush(); + + expect(await collect(SCOPE, KEY)).toEqual([{ n: 1 }]); + expect(attempts).toBe(3); + }); + + it('close() makes a final recovery attempt and persists buffered records', async () => { + const failure = diskFull(); + let attempts = 0; + const originalAppend = storage.append.bind(storage); + storage.append = async (...args) => { + attempts++; + if (attempts === 1) throw failure; + return originalAppend(...args); + }; + let reportFailure!: () => void; + const reportedFailure = new Promise((resolve) => { + reportFailure = resolve; + }); + + record.append(SCOPE, KEY, { n: 1 }, { onError: reportFailure }); + record.append(SCOPE, KEY, { n: 2 }); + await reportedFailure; + + await record.close(); + + expect(await collect(SCOPE, KEY)).toEqual([{ n: 1 }, { n: 2 }]); + expect(attempts).toBe(2); + }); + + it('does not enter recovery for non-retryable storage errors', async () => { + const failure = new StorageError( + StorageErrors.codes.STORAGE_PERMISSION_DENIED, + 'storage append failed: permission denied', + ); + let attempts = 0; + const originalAppend = storage.append.bind(storage); + storage.append = async (...args) => { + attempts++; + if (attempts === 1) throw failure; + return originalAppend(...args); + }; + let reportFailure!: () => void; + const reportedFailure = new Promise((resolve) => { + reportFailure = resolve; + }); + + record.append(SCOPE, KEY, { n: 1 }, { onError: reportFailure }); + await reportedFailure; + + await vi.advanceTimersByTimeAsync(5000); + await settle(); + expect(attempts).toBe(1); + await expect(record.flush()).rejects.toBe(failure); + }); + + it('also recovers transient io failures, not only disk-full', async () => { + const failure = new StorageError( + StorageErrors.codes.STORAGE_IO_FAILED, + 'storage append failed: unrecognized I/O error', + ); + let attempts = 0; + const originalAppend = storage.append.bind(storage); + storage.append = async (...args) => { + attempts++; + if (attempts === 1) throw failure; + return originalAppend(...args); + }; + let reportFailure!: () => void; + const reportedFailure = new Promise((resolve) => { + reportFailure = resolve; + }); + const recovered = watchRecovery(); + + record.append(SCOPE, KEY, { n: 1 }, { onError: reportFailure }); + await reportedFailure; + + await vi.advanceTimersByTimeAsync(1000); + await recovered; + await record.flush(); + + expect(await collect(SCOPE, KEY)).toEqual([{ n: 1 }]); + expect(attempts).toBe(2); + }); + + it('retirement performs a final recovery attempt before dropping the log', async () => { + const failure = diskFull(); + let attempts = 0; + const originalAppend = storage.append.bind(storage); + storage.append = async (...args) => { + attempts++; + if (attempts === 1) throw failure; + return originalAppend(...args); + }; + let reportFailure!: () => void; + const reportedFailure = new Promise((resolve) => { + reportFailure = resolve; + }); + + const owner = record.acquire(SCOPE, KEY); + record.append(SCOPE, KEY, { n: 1 }, { onError: reportFailure }); + await reportedFailure; + + owner.dispose(); + await record.drainRetirements(); + + expect(await collect(SCOPE, KEY)).toEqual([{ n: 1 }]); + expect(attempts).toBe(2); + }); + + it('rewrite cancels a pending recovery and keeps the buffered records', async () => { + const failure = diskFull(); + let attempts = 0; + const originalAppend = storage.append.bind(storage); + storage.append = async (...args) => { + attempts++; + if (attempts === 1) throw failure; + return originalAppend(...args); + }; + let reportFailure!: () => void; + const reportedFailure = new Promise((resolve) => { + reportFailure = resolve; + }); + + record.append(SCOPE, KEY, { n: 1 }, { onError: reportFailure }); + await reportedFailure; + expect(attempts).toBe(1); + + await record.rewrite(SCOPE, KEY, [{ n: 9 }]); + expect(await collect(SCOPE, KEY)).toEqual([{ n: 9 }, { n: 1 }]); + + await vi.advanceTimersByTimeAsync(60_000); + await settle(); + expect(attempts).toBe(2); + }); + + it('recovers a failed rewrite once the disk has space again', async () => { + const failure = diskFull(); + record.append(SCOPE, KEY, { n: 1 }); + await record.flush(); + + let writeAttempts = 0; + const originalWrite = storage.write.bind(storage); + storage.write = async (...args) => { + writeAttempts++; + if (writeAttempts === 1) throw failure; + return originalWrite(...args); + }; + const recovered = watchRecovery(); + + await expect(record.rewrite(SCOPE, KEY, [{ n: 9 }])).rejects.toBe(failure); + await expect(record.flush()).rejects.toBe(failure); + + await vi.advanceTimersByTimeAsync(1000); + await recovered; + + record.append(SCOPE, KEY, { n: 2 }); + await record.flush(); + expect(await collect(SCOPE, KEY)).toEqual([{ n: 1 }, { n: 2 }]); + }); + }); }); diff --git a/packages/agent-core-v2/test/session/agentLifecycle/agentLifecycle.test.ts b/packages/agent-core-v2/test/session/agentLifecycle/agentLifecycle.test.ts index 42a8678c499..ad2a30353ff 100644 --- a/packages/agent-core-v2/test/session/agentLifecycle/agentLifecycle.test.ts +++ b/packages/agent-core-v2/test/session/agentLifecycle/agentLifecycle.test.ts @@ -158,6 +158,7 @@ function recordingAppendLog(initial: readonly WireRecord[] = []): { const store: IAppendLogStore = { _serviceBrand: undefined, onDidWrite: Event.None as IAppendLogStore['onDidWrite'], + onDidRecover: Event.None as IAppendLogStore['onDidRecover'], append: (_scope: string, _key: string, record: R) => { const persisted = record as unknown as WireRecord; records.push(persisted); diff --git a/packages/agent-core-v2/test/session/subagent/forkParity.test.ts b/packages/agent-core-v2/test/session/subagent/forkParity.test.ts index 9858f64402c..8db37ce37f3 100644 --- a/packages/agent-core-v2/test/session/subagent/forkParity.test.ts +++ b/packages/agent-core-v2/test/session/subagent/forkParity.test.ts @@ -43,6 +43,7 @@ class ScopedAppendLogStore implements IAppendLogStore { declare readonly _serviceBrand: undefined; private readonly logs = new Map(); readonly onDidWrite: IAppendLogStore['onDidWrite'] = Event.None as IAppendLogStore['onDidWrite']; + readonly onDidRecover: IAppendLogStore['onDidRecover'] = Event.None as IAppendLogStore['onDidRecover']; recordsFor(scope: string, key: string): WireRecord[] { return structuredClone(this.logs.get(`${scope}/${key}`) ?? []); diff --git a/packages/agent-core-v2/test/wire/stubs.ts b/packages/agent-core-v2/test/wire/stubs.ts index 886e2b806b8..1fb2da12d19 100644 --- a/packages/agent-core-v2/test/wire/stubs.ts +++ b/packages/agent-core-v2/test/wire/stubs.ts @@ -35,6 +35,7 @@ interface TestAgentWireDependencies { const noopLog: IAppendLogStore = { _serviceBrand: undefined, onDidWrite: Event.None as IAppendLogStore['onDidWrite'], + onDidRecover: Event.None as IAppendLogStore['onDidRecover'], append: () => {}, read: async function* () {}, rewrite: async () => {}, @@ -184,6 +185,7 @@ export function recordingWireLog( return { _serviceBrand: undefined, onDidWrite: Event.None as IAppendLogStore['onDidWrite'], + onDidRecover: Event.None as IAppendLogStore['onDidRecover'], append: (_scope, _key, record) => { records.push(record as WireRecord); onAppend?.(record as WireRecord); diff --git a/packages/kap-server/src/transport/ws/v1/sessionEventJournal.ts b/packages/kap-server/src/transport/ws/v1/sessionEventJournal.ts index 83d41e297aa..503ae159f05 100644 --- a/packages/kap-server/src/transport/ws/v1/sessionEventJournal.ts +++ b/packages/kap-server/src/transport/ws/v1/sessionEventJournal.ts @@ -4,6 +4,8 @@ import { dirname, join } from 'node:path'; import { ulid } from 'ulid'; const JOURNAL_VERSION = 1; +const RETRY_BASE_DELAY_MS = 500; +const RETRY_MAX_DELAY_MS = 30_000; export interface EventEnvelope { readonly type: string; @@ -45,6 +47,8 @@ export class SessionEventJournal { private _seq: number; private pendingLines: string[] = []; private flushPromise: Promise | undefined; + private retryTimer: ReturnType | undefined; + private consecutiveFailures = 0; private headerPending: boolean; private closed = false; @@ -129,13 +133,23 @@ export class SessionEventJournal { } async flush(): Promise { - while (this.flushPromise !== undefined || this.pendingLines.length > 0) { - if (this.flushPromise === undefined) { - this.flushPromise = this.flushOnce().finally(() => { - this.flushPromise = undefined; - }); + if (this.retryTimer !== undefined) { + clearTimeout(this.retryTimer); + this.retryTimer = undefined; + } + for (;;) { + let inFlight = this.flushPromise; + if (inFlight === undefined) { + if (this.pendingLines.length === 0) return; + this.scheduleFlush(); + inFlight = this.flushPromise; + if (inFlight === undefined) return; } - await this.flushPromise; + const failuresBefore = this.consecutiveFailures; + await inFlight; + if (this.flushPromise !== undefined) continue; + if (this.pendingLines.length === 0) return; + if (this.consecutiveFailures !== failuresBefore) return; } } @@ -145,16 +159,30 @@ export class SessionEventJournal { } private scheduleFlush(): void { - if (this.flushPromise !== undefined) return; + if (this.flushPromise !== undefined || this.retryTimer !== undefined) return; this.flushPromise = this.flushOnce().finally(() => { this.flushPromise = undefined; - if (this.pendingLines.length > 0) this.scheduleFlush(); + if (this.pendingLines.length === 0) return; + if (this.consecutiveFailures === 0) { + this.scheduleFlush(); + return; + } + const delay = Math.min( + RETRY_BASE_DELAY_MS * 2 ** (this.consecutiveFailures - 1), + RETRY_MAX_DELAY_MS, + ); + this.retryTimer = setTimeout(() => { + this.retryTimer = undefined; + this.scheduleFlush(); + }, delay); + this.retryTimer.unref(); }); } private async flushOnce(): Promise { const lines: string[] = []; - if (this.headerPending) { + const headerIncluded = this.headerPending; + if (headerIncluded) { const header: JournalHeaderLine = { kind: 'journal_header', version: JOURNAL_VERSION, @@ -162,7 +190,6 @@ export class SessionEventJournal { created_at: Date.now(), }; lines.push(JSON.stringify(header)); - this.headerPending = false; } lines.push(...this.pendingLines); this.pendingLines = []; @@ -170,10 +197,14 @@ export class SessionEventJournal { try { await mkdir(dirname(this.filePath), { recursive: true }); await appendFile(this.filePath, lines.join('\n') + '\n', 'utf8'); + if (headerIncluded) this.headerPending = false; + this.consecutiveFailures = 0; } catch (error) { + this.pendingLines = (headerIncluded ? lines.slice(1) : lines).concat(this.pendingLines); + this.consecutiveFailures++; this.logger.warn( { filePath: this.filePath, err: String(error) }, - 'event journal write failed; events remain live-only this round', + 'event journal write failed; events stay buffered and will be retried', ); } } diff --git a/packages/kap-server/test/sessionEventJournal.test.ts b/packages/kap-server/test/sessionEventJournal.test.ts index e5b2608940c..f348227e5b0 100644 --- a/packages/kap-server/test/sessionEventJournal.test.ts +++ b/packages/kap-server/test/sessionEventJournal.test.ts @@ -2,7 +2,25 @@ import { mkdtemp, readFile, rm, writeFile } from 'node:fs/promises'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; -import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; + +const fsMock = vi.hoisted(() => ({ appendFailures: 0 })); + +vi.mock('node:fs/promises', async (importActual) => { + const actual = await importActual(); + return { + ...actual, + appendFile: async (...args: Parameters) => { + if (fsMock.appendFailures > 0) { + fsMock.appendFailures--; + const error = new Error('ENOSPC: no space left on device, write') as NodeJS.ErrnoException; + error.code = 'ENOSPC'; + throw error; + } + return actual.appendFile(...args); + }, + }; +}); import { type EventEnvelope, @@ -111,4 +129,61 @@ describe('SessionEventJournal', () => { expect(lines).toBe(13); await j.close(); }); + + it('keeps buffered lines after a write failure and retries with backoff', async () => { + const j = await SessionEventJournal.open(filePath); + fsMock.appendFailures = 2; + try { + j.append(j.nextSeq(), envelope(1)); + j.append(j.nextSeq(), envelope(2)); + + await new Promise((resolve) => setTimeout(resolve, 50)); + expect(await readFile(filePath, 'utf8').catch(() => '')).toBe(''); + + const deadline = Date.now() + 5000; + let content = ''; + while (Date.now() < deadline) { + content = await readFile(filePath, 'utf8').catch(() => ''); + if (content.trim().split('\n').filter(Boolean).length >= 3) break; + await new Promise((resolve) => setTimeout(resolve, 20)); + } + const lines = content.trim().split('\n').filter(Boolean); + expect(lines).toHaveLength(3); + expect(lines[0]).toContain('"journal_header"'); + expect(lines[1]).toContain('"seq":1'); + expect(lines[2]).toContain('"seq":2'); + expect(fsMock.appendFailures).toBe(0); + } finally { + fsMock.appendFailures = 0; + await j.close(); + } + }); + + it('keeps the journal header pending when the first write fails, preserving the epoch', async () => { + const j = await SessionEventJournal.open(filePath); + const epoch = j.epoch; + fsMock.appendFailures = 1; + try { + j.append(j.nextSeq(), envelope(1)); + const deadline = Date.now() + 5000; + let content = ''; + while (Date.now() < deadline) { + content = await readFile(filePath, 'utf8').catch(() => ''); + if (content.includes('"journal_header"') && content.includes('"seq":1')) break; + await new Promise((resolve) => setTimeout(resolve, 20)); + } + const lines = content.trim().split('\n'); + expect(lines).toHaveLength(2); + expect(lines[0]).toContain('"journal_header"'); + expect(lines[1]).toContain('"seq":1'); + } finally { + fsMock.appendFailures = 0; + await j.close(); + } + + const reopened = await SessionEventJournal.open(filePath); + expect(reopened.epoch).toBe(epoch); + expect(reopened.seq).toBe(1); + await reopened.close(); + }); }); From 958dcbc00328ac8748203243614cf08fb7885ee5 Mon Sep 17 00:00:00 2001 From: LunarFeller Date: Thu, 10 Sep 2026 15:06:44 +0800 Subject: [PATCH 2/3] fix(kap-server): stop event journal retries after close A retry armed by a failed write could fire after the journal was closed and recreate a deleted session's journal file. Retry scheduling is now gated on the journal still being open. --- .../src/transport/ws/v1/sessionEventJournal.ts | 1 + packages/kap-server/test/sessionEventJournal.test.ts | 11 +++++++++++ 2 files changed, 12 insertions(+) diff --git a/packages/kap-server/src/transport/ws/v1/sessionEventJournal.ts b/packages/kap-server/src/transport/ws/v1/sessionEventJournal.ts index 503ae159f05..926c61c61e9 100644 --- a/packages/kap-server/src/transport/ws/v1/sessionEventJournal.ts +++ b/packages/kap-server/src/transport/ws/v1/sessionEventJournal.ts @@ -163,6 +163,7 @@ export class SessionEventJournal { this.flushPromise = this.flushOnce().finally(() => { this.flushPromise = undefined; if (this.pendingLines.length === 0) return; + if (this.closed) return; if (this.consecutiveFailures === 0) { this.scheduleFlush(); return; diff --git a/packages/kap-server/test/sessionEventJournal.test.ts b/packages/kap-server/test/sessionEventJournal.test.ts index f348227e5b0..8e74588cc92 100644 --- a/packages/kap-server/test/sessionEventJournal.test.ts +++ b/packages/kap-server/test/sessionEventJournal.test.ts @@ -186,4 +186,15 @@ describe('SessionEventJournal', () => { expect(reopened.seq).toBe(1); await reopened.close(); }); + + it('does not write buffered lines after close', async () => { + const j = await SessionEventJournal.open(filePath); + fsMock.appendFailures = 100; + j.append(j.nextSeq(), envelope(1)); + await new Promise((resolve) => setTimeout(resolve, 50)); + await j.close(); + fsMock.appendFailures = 0; + await new Promise((resolve) => setTimeout(resolve, 1200)); + expect(await readFile(filePath, 'utf8').catch(() => '')).toBe(''); + }); }); From dff7d142d667d0902c0cfd59fa9f5f84992d1a86 Mon Sep 17 00:00:00 2001 From: LunarFeller Date: Thu, 10 Sep 2026 15:06:52 +0800 Subject: [PATCH 3/3] fix(agent-core-v2): classify recovery failures by the latest storage error A recovery attempt that failed with a different, permanent error (e.g. permission denied) used to be classified by the original disk-full error and retried forever. The latest failure now replaces the stored one. --- .../backends/node-fs/appendLogStore.ts | 4 +-- .../backends/node-fs/appendLogStore.test.ts | 31 +++++++++++++++++++ 2 files changed, 33 insertions(+), 2 deletions(-) diff --git a/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts b/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts index 88078bddd4f..61147300653 100644 --- a/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts +++ b/packages/agent-core-v2/src/persistence/backends/node-fs/appendLogStore.ts @@ -443,12 +443,12 @@ export class AppendLogStore extends Disposable implements IAppendLogStore { wrote = true; if (wroteBox !== undefined) wroteBox.value = true; } catch (error) { - const failure = (state.storageFailure ??= { error }); + state.storageFailure = { error }; if (state.recovery === undefined && isRecoverableStorageError(error)) { state.recovery = { failedBatch: batch, attempts: 0, timer: undefined }; this.scheduleRecovery(scope, key, state); } - throw failure.error; + throw error; } if (state.cutoverEpoch !== cutoverEpoch) return wrote; state.pending.splice(0, batch.length); diff --git a/packages/agent-core-v2/test/persistence/backends/node-fs/appendLogStore.test.ts b/packages/agent-core-v2/test/persistence/backends/node-fs/appendLogStore.test.ts index ba12f7c0773..4a85026c5dd 100644 --- a/packages/agent-core-v2/test/persistence/backends/node-fs/appendLogStore.test.ts +++ b/packages/agent-core-v2/test/persistence/backends/node-fs/appendLogStore.test.ts @@ -1145,5 +1145,36 @@ describe('AppendLogStore', () => { await record.flush(); expect(await collect(SCOPE, KEY)).toEqual([{ n: 1 }, { n: 2 }]); }); + + it('becomes sticky when a recovery attempt hits a non-retryable error', async () => { + const permanent = new StorageError( + StorageErrors.codes.STORAGE_PERMISSION_DENIED, + 'storage append failed: permission denied', + ); + const failure = diskFull(); + let attempts = 0; + const originalAppend = storage.append.bind(storage); + storage.append = async (...args) => { + attempts++; + if (attempts === 1) throw failure; + if (attempts === 2) throw permanent; + return originalAppend(...args); + }; + let reportFailure!: () => void; + const reportedFailure = new Promise((resolve) => { + reportFailure = resolve; + }); + + record.append(SCOPE, KEY, { n: 1 }, { onError: reportFailure }); + await reportedFailure; + + await vi.advanceTimersByTimeAsync(1000); + await settle(); + await vi.advanceTimersByTimeAsync(60_000); + await settle(); + + expect(attempts).toBe(2); + await expect(record.flush()).rejects.toBe(permanent); + }); }); });