diff --git a/src/main/agentActivity/AgentActivityStore.test.ts b/src/main/agentActivity/AgentActivityStore.test.ts index a4c1eab99..d12ea96c9 100644 --- a/src/main/agentActivity/AgentActivityStore.test.ts +++ b/src/main/agentActivity/AgentActivityStore.test.ts @@ -1,4 +1,4 @@ -import { appendFile, mkdtemp, readFile, rm } from 'node:fs/promises' +import { appendFile, chmod, mkdir, mkdtemp, readFile, readdir, rm, writeFile } from 'node:fs/promises' import { tmpdir } from 'node:os' import { join } from 'node:path' @@ -124,3 +124,224 @@ describe('AgentActivityStore', () => { expect(keys.size).toBe(1) }) }) + +// #1303: the context id was cached BEFORE its context line was written. One +// failed append (ENOSPC, EIO) then left every later interval for that agent +// this month pointing at a context line that never reached disk, and +// readIntervals dropped each one silently. +describe('a failed context write', () => { + it('does not orphan the agent\'s later intervals', async () => { + const store = new AgentActivityStore(dir) + const internal = store as unknown as { appendLines: (file: string, lines: string[]) => Promise } + const realAppend = internal.appendLines.bind(store) + let failNext = true + internal.appendLines = async (file, lines) => { + if (failNext) { failNext = false; throw Object.assign(new Error('no space left'), { code: 'ENOSPC' }) } + return realAppend(file, lines) + } + const start = Date.parse('2026-09-01T09:00:00Z') + await expect(store.appendInterval({ context, startedAt: start, endedAt: start + HOUR })).rejects.toThrow('no space left') + await store.appendInterval({ context, startedAt: start + 2 * HOUR, endedAt: start + 3 * HOUR }) + const read = await new AgentActivityStore(dir).readIntervals(start, start + 4 * HOUR) + expect(read.map(interval => interval.startedAt)).toEqual([start + 2 * HOUR]) + }) +}) + +// #1414 review a+b: a PARTIAL write (A's context line lands, then the append +// fails before its interval line) left id 1 on disk for A. A was not cached, so +// B's next context also took id 1; after a restart a later A interval reused +// id 1 and read back as B's time. +describe('a partially written context', () => { + it('never lets another agent reuse its id', async () => { + const store = new AgentActivityStore(dir) + const internal = store as unknown as { appendLines: (file: string, lines: string[]) => Promise } + const realAppend = internal.appendLines.bind(store) + let partial = true + internal.appendLines = async (file, lines) => { + if (partial) { + partial = false + await appendFile(file, lines[0] + '\n') + throw Object.assign(new Error('no space left'), { code: 'ENOSPC' }) + } + return realAppend(file, lines) + } + const start = Date.parse('2026-09-01T09:00:00Z') + const a = { ...context, agentKey: 'A', label: 'A' } + const b = { ...context, agentKey: 'B', label: 'B' } + await expect(store.appendInterval({ context: a, startedAt: start, endedAt: start + HOUR })).rejects.toThrow('no space left') + await store.appendInterval({ context: b, startedAt: start + HOUR, endedAt: start + 2 * HOUR }) + const restarted = new AgentActivityStore(dir) + await restarted.appendInterval({ context: a, startedAt: start + 2 * HOUR, endedAt: start + 3 * HOUR }) + const read = await new AgentActivityStore(dir).readIntervals(start, start + 4 * HOUR) + expect(read.map(interval => interval.context.agentKey)).toEqual(['B', 'A']) + }) +}) + +// #1414 review a round 2 (q115, "unknown is never empty"): an existing month +// file that cannot be READ was treated as absent, so a restarted store started +// ids at 1 and gave a second agent the id the first agent's lines already use. +// Once readable again, the first agent's later hours read back as the second +// agent's. Only ENOENT means "no file yet"; any other failure refuses the +// append, and the bytes stay as they were. +describe('an unreadable month file', () => { + it('refuses the append instead of restarting ids, and ids continue once it is readable', async () => { + const start = Date.parse('2026-09-01T09:00:00Z') + const a = { ...context, agentKey: 'A', label: 'A' } + const b = { ...context, agentKey: 'B', label: 'B' } + await new AgentActivityStore(dir).appendInterval({ context: a, startedAt: start, endedAt: start + HOUR }) + const file = join(dir, '2026-09.jsonl') + const before = await readFile(file) + await chmod(file, 0o200) + try { + await expect(new AgentActivityStore(dir).appendInterval({ context: b, startedAt: start + HOUR, endedAt: start + 2 * HOUR })).rejects.toThrow() + } finally { + await chmod(file, 0o600) + } + expect(await readFile(file)).toEqual(before) + const restarted = new AgentActivityStore(dir) + await restarted.appendInterval({ context: b, startedAt: start + HOUR, endedAt: start + 2 * HOUR }) + await restarted.appendInterval({ context: a, startedAt: start + 2 * HOUR, endedAt: start + 3 * HOUR }) + const read = await new AgentActivityStore(dir).readIntervals(start, start + 4 * HOUR) + expect(read.map(interval => interval.context.agentKey)).toEqual(['A', 'B', 'A']) + }) +}) + +// #1414 review b round 2 (test gap): a failed write that wrote NOTHING still +// consumed its id, and a restart must continue from the highest id on disk, +// not from the count of contexts (which would reissue a live id). +describe('an id gap left by a failed write', () => { + it('is never filled by a later context after a restart', async () => { + const store = new AgentActivityStore(dir) + const internal = store as unknown as { appendLines: (file: string, lines: string[]) => Promise } + const realAppend = internal.appendLines.bind(store) + let failNext = true + internal.appendLines = async (file, lines) => { + if (failNext) { failNext = false; throw Object.assign(new Error('no space left'), { code: 'ENOSPC' }) } + return realAppend(file, lines) + } + const start = Date.parse('2026-09-01T09:00:00Z') + const a = { ...context, agentKey: 'A', label: 'A' } + const b = { ...context, agentKey: 'B', label: 'B' } + const c = { ...context, agentKey: 'C', label: 'C' } + await expect(store.appendInterval({ context: a, startedAt: start, endedAt: start + HOUR })).rejects.toThrow('no space left') + await store.appendInterval({ context: b, startedAt: start + HOUR, endedAt: start + 2 * HOUR }) + const restarted = new AgentActivityStore(dir) + await restarted.appendInterval({ context: c, startedAt: start + 2 * HOUR, endedAt: start + 3 * HOUR }) + await restarted.appendInterval({ context: b, startedAt: start + 3 * HOUR, endedAt: start + 4 * HOUR }) + const read = await new AgentActivityStore(dir).readIntervals(start, start + 5 * HOUR) + expect(read.map(interval => interval.context.agentKey)).toEqual(['B', 'C', 'B']) + }) +}) + +// #1414 review a round 3 (q115): recovery treated an UNREADABLE open.json as +// "no open file" and overwrote it with an empty snapshot, so the pending +// interval was lost for good. Now an unreadable (or corrupt) snapshot is moved +// aside, bytes intact, and a later start recovers it once it can be read. +describe('an unreadable open-interval snapshot', () => { + it('is set aside instead of overwritten, and recovered once readable', async () => { + const start = Date.parse('2026-09-01T09:00:00Z') + const lastTouch = start + 2 * HOUR + await new AgentActivityStore(dir).writeOpen([{ sessionId: 'session-1', context, startedAt: start }], lastTouch) + const file = join(dir, 'open.json') + const before = await readFile(file) + await chmod(file, 0o200) + expect(await new AgentActivityStore(dir).recoverOpenIntervals(start + 10 * HOUR)).toBe(0) + const aside = (await readdir(dir)).filter(name => name.startsWith('open.json.unrecovered-')) + expect(aside).toHaveLength(1) + await chmod(join(dir, aside[0]!), 0o600) + expect(await readFile(join(dir, aside[0]!))).toEqual(before) + // Readable again: the next start recovers it and removes the set-aside copy. + expect(await new AgentActivityStore(dir).recoverOpenIntervals(start + 11 * HOUR)).toBe(1) + expect((await readdir(dir)).filter(name => name.startsWith('open.json.unrecovered-'))).toEqual([]) + const read = await new AgentActivityStore(dir).readIntervals(start, start + 12 * HOUR) + expect(read.map(interval => [interval.startedAt, interval.endedAt])).toEqual([[start, lastTouch]]) + }) +}) + +// #1414 review c: a recovery that fails PARTWAY (the second entry's append +// fails) set the whole snapshot aside, so the next start re-appended the +// entry already recovered. Only the entries not yet recovered are set aside. +describe('a recovery that fails partway', () => { + it('sets aside only what was not recovered, so nothing is counted twice', async () => { + const start = Date.parse('2026-09-01T09:00:00Z') + const lastTouch = start + 2 * HOUR + const a = { ...context, agentKey: 'A', label: 'A' } + const b = { ...context, agentKey: 'B', label: 'B' } + await new AgentActivityStore(dir).writeOpen([ + { sessionId: 'session-a', context: a, startedAt: start }, + { sessionId: 'session-b', context: b, startedAt: start + HOUR }, + ], lastTouch) + const store = new AgentActivityStore(dir) + const internal = store as unknown as { appendLines: (file: string, lines: string[]) => Promise } + const realAppend = internal.appendLines.bind(store) + let calls = 0 + internal.appendLines = async (file, lines) => { + calls += 1 + if (calls === 2) throw Object.assign(new Error('no space left'), { code: 'ENOSPC' }) + return realAppend(file, lines) + } + expect(await store.recoverOpenIntervals(start + 10 * HOUR)).toBe(1) + expect(await new AgentActivityStore(dir).recoverOpenIntervals(start + 11 * HOUR)).toBe(1) + const read = await new AgentActivityStore(dir).readIntervals(start, start + 12 * HOUR) + expect(read.map(interval => interval.context.agentKey).sort()).toEqual(['A', 'B']) + }) + + // #1414 review c (test gap): a set-aside copy that still cannot be read is + // KEPT for a later start, never deleted. + it('keeps a set-aside copy that still cannot be read', async () => { + const aside = join(dir, 'open.json.unrecovered-1000') + await mkdir(dir, { recursive: true }) + await writeFile(aside, '{"aliveAt":1,"open":[]}') + await chmod(aside, 0o200) + try { + await new AgentActivityStore(dir).recoverOpenIntervals(Date.parse('2026-09-01T12:00:00Z')) + expect((await readdir(dir)).filter(name => name.startsWith('open.json.unrecovered-'))).toEqual(['open.json.unrecovered-1000']) + } finally { + await chmod(aside, 0o600) + } + }) +}) + +// B6 check (q115): an aliases file that exists but cannot be READ was treated +// as absent twice over. The tail repair skipped its newline check, so a new +// edge was glued onto a torn last line and lost; and the aliases were cached +// as empty. Only ENOENT is "no file": otherwise the append is refused, the +// bytes are untouched, and nothing is cached, so a later call reads again. +describe('an unreadable aliases file', () => { + it('refuses the append instead of gluing onto an unseen tail, and works once readable', async () => { + const file = join(dir, 'aliases.jsonl') + await mkdir(dir, { recursive: true }) + await writeFile(file, '{"f":"A","t":"B"}') + const before = await readFile(file) + await chmod(file, 0o200) + try { + await expect(new AgentActivityStore(dir).appendAliases([['B', 'C']])).rejects.toThrow() + } finally { + await chmod(file, 0o600) + } + expect(await readFile(file)).toEqual(before) + await new AgentActivityStore(dir).appendAliases([['B', 'C']]) + const lines = (await readFile(file, 'utf8')).trim().split('\n').map(line => JSON.parse(line) as { f: string; t: string }) + expect(lines).toEqual([{ f: 'A', t: 'B' }, { f: 'B', t: 'C' }]) + }) + + // The READ path has one guard: `loadAliases`' ENOENT-only catch (B6 check, + // 2110). The append test above cannot pin it, because the tail repair + // refuses that append on its own. Here the SAME store instance reads while + // the file is unreadable: with a catch-all, the empty map would be cached + // and A and B would stay split for the life of the store. + it('does not cache an unreadable file as empty: the same store groups once it is readable', async () => { + const file = join(dir, 'aliases.jsonl') + await mkdir(dir, { recursive: true }) + await writeFile(file, '{"f":"A","t":"B"}\n') + const store = new AgentActivityStore(dir) + await chmod(file, 0o200) + try { + await expect(store.agentKeyGrouping()).rejects.toThrow() + } finally { + await chmod(file, 0o600) + } + const group = await store.agentKeyGrouping() + expect(group('A')).toBe(group('B')) + }) +}) diff --git a/src/main/agentActivity/AgentActivityStore.ts b/src/main/agentActivity/AgentActivityStore.ts index 49d69572d..0333a737d 100644 --- a/src/main/agentActivity/AgentActivityStore.ts +++ b/src/main/agentActivity/AgentActivityStore.ts @@ -1,4 +1,4 @@ -import { appendFile, mkdir, open, readFile, readdir, rename, writeFile } from 'node:fs/promises' +import { appendFile, mkdir, open, readFile, readdir, rename, rm, writeFile } from 'node:fs/promises' import { join } from 'node:path' import type { SystemSuspension } from '@shared/types/systemSuspension.js' @@ -107,6 +107,16 @@ function parseJsonLines(text: string): Record[] { return out } +type SetAsideSnapshot = { aliveAt: number | null; open: unknown[] } + +/** An open-interval snapshot recovered only partway: `remaining` starts at + * the entry whose append failed (#1414 review c). */ +class PartialRecovery extends Error { + constructor(readonly recovered: number, readonly remaining: SetAsideSnapshot) { + super('open-interval recovery stopped partway') + } +} + export class AgentActivityStore { private tail: Promise = Promise.resolve() /** aliases.jsonl in memory, loaded on first use. */ @@ -115,6 +125,15 @@ export class AgentActivityStore { private readonly cleanTails = new Set() /** Context ids already written to each month file this process has touched. */ private readonly monthContexts = new Map>() + /** + * The next context id to mint per month (#1414 review). Minting always + * advances it, even when the write then fails, so an id is never issued + * twice: a PARTIAL write can leave a context line on disk for an id whose + * mapping was never cached, and reusing that id for another agent made the + * later reads attribute one agent's time to the other. Loaded as the + * file's highest id + 1. + */ + private readonly monthNextId = new Map() constructor(private readonly dir: string) {} @@ -132,16 +151,26 @@ export class AgentActivityStore { const known = this.monthContexts.get(month) if (known) return known const ids = new Map() + let highest = 0 try { for (const line of parseJsonLines(await readFile(join(this.dir, `${month}.jsonl`), 'utf8'))) { if (line.t !== 'c' || !isNumber(line.c)) continue + highest = Math.max(highest, line.c) const context = parseContext(line) if (context) ids.set(contextKey(context), line.c) } - } catch { - // No file yet for this month. + } catch (error) { + // Only a missing file means "no contexts yet" (#1414 review a round 2, + // q115 "unknown is never empty"). A file that exists but cannot be read + // was treated as empty: ids restarted at 1, a second agent got the id + // the first agent's lines already use, and once readable the first + // agent's later hours read back as the second's. Refuse the append + // instead (the interval is lost, as for any failed write); nothing is + // cached, so the next append reads again. + if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error } this.monthContexts.set(month, ids) + this.monthNextId.set(month, highest + 1) return ids } @@ -154,15 +183,28 @@ export class AgentActivityStore { const key = contextKey(interval.context) const lines: string[] = [] let id = ids.get(key) + const isNewContext = id === undefined if (id === undefined) { - id = ids.size + 1 - ids.set(key, id) + // contextsFor always sets the month's next id. No `ids.size + 1` + // fallback: counting contexts reissues an id after any gap (#1414 + // review c; the gap test pins it). + id = this.monthNextId.get(month)! + this.monthNextId.set(month, id + 1) const contextLine: ContextLine = { t: 'c', c: id, ...interval.context } lines.push(JSON.stringify(contextLine)) } const intervalLine: IntervalLine = { t: 'i', c: id, s: interval.startedAt, e: interval.endedAt } lines.push(JSON.stringify(intervalLine)) await this.appendLines(join(this.dir, `${month}.jsonl`), lines) + // Cache the id only once its context line is on disk (#1303). Caching it + // first meant one failed append (ENOSPC, EIO) left every later interval + // for this agent this month pointing at a context line that never + // landed, and readIntervals drops an interval with no context. Not + // caching on failure means the next interval for this agent mints a + // NEW id (monthNextId already advanced) and writes its context line + // again. The failed id is burned: if its line landed partially, it + // still names this agent, and no other agent is ever given that id. + if (isNewContext) ids.set(key, id) }) } @@ -175,8 +217,12 @@ export class AgentActivityStore { aliases.set(line.f, line.t) } } - } catch { - // No aliases yet. + } catch (error) { + // Only a missing file means "no aliases yet" (B6 check, q115). An + // unreadable one was cached as empty, so every read grouped agents + // wrongly until restart. Throw instead; nothing is cached, so a later + // call reads again. + if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error } this.aliases = aliases return aliases @@ -239,8 +285,12 @@ export class AgentActivityStore { } finally { await handle.close() } - } catch { - // No file yet: nothing to repair. + } catch (error) { + // Only a missing file has no tail to repair (B6 check, q115). A file + // that cannot be read (write-only, EIO) has an unseen tail: appending + // blindly glued the new line onto a torn last line and lost it. + // Refuse the append instead. + if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error } } await appendFile(path, text) @@ -310,32 +360,104 @@ export class AgentActivityStore { /** Close every interval a previous run left open at that run's last touch. * An unclean shutdown therefore contributes at most one touch period of time - * that was not observed, never the hours until the next launch. */ + * that was not observed, never the hours until the next launch. + * + * "Unknown is never empty" (#1414 review a round 3, q115): recovery used to + * treat ANY failure (unreadable, corrupt, a failed append midway) as "no + * open file" and then overwrite open.json with an empty snapshot, losing the + * pending interval for good. Now only a missing file (ENOENT) is empty. + * Anything else is moved aside to `open.json.unrecovered-