mirror of
https://github.com/stablyai/orca.git
synced 2026-10-01 08:01:56 +00:00
fix(native-chat): settle a roster the dying host never got to sweep
`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.
This commit is contained in:
@@ -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<unknown>
|
||||
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<unknown>
|
||||
highestFence: () => number
|
||||
readOnly: () => boolean
|
||||
},
|
||||
loaded: JournalLoad
|
||||
): Promise<void> {
|
||||
if (input.readOnly() || loaded.corrupt) {
|
||||
return
|
||||
}
|
||||
for (const revision of staleSubagentRosterRevisions(loaded.state.items.values())) {
|
||||
await input.appendItem(revision.identity, revision.body, input.highestFence())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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<Parameters<typeof openAgentSessionJournal>[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)
|
||||
})
|
||||
})
|
||||
@@ -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<AgentJournalRenderItem>
|
||||
): 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
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -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'
|
||||
|
||||
Reference in New Issue
Block a user