From 09531c620c1f9ec4d193def2bab39f3eee3ea0fa Mon Sep 17 00:00:00 2001 From: Merge Sim Date: Mon, 7 Sep 2026 02:18:08 -0700 Subject: [PATCH] fix(native-chat): settle a roster the dying host never got to sweep MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `settleSession` only fires when the provider goes away while this process is alive. If the host itself dies, nothing sweeps and nothing reconciles on restore, so a `subagent-group` row persisted as `working` claimed live children forever — the mirror of the defect the previous commit fixed, and the same `ssh-execution-boundary.md` violation in the other direction. Reconciled host-side, at journal open, not in the renderer: mobile shows only the durable text twin and reconciles nothing, so a renderer-only fix would leave it claiming live children indefinitely. Opening the journal is also the one moment a host can honestly say the previous writer is gone. - `staleSubagentRosterRevisions` rewrites every child still reading `working` to `unverifiable` and regenerates the twin from the same summary, so the block and the sentence cannot disagree. - No terminal timestamp: the child stopped being observable at an unknown moment, and stamping the reopen would report the downtime as its run length. - Revises in place under the parsed identity, so a reopen upserts the row rather than appending a duplicate, and a second reopen writes nothing. - Skipped on a corrupt load: that journal is still owed a rebuild from provider history, and content past the repair's free sequence retires the demand. Reconciles journal ROWS, not roster state — the producer's in-process group map is untouched, so the roster's known seeding limitation is unchanged, as is `canReplaceSubagentState`: `unverifiable` still does not latch. --- .../journal-store-open.ts | 45 +++- .../journal-store-restore.ts | 3 +- .../journal-subagent-liveness.test.ts | 202 ++++++++++++++++++ .../journal-subagent-liveness.ts | 101 +++++++++ .../native-chat/NativeChatSubagentRun.tsx | 17 +- 5 files changed, 352 insertions(+), 16 deletions(-) create mode 100644 src/main/native-chat/agent-session-journal/journal-subagent-liveness.test.ts create mode 100644 src/main/native-chat/agent-session-journal/journal-subagent-liveness.ts diff --git a/src/main/native-chat/agent-session-journal/journal-store-open.ts b/src/main/native-chat/agent-session-journal/journal-store-open.ts index 721e5f4ba7f..7b5b6d0dff8 100644 --- a/src/main/native-chat/agent-session-journal/journal-store-open.ts +++ b/src/main/native-chat/agent-session-journal/journal-store-open.ts @@ -1,4 +1,8 @@ import { mkdir } from 'node:fs/promises' +import type { + AgentJournalItemBody, + AgentJournalItemIdentity +} from '../../../shared/agent-session-journal-types' import type { AgentType } from '../../../shared/agent-status-types' import { findJournalFileFormatRemnant, @@ -6,6 +10,7 @@ import { } from './journal-file-format-remnant' import type { JournalLoad } from './journal-open' import { journalRepairDisclosure, type JournalRepairDisclosure } from './journal-repair-disclosure' +import { staleSubagentRosterRevisions } from './journal-subagent-liveness' /** What any of this file's disclosures hands the store — a repair's, or the * pre-SQLite notice's. Same shape, and neither is only a repair. */ @@ -36,9 +41,9 @@ export async function openJournalStoreState(input: { adopt: (loaded: JournalLoad) => void /** Republishes an anchor row for an epoch a repair emptied. */ publishRepairEpoch: () => void - appendDisclosure: ( - identity: JournalRepairDisclosure['identity'], - body: JournalRepairDisclosure['body'], + appendItem: ( + identity: AgentJournalItemIdentity, + body: AgentJournalItemBody, fence: number ) => Promise agent: AgentType @@ -68,8 +73,9 @@ export async function openJournalStoreState(input: { } if (input.malformedRows() > 0 && !input.readOnly()) { const disclosure = journalRepairDisclosure({ malformedRows: input.malformedRows() }) - await input.appendDisclosure(disclosure.identity, disclosure.body, input.highestFence()) + await input.appendItem(disclosure.identity, disclosure.body, input.highestFence()) } + await settleStaleSubagentRosters(input, loaded) // Founding the epoch and appending the row are two transactions, and a // committed epoch sends every later open down this branch instead. Anything // that interrupts between them — a quit during startup restore, a failed @@ -92,7 +98,7 @@ export async function openJournalStoreState(input: { async function discloseFileFormatRemnant(input: { journalDir: string agent: AgentType - appendDisclosure: ( + appendItem: ( identity: JournalDisclosure['identity'], body: JournalDisclosure['body'], fence: number @@ -108,5 +114,32 @@ async function discloseFileFormatRemnant(input: { return } const disclosure = journalFileFormatRemnantDisclosure({ transcriptPath, agent: input.agent }) - await input.appendDisclosure(disclosure.identity, disclosure.body, input.highestFence()) + await input.appendItem(disclosure.identity, disclosure.body, input.highestFence()) +} + +/** + * Retires a `working` subagent roster the previous host never got to settle. + * + * Skipped on a corrupt load: that journal is still owed a rebuild from provider + * history, and content written past the repair's free sequence retires the + * demand for it. + */ +async function settleStaleSubagentRosters( + input: { + appendItem: ( + identity: AgentJournalItemIdentity, + body: AgentJournalItemBody, + fence: number + ) => Promise + highestFence: () => number + readOnly: () => boolean + }, + loaded: JournalLoad +): Promise { + if (input.readOnly() || loaded.corrupt) { + return + } + for (const revision of staleSubagentRosterRevisions(loaded.state.items.values())) { + await input.appendItem(revision.identity, revision.body, input.highestFence()) + } } diff --git a/src/main/native-chat/agent-session-journal/journal-store-restore.ts b/src/main/native-chat/agent-session-journal/journal-store-restore.ts index 3a5d3c7ac6e..fc69dd339d8 100644 --- a/src/main/native-chat/agent-session-journal/journal-store-restore.ts +++ b/src/main/native-chat/agent-session-journal/journal-store-restore.ts @@ -39,8 +39,7 @@ export function restoreJournalStore( publishRepairEpoch: () => collaborators.epochController.start('unreconcilable_prefix', host.state().highestFence), adopt: host.adopt, - appendDisclosure: (identity, body, fence) => - host.journal().appendItem(identity, body, { fence }), + appendItem: (identity, body, fence) => host.journal().appendItem(identity, body, { fence }), agent: host.identity.agent, highestFence: () => host.state().highestFence, malformedRows: host.malformedRows, diff --git a/src/main/native-chat/agent-session-journal/journal-subagent-liveness.test.ts b/src/main/native-chat/agent-session-journal/journal-subagent-liveness.test.ts new file mode 100644 index 00000000000..9d6887de61e --- /dev/null +++ b/src/main/native-chat/agent-session-journal/journal-subagent-liveness.test.ts @@ -0,0 +1,202 @@ +import { mkdtemp, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, beforeEach, describe, expect, it } from 'vitest' +import type { + AgentJournalRenderItem, + AgentSessionJournalIdentity +} from '../../../shared/agent-session-journal-types' +import { agentJournalItemKey } from '../../../shared/agent-session-journal-item-key' +import { isSubagentGroupBlock } from '../../../shared/native-chat-types' +import type { NativeChatSubagentEntry } from '../../../shared/native-chat-types' +import { + codexSubagentGroupBody, + codexSubagentGroupIdentity +} from '../../codex/codex-subagent-roster' +import type { openAgentSessionJournal } from './journal-store-factory' +import { createTrackedJournalOpener } from './journal-store-test-open' +import { staleSubagentRosterRevisions } from './journal-subagent-liveness' + +const IDENTITY: AgentSessionJournalIdentity = { + sessionId: 'session-1', + workspaceId: 'ws-1', + hostId: 'host-1', + agent: 'codex', + providerHandle: { kind: 'codex', threadId: 'thread-1' } +} + +const GROUP_ID = 'thread-1:turn-1' + +let root: string +let clock = 1_000 + +function tick(): number { + clock += 1 + return clock +} + +const journals = createTrackedJournalOpener() + +async function open(overrides: Partial[0]> = {}) { + return journals.open({ + identity: IDENTITY, + journalDir: root, + now: tick, + mintEpoch: () => `epoch-${clock}`, + ...overrides + }) +} + +/** The row as the producer writes it: the structured block plus its twin. */ +function rosterRow(agents: NativeChatSubagentEntry[]) { + return { + identity: codexSubagentGroupIdentity(GROUP_ID), + body: codexSubagentGroupBody(GROUP_ID, agents) + } +} + +function renderItem(agents: NativeChatSubagentEntry[]): AgentJournalRenderItem { + const row = rosterRow(agents) + return { + itemId: agentJournalItemKey(row.identity), + revision: 1, + body: row.body, + sequence: 2, + observedAt: 1 + } +} + +function rosterOf(body: AgentJournalRenderItem['body']): NativeChatSubagentEntry[] { + return body.kind === 'message' ? (body.blocks.find(isSubagentGroupBlock)?.agents ?? []) : [] +} + +function twinOf(body: AgentJournalRenderItem['body']): string | undefined { + return body.kind === 'message' + ? body.blocks.find((block) => block.type === 'text')?.text + : undefined +} + +beforeEach(async () => { + root = await mkdtemp(join(tmpdir(), 'orca-journal-subagents-')) + clock = 1_000 +}) + +afterEach(async () => { + await journals.closeAll() + await rm(root, { recursive: true, force: true }) +}) + +describe('staleSubagentRosterRevisions', () => { + it('settles a child the previous host left working, and moves the twin with it', () => { + const revisions = staleSubagentRosterRevisions([ + renderItem([ + { id: 'a', label: 'read_readme', state: 'working', startedAt: 10 }, + { id: 'b', label: 'read_package', state: 'completed', startedAt: 10, settledAt: 20 } + ]) + ]) + + expect(revisions).toHaveLength(1) + expect(rosterOf(revisions[0]!.body)).toMatchObject([ + { id: 'a', state: 'unverifiable' }, + { id: 'b', state: 'completed' } + ]) + // Mobile reads only this sentence, so it may not go on saying `Kicked off`. + expect(twinOf(revisions[0]!.body)).toBe('Ran 2 subagents (1 unverifiable)') + }) + + // The child stopped being observable at an unknown moment. A stamp taken now + // would report the time the app was down as how long the child ran. + it('records no terminal timestamp for a child whose run length is unknown', () => { + const revisions = staleSubagentRosterRevisions([ + renderItem([{ id: 'a', label: 'read', state: 'working', startedAt: 10 }]) + ]) + + expect(rosterOf(revisions[0]!.body)[0]).not.toHaveProperty('settledAt') + }) + + it('owes nothing for a roster whose children all settled', () => { + expect( + staleSubagentRosterRevisions([ + renderItem([{ id: 'a', label: 'read', state: 'completed', settledAt: 20 }]) + ]) + ).toEqual([]) + }) + + it('leaves rows that carry no roster alone', () => { + expect( + staleSubagentRosterRevisions([ + { + itemId: 'orca:plain', + revision: 1, + body: { kind: 'message', role: 'assistant', blocks: [{ type: 'text', text: 'hi' }] }, + sequence: 2, + observedAt: 1 + } + ]) + ).toEqual([]) + }) + + // Appending under a fresh identity would add a second row rather than revise + // the one on disk, so an unaddressable key is left exactly as it is. + it('skips a row whose key cannot be parsed back to its identity', () => { + expect( + staleSubagentRosterRevisions([ + { ...renderItem([{ id: 'a', label: 'r', state: 'working' }]), itemId: 'not-a-key' } + ]) + ).toEqual([]) + }) +}) + +describe('journal reopen after the writing host is gone', () => { + it('settles a persisted working roster to unverifiable, while the live row still reads working', async () => { + const live = await open() + const row = rosterRow([ + { id: 'a', label: 'read_readme', state: 'working', startedAt: 10 }, + { id: 'b', label: 'read_package', state: 'working', startedAt: 10 } + ]) + await live.appendItem(row.identity, row.body, { fence: 0 }) + + // Still the writing host: it can see the children, so the row says so. + const beforeRestart = live.snapshot().items.at(-1)! + expect(rosterOf(beforeRestart.body)).toMatchObject([{ state: 'working' }, { state: 'working' }]) + expect(twinOf(beforeRestart.body)).toBe('Kicked off 2 subagents') + + // The host dies without ever settling them — no `ended`, so no session sweep. + await live.close() + + const reopened = await open() + const afterRestart = reopened.snapshot().items.at(-1)! + expect(afterRestart.itemId).toBe(beforeRestart.itemId) + expect(rosterOf(afterRestart.body)).toMatchObject([ + { id: 'a', state: 'unverifiable' }, + { id: 'b', state: 'unverifiable' } + ]) + expect(twinOf(afterRestart.body)).toBe('Ran 2 subagents (2 unverifiable)') + }) + + it('revises the row in place rather than appending a second one', async () => { + const live = await open() + const row = rosterRow([{ id: 'a', label: 'read', state: 'working', startedAt: 10 }]) + await live.appendItem(row.identity, row.body, { fence: 0 }) + const before = live.snapshot().items.length + await live.close() + + const reopened = await open() + expect(reopened.snapshot().items).toHaveLength(before) + expect(reopened.snapshot().items.at(-1)?.revision).toBe(2) + }) + + it('writes nothing on a second reopen once every child is settled', async () => { + const live = await open() + const row = rosterRow([{ id: 'a', label: 'read', state: 'working', startedAt: 10 }]) + await live.appendItem(row.identity, row.body, { fence: 0 }) + await live.close() + + const once = await open() + const revision = once.snapshot().items.at(-1)?.revision + await once.close() + + const twice = await open() + expect(twice.snapshot().items.at(-1)?.revision).toBe(revision) + }) +}) diff --git a/src/main/native-chat/agent-session-journal/journal-subagent-liveness.ts b/src/main/native-chat/agent-session-journal/journal-subagent-liveness.ts new file mode 100644 index 00000000000..9b2724e9d1d --- /dev/null +++ b/src/main/native-chat/agent-session-journal/journal-subagent-liveness.ts @@ -0,0 +1,101 @@ +// A roster row left claiming live children by a host that is gone. +// +// The writing host revises its `subagent-group` rows in place while it can see +// the children, and sweeps whatever is still `working` when the provider goes +// away. A host that DIED — crash, quit, force-restart — does neither: its last +// revision goes on saying `working`, and nothing replays those children, so no +// later event can ever settle them. Opening the journal is the one moment a new +// host can state the truth about the old one: contact was lost. That is +// `unverifiable`, never a synthesized exit — see +// `docs/reference/ssh-execution-boundary.md`. +// +// Reconciles JOURNAL ROWS, not roster state: nothing here seeds the producer's +// in-process group map, so the roster's known limitation is untouched. + +import { + agentJournalItemKey, + parseAgentJournalItemKey +} from '../../../shared/agent-session-journal-item-key' +import type { + AgentJournalItemBody, + AgentJournalItemIdentity, + AgentJournalRenderItem +} from '../../../shared/agent-session-journal-types' +import { + isSubagentGroupFallbackText, + normalizeSubagentState, + subagentGroupFallbackText +} from '../../../shared/native-chat-subagent-summary' +import { + isSubagentGroupBlock, + type NativeChatBlock, + type NativeChatSubagentGroupBlock +} from '../../../shared/native-chat-types' + +export type JournalSubagentLivenessRevision = { + identity: AgentJournalItemIdentity + body: AgentJournalItemBody +} + +/** The revisions a reopened journal owes: one per row still claiming a live + * child. Empty — the common case — when nothing was left mid-flight. */ +export function staleSubagentRosterRevisions( + items: Iterable +): JournalSubagentLivenessRevision[] { + const revisions: JournalSubagentLivenessRevision[] = [] + for (const item of items) { + const body = item.body + if (body.kind !== 'message' || !body.blocks.some(hasWorkingChild)) { + continue + } + // A key that will not parse cannot be re-addressed, and appending under a + // fresh identity would duplicate the row rather than revise it. + const identity = parseAgentJournalItemKey(item.itemId) + if (!identity || agentJournalItemKey(identity) !== item.itemId) { + continue + } + revisions.push({ identity, body: { ...body, blocks: settleBlocks(body.blocks) } }) + } + return revisions +} + +function hasWorkingChild(block: NativeChatBlock): boolean { + return ( + isSubagentGroupBlock(block) && + block.agents.some((agent) => normalizeSubagentState(agent.state) === 'working') + ) +} + +/** No `settledAt`: the child stopped being observable at an unknown moment, and + * stamping the reopen would report the time the app was down as how long it + * ran. Readers already draw an unverifiable child with no stamp as having no + * known run length. */ +function settleBlocks(blocks: readonly NativeChatBlock[]): NativeChatBlock[] { + const settled = blocks.map((block) => + hasWorkingChild(block) ? settleGroup(block as NativeChatSubagentGroupBlock) : block + ) + const rosters = settled.filter(isSubagentGroupBlock) + const only = rosters.length === 1 ? rosters[0] : undefined + if (!only) { + return settled + } + // The plain-text twin is all a client without the block type ever shows, so it + // has to move with the block or the two would disagree about the same row. + const twin = subagentGroupFallbackText(only.agents) + return settled.map((block) => + block.type === 'text' && isSubagentGroupFallbackText(block.text) + ? { ...block, text: twin } + : block + ) +} + +function settleGroup(block: NativeChatSubagentGroupBlock): NativeChatSubagentGroupBlock { + return { + ...block, + agents: block.agents.map((agent) => + normalizeSubagentState(agent.state) === 'working' + ? { ...agent, state: 'unverifiable' as const } + : agent + ) + } +} diff --git a/src/renderer/src/components/native-chat/NativeChatSubagentRun.tsx b/src/renderer/src/components/native-chat/NativeChatSubagentRun.tsx index a19fa544a5a..af3bf468011 100644 --- a/src/renderer/src/components/native-chat/NativeChatSubagentRun.tsx +++ b/src/renderer/src/components/native-chat/NativeChatSubagentRun.tsx @@ -153,8 +153,9 @@ function SubagentElapsed({ * consulted: `spawn_agent` children outlive the turn that spawned them and keep * reporting into this group long after a newer turn opened, so a turn boundary * is a fact about the turn and never evidence that contact with a child was - * lost. Only the writing host can say that, and it does — see - * `CodexSubagentRoster.settleSession`. */ + * lost. Only a host can say that, and one does: `CodexSubagentRoster.settleSession` + * when the provider goes away, and `staleSubagentRosterRevisions` on the next + * journal open when the host itself died mid-flight. */ export function NativeChatSubagentRun({ block }: { @@ -189,12 +190,12 @@ export function NativeChatSubagentRun({ const alertState = working ? summary.adverseState : null const alert = alertState === null ? null : subagentStateLabel(alertState, summary.adverseCount, summary.total) - // A child can read `unverifiable` with no terminal timestamp — a state a newer - // build wrote that this one cannot name. Its run length is unknown, and - // measuring it to `now` would report the time since we lost sight of it as how - // long it ran, on a row that is not even counting. A sibling's stamp is no - // better: in a mixed group it would present that sibling's duration as the - // group's while a child's fate is still unknown. + // A child settled by the reopen reads `unverifiable` with no terminal stamp: + // it stopped being observable at an unknown moment. Measuring to `now` would + // report the time since the host died as how long the child ran, on a row that + // is not even counting. A sibling's stamp is no better: in a mixed group it + // would present that sibling's duration as the group's while a child's fate is + // still unknown. const runLengthUnknown = agents.some( (agent) => normalizeSubagentState(agent.state) === 'unverifiable' && typeof agent.settledAt !== 'number'