kimi-code / packages /agent-core-v2 /test /wire /wireService.test.ts
SaylorTwift's picture
SaylorTwift HF Staff
Add files using upload-large-folder tool
f0634fb verified
Raw History Blame Contribute Delete
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',
]);
});
});