Download packages/agent-core-v2/test/wire/wireService.test.ts from SaylorTwift/kimi-code: direct link, hf CLI and curl.
- Browser
- Download file 41 kB
-
https://huggingface.co/SaylorTwift/kimi-code/resolve/main/packages/agent-core-v2/test/wire/wireService.test.ts
- Command line
-
hf download hf://SaylorTwift/kimi-code/packages/agent-core-v2/test/wire/wireService.test.ts
-
curl -L -o wireService.test.ts https://huggingface.co/SaylorTwift/kimi-code/resolve/main/packages/agent-core-v2/test/wire/wireService.test.ts
41 kB
| import { afterEach, beforeEach, describe, expect, it } from 'vitest'; | |
| import { DisposableStore } from '#/_base/di/lifecycle'; | |
| import { SyncDescriptor } from '#/_base/di/descriptors'; | |
| import { TestInstantiationService } from '#/_base/di/test'; | |
| import { resetUnexpectedErrorHandler, setUnexpectedErrorHandler } from '#/_base/errors/unexpectedError'; | |
| import { ILogService } from '#/_base/log/log'; | |
| import { IAgentBlobService } from '#/agent/blob/agentBlobService'; | |
| import { ITelemetryService, noopTelemetryService } from '#/app/telemetry/telemetry'; | |
| import type { ContentPart } from '#human/llm/message'; | |
| import { AppendLogStore } from '#/persistence/backends/node-fs/appendLogStore'; | |
| import { InMemoryStorageService } from '#/persistence/backends/memory/inMemoryStorageService'; | |
| import { IAppendLogStore } from '#/persistence/interface/appendLogStore'; | |
| import { IFileSystemStorageService, StorageError, StorageErrors } from '#/persistence/interface/storage'; | |
| import { WIRE_PROTOCOL_VERSION } from '#/wire/migration/migration'; | |
| import { wireJournalBackupKey } from '#/wire/repair'; | |
| import { WireError, WireErrors } from '#/wire/errors'; | |
| import { IWireService } from '#/wire/wire'; | |
| import { AGENT_WIRE_RECORD_KEY, type WireRecord } from '#/wire/record'; | |
| import { recordingWireLog, registerTestAgentWire, testWireScope, noopLogger } from './stubs'; | |
| const SCOPE = 'wire'; | |
| const KEY = 'journal-test'; | |
| let disposables: DisposableStore; | |
| let ix: TestInstantiationService; | |
| let wire: IWireService; | |
| let log: IAppendLogStore; | |
| let storage: InMemoryStorageService; | |
| beforeEach(() => { | |
| disposables = new DisposableStore(); | |
| ix = disposables.add(new TestInstantiationService()); | |
| storage = new InMemoryStorageService(); | |
| ix.stub(IFileSystemStorageService, storage); | |
| ix.set(IAppendLogStore, new SyncDescriptor(AppendLogStore)); | |
| log = ix.get(IAppendLogStore); | |
| wire = registerTestAgentWire(ix, testWireScope(SCOPE, KEY), { | |
| log, | |
| storage, | |
| logger: noopLogger, | |
| telemetry: noopTelemetryService, | |
| }); | |
| }); | |
| afterEach(() => disposables.dispose()); | |
| async function readRecords( | |
| target: IAppendLogStore = log, | |
| scope = SCOPE, | |
| key = KEY, | |
| ): Promise<WireRecord[]> { | |
| const out: WireRecord[] = []; | |
| for await (const record of target.read<WireRecord>(testWireScope(scope, key), AGENT_WIRE_RECORD_KEY)) { | |
| out.push(record); | |
| } | |
| return out; | |
| } | |
| async function collect(journal: AsyncIterable<WireRecord>): Promise<WireRecord[]> { | |
| const out: WireRecord[] = []; | |
| for await (const record of journal) { | |
| out.push(record); | |
| } | |
| return out; | |
| } | |
| function wireOverLog( | |
| stubLog: IAppendLogStore, | |
| key: string, | |
| dependencies: { blob?: IAgentBlobService; storage?: IFileSystemStorageService; telemetry?: ITelemetryService } = {}, | |
| ): IWireService { | |
| const stubIx = disposables.add(new TestInstantiationService()); | |
| return registerTestAgentWire(stubIx, testWireScope(SCOPE, key), { log: stubLog, ...dependencies }); | |
| } | |
| describe('WireService seal', () => { | |
| it('writes the metadata envelope once and ignores repeated calls', async () => { | |
| await wire.seal(); | |
| await wire.seal(); | |
| expect(await readRecords()).toEqual([ | |
| { | |
| type: 'metadata', | |
| protocol_version: WIRE_PROTOCOL_VERSION, | |
| created_at: expect.any(Number), | |
| }, | |
| ]); | |
| }); | |
| it('does not seal a journal that already has records', async () => { | |
| log.append(testWireScope(SCOPE, KEY), AGENT_WIRE_RECORD_KEY, { | |
| type: 'wire.test.existing', | |
| time: 1, | |
| }); | |
| await log.flush(); | |
| await wire.seal(); | |
| expect(await readRecords()).toEqual([{ type: 'wire.test.existing', time: 1 }]); | |
| }); | |
| }); | |
| describe('WireService appendRecord', () => { | |
| it('appends flat records without a dehydrator', async () => { | |
| wire.appendRecord({ type: 'wire.test.append', value: 1, time: 10 }); | |
| wire.appendRecord({ type: 'wire.test.append', value: 2, time: 11 }); | |
| expect(await readRecords()).toEqual([ | |
| { type: 'wire.test.append', value: 1, time: 10 }, | |
| { type: 'wire.test.append', value: 2, time: 11 }, | |
| ]); | |
| }); | |
| it('runs records through the dehydrate queue in append order', async () => { | |
| const order: string[] = []; | |
| wire.appendRecord({ type: 'wire.test.a', time: 1 }, async (record) => { | |
| order.push('a'); | |
| return { ...record, dehydrated: true }; | |
| }); | |
| wire.appendRecord({ type: 'wire.test.b', time: 2 }, async (record) => { | |
| order.push('b'); | |
| return record; | |
| }); | |
| await wire.flush(); | |
| expect(order).toEqual(['a', 'b']); | |
| expect(await readRecords()).toEqual([ | |
| { type: 'wire.test.a', time: 1, dehydrated: true }, | |
| { type: 'wire.test.b', time: 2 }, | |
| ]); | |
| }); | |
| it('queues a plain append behind a pending dehydrate', async () => { | |
| const records: WireRecord[] = []; | |
| const queued = wireOverLog(recordingWireLog(records), 'queued'); | |
| let release!: () => void; | |
| const gate = new Promise<void>((resolve) => { | |
| release = resolve; | |
| }); | |
| queued.appendRecord({ type: 'wire.test.gated', time: 1 }, async (record) => { | |
| await gate; | |
| return record; | |
| }); | |
| queued.appendRecord({ type: 'wire.test.plain', time: 2 }); | |
| await new Promise((resolve) => setImmediate(resolve)); | |
| expect(records).toEqual([]); | |
| release(); | |
| await queued.flush(); | |
| expect(records).toEqual([ | |
| { type: 'wire.test.gated', time: 1 }, | |
| { type: 'wire.test.plain', time: 2 }, | |
| ]); | |
| }); | |
| it('hands the dehydrator a blob offload transform backed by the blob service', async () => { | |
| const offloaded: unknown[][] = []; | |
| const blob: IAgentBlobService = { | |
| _serviceBrand: undefined, | |
| offloadParts: async (parts) => { | |
| offloaded.push([...parts]); | |
| return parts.map((part) => ({ type: 'blob_ref', part })) as unknown as ContentPart[]; | |
| }, | |
| loadParts: async (parts) => parts, | |
| isBlobRef: () => false, | |
| }; | |
| const records: WireRecord[] = []; | |
| const withBlob = wireOverLog(recordingWireLog(records), 'blob', { blob }); | |
| withBlob.appendRecord( | |
| { type: 'wire.test.blob', parts: [{ type: 'text', text: 'x' }], time: 1 }, | |
| async (record, transform) => ({ | |
| ...record, | |
| parts: await transform(record['parts'] as readonly unknown[]), | |
| }), | |
| ); | |
| await withBlob.flush(); | |
| expect(offloaded).toEqual([[{ type: 'text', text: 'x' }]]); | |
| expect(records).toEqual([ | |
| { | |
| type: 'wire.test.blob', | |
| parts: [{ type: 'blob_ref', part: { type: 'text', text: 'x' } }], | |
| time: 1, | |
| }, | |
| ]); | |
| }); | |
| it('reports a synchronous append failure through onUnexpectedError instead of throwing', () => { | |
| const expected = new Error('append exploded'); | |
| const failing = recordingWireLog([]); | |
| failing.append = () => { | |
| throw expected; | |
| }; | |
| const stub = wireOverLog(failing, 'failing'); | |
| const unexpected: unknown[] = []; | |
| setUnexpectedErrorHandler((error) => unexpected.push(error)); | |
| try { | |
| stub.appendRecord({ type: 'wire.test.fail', time: 1 }); | |
| expect(unexpected).toEqual([expected]); | |
| } finally { | |
| resetUnexpectedErrorHandler(); | |
| } | |
| }); | |
| it('reports a dehydrate failure and keeps the queue usable for later appends', async () => { | |
| const expected = new Error('dehydrate exploded'); | |
| const records: WireRecord[] = []; | |
| const stub = wireOverLog(recordingWireLog(records), 'dehydrate-fail'); | |
| const unexpected: unknown[] = []; | |
| setUnexpectedErrorHandler((error) => unexpected.push(error)); | |
| try { | |
| stub.appendRecord({ type: 'wire.test.bad', time: 1 }, async () => { | |
| throw expected; | |
| }); | |
| stub.appendRecord({ type: 'wire.test.good', time: 2 }); | |
| await stub.flush(); | |
| expect(unexpected).toEqual([expected]); | |
| expect(records).toEqual([{ type: 'wire.test.good', time: 2 }]); | |
| } finally { | |
| resetUnexpectedErrorHandler(); | |
| } | |
| }); | |
| }); | |
| describe('WireService readJournal', () => { | |
| it('normalizes legacy plan revision paths and rewrites them as keys', async () => { | |
| const telemetryRecords: { event: string; properties: unknown }[] = []; | |
| const telemetry = { | |
| ...noopTelemetryService, | |
| track2: (event: string, properties: unknown) => telemetryRecords.push({ event, properties }), | |
| } as unknown as ITelemetryService; | |
| const seeded: WireRecord[] = [ | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| { type: 'wire.test.before', value: 1, time: 2 }, | |
| { | |
| type: 'plan.revision', | |
| id: 'plan-1', | |
| version: 1, | |
| path: 'sessions/source/session-1/agents/test-agent/plan/plan-1/v1.md', | |
| sha256: 'sha', | |
| bytes: 3, | |
| time: 2, | |
| }, | |
| { type: 'wire.test.after', value: 2, time: 3 }, | |
| ]; | |
| const stubLog = recordingWireLog(seeded); | |
| const stub = wireOverLog(stubLog, 'legacy-plan', { telemetry }); | |
| expect(await collect(stub.readJournal())).toEqual([ | |
| seeded[0], | |
| { type: 'wire.test.before', value: 1, time: 2 }, | |
| { | |
| type: 'plan.revision', | |
| id: 'plan-1', | |
| version: 1, | |
| key: 'plan/plan-1/v1.md', | |
| sha256: 'sha', | |
| bytes: 3, | |
| time: 2, | |
| }, | |
| { type: 'wire.test.after', value: 2, time: 3 }, | |
| ]); | |
| expect(telemetryRecords).toEqual([ | |
| { | |
| event: 'wire_plan_revision_migrated', | |
| properties: { | |
| record_type: 'plan.revision', | |
| legacy_field: 'path', | |
| migration_outcome: 'migrated', | |
| }, | |
| }, | |
| ]); | |
| expect(await readRecords(stubLog, SCOPE, 'legacy-plan')).toEqual([ | |
| seeded[0], | |
| { type: 'wire.test.before', value: 1, time: 2 }, | |
| { | |
| type: 'plan.revision', | |
| id: 'plan-1', | |
| version: 1, | |
| key: 'plan/plan-1/v1.md', | |
| sha256: 'sha', | |
| bytes: 3, | |
| time: 2, | |
| }, | |
| { type: 'wire.test.after', value: 2, time: 3 }, | |
| ]); | |
| }); | |
| it('skips unsafe legacy plan revision paths and reports the migration outcome', async () => { | |
| const telemetryRecords: { event: string; properties: unknown }[] = []; | |
| const telemetry = { | |
| ...noopTelemetryService, | |
| track2: (event: string, properties: unknown) => telemetryRecords.push({ event, properties }), | |
| } as unknown as ITelemetryService; | |
| const stub = wireOverLog( | |
| recordingWireLog([ | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| { | |
| type: 'plan.revision', | |
| id: 'plan-1', | |
| version: 1, | |
| path: 'sessions/source/session-1/agents/other-agent/plan/plan-1/v1.md', | |
| sha256: 'sha', | |
| bytes: 3, | |
| }, | |
| ]), | |
| 'unsafe-legacy-plan', | |
| { telemetry }, | |
| ); | |
| expect(await collect(stub.readJournal())).toEqual([ | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| ]); | |
| expect(telemetryRecords).toContainEqual({ | |
| event: 'wire_plan_revision_migrated', | |
| properties: { | |
| record_type: 'plan.revision', | |
| legacy_field: 'path', | |
| migration_outcome: 'skipped', | |
| }, | |
| }); | |
| }); | |
| it('bootstraps the metadata envelope onto an empty journal', async () => { | |
| expect(await collect(wire.readJournal())).toEqual([]); | |
| expect(await readRecords()).toEqual([ | |
| { | |
| type: 'metadata', | |
| protocol_version: WIRE_PROTOCOL_VERSION, | |
| created_at: expect.any(Number), | |
| }, | |
| ]); | |
| }); | |
| it('heals an envelope-less legacy journal through the v1.4 to v1.5 migration', async () => { | |
| log.append(testWireScope(SCOPE, KEY), AGENT_WIRE_RECORD_KEY, { | |
| type: 'goal.create', | |
| goalId: 'g1', | |
| objective: 'legacy', | |
| time: 7, | |
| }); | |
| await log.flush(); | |
| const yielded = await collect(wire.readJournal()); | |
| expect(yielded).toEqual([ | |
| { | |
| type: 'goal.create', | |
| goalId: 'g1', | |
| objective: 'legacy', | |
| time: 7, | |
| wallClockResumedAt: 7, | |
| }, | |
| ]); | |
| expect(await readRecords()).toEqual([ | |
| { | |
| type: 'metadata', | |
| protocol_version: WIRE_PROTOCOL_VERSION, | |
| created_at: expect.any(Number), | |
| }, | |
| ...yielded, | |
| ]); | |
| }); | |
| it('migrates a v1.4 journal and rewrites it at the current protocol version', async () => { | |
| log.append(testWireScope(SCOPE, KEY), AGENT_WIRE_RECORD_KEY, { | |
| type: 'metadata', | |
| protocol_version: '1.4', | |
| created_at: 1, | |
| }); | |
| log.append(testWireScope(SCOPE, KEY), AGENT_WIRE_RECORD_KEY, { | |
| type: 'goal.create', | |
| goalId: 'g1', | |
| time: 9, | |
| }); | |
| await log.flush(); | |
| const yielded = await collect(wire.readJournal()); | |
| expect(yielded).toEqual([ | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| { type: 'goal.create', goalId: 'g1', time: 9, wallClockResumedAt: 9 }, | |
| ]); | |
| expect(await readRecords()).toEqual(yielded); | |
| }); | |
| it('reads a current-version journal without rewriting it', async () => { | |
| const seeded: WireRecord[] = [ | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| { type: 'wire.test.current', value: 1, time: 2 }, | |
| ]; | |
| let rewrites = 0; | |
| const counting = recordingWireLog(seeded); | |
| const rewrite = counting.rewrite.bind(counting); | |
| counting.rewrite = async (scope, key, next) => { | |
| rewrites += 1; | |
| return rewrite(scope, key, next); | |
| }; | |
| const stub = wireOverLog(counting, 'current'); | |
| expect(await collect(stub.readJournal())).toEqual(seeded); | |
| expect(rewrites).toBe(0); | |
| const journal = stub; | |
| expect(await collect(journal.readRaw())).toEqual(seeded); | |
| expect(await collect(journal.read())).toEqual(seeded); | |
| expect(journal.journalRef).toEqual({ tree: testWireScope(SCOPE, 'current'), branch: 'main' }); | |
| expect(journal.branches()).toEqual(['main']); | |
| expect(journal.nextSeq()).toBe(3); | |
| expect(rewrites).toBe(0); | |
| }); | |
| it('reads a newer-version journal without stamping or rewriting it', async () => { | |
| const seeded: WireRecord[] = [ | |
| { type: 'metadata', protocol_version: '9.9', created_at: 1 }, | |
| { type: 'wire.test.newer', value: 1, time: 2 }, | |
| ]; | |
| let rewrites = 0; | |
| const counting = recordingWireLog(seeded); | |
| const rewrite = counting.rewrite.bind(counting); | |
| counting.rewrite = async (scope, key, next) => { | |
| rewrites += 1; | |
| return rewrite(scope, key, next); | |
| }; | |
| const stub = wireOverLog(counting, 'newer'); | |
| expect(await collect(stub.readJournal())).toEqual(seeded); | |
| expect(rewrites).toBe(0); | |
| }); | |
| it('leaves legacy plan revision paths untouched in a newer-version journal', async () => { | |
| const telemetryRecords: { event: string; properties: unknown }[] = []; | |
| const telemetry = { | |
| ...noopTelemetryService, | |
| track2: (event: string, properties: unknown) => telemetryRecords.push({ event, properties }), | |
| } as unknown as ITelemetryService; | |
| const seeded: WireRecord[] = [ | |
| { type: 'metadata', protocol_version: '9.9', created_at: 1 }, | |
| { | |
| type: 'plan.revision', | |
| id: 'plan-1', | |
| version: 1, | |
| path: 'sessions/source/session-1/agents/test-agent/plan/plan-1/v1.md', | |
| sha256: 'sha', | |
| bytes: 3, | |
| time: 2, | |
| }, | |
| ]; | |
| let rewrites = 0; | |
| const counting = recordingWireLog(seeded); | |
| const rewrite = counting.rewrite.bind(counting); | |
| counting.rewrite = async (scope, key, next) => { | |
| rewrites += 1; | |
| return rewrite(scope, key, next); | |
| }; | |
| const stub = wireOverLog(counting, 'newer-legacy-plan', { telemetry }); | |
| expect(await collect(stub.readJournal())).toEqual(seeded); | |
| expect(rewrites).toBe(0); | |
| expect(telemetryRecords).toEqual([]); | |
| }); | |
| it('rejects a malformed metadata envelope as corrupted storage', async () => { | |
| const stub = wireOverLog( | |
| recordingWireLog([{ type: 'metadata' }]), | |
| 'malformed-metadata', | |
| ); | |
| const failure = await collect(stub.readJournal()).catch((error: unknown) => error); | |
| expect(failure).toBeInstanceOf(StorageError); | |
| expect(failure).toMatchObject({ | |
| code: StorageErrors.codes.STORAGE_CORRUPTED, | |
| }); | |
| }); | |
| it('skips malformed lines and reports them through onUnexpectedError', async () => { | |
| const seeded: WireRecord[] = [ | |
| 'garbage' as unknown as WireRecord, | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| 42 as unknown as WireRecord, | |
| { type: 'wire.test.ok', time: 3 }, | |
| ]; | |
| const stub = wireOverLog(recordingWireLog(seeded), 'malformed-lines'); | |
| const unexpected: unknown[] = []; | |
| setUnexpectedErrorHandler((error) => unexpected.push(error)); | |
| try { | |
| const yielded = await collect(stub.readJournal()); | |
| expect(yielded).toEqual([ | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| { type: 'wire.test.ok', time: 3 }, | |
| ]); | |
| expect(unexpected).toHaveLength(2); | |
| expect(unexpected[0]).toMatchObject({ | |
| code: 'wire.unknown_record', | |
| details: { type: undefined, index: 0 }, | |
| }); | |
| expect(unexpected[1]).toMatchObject({ | |
| code: 'wire.unknown_record', | |
| details: { type: undefined, index: 1 }, | |
| }); | |
| } finally { | |
| resetUnexpectedErrorHandler(); | |
| } | |
| }); | |
| it('throws when the journal version has no migration path', async () => { | |
| const stub = wireOverLog( | |
| recordingWireLog([{ type: 'metadata', protocol_version: '0.9', created_at: 1 }]), | |
| 'no-migration', | |
| ); | |
| await expect(collect(stub.readJournal())).rejects.toThrow( | |
| 'Missing wire migration for version 0.9', | |
| ); | |
| }); | |
| }); | |
| describe('WireService corruption repair', () => { | |
| const enc = new TextEncoder(); | |
| const dec = new TextDecoder(); | |
| const BACKUP_KEY = wireJournalBackupKey(AGENT_WIRE_RECORD_KEY); | |
| interface RepairCapture { | |
| readonly warnings: Array<{ message: string; payload?: unknown }>; | |
| readonly events: Array<{ name: string; payload: unknown }>; | |
| } | |
| function wireWithCapture(key: string, capture: RepairCapture): IWireService { | |
| const logger: ILogService = { | |
| _serviceBrand: undefined, | |
| level: 'off', | |
| error: () => {}, | |
| warn: (message, payload) => capture.warnings.push({ message, payload }), | |
| info: () => {}, | |
| debug: () => {}, | |
| child: () => logger, | |
| setLevel: () => {}, | |
| flush: async () => {}, | |
| }; | |
| const telemetry: ITelemetryService = { | |
| ...noopTelemetryService, | |
| track2: ((name: string, payload: unknown) => { | |
| capture.events.push({ name, payload }); | |
| }) as ITelemetryService['track2'], | |
| }; | |
| const localIx = disposables.add(new TestInstantiationService()); | |
| return registerTestAgentWire(localIx, testWireScope(SCOPE, key), { | |
| log, | |
| storage, | |
| logger, | |
| telemetry, | |
| }); | |
| } | |
| async function rawBytes(key = AGENT_WIRE_RECORD_KEY): Promise<string | undefined> { | |
| const bytes = await storage.read(testWireScope(SCOPE, KEY), key); | |
| return bytes === undefined ? undefined : dec.decode(bytes); | |
| } | |
| async function seedCorrupt(raw: string): Promise<void> { | |
| await storage.write(testWireScope(SCOPE, KEY), AGENT_WIRE_RECORD_KEY, enc.encode(raw)); | |
| } | |
| function currentMetadata(createdAt = 1): string { | |
| return JSON.stringify({ | |
| type: 'metadata', | |
| protocol_version: WIRE_PROTOCOL_VERSION, | |
| created_at: createdAt, | |
| }); | |
| } | |
| it('heals a torn final line, keeps the valid prefix, and backs up the original bytes', async () => { | |
| const valid = `${currentMetadata()}\n${JSON.stringify({ type: 'wire.test.ok', time: 3 })}\n`; | |
| const torn = `${JSON.stringify({ type: 'wire.test.torn', time: 4 }).slice(0, 12)}`; | |
| await seedCorrupt(valid + torn); | |
| const yielded = await collect(wire.readJournal()); | |
| expect(yielded).toEqual([ | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| { type: 'wire.test.ok', time: 3 }, | |
| ]); | |
| expect(await rawBytes()).toBe(valid); | |
| expect(await rawBytes(BACKUP_KEY)).toBe(valid + torn); | |
| expect(await collect(wire.readJournal())).toEqual(yielded); | |
| }); | |
| it('truncates at a corrupted middle line, dropping later valid lines', async () => { | |
| const edge = { | |
| type: 'agent.switched', | |
| agentId: 'test-agent', | |
| branch: 'b1', | |
| reason: 'undo', | |
| base: { branch: 'main', line: 2 }, | |
| turns: 1, | |
| legacyUndoLine: 4, | |
| time: 2, | |
| }; | |
| const prefix = `${currentMetadata()}\n${JSON.stringify({ type: 'wire.test.a', time: 1 })}\n${JSON.stringify(edge)}\n`; | |
| const raw = `${prefix}GARBAGE\n${JSON.stringify({ type: 'wire.test.b', time: 2 })}\n`; | |
| await seedCorrupt(raw); | |
| const yielded = await collect(wire.readJournal()); | |
| expect(yielded).toEqual([ | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| { type: 'wire.test.a', time: 1 }, | |
| edge, | |
| ]); | |
| expect(await rawBytes()).toBe(prefix); | |
| expect(await rawBytes(BACKUP_KEY)).toBe(raw); | |
| const journal = wire; | |
| expect(await collect(journal.read())).toEqual([ | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| { type: 'wire.test.a', time: 1 }, | |
| ]); | |
| expect(journal.journalRef.branch).toBe('b1'); | |
| }); | |
| it('keeps the first backup when the journal corrupts again after a repair', async () => { | |
| const first = `${currentMetadata()}\nGARBAGE-1\n`; | |
| await seedCorrupt(first); | |
| await collect(wire.readJournal()); | |
| const second = `${currentMetadata()}\nGARBAGE-2\n`; | |
| await seedCorrupt(second); | |
| await collect(wire.readJournal()); | |
| expect(await rawBytes(BACKUP_KEY)).toBe(first); | |
| }); | |
| it('repairs through the migration rewrite path when corruption meets an old version', async () => { | |
| const legacy = `${JSON.stringify({ type: 'metadata', protocol_version: '1.4', created_at: 1 })}\n`; | |
| const raw = `${legacy}${JSON.stringify({ type: 'wire.test.legacy', time: 9 })}\nGARBAGE\n`; | |
| await seedCorrupt(raw); | |
| const yielded = await collect(wire.readJournal()); | |
| expect(yielded).toEqual([ | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| { type: 'wire.test.legacy', time: 9 }, | |
| ]); | |
| expect(await rawBytes()).toBe( | |
| `${currentMetadata()}\n${JSON.stringify({ type: 'wire.test.legacy', time: 9 })}\n`, | |
| ); | |
| expect(await rawBytes(BACKUP_KEY)).toBe(raw); | |
| }); | |
| it('repairs a corrupted journal that also carries a migratable legacy plan revision', async () => { | |
| const legacyPlan = JSON.stringify({ | |
| type: 'plan.revision', | |
| id: 'plan-1', | |
| version: 1, | |
| path: 'sessions/source/session-1/agents/test-agent/plan/plan-1/v1.md', | |
| sha256: 'sha', | |
| bytes: 3, | |
| time: 2, | |
| }); | |
| const raw = `${currentMetadata()}\n${legacyPlan}\nGARBAGE\n${JSON.stringify({ type: 'wire.test.dropped', time: 3 })}\n`; | |
| await seedCorrupt(raw); | |
| const yielded = await collect(wire.readJournal()); | |
| expect(yielded).toEqual([ | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| { | |
| type: 'plan.revision', | |
| id: 'plan-1', | |
| version: 1, | |
| sha256: 'sha', | |
| bytes: 3, | |
| time: 2, | |
| key: 'plan/plan-1/v1.md', | |
| }, | |
| ]); | |
| expect(await rawBytes()).toBe( | |
| `${currentMetadata()}\n${JSON.stringify({ | |
| type: 'plan.revision', | |
| id: 'plan-1', | |
| version: 1, | |
| sha256: 'sha', | |
| bytes: 3, | |
| time: 2, | |
| key: 'plan/plan-1/v1.md', | |
| })}\n`, | |
| ); | |
| expect(await rawBytes(BACKUP_KEY)).toBe(raw); | |
| }); | |
| it('reports a corrupted middle line through a warn log and the wire_repair event', async () => { | |
| const capture: RepairCapture = { warnings: [], events: [] }; | |
| const svc = wireWithCapture(KEY, capture); | |
| const prefix = `${currentMetadata()}\n${JSON.stringify({ type: 'wire.test.a', time: 1 })}\n`; | |
| const raw = `${prefix}GARBAGE\n${JSON.stringify({ type: 'wire.test.b', time: 2 })}\n`; | |
| await seedCorrupt(raw); | |
| await collect(svc.readJournal()); | |
| expect(capture.warnings).toHaveLength(1); | |
| expect(capture.warnings[0]!.message).toBe( | |
| 'corrupted wire journal truncated to its valid prefix', | |
| ); | |
| expect(capture.warnings[0]!.payload).toMatchObject({ | |
| lineNumber: 3, | |
| reason: 'corrupted', | |
| outcome: 'repaired', | |
| droppedCount: 2, | |
| backupCreated: true, | |
| }); | |
| expect(capture.events).toEqual([ | |
| { | |
| name: 'wire_repair', | |
| payload: { | |
| kind: 'corrupted', | |
| outcome: 'repaired', | |
| dropped_count: 2, | |
| backup_created: true, | |
| }, | |
| }, | |
| ]); | |
| }); | |
| it('reports a torn tail as truncation through the wire_repair event', async () => { | |
| const capture: RepairCapture = { warnings: [], events: [] }; | |
| const svc = wireWithCapture(KEY, capture); | |
| const raw = `${currentMetadata()}\n${JSON.stringify({ type: 'wire.test.a' }).slice(0, 10)}`; | |
| await seedCorrupt(raw); | |
| await collect(svc.readJournal()); | |
| expect(capture.events).toEqual([ | |
| { | |
| name: 'wire_repair', | |
| payload: { | |
| kind: 'truncated', | |
| outcome: 'repaired', | |
| dropped_count: 1, | |
| backup_created: true, | |
| }, | |
| }, | |
| ]); | |
| }); | |
| it('keeps restoring from the valid prefix when the on-disk repair itself fails', async () => { | |
| const capture: RepairCapture = { warnings: [], events: [] }; | |
| const svc = wireWithCapture(KEY, capture); | |
| const prefix = `${currentMetadata()}\n`; | |
| const raw = `${prefix}GARBAGE\n`; | |
| await seedCorrupt(raw); | |
| const originalWrite = storage.write.bind(storage); | |
| storage.write = async (scope, key, data, options) => { | |
| if (key === AGENT_WIRE_RECORD_KEY) throw new Error('disk full'); | |
| return originalWrite(scope, key, data, options); | |
| }; | |
| const yielded = await collect(svc.readJournal()); | |
| expect(yielded).toEqual([ | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| ]); | |
| expect(capture.events).toEqual([ | |
| { | |
| name: 'wire_repair', | |
| payload: { | |
| kind: 'corrupted', | |
| outcome: 'failed', | |
| dropped_count: 1, | |
| backup_created: true, | |
| }, | |
| }, | |
| ]); | |
| }); | |
| it('retries a failed repair before the next append and heals the journal first', async () => { | |
| const capture: RepairCapture = { warnings: [], events: [] }; | |
| const svc = wireWithCapture(KEY, capture); | |
| const prefix = `${currentMetadata()}\n`; | |
| const raw = `${prefix}GARBAGE\n`; | |
| await seedCorrupt(raw); | |
| const originalWrite = storage.write.bind(storage); | |
| storage.write = async (scope, key, data, options) => { | |
| if (key === AGENT_WIRE_RECORD_KEY) throw new Error('disk full'); | |
| return originalWrite(scope, key, data, options); | |
| }; | |
| await collect(svc.readJournal()); | |
| storage.write = originalWrite; | |
| svc.appendRecord({ type: 'wire.test.new', time: 7 }); | |
| await svc.flush(); | |
| expect(await rawBytes()).toBe(`${prefix}${JSON.stringify({ type: 'wire.test.new', time: 7 })}\n`); | |
| expect(await rawBytes(BACKUP_KEY)).toBe(raw); | |
| expect(await collect(svc.readJournal())).toEqual([ | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| { type: 'wire.test.new', time: 7 }, | |
| ]); | |
| expect(capture.events).toEqual([ | |
| { | |
| name: 'wire_repair', | |
| payload: { | |
| kind: 'corrupted', | |
| outcome: 'failed', | |
| dropped_count: 1, | |
| backup_created: true, | |
| }, | |
| }, | |
| { | |
| name: 'wire_repair', | |
| payload: { | |
| kind: 'corrupted', | |
| outcome: 'repaired', | |
| dropped_count: 1, | |
| backup_created: false, | |
| }, | |
| }, | |
| ]); | |
| }); | |
| it('refuses to append behind the corrupted tail while a failed repair keeps failing', async () => { | |
| const capture: RepairCapture = { warnings: [], events: [] }; | |
| const svc = wireWithCapture(KEY, capture); | |
| const prefix = `${currentMetadata()}\n`; | |
| const raw = `${prefix}GARBAGE\n`; | |
| await seedCorrupt(raw); | |
| const originalWrite = storage.write.bind(storage); | |
| storage.write = async (scope, key, data, options) => { | |
| if (key === AGENT_WIRE_RECORD_KEY) throw new Error('disk full'); | |
| return originalWrite(scope, key, data, options); | |
| }; | |
| await collect(svc.readJournal()); | |
| const unexpected: unknown[] = []; | |
| setUnexpectedErrorHandler((error) => unexpected.push(error)); | |
| try { | |
| svc.appendRecord({ type: 'wire.test.doomed', time: 8 }); | |
| await expect(svc.flush()).rejects.toThrow('Wire journal repair did not complete'); | |
| expect(await rawBytes()).toBe(raw); | |
| expect(unexpected).toHaveLength(1); | |
| expect(unexpected[0]).toBeInstanceOf(WireError); | |
| expect((unexpected[0] as WireError).code).toBe(WireErrors.codes.RECORDS_WRITE_FAILED); | |
| expect(capture.events).toEqual([ | |
| { | |
| name: 'wire_repair', | |
| payload: { | |
| kind: 'corrupted', | |
| outcome: 'failed', | |
| dropped_count: 1, | |
| backup_created: true, | |
| }, | |
| }, | |
| { | |
| name: 'wire_repair', | |
| payload: { | |
| kind: 'corrupted', | |
| outcome: 'failed', | |
| dropped_count: 1, | |
| backup_created: false, | |
| }, | |
| }, | |
| ]); | |
| } finally { | |
| resetUnexpectedErrorHandler(); | |
| } | |
| }); | |
| it('surfaces the discarded record to flush callers when the repair never reaches the rewrite', async () => { | |
| const capture: RepairCapture = { warnings: [], events: [] }; | |
| const svc = wireWithCapture(KEY, capture); | |
| const prefix = `${currentMetadata()}\n`; | |
| const raw = `${prefix}GARBAGE\n`; | |
| await seedCorrupt(raw); | |
| const originalWrite = storage.write.bind(storage); | |
| storage.write = async (scope, key, data, options) => { | |
| if (key === BACKUP_KEY) throw new Error('disk full'); | |
| return originalWrite(scope, key, data, options); | |
| }; | |
| await collect(svc.readJournal()); | |
| const unexpected: unknown[] = []; | |
| setUnexpectedErrorHandler((error) => unexpected.push(error)); | |
| try { | |
| svc.appendRecord({ type: 'wire.test.doomed', time: 9 }); | |
| await expect(svc.flush()).rejects.toThrow('Wire journal repair did not complete'); | |
| expect(await rawBytes()).toBe(raw); | |
| expect(unexpected).toHaveLength(1); | |
| expect(unexpected[0]).toBeInstanceOf(WireError); | |
| expect((unexpected[0] as WireError).code).toBe(WireErrors.codes.RECORDS_WRITE_FAILED); | |
| expect(capture.events).toEqual([ | |
| { | |
| name: 'wire_repair', | |
| payload: { | |
| kind: 'corrupted', | |
| outcome: 'failed', | |
| dropped_count: 1, | |
| backup_created: false, | |
| }, | |
| }, | |
| { | |
| name: 'wire_repair', | |
| payload: { | |
| kind: 'corrupted', | |
| outcome: 'failed', | |
| dropped_count: 1, | |
| backup_created: false, | |
| }, | |
| }, | |
| ]); | |
| } finally { | |
| resetUnexpectedErrorHandler(); | |
| } | |
| }); | |
| }); | |
| describe('WireService flush', () => { | |
| it('drains the dehydrate queue before resolving', async () => { | |
| const records: WireRecord[] = []; | |
| const stub = wireOverLog(recordingWireLog(records), 'flush'); | |
| let release!: () => void; | |
| const gate = new Promise<void>((resolve) => { | |
| release = resolve; | |
| }); | |
| stub.appendRecord({ type: 'wire.test.gated', time: 1 }, async (record) => { | |
| await gate; | |
| return record; | |
| }); | |
| await new Promise((resolve) => setImmediate(resolve)); | |
| expect(records).toEqual([]); | |
| let flushed = false; | |
| const flushPromise = stub.flush().then(() => { | |
| flushed = true; | |
| }); | |
| await new Promise((resolve) => setImmediate(resolve)); | |
| expect(flushed).toBe(false); | |
| release(); | |
| await flushPromise; | |
| expect(flushed).toBe(true); | |
| expect(records).toEqual([{ type: 'wire.test.gated', time: 1 }]); | |
| }); | |
| }); | |
| describe('WireService journal location', () => { | |
| it('counts journal lines across appends and a fresh read-back', async () => { | |
| await wire.seal(); | |
| wire.appendRecord({ type: 'wire.test.one', time: 1 }); | |
| wire.appendRecord({ type: 'wire.test.two', time: 2 }); | |
| await wire.flush(); | |
| expect(wire.lineCount()).toBe(3); | |
| const reopened = wireOverLog(log, KEY); | |
| expect(reopened.lineCount()).toBe(0); | |
| await collect(reopened.readJournal()); | |
| expect(reopened.lineCount()).toBe(3); | |
| reopened.appendRecord({ type: 'wire.test.three', time: 3 }); | |
| await reopened.flush(); | |
| expect(reopened.lineCount()).toBe(4); | |
| }); | |
| it('reports no journal path when the storage has no on-disk location', () => { | |
| expect(wire.journalPath()).toBeUndefined(); | |
| }); | |
| it('tracks the latest context.clear line across appends and reads', async () => { | |
| await wire.seal(); | |
| wire.appendRecord({ type: 'wire.test.one', time: 1 }); | |
| expect(wire.lastContextClearLine()).toBeUndefined(); | |
| wire.appendRecord({ type: 'context.clear', time: 2 }); | |
| wire.appendRecord({ type: 'wire.test.two', time: 3 }); | |
| await wire.flush(); | |
| expect(wire.lastContextClearLine()).toBe(3); | |
| const reopened = wireOverLog(log, KEY); | |
| expect(reopened.lastContextClearLine()).toBeUndefined(); | |
| await collect(reopened.readJournal()); | |
| expect(reopened.lastContextClearLine()).toBe(3); | |
| }); | |
| it('reports the on-disk journal path resolved by the storage layer', () => { | |
| const locatedStorage: IFileSystemStorageService = Object.assign(Object.create(storage), { | |
| pathFor: (scope: string, key: string) => `/home/user/.kimi-code/${scope}/${key}`, | |
| }) as IFileSystemStorageService; | |
| const located = wireOverLog(log, 'located', { storage: locatedStorage }); | |
| expect(located.journalPath()).toBe( | |
| `/home/user/.kimi-code/${testWireScope(SCOPE, 'located')}/${AGENT_WIRE_RECORD_KEY}`, | |
| ); | |
| }); | |
| }); | |
| describe('WireService tree projection', () => { | |
| function userPrompt(text: string, id: string) { | |
| return { | |
| role: 'user', | |
| content: [{ type: 'text', text }], | |
| toolCalls: [], | |
| origin: { kind: 'user' }, | |
| id, | |
| }; | |
| } | |
| it('writes the undo switch triple with paired physical line numbers after re-reading a rewritten journal', async () => { | |
| const enc = new TextEncoder(); | |
| const seed = [ | |
| JSON.stringify({ type: 'context.append_message', message: userPrompt('first', 'p1'), time: 1 }), | |
| JSON.stringify({ | |
| type: 'context.append_message', | |
| message: { role: 'assistant', content: [{ type: 'text', text: 'answer' }], toolCalls: [] }, | |
| time: 2, | |
| }), | |
| JSON.stringify({ type: 'context.append_message', message: userPrompt('second', 'p2'), time: 3 }), | |
| ]; | |
| await storage.write(testWireScope(SCOPE, KEY), AGENT_WIRE_RECORD_KEY, enc.encode(`${seed.join('\n')}\n`)); | |
| const switched = await wire.switchBranch({ turns: 1 }); | |
| expect(switched).toEqual({ | |
| branch: 'b1', | |
| base: { branch: 'main', line: 3 }, | |
| edgeLine: 5, | |
| forkLine: 3, | |
| }); | |
| const onDisk = await readRecords(); | |
| expect(onDisk.slice(4)).toEqual([ | |
| { | |
| type: 'agent.switched', | |
| agentId: 'test-agent', | |
| branch: 'b1', | |
| reason: 'undo', | |
| base: { branch: 'main', line: 3 }, | |
| turns: 1, | |
| legacyUndoLine: 6, | |
| time: expect.any(Number), | |
| }, | |
| { type: 'context.undo', agentId: 'test-agent', count: 1, time: expect.any(Number) }, | |
| { type: 'context.undone', agentId: 'test-agent', turns: 1, time: expect.any(Number) }, | |
| ]); | |
| expect(wire.lineCount()).toBe(7); | |
| expect(wire.journalRef).toEqual({ tree: testWireScope(SCOPE, KEY), branch: 'b1' }); | |
| expect(wire.branches()).toEqual(['main', 'b1']); | |
| expect(wire.nextSeq()).toBe(8); | |
| expect(await collect(wire.read())).toEqual([ | |
| ...onDisk.slice(0, 3), | |
| onDisk[5]!, | |
| onDisk[6]!, | |
| ]); | |
| expect(await collect(wire.readRaw())).toEqual(onDisk); | |
| await expect(wire.switchBranch({ turns: 2 })).rejects.toMatchObject({ | |
| reason: 'insufficient', | |
| }); | |
| expect(wire.lineCount()).toBe(7); | |
| wire.appendRecord({ | |
| type: 'context.apply_compaction', | |
| summary: 'compacted', | |
| compactedCount: 2, | |
| tokensBefore: 10, | |
| time: 4, | |
| }); | |
| await wire.flush(); | |
| await expect(wire.switchBranch({ turns: 1 })).rejects.toMatchObject({ | |
| reason: 'compaction_boundary', | |
| }); | |
| expect(wire.lineCount()).toBe(8); | |
| expect(await readRecords()).toHaveLength(8); | |
| }); | |
| it('walks nested switch edges to compute the active chain', async () => { | |
| const seeded: WireRecord[] = [ | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| { type: 'wire.test.one', time: 1 }, | |
| { type: 'wire.test.two', time: 2 }, | |
| { | |
| type: 'agent.switched', | |
| agentId: 'test-agent', | |
| branch: 'b1', | |
| reason: 'undo', | |
| base: { branch: 'main', line: 2 }, | |
| turns: 1, | |
| legacyUndoLine: 5, | |
| time: 3, | |
| }, | |
| { type: 'context.undo', agentId: 'test-agent', count: 1, time: 3 }, | |
| { type: 'context.undone', agentId: 'test-agent', turns: 1, time: 3 }, | |
| { type: 'wire.test.three', time: 4 }, | |
| { | |
| type: 'agent.switched', | |
| agentId: 'test-agent', | |
| branch: 'b2', | |
| reason: 'undo', | |
| base: { branch: 'b1', line: 7 }, | |
| turns: 1, | |
| legacyUndoLine: 9, | |
| time: 5, | |
| }, | |
| { type: 'context.undo', agentId: 'test-agent', count: 1, time: 5 }, | |
| { type: 'context.undone', agentId: 'test-agent', turns: 1, time: 5 }, | |
| { type: 'wire.test.four', time: 6 }, | |
| ]; | |
| const stub = wireOverLog(recordingWireLog(seeded), 'tree'); | |
| const journal = stub; | |
| expect(await collect(journal.readRaw())).toEqual(seeded); | |
| expect(await collect(journal.read())).toEqual([ | |
| seeded[0]!, | |
| seeded[1]!, | |
| seeded[4]!, | |
| seeded[5]!, | |
| seeded[6]!, | |
| seeded[8]!, | |
| seeded[9]!, | |
| seeded[10]!, | |
| ]); | |
| expect(journal.journalRef.branch).toBe('b2'); | |
| expect(journal.branches()).toEqual(['main', 'b1', 'b2']); | |
| expect(journal.nextSeq()).toBe(12); | |
| }); | |
| it('retries the branch switch on read interference, then rejects when the switch triple is interleaved', async () => { | |
| const records: WireRecord[] = [ | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| { type: 'context.append_message', message: userPrompt('first', 'p1'), time: 1 }, | |
| { | |
| type: 'context.append_message', | |
| message: { role: 'assistant', content: [{ type: 'text', text: 'answer' }], toolCalls: [] }, | |
| time: 2, | |
| }, | |
| ]; | |
| let concurrent: IWireService | undefined; | |
| let interferences = 0; | |
| let maxInterferences = 1; | |
| const log = recordingWireLog(records); | |
| log.read = function <R>(): AsyncIterable<R> { | |
| const snapshot = [...records]; | |
| const interfere = interferences < maxInterferences; | |
| if (interfere) interferences += 1; | |
| return (async function* () { | |
| try { | |
| for (const record of snapshot) { | |
| yield record as R; | |
| } | |
| } finally { | |
| if (interfere) { | |
| concurrent!.appendRecord({ type: 'wire.test.concurrent', time: 9 }); | |
| } | |
| } | |
| })(); | |
| }; | |
| concurrent = wireOverLog(log, 'concurrent'); | |
| const switched = await concurrent.switchBranch({ turns: 1 }); | |
| expect(switched).toMatchObject({ branch: 'b1', base: { branch: 'main', line: 1 }, forkLine: 1 }); | |
| expect(records[3]).toEqual({ type: 'wire.test.concurrent', time: 9 }); | |
| expect(records[4]?.type).toBe('agent.switched'); | |
| interferences = 0; | |
| maxInterferences = Number.MAX_SAFE_INTEGER; | |
| await expect(concurrent.switchBranch({ turns: 1 })).rejects.toMatchObject({ | |
| code: WireErrors.codes.RECORDS_WRITE_FAILED, | |
| }); | |
| const tailRecords: WireRecord[] = [ | |
| { type: 'metadata', protocol_version: WIRE_PROTOCOL_VERSION, created_at: 1 }, | |
| { type: 'context.append_message', message: userPrompt('first', 'p1'), time: 1 }, | |
| { | |
| type: 'context.append_message', | |
| message: { role: 'assistant', content: [{ type: 'text', text: 'answer' }], toolCalls: [] }, | |
| time: 2, | |
| }, | |
| ]; | |
| const tailLog = recordingWireLog(tailRecords); | |
| const baseAppend = tailLog.append.bind(tailLog); | |
| tailLog.append = (scope, key, record, options) => { | |
| if ((record as WireRecord).type === 'context.undo') { | |
| baseAppend(scope, key, { type: 'wire.test.interleaved', time: 9 }, options); | |
| } | |
| baseAppend(scope, key, record, options); | |
| }; | |
| const tail = wireOverLog(tailLog, 'tail'); | |
| await expect(tail.switchBranch({ turns: 1 })).rejects.toMatchObject({ | |
| code: WireErrors.codes.RECORDS_WRITE_FAILED, | |
| }); | |
| expect(tailRecords.map((record) => record.type)).toEqual([ | |
| 'metadata', | |
| 'context.append_message', | |
| 'context.append_message', | |
| 'agent.switched', | |
| 'wire.test.interleaved', | |
| 'context.undo', | |
| 'context.undone', | |
| ]); | |
| }); | |
| }); | |