fix(native-chat): make turn forking work, recoverable and bounded

The fork feature was inert. `completedTurnEnd` required a journal row whose
`turnLifecycle.state` is `completed`, and nothing writes one: settlement
TOMBSTONES the running row instead. Every turn therefore looked live and every
fork refused `busy`. Completion is now read as "not the active turn", off the
same projection the chat view reads, and the fixtures that hand-built the
impossible row are replaced with journals the real producers emit.

- Read completion as the absence of a running lifecycle row, via
  `activeStructuredAgentSessionTurnId`; Codex matches by turn id, Claude by the
  turn being the newest one.
- Carry `legacy:` rows into the child instead of throwing. Adoption imports the
  whole prior transcript under that namespace, so forking any turn taken after
  adopting a conversation refused `unsupported`. Parent turn-lifecycle rows are
  dropped with the other host operation-lifecycle rows.
- Give fork a refusal terminal. A failure that provably preceded any provider
  session now settles to `refused` and clears `retained`, instead of stranding
  the record at `attempted` where every later attach threw
  `agent_session_operation_unknown` and the dead prefix was rewritten on every
  lease renewal. An ambiguous outcome still keeps the guard and never makes a
  second provider session.
- Settle an existing epoch rather than replacing the child journal twice when a
  crash lands between the journal transaction and the completion transition.
- Evict the oldest idle client attempt at the cap instead of rejecting, which
  wedged forking for every session, tab and worktree until a restart, and report
  each typed refusal in its own words instead of collapsing them into "wait for
  the conversation to finish". The pre-request capability error now surfaces as
  itself rather than as an unconfirmed outcome.
- Stop a failed fork-support probe from blanking the slash-command menu and the
  per-turn option refresh on a remote target.
- Scan eligibility in one linear pass and skip it entirely while the action is
  unavailable: 2000 items went from 1.43ms to 0.05ms, and a single long turn
  from 16.2ms to 0.09ms.
- Drop six orphan i18n keys a mis-resolved merge resurrected, which failed
  `verify:localization-runtime-catalog`.
This commit is contained in:
Merge Sim
2026-09-09 18:40:12 -07:00
parent c0956f388d
commit 6a79bfb958
20 changed files with 796 additions and 164 deletions
@@ -0,0 +1,130 @@
import { describe, expect, it } from 'vitest'
import { agentJournalItemKey } from '../../shared/agent-session-journal-item-key'
import type {
AgentJournalItemBody,
AgentJournalItemIdentity,
AgentJournalRenderItem
} from '../../shared/agent-session-journal-types'
import {
selectAgentSessionPrefix,
structuredForkEligibleItems
} from '../../shared/agent-session-prefix'
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import { createCodexJournalTranslator } from './codex-structured-journal-translation'
import type { CodexStructuredSessionEvent } from './codex-structured-session-adapter'
const SESSION_ID = 'session-1'
const THREAD_ID = 'thread-abc'
const TURN_ID = 'turn-1'
const HANDLE = { provider: 'codex', threadId: THREAD_ID } as const
/** Reduces the sink's append/tombstone stream the way the journal store does, so the predicate
* under test sees what the real producer actually leaves behind. */
function reducingSink(): {
sink: StructuredAgentSessionEventSink
items: () => AgentJournalRenderItem[]
} {
const rows = new Map<string, AgentJournalItemBody>()
return {
sink: {
appendItem: (identity: AgentJournalItemIdentity, body: AgentJournalItemBody) => {
rows.set(agentJournalItemKey(identity), body)
},
appendTombstone: (identity: AgentJournalItemIdentity) => {
rows.delete(agentJournalItemKey(identity))
},
publish: () => {}
},
items: () =>
[...rows].map(([itemId, body], index) => ({
itemId,
body,
sequence: index,
observedAt: 1
})) as AgentJournalRenderItem[]
}
}
function notification(method: string, params: unknown): CodexStructuredSessionEvent {
return { type: 'notification', sessionId: SESSION_ID, threadId: THREAD_ID, method, params }
}
/** Drives the real translator through one turn. Stopping before `turn/completed` leaves the
* journal in the shape a live turn really has. */
function journalAfterTurn(settled: boolean, turnId: string = TURN_ID): AgentJournalRenderItem[] {
const tap = reducingSink()
const translator = createCodexJournalTranslator({
sink: tap.sink,
primaryThreadId: () => THREAD_ID
})
// Live user bubbles come from the submission, not the provider echo.
tap.sink.appendItem(
{ provider: 'orca', clientMessageId: `prompt-${turnId}` },
{ kind: 'message', role: 'user', blocks: [{ type: 'text', text: 'Ask' }] }
)
translator.handle(notification('turn/started', { turn: { id: turnId } }))
translator.handle(
notification('item/completed', {
item: { type: 'agentMessage', id: `answer-${turnId}`, text: 'Answered' }
})
)
if (settled) {
translator.handle(notification('turn/completed', { turn: { id: turnId } }))
}
return tap.items()
}
function assistantId(items: readonly AgentJournalRenderItem[]): string {
const answer = items.find(
(item) => item.body.kind === 'message' && item.body.role === 'assistant'
)
expect(answer).toBeDefined()
return answer!.itemId
}
describe('forking a turn the real Codex producer settled', () => {
it('leaves no lifecycle row behind once the turn completes', () => {
const items = journalAfterTurn(true)
// The producer contract this feature must read: settlement TOMBSTONES the row. Nothing in the
// product ever writes `turnLifecycle.state: 'completed'`, so a fixture that does is not a
// model of any journal Orca can produce.
expect(items.some((item) => item.body.kind === 'status' && item.body.turnLifecycle)).toBe(false)
expect(items.map((item) => item.body.kind)).toEqual(['message', 'message'])
})
it('offers the settled turn as forkable and retains it inclusively', () => {
const items = journalAfterTurn(true)
const itemId = assistantId(items)
expect(structuredForkEligibleItems(items).has(itemId)).toBe(true)
const selected = selectAgentSessionPrefix({
items,
itemId,
handle: HANDLE,
boundary: 'through'
})
expect(selected).toMatchObject({ ok: true, throughId: TURN_ID })
expect(selected.ok && selected.retained.map((item) => item.itemId)).toEqual(
items.map((item) => item.itemId)
)
})
it('refuses the same turn while it is still running', () => {
const items = journalAfterTurn(false)
const itemId = assistantId(items)
expect(
items.some(
(item) => item.body.kind === 'status' && item.body.turnLifecycle?.state === 'running'
)
).toBe(true)
expect(structuredForkEligibleItems(items).has(itemId)).toBe(false)
expect(
selectAgentSessionPrefix({ items, itemId, handle: HANDLE, boundary: 'through' })
).toMatchObject({ ok: false, reason: 'busy' })
})
it('keeps a settled earlier turn forkable while a later turn runs', () => {
const settled = journalAfterTurn(true)
const items = [...settled, ...journalAfterTurn(false, 'turn-2')]
expect(structuredForkEligibleItems(items)).toEqual(new Set([assistantId(settled)]))
})
})
@@ -1,6 +1,7 @@
import {
beginStructuredForkAttempt,
proveStructuredForkAcquisition
proveStructuredForkAcquisition,
refuseStructuredForkAttempt
} from './structured-agent-session-fork-lifecycle'
import { isDeepStrictEqual } from 'node:util'
import { claudeRewindAcquisitionProofs } from './structured-rewind-claude-proof'
@@ -80,8 +81,16 @@ export async function acquireOwner(
}
} catch (error) {
if (isAgentSessionPreSpawnError(error)) {
// The launch never resolved, so no provider session can exist: settle the attempt rather
// than strand it. A failure past this point stays ambiguous and keeps the guard.
await refuseStructuredForkAttempt(input.store, record, describePreSpawnRefusal(error))
throw error
}
return rethrowAfterAgentSessionAcquisitionCleanup(input.adapter, record.sessionId, error)
}
}
function describePreSpawnRefusal(error: unknown): string {
const cause = error instanceof Error ? (error.cause ?? error) : error
return cause instanceof Error && cause.message ? cause.message : 'agent_session_pre_spawn'
}
@@ -1,7 +1,7 @@
import { mkdtemp, rm } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterEach, describe, expect, it } from 'vitest'
import { afterEach, describe, expect, it, vi } from 'vitest'
import { AgentSessionRecordStore } from '../../runtime/agent-session-record-store'
import { AgentSessionJournal } from '../agent-session-journal/journal-store'
import type { AgentSessionForkRecord } from '../../../shared/agent-session-fork'
@@ -9,7 +9,9 @@ import { isAgentSessionRecord } from '../../../shared/agent-session-record'
import {
beginStructuredForkAttempt,
proveStructuredForkAcquisition,
publishStructuredForkJournal
publishStructuredForkJournal,
refuseStructuredForkAttempt,
restartRefusedStructuredFork
} from './structured-agent-session-fork-lifecycle'
const NOW = 1_800_000_000_000
@@ -198,4 +200,109 @@ describe('durable fork lifecycle', () => {
})
).toThrow('agent_session_provider_handle_invalid')
})
it('settles an attempt that never reached the provider and drops its stranded prefix', async () => {
const { store, record, fork } = await prepare()
await beginStructuredForkAttempt(store, record)
await refuseStructuredForkAttempt(store, record, 'managed account is switching')
// The store re-serializes in full on every lease renewal, so the dead prefix must not survive.
expect(store.getRecord(record.sessionId)?.fork).toMatchObject({
phase: 'refused',
reason: 'managed account is switching',
retained: []
})
expect(isAgentSessionRecord(store.getRecord(record.sessionId))).toBe(true)
// A refused record is terminal: it is re-armed deliberately, never resumed in place.
await expect(
beginStructuredForkAttempt(store, store.getRecord(record.sessionId)!)
).rejects.toThrow('agent_session_operation_unknown')
await restartRefusedStructuredFork(store, record.sessionId, fork)
expect(
await beginStructuredForkAttempt(store, store.getRecord(record.sessionId)!)
).toMatchObject({ throughId: 'turn' })
})
it('leaves an ambiguous outcome under the guard rather than settling it', async () => {
const { store, record } = await prepare()
await beginStructuredForkAttempt(store, record)
await store.commitProcessIdentity({
sessionId: record.sessionId,
fence: 1,
process: { hostId: 'local', pid: 123, processStartTimeMs: NOW, spawnToken: 'child-token' },
now: NOW
})
const proved = await store.proveOwner({
sessionId: record.sessionId,
fence: 1,
now: NOW,
link: proveStructuredForkAcquisition(record, {
linkId: 'fork-link',
handle: { provider: 'codex', threadId: 'child' },
origin: 'created',
mintedAtFence: 1,
observedAt: NOW
})
})
expect(proved.fork?.phase).toBe('provider-succeeded')
// A provider session may already exist here; settling would license a second one.
await refuseStructuredForkAttempt(store, proved, 'too late')
expect(store.getRecord(record.sessionId)?.fork?.phase).toBe('provider-succeeded')
})
it('settles its existing epoch after a crash instead of replacing the child journal twice', async () => {
const { root, store, record } = await prepare()
await beginStructuredForkAttempt(store, record)
await store.commitProcessIdentity({
sessionId: record.sessionId,
fence: 1,
process: { hostId: 'local', pid: 123, processStartTimeMs: NOW, spawnToken: 'child-token' },
now: NOW
})
const proved = await store.proveOwner({
sessionId: record.sessionId,
fence: 1,
now: NOW,
link: proveStructuredForkAcquisition(record, {
linkId: 'fork-link',
handle: { provider: 'codex', threadId: 'child' },
origin: 'created',
mintedAtFence: 1,
observedAt: NOW
})
})
const openJournal = async () => {
const journal = new AgentSessionJournal({
journalDir: join(root, 'journal'),
identity: {
sessionId: record.sessionId,
workspaceId: 'workspace',
hostId: 'local',
agent: 'codex',
providerHandle: { kind: 'codex', threadId: 'child' }
}
})
journals.push(journal)
await journal.open()
return journal
}
// Crash in the window between the journal transaction and the completion transition.
const transition = vi
.spyOn(store, 'transitionHandoff')
.mockRejectedValueOnce(new Error('host exited'))
await expect(publishStructuredForkJournal(store, proved, await openJournal())).rejects.toThrow(
'host exited'
)
transition.mockRestore()
expect(store.getRecord(record.sessionId)?.fork?.phase).toBe('provider-succeeded')
const reopened = await openJournal()
const epoch = reopened.cursor().epoch
await publishStructuredForkJournal(store, store.getRecord(record.sessionId)!, reopened)
expect(reopened.cursor().epoch).toBe(epoch)
expect(reopened.snapshot().items.map((item) => item.itemId)).toEqual(['codex:child:turn:0'])
expect(store.getRecord(record.sessionId)).toMatchObject({
schemaVersion: 2,
fork: { phase: 'completed', retained: [] }
})
})
})
@@ -1,5 +1,10 @@
import { isDeepStrictEqual } from 'node:util'
import { agentSessionForkAnchor } from '../../../shared/agent-session-fork'
import type { AgentSessionForkTarget } from '../../../shared/agent-session-fork'
import type {
AgentSessionForkRecord,
AgentSessionForkTarget
} from '../../../shared/agent-session-fork'
import { agentJournalItemKey } from '../../../shared/agent-session-journal-item-key'
import { isAdmissibleAgentJournalItemBody } from '../../../shared/agent-session-journal-schemas'
import {
agentSessionProviderHandleKey,
@@ -39,6 +44,56 @@ export async function beginStructuredForkAttempt(
}
}
/**
* Settle an attempt that provably never reached the provider.
*
* Without a terminal for this case the record stays `attempted` forever: every later attach throws
* `agent_session_operation_unknown` and replay deliberately skips the phase, so the child is
* wedged. Clearing `retained` matters just as much — the store re-serializes in full on every
* lease renewal, so a stranded prefix is rewritten every few seconds for the life of the record.
*
* Only a proven pre-provider failure may settle here. An ambiguous one must keep `attempted` and
* refuse a second attempt rather than risk a second provider session.
*/
export async function refuseStructuredForkAttempt(
store: AgentSessionRecordStore,
record: AgentSessionRecord,
reason: string
): Promise<void> {
const fork = store.getRecord(record.sessionId)?.fork
if (fork?.phase !== 'attempted') {
return
}
await store.transitionHandoff(record.sessionId, (current) => {
if (
current.lease.runtimeFence !== record.lease.runtimeFence ||
current.fork?.phase !== 'attempted'
) {
throw new Error('agent_session_checkpoint_stale')
}
return {
...current,
fork: { ...current.fork, phase: 'refused', reason: reason.slice(0, 512), retained: [] }
}
})
}
/** Re-arm a refused fork. The retry re-derives the prefix from the parent, because settling the
* refusal dropped the stranded copy. */
export function restartRefusedStructuredFork(
store: AgentSessionRecordStore,
sessionId: string,
fork: AgentSessionForkRecord
): Promise<unknown> {
return store.transitionHandoff(sessionId, (current) => {
if (current.fork?.phase !== 'refused') {
throw new Error('agent_session_checkpoint_stale')
}
// Merged, not replaced: the replay that preceded this stamped the recovery operation ids.
return { ...current, fork: { ...current.fork, ...fork, phase: 'prepared' } }
})
}
export function proveStructuredForkAcquisition(
record: AgentSessionRecord,
acquired: AgentSessionProviderHandleLink
@@ -83,11 +138,23 @@ export async function publishStructuredForkJournal(
}
return { ...item, body: item.body }
})
await journal.replaceEpochItems(
'handle_forked',
record.lease.runtimeFence,
forkJournalSeed(items, fork.source, head.handle)
)
const seed = forkJournalSeed(items, fork.source, head.handle)
// A crash after the journal transaction must settle its existing epoch, not replace it twice.
// The child journal is opened and seeded before its event sink is bound, so on re-entry it
// holds exactly the seed; anything else means the child moved on and must not be overwritten.
const landed = journal.snapshot().items.map(({ itemId, body }) => ({ itemId, body }))
if (landed.length > 0) {
if (
!isDeepStrictEqual(
landed,
seed.map(({ identity, body }) => ({ itemId: agentJournalItemKey(identity), body }))
)
) {
throw new Error('agent_session_operation_invalid')
}
} else {
await journal.replaceEpochItems('handle_forked', record.lease.runtimeFence, seed)
}
await store.transitionHandoff(record.sessionId, (current) => {
if (
current.lease.runtimeFence !== record.lease.runtimeFence ||
@@ -9,9 +9,10 @@ import type { AgentSessionForkSource } from '../../../shared/agent-session-fork'
import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope'
import { AgentSessionRecordStore } from '../../runtime/agent-session-record-store'
import { AgentSessionJournal } from '../agent-session-journal/journal-store'
import type {
StructuredAgentSessionAdapter,
StructuredAgentSessionAcquireInput
import {
AgentSessionPreSpawnError,
type StructuredAgentSessionAdapter,
type StructuredAgentSessionAcquireInput
} from './structured-agent-session-adapter'
import { StructuredAgentSessionHost } from './structured-agent-session-host'
import {
@@ -61,6 +62,8 @@ async function setup(provider: 'claude' | 'codex' = 'codex') {
sessionId: isChild ? 'child-thread' : 'parent-thread',
leafUuid: forked ? input.fork!.throughId : 'turn-2-1'
} as const)
// Settlement TOMBSTONES a turn's lifecycle row, so a finished turn leaves none behind. A
// fixture that appends `turnLifecycle.state: 'completed'` models no journal Orca can produce.
if (!isChild) {
for (const turn of ['turn-1', 'turn-2']) {
input.events?.appendItem(identity(turn, 0), {
@@ -73,10 +76,6 @@ async function setup(provider: 'claude' | 'codex' = 'codex') {
role: 'assistant',
blocks: [{ type: 'text', text: `${turn} answer` }]
})
input.events?.appendItem(
{ provider: 'orca', clientMessageId: `${turn}-done` },
{ kind: 'status', text: 'Completed', turnLifecycle: { turnId: turn, state: 'completed' } }
)
}
}
return {
@@ -104,10 +103,6 @@ async function setup(provider: 'claude' | 'codex' = 'codex') {
role: 'assistant',
blocks: [{ type: 'text', text: 'Third answer' }]
})
input.events?.appendItem(
{ provider: 'orca', clientMessageId: 'turn-3-done' },
{ kind: 'status', text: 'Completed', turnLifecycle: { turnId: 'turn-3', state: 'completed' } }
)
return { state: 'accepted', providerIdentity: identity('turn-3', 0) }
})
const adapter: StructuredAgentSessionAdapter = {
@@ -256,10 +251,28 @@ describe('fork from a structured turn', () => {
expect(store.getRecord('child-session')?.fork?.phase).toBe('attempted')
expect(store.getVisibleSessionTabIndex().sessionIds).not.toContain('child-session')
expect(host.hasSession('child-session')).toBe(false)
await host.fork(caller, child, source)
// The ambiguity guard: the provider may hold a child, so the retry must never make a second.
expect(await host.fork(caller, child, source)).toMatchObject({ ok: false })
expect(acquire).toHaveBeenCalledTimes(2)
})
it('recovers a fork whose launch failed before any provider session existed', async () => {
const { host, store, child, source, acquire } = await setup()
acquire.mockImplementationOnce(async () => {
throw new AgentSessionPreSpawnError(new Error('managed account is switching'))
})
await expect(host.fork(caller, child, source)).rejects.toThrow('managed account is switching')
// Settled, not stranded — and the dead prefix is dropped rather than rewritten on every
// lease renewal for the life of the record.
expect(store.getRecord('child-session')?.fork).toMatchObject({
phase: 'refused',
retained: []
})
expect(await host.fork(caller, child, source)).toMatchObject({ ok: true })
expect(host.journalSnapshot('child-session').items).not.toHaveLength(0)
expect(acquire.mock.calls.filter(([input]) => input.fork)).toHaveLength(2)
})
it('resumes the proved child when journal publication fails instead of forking again', async () => {
const { host, child, source, acquire, store } = await setup()
const replace = vi
@@ -2,6 +2,7 @@ import { forkJournalSeed } from './structured-fork-journal-seed'
import { restoreRewindJournalBody } from './structured-rewind-journal-body'
import type { AgentSessionRewindReason } from '../../../shared/agent-session-rewind'
import { prepareStructuredForkReplay } from './structured-agent-session-fork-replay'
import { restartRefusedStructuredFork } from './structured-agent-session-fork-lifecycle'
import type { AgentSessionForkSource } from '../../../shared/agent-session-fork'
import { agentSessionLeaseAdmitsWriter } from '../../../shared/agent-session-lease-adjudication'
import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope'
@@ -65,7 +66,9 @@ export function forkStructuredAgentSession(
) {
return refuse('invalid-target')
}
let fork = prior
// A refused attempt proved no provider session exists, so the retry re-derives the prefix
// instead of replaying a record whose retained copy was dropped when it settled.
let fork = prior?.phase === 'refused' ? undefined : prior
if (!fork) {
const support = context.deps.adapter.forkSupport?.(source.sessionId)
if (!support?.supported) {
@@ -140,10 +143,15 @@ export function forkStructuredAgentSession(
fields: attachFingerprintFields(attach)
})
}
// Ordered: replay reads the settled `refused` record to mint the recovery envelope this
// existing child record requires, and only then is the record re-armed for a fresh attempt.
const replay = await prepareStructuredForkReplay(context, caller, attach)
if (replay.result) {
return replay.result
}
if (prior?.phase === 'refused') {
await restartRefusedStructuredFork(store, params.envelope.sessionId, fork)
}
return attachStructuredAgentSession(attachContext, caller.callerKey, replay.params)
})
}
@@ -85,13 +85,62 @@ describe('fork journal identities', () => {
expect(() =>
forkJournalIdentity({ provider: 'claude', sessionId: 'parent', uuid: 'id' }, source, target)
).toThrow('agent_session_identity_required')
expect(() =>
forkJournalIdentity(
{ provider: 'legacy', agent: 'codex', sessionId: 'parent', recordId: '1' },
source,
target
)
).toThrow('agent_session_identity_required')
})
it('carries an adopted transcript into the child and continues it on the new thread', () => {
// Adoption imports the whole prior conversation under `legacy:` keys, so a turn taken after
// adopting is a mix of legacy history and Codex echoes. Both belong in the child.
const imported = ['msg_1', 'msg_2'].map((recordId) => ({
itemId: agentJournalItemKey({
provider: 'legacy',
agent: 'codex',
sessionId: 'parent',
recordId
}),
body,
observedAt: 1
}))
const seed = forkJournalSeed(
[
...imported,
{ itemId: 'codex:parent:turn-a:0', body, observedAt: 2 },
{ itemId: 'codex:parent:turn-a:1', body, observedAt: 3 }
],
source,
target
)
expect(seed.map((item) => agentJournalItemKey(item.identity))).toEqual([
'legacy:codex:parent:msg_1',
'legacy:codex:parent:msg_2',
'codex:child:turn-a:0',
'codex:child:turn-a:1'
])
})
it('leaves the parent turn-lifecycle row behind rather than seeding a phantom running turn', () => {
const lifecycle: AgentJournalItemBody = {
kind: 'status',
text: 'Codex is working…',
turnLifecycle: { turnId: 'turn-a', state: 'running' }
}
const seed = forkJournalSeed(
[
{
itemId: agentJournalItemKey({
provider: 'legacy',
agent: 'codex',
sessionId: 'parent',
recordId: 'turn-lifecycle:turn-a'
}),
body: lifecycle,
observedAt: 1
},
{ itemId: 'codex:parent:turn-a:1', body, observedAt: 2 }
],
source,
target
)
expect(seed.map((item) => agentJournalItemKey(item.identity))).toEqual(['codex:child:turn-a:1'])
})
it('does not retain source submission identities or prompt authority', () => {
@@ -32,6 +32,12 @@ export function forkJournalIdentity(
ordinal: identity.ordinal
}
}
// Import-scoped bridge-era rows: no provider echo can ever reconcile against one, and the
// session id already names the provider conversation the child forks from, so it carries over
// verbatim rather than being minted into a namespace the child cannot reproduce.
if (identity.provider === 'legacy') {
return identity
}
if (identity.provider === 'claude' && target.provider === 'claude') {
// Forked files rewrite sessionId; retained UUIDs still belong to the parent namespace.
return {
@@ -57,8 +63,9 @@ export function forkJournalSeed(
if (!identity) {
throw new Error('agent_session_identity_required')
}
// Host status and submission identities belong to the parent's operation lifecycle.
if (identity.provider === 'orca') {
// Host status and submission identities belong to the parent's operation lifecycle. So does a
// turn-lifecycle row, whichever namespace it was keyed in — the child opens its own.
if (identity.provider === 'orca' || (item.body.kind === 'status' && item.body.turnLifecycle)) {
return []
}
const rekeyed = forkJournalIdentity(identity, source, target)
@@ -244,7 +244,6 @@ export function NativeChatMessageList({
onScroll={handleScroll}
className="scrollbar-sleek h-full overflow-y-auto [scrollbar-gutter:stable_both-edges] px-3 pt-10 pb-4 sm:px-4"
>
<div
ref={contentRef}
// Why: matches composer column (max-w-4xl) with 5px horizontal inset
@@ -314,7 +313,13 @@ export function NativeChatMessageList({
/>
)}
{forkAction?.eligibleIds.has(message.id) ? (
<Button variant="ghost" size="sm" className="self-start text-muted-foreground" disabled={forkAction.pending} onClick={() => forkAction.onFork(message.id)}>
<Button
variant="ghost"
size="sm"
className="self-start text-muted-foreground"
disabled={forkAction.pending}
onClick={() => forkAction.onFork(message.id)}
>
<GitFork className="size-3.5" />
{translate('components.native-chat.forkFromTurn', 'Fork from this turn')}
</Button>
@@ -355,7 +360,6 @@ export function NativeChatMessageList({
) : null}
{!showTurnStatus && showTypingIndicator ? <NativeChatTypingIndicatorRow /> : null}
</div>
</div>
{showJump ? (
<button
@@ -1,6 +1,12 @@
import { beforeEach, describe, expect, it, vi } from 'vitest'
const { call } = vi.hoisted(() => ({ call: vi.fn() }))
vi.mock('@/runtime/structured-agent-session-client', () => ({ callStructuredAgentSession: call }))
const { call, CapabilityError } = vi.hoisted(() => ({
call: vi.fn(),
CapabilityError: class extends Error {}
}))
vi.mock('@/runtime/structured-agent-session-client', () => ({
callStructuredAgentSession: call,
StructuredAgentSessionCapabilityError: CapabilityError
}))
vi.mock('@/i18n/i18n', () => ({ translate: (_key: string, fallback: string) => fallback }))
import { forkStructuredSessionFromTurn } from './structured-agent-session-fork-command'
@@ -68,10 +74,53 @@ describe('fork create intent replay', () => {
call
.mockResolvedValueOnce({ ok: false, refusal: { forkReason: 'busy' } })
.mockResolvedValueOnce({ ok: true })
await expect(forkStructuredSessionFromTurn(args)).rejects.toThrow('Wait for the conversation')
await expect(forkStructuredSessionFromTurn(args)).rejects.toThrow('Wait for this conversation')
await forkStructuredSessionFromTurn(args)
expect(call.mock.calls[1]?.[2].envelope.sessionId).not.toBe(
call.mock.calls[0]?.[2].envelope.sessionId
)
})
it.each([
['unsupported', 'cannot be forked'],
['history-limit', 'too much history'],
['stale-epoch', 'moved on'],
['invalid-target', 'Could not fork this turn.'],
['outcome-unknown', 'could not be confirmed']
])('reports a %s refusal in its own words', async (forkReason, message) => {
call.mockResolvedValue({ ok: false, refusal: { forkReason } })
await expect(forkStructuredSessionFromTurn(input())).rejects.toThrow(message)
})
it('surfaces the pre-request capability guard instead of an unconfirmed outcome', async () => {
const args = input()
call
.mockRejectedValueOnce(new CapabilityError('Forking requires a newer Orca server.'))
.mockResolvedValueOnce({ ok: true })
await expect(forkStructuredSessionFromTurn(args)).rejects.toThrow('newer Orca server')
// Nothing was sent, so the retired attempt must not replay the abandoned child id.
await forkStructuredSessionFromTurn(args)
expect(call.mock.calls[1]?.[2].envelope.sessionId).not.toBe(
call.mock.calls[0]?.[2].envelope.sessionId
)
})
it('evicts the oldest unconfirmed attempt instead of wedging every later fork', async () => {
call.mockRejectedValue(new Error('disconnected'))
const first = input()
await expect(forkStructuredSessionFromTurn(first)).rejects.toThrow('could not be confirmed')
for (let index = 0; index < 200; index += 1) {
await expect(forkStructuredSessionFromTurn(input())).rejects.toThrow('could not be confirmed')
}
call.mockReset()
call.mockResolvedValue({ ok: true })
// The 202nd fork in this app session must still reach the host.
await forkStructuredSessionFromTurn(input())
expect(call).toHaveBeenCalledTimes(1)
// The evicted entry no longer replays, so the retry mints a fresh child id.
await forkStructuredSessionFromTurn(first)
expect(call.mock.calls[1]?.[2].envelope.sessionId).not.toBe(
call.mock.calls[0]?.[2].envelope.sessionId
)
})
})
@@ -9,13 +9,16 @@ import type {
AgentSessionMutationResult
} from '../../../../shared/agent-session-wire'
import type { RuntimeClientTarget } from '@/runtime/runtime-rpc-client'
import { callStructuredAgentSession } from '@/runtime/structured-agent-session-client'
import {
callStructuredAgentSession,
StructuredAgentSessionCapabilityError
} from '@/runtime/structured-agent-session-client'
import { translate } from '@/i18n/i18n'
const attempts = new Map<
string,
{ params: StructuredAgentSessionCreateParams; running?: Promise<void> }
>()
type ForkAttempt = { params: StructuredAgentSessionCreateParams; running?: Promise<void> }
const attempts = new Map<string, ForkAttempt>()
const MAX_TRACKED_ATTEMPTS = 128
export function forkStructuredSessionFromTurn(input: {
target: RuntimeClientTarget
@@ -37,16 +40,6 @@ export function forkStructuredSessionFromTurn(input: {
return attempt.running
}
if (!attempt) {
if (attempts.size >= 128) {
return Promise.reject(
new Error(
translate(
'components.native-chat.forkUnconfirmed',
'A fork could not be confirmed. Retry the same turn to check its outcome.'
)
)
)
}
attempt = {
params: structuredAgentSessionCreateParams({
sessionId: createStructuredAgentSessionId(input.agent, () => crypto.randomUUID()),
@@ -56,7 +49,7 @@ export function forkStructuredSessionFromTurn(input: {
randomUuid: () => crypto.randomUUID()
})
}
attempts.set(key, attempt)
track(key, attempt)
}
const current = attempt
current.running = callStructuredAgentSession<
@@ -65,35 +58,18 @@ export function forkStructuredSessionFromTurn(input: {
.then(
(result) => {
if (!result.ok) {
if (
result.refusal.forkReason &&
result.refusal.forkReason !== 'outcome-unknown' &&
result.refusal.forkReason !== 'proof-mismatch'
) {
attempts.delete(key)
throw new Error(
translate(
'components.native-chat.forkFailed',
'Could not fork this turn. Wait for the conversation to finish and try again.'
)
)
}
throw new Error(
translate(
'components.native-chat.forkUnconfirmed',
'A fork could not be confirmed. Retry the same turn to check its outcome.'
)
)
throw refusalError(key, result.refusal.forkReason)
}
attempts.delete(key)
},
() => {
throw new Error(
translate(
'components.native-chat.forkUnconfirmed',
'A fork could not be confirmed. Retry the same turn to check its outcome.'
)
)
(error: unknown) => {
// The capability guard runs before the request leaves the client, so no child exists: the
// attempt is retired and the real reason reaches the user.
if (error instanceof StructuredAgentSessionCapabilityError) {
attempts.delete(key)
throw error
}
throw new Error(unconfirmed())
}
)
.finally(() => {
@@ -101,3 +77,67 @@ export function forkStructuredSessionFromTurn(input: {
})
return current.running
}
/** Bound the table by EVICTING the oldest idle entry. Refusing at the cap instead wedged forking
* app-wide — every session, tab and worktree — until a restart, reported as an unconfirmed fork. */
function track(key: string, attempt: ForkAttempt): void {
while (attempts.size >= MAX_TRACKED_ATTEMPTS) {
let evicted = false
for (const [candidate, entry] of attempts) {
if (!entry.running) {
attempts.delete(candidate)
evicted = true
break
}
}
if (!evicted) {
break
}
}
attempts.set(key, attempt)
}
/** A settled refusal proves the host minted no provider session, so the child id is retired and a
* retry starts clean. An unknown or mismatched outcome must reuse it to adjudicate the original. */
function refusalError(key: string, reason: string | undefined): Error {
if (reason === undefined || reason === 'outcome-unknown' || reason === 'proof-mismatch') {
return new Error(unconfirmed())
}
attempts.delete(key)
return new Error(refusalMessage(reason))
}
function refusalMessage(reason: string): string {
switch (reason) {
case 'busy':
return translate(
'components.native-chat.forkBusy',
'Wait for this conversation to finish, then fork the turn.'
)
case 'unsupported':
case 'history-not-paginated':
return translate(
'components.native-chat.forkUnsupported',
'This conversation cannot be forked.'
)
case 'history-limit':
return translate(
'components.native-chat.forkHistoryLimit',
'This turn carries too much history to fork.'
)
case 'stale-epoch':
return translate(
'components.native-chat.forkStaleEpoch',
'This conversation moved on. Reopen it and fork the turn again.'
)
default:
return translate('components.native-chat.forkFailed', 'Could not fork this turn.')
}
}
function unconfirmed(): string {
return translate(
'components.native-chat.forkUnconfirmed',
'A fork could not be confirmed. Retry the same turn to check its outcome.'
)
}
@@ -5,6 +5,7 @@ import { beforeEach, describe, expect, it, vi } from 'vitest'
const mocks = vi.hoisted(() => ({
call: vi.fn(),
forkAvailable: vi.fn(),
operationId: vi.fn(),
enqueueSettingsWrite: vi.fn()
}))
@@ -12,7 +13,8 @@ let fence = 3
let sessionCommands: { name: string; kind: 'command' | 'skill' }[] | undefined
vi.mock('@/runtime/structured-agent-session-client', () => ({
callStructuredAgentSession: mocks.call
callStructuredAgentSession: mocks.call,
structuredAgentSessionForkAvailable: mocks.forkAvailable
}))
vi.mock('./native-chat-session-option-settings-write', () => ({
@@ -68,6 +70,7 @@ function seededByNextLaunch(): Record<string, string> | undefined {
}
const LOCAL_TARGET = { kind: 'local' } as const
const REMOTE_TARGET = { kind: 'environment', environmentId: 'remote-host' } as const
const OPTIONS = {
models: [
@@ -106,6 +109,7 @@ describe('useStructuredAgentSession options', () => {
mocks.call.mockImplementation((_target, method) =>
method === 'agentSession.options' ? Promise.resolve(OPTIONS) : Promise.resolve(null)
)
mocks.forkAvailable.mockReset().mockResolvedValue(true)
})
it('applies provider-reconciled values after a model change', async () => {
@@ -523,3 +527,52 @@ describe('session command catalog stream', () => {
).toHaveLength(0)
})
})
describe('fork support probe isolation', () => {
const COMMANDS = [{ name: 'compact', kind: 'command' as const }]
beforeEach(() => {
vi.clearAllMocks()
fence = 3
mocks.operationId.mockReset().mockReturnValue('operation-1')
mocks.call.mockImplementation((_target: unknown, method: string) =>
method === 'agentSession.options'
? Promise.resolve({ ...OPTIONS, conversationCommands: COMMANDS, fork: { supported: true } })
: Promise.resolve(null)
)
})
it('keeps conversation commands and options when the remote fork probe fails', async () => {
// The probe is a second RPC to the remote host. A transient failure used to blank the
// slash-command menu and drop this turn's option refresh with it.
mocks.forkAvailable.mockRejectedValue(new Error('unreachable'))
const { result } = renderHook(() =>
useStructuredAgentSession({
sessionId: 'session-1',
target: REMOTE_TARGET,
agent: 'codex',
isVisible: true
})
)
await waitFor(() => expect(result.current.conversationCommands).toEqual(COMMANDS))
expect(result.current.forkSupported).toBe(false)
// The per-turn option refresh rode on the same `.then`; it must still land.
expect(result.current.optionSnapshot.find((entry) => entry.id === 'model')?.kind).toMatchObject(
{ currentValue: 'gpt-live' }
)
})
it('reports fork support when the probe answers', async () => {
mocks.forkAvailable.mockResolvedValue(true)
const { result } = renderHook(() =>
useStructuredAgentSession({
sessionId: 'session-1',
target: REMOTE_TARGET,
agent: 'codex',
isVisible: true
})
)
await waitFor(() => expect(result.current.forkSupported).toBe(true))
expect(result.current.conversationCommands).toEqual(COMMANDS)
})
})
@@ -97,8 +97,11 @@ export function useStructuredAgentSession(args: {
sessionId
})
.then(async (result) => {
// The probe is a separate RPC to the remote host. Its failure may only cost the fork
// affordance — the slash-command menu and this turn's options come from `result`.
const forkSupported =
result.fork?.supported === true && (await structuredAgentSessionForkAvailable(target))
result.fork?.supported === true &&
(await structuredAgentSessionForkAvailable(target).catch(() => false))
if (!stale) {
setConversationSupport({
sessionId,
@@ -0,0 +1,53 @@
// @vitest-environment happy-dom
import { renderHook } from '@testing-library/react'
import { describe, expect, it, vi } from 'vitest'
const { eligible } = vi.hoisted(() => ({ eligible: vi.fn(() => new Set<string>(['assistant-1'])) }))
vi.mock('../../../../shared/agent-session-prefix', () => ({
structuredForkEligibleItems: eligible
}))
vi.mock('./structured-agent-session-fork-command', () => ({
forkStructuredSessionFromTurn: vi.fn()
}))
import { useStructuredForkAction } from './use-structured-fork-action'
type Props = Parameters<typeof useStructuredForkAction>[0]
type Controller = Parameters<typeof useStructuredForkAction>[1]
const props = { agent: 'codex', target: { kind: 'local' } } as unknown as Props
function controller(isWorking: boolean): Controller {
return {
forkSupported: true,
forkSource: { sessionId: 'parent', expectedEpoch: 'epoch', expectedRuntimeFence: 1 },
isWorking,
journalItems: [{ itemId: 'assistant-1' }]
} as unknown as Controller
}
describe('fork action eligibility work', () => {
it('does no eligibility scan while a turn is live', () => {
eligible.mockClear()
// The action is unavailable mid-turn, but a live turn emits a journal delta per frame and the
// memo runs before the early return — so an ungated memo rescans the transcript for nothing.
const { result, rerender } = renderHook(() =>
useStructuredForkAction(props, controller(true), 'worktree', () => {})
)
for (let index = 0; index < 5; index += 1) {
rerender()
}
expect(result.current).toBeUndefined()
expect(eligible).not.toHaveBeenCalled()
})
it('scans once the turn settles and the action becomes available', () => {
eligible.mockClear()
const { result } = renderHook(() =>
useStructuredForkAction(props, controller(false), 'worktree', () => {})
)
expect(eligible).toHaveBeenCalled()
expect(result.current?.eligibleIds.has('assistant-1')).toBe(true)
})
})
@@ -5,6 +5,8 @@ import type { useStructuredAgentSession } from './use-structured-agent-session'
import { forkStructuredSessionFromTurn } from './structured-agent-session-fork-command'
import { toRuntimeWorktreeSelector } from '@/runtime/runtime-worktree-selector'
const NO_ELIGIBLE_ITEMS: ReadonlySet<string> = new Set()
export function useStructuredForkAction(
props: Omit<NativeChatStructuredViewProps, 'mode'>,
controller: ReturnType<typeof useStructuredAgentSession>,
@@ -12,18 +14,22 @@ export function useStructuredForkAction(
onError: (message: string) => void
) {
const [pending, setPending] = useState(false)
const eligibleIds = useMemo(
() => structuredForkEligibleItems(controller.journalItems ?? []),
[controller.journalItems]
)
const agent = props.agent === 'claude' ? 'claude' : props.agent === 'codex' ? 'codex' : undefined
if (
!controller.forkSupported ||
!controller.forkSource ||
!worktreeId ||
controller.isWorking ||
!agent
) {
const enabled = Boolean(
controller.forkSupported &&
controller.forkSource &&
worktreeId &&
!controller.isWorking &&
agent
)
// Hooks cannot be skipped, so the unavailable case is gated inside the memo instead: a live turn
// emits a journal delta per frame and every one of them would rescan for a discarded result.
const eligibleIds = useMemo(
() =>
enabled ? structuredForkEligibleItems(controller.journalItems ?? []) : NO_ELIGIBLE_ITEMS,
[enabled, controller.journalItems]
)
if (!enabled || !controller.forkSource || !worktreeId || !agent) {
return undefined
}
return {
+7 -9
View File
@@ -5298,11 +5298,7 @@
"5c9c7c16aa": "Add a project to create workspaces",
"a30e34eb5c": "Close workspace board",
"views": "Sidebar view",
"createMenu": "Create",
"addProject": "Add project",
"25a95899c9": "Add Project",
"ca6f729da2": "New workspace ({{value0}})",
"moreActions": "More workspace actions"
"addProject": "Add project"
},
"SidebarNav": {
"80611a8b10": "Search",
@@ -16268,9 +16264,7 @@
"permission": "Needs attention"
},
"filtersSection": "Filters",
"viewSection": "View",
"compactModeDescription": "Shows shorter thread rows with one-line titles and two-line status messages.",
"unreadOnlyDescription": "Filters the activity list to show only threads with unread updates."
"viewSection": "View"
},
"clearCompleted": {
"clearedOne": "Cleared 1 completed agent",
@@ -17260,7 +17254,11 @@
},
"forkServerUpdateRequired": "Forking requires a newer Orca server. Update the server and try again.",
"forkFromTurn": "Fork from this turn",
"forkFailed": "Could not fork this turn. Wait for the conversation to finish and try again.",
"forkFailed": "Could not fork this turn.",
"forkBusy": "Wait for this conversation to finish, then fork the turn.",
"forkUnsupported": "This conversation cannot be forked.",
"forkHistoryLimit": "This turn carries too much history to fork.",
"forkStaleEpoch": "This conversation moved on. Reopen it and fork the turn again.",
"forkUnconfirmed": "A fork could not be confirmed. Retry the same turn to check its outcome."
},
"tab": {
@@ -15,6 +15,10 @@ import {
type RuntimeClientTarget
} from './runtime-rpc-client'
/** The paired runtime is too old for this method. Thrown BEFORE any request leaves the client, so
* a caller may report the real reason and retire the attempt: no session was created. */
export class StructuredAgentSessionCapabilityError extends Error {}
export async function callStructuredAgentSession<TResult>(
target: RuntimeClientTarget,
method: string,
@@ -28,7 +32,9 @@ export async function callStructuredAgentSession<TResult>(
AGENT_SESSION_REWIND_RUNTIME_CAPABILITY
))
) {
throw new Error('Rewinding requires a newer Orca server. Update the server and try again.')
throw new StructuredAgentSessionCapabilityError(
'Rewinding requires a newer Orca server. Update the server and try again.'
)
}
if (
method === 'agentSession.create' &&
@@ -37,7 +43,7 @@ export async function callStructuredAgentSession<TResult>(
'forkFrom' in params &&
!(await structuredAgentSessionForkAvailable(target))
) {
throw new Error(
throw new StructuredAgentSessionCapabilityError(
translate(
'components.native-chat.forkServerUpdateRequired',
'Forking requires a newer Orca server. Update the server and try again.'
+5 -1
View File
@@ -33,7 +33,11 @@ export const AgentSessionForkRecordSchema = z.object({
expectedRuntimeFence: z.number().int().positive(),
source: z.custom<AgentSessionProviderHandle>(isAgentSessionProviderHandle),
throughId: Key,
phase: z.enum(['prepared', 'attempted', 'provider-succeeded', 'completed']),
/** `attempted` is the ambiguity guard: the provider may or may not hold a child, so it never
* retries. `refused` is its terminal counterpart for a failure that provably preceded any
* provider session, and is recoverable. */
phase: z.enum(['prepared', 'attempted', 'provider-succeeded', 'completed', 'refused']),
reason: z.string().min(1).max(512).optional(),
retained: z
.array(
z.object({
+39 -37
View File
@@ -7,32 +7,35 @@ import {
agentSessionPrefixWithinBounds
} from './agent-session-prefix-bounds'
function items(provider: 'claude' | 'codex'): AgentJournalRenderItem[] {
return ['a', 'b'].flatMap(
(turnId, turn) =>
[
{
itemId:
provider === 'codex' ? `codex:parent:${turnId}:0` : `claude:parent:${turnId}-prompt`,
body: { kind: 'message', role: 'user', blocks: [] },
sequence: turn * 3,
observedAt: 1
},
{
itemId:
provider === 'codex' ? `codex:parent:${turnId}:1` : `claude:parent:${turnId}-answer`,
body: { kind: 'message', role: 'assistant', blocks: [] },
sequence: turn * 3 + 1,
observedAt: 1
},
{
itemId: `orca:${turnId}`,
body: { kind: 'status', text: 'Done', turnLifecycle: { turnId, state: 'completed' } },
sequence: turn * 3 + 2,
observedAt: 1
}
] as AgentJournalRenderItem[]
)
/** Models a real journal: settlement TOMBSTONES a turn's lifecycle row, so a finished turn leaves
* none behind. Only a live turn has one, appended at turn start — before its answer. */
function items(provider: 'claude' | 'codex', running?: string): AgentJournalRenderItem[] {
return ['a', 'b'].flatMap((turnId, turn) => {
const rows: AgentJournalRenderItem[] = [
{
itemId:
provider === 'codex' ? `codex:parent:${turnId}:0` : `claude:parent:${turnId}-prompt`,
body: { kind: 'message', role: 'user', blocks: [] },
sequence: turn * 3,
observedAt: 1
} as AgentJournalRenderItem
]
if (running === turnId) {
rows.push({
itemId: `legacy:${provider}:parent:turn-lifecycle%3A${turnId}`,
body: { kind: 'status', text: 'Working', turnLifecycle: { turnId, state: 'running' } },
sequence: turn * 3 + 1,
observedAt: 1
} as AgentJournalRenderItem)
}
rows.push({
itemId: provider === 'codex' ? `codex:parent:${turnId}:1` : `claude:parent:${turnId}-answer`,
body: { kind: 'message', role: 'assistant', blocks: [] },
sequence: turn * 3 + 2,
observedAt: 1
} as AgentJournalRenderItem)
return rows
})
}
describe('bounded conversation prefix', () => {
@@ -53,11 +56,11 @@ describe('bounded conversation prefix', () => {
expect(result).toMatchObject({ ok: true, throughId: provider === 'codex' ? 'a' : 'a-answer' })
if (result.ok) {
expect(result.retained.map((item) => item.itemId)).toEqual(
history.slice(0, 3).map((item) => item.itemId)
history.slice(0, 2).map((item) => item.itemId)
)
}
expect(structuredForkEligibleItems(history)).toEqual(
new Set([history[1]!.itemId, history[4]!.itemId])
new Set([history[1]!.itemId, history[3]!.itemId])
)
}
)
@@ -71,13 +74,13 @@ describe('bounded conversation prefix', () => {
: ({ provider, sessionId: 'parent', leafUuid: 'b-answer' } as const)
const result = selectAgentSessionPrefix({
items: history,
itemId: history[4]!.itemId,
itemId: history[3]!.itemId,
handle,
boundary: 'before'
})
expect(result.ok).toBe(true)
if (result.ok) {
expect(result.retained).toHaveLength(provider === 'codex' ? 3 : 4)
expect(result.retained).toHaveLength(provider === 'codex' ? 2 : 3)
}
}
})
@@ -97,13 +100,12 @@ describe('bounded conversation prefix', () => {
expect(
selectAgentSessionPrefix({ ...args, handle: { provider: 'codex', threadId: 'foreign' } })
).toMatchObject({ ok: false, reason: 'invalid-target' })
history[2]!.body = {
kind: 'status',
text: 'Working',
turnLifecycle: { turnId: 'a', state: 'running' }
}
expect(selectAgentSessionPrefix(args)).toMatchObject({ ok: false, reason: 'busy' })
expect(structuredForkEligibleItems(history).has(history[1]!.itemId)).toBe(false)
// A live turn: the running row is the one the producer actually leaves in the journal.
const live = items('codex', 'b')
expect(
selectAgentSessionPrefix({ ...args, items: live, itemId: live[4]!.itemId })
).toMatchObject({ ok: false, reason: 'busy' })
expect(structuredForkEligibleItems(live)).toEqual(new Set([live[1]!.itemId]))
})
it('inherits the retained entry and UTF-8 byte bounds', () => {
+47 -23
View File
@@ -2,10 +2,15 @@ import {
agentJournalSubmissionKey,
parseAgentJournalItemKey
} from './agent-session-journal-item-key'
import type { AgentJournalRenderItem, AgentJournalSnapshot } from './agent-session-journal-types'
import type {
AgentJournalItemIdentity,
AgentJournalRenderItem,
AgentJournalSnapshot
} from './agent-session-journal-types'
import type { AgentSessionProviderHandle } from './agent-session-provider-handle'
import type { AgentSessionRewindReason, AgentSessionRewindRecord } from './agent-session-rewind'
import { agentSessionPrefixWithinBounds } from './agent-session-prefix-bounds'
import { activeStructuredAgentSessionTurnId } from './structured-agent-session-projection'
type PrefixSelection =
| { ok: false; reason: AgentSessionRewindReason }
@@ -125,38 +130,57 @@ export function selectAgentSessionPrefix(
}
}
/** End of the turn containing `selected`, or null while that turn is still live.
*
* Settlement TOMBSTONES a turn's lifecycle row rather than rewriting it to `completed`, so a
* finished turn leaves no row behind and only the live turn still has one. Completion is
* therefore read as "not the active turn", off the same projection the chat view reads. */
function completedTurnEnd(
items: readonly AgentJournalRenderItem[],
selected: number,
providerKey: (id: string) => string
): number | null {
const key = parseAgentJournalItemKey(providerKey(items[selected]!.itemId))
const nextPrompt = items.findIndex(
(item, index) => index > selected && item.body.kind === 'message' && item.body.role === 'user'
)
const end = nextPrompt === -1 ? items.length : nextPrompt
const lifecycle = items
.slice(selected, end)
.findLast(
(item) =>
item.body.kind === 'status' &&
item.body.turnLifecycle &&
(key?.provider !== 'codex' || item.body.turnLifecycle.turnId === key.turnId)
)
return lifecycle?.body.kind === 'status' && lifecycle.body.turnLifecycle?.state === 'completed'
? end
: null
const key = parseAgentJournalItemKey(providerKey(items[selected]!.itemId))
return liveTurn(activeStructuredAgentSessionTurnId(items), key, nextPrompt === -1)
? null
: nextPrompt === -1
? items.length
: nextPrompt
}
/** Codex names the live turn in its own item keys; Claude does not, and a running turn is always
* the newest one, so a Claude row is live exactly when no later prompt bounds it. */
function liveTurn(
activeTurnId: string | null,
key: AgentJournalItemIdentity | null,
isLastTurn: boolean
): boolean {
if (activeTurnId === null) {
return false
}
return key?.provider === 'codex' ? key.turnId === activeTurnId : isLastTurn
}
export function structuredForkEligibleItems(items: readonly AgentJournalRenderItem[]): Set<string> {
return new Set(
items
.filter(
(item, index) =>
item.body.kind === 'message' &&
item.body.role === 'assistant' &&
completedTurnEnd(items, index, (id) => id) !== null
)
.map((item) => item.itemId)
const activeTurnId = activeStructuredAgentSessionTurnId(items)
const lastPrompt = items.findLastIndex(
(item) => item.body.kind === 'message' && item.body.role === 'user'
)
const eligible = new Set<string>()
items.forEach((item, index) => {
if (item.body.kind !== 'message' || item.body.role !== 'assistant') {
return
}
// Parsing every key is wasted work on the idle journal a fork is actually taken from.
if (
activeTurnId === null ||
!liveTurn(activeTurnId, parseAgentJournalItemKey(item.itemId), index > lastPrompt)
) {
eligible.add(item.itemId)
}
})
return eligible
}