mirror of
https://github.com/stablyai/orca.git
synced 2026-10-01 00:02:10 +00:00
feat(native-chat): Codex sessions write their subagents into the host status store (#22553)
* refactor(native-chat): the Codex acquire names its turn-boundary methods as a set Behavior-neutral: the same two methods stamp receipt time. Keeps the file under the size limit once the child-work sink lands. * feat(native-chat): Codex sessions write their subagents into the host status store A Codex child thread and each persistent command become host child records, fed through the same delivery, ingest and reducer the Claude lane uses. The child's own turn decides it: turn start is live, turn completion settles it with the outcome Codex reports, and a follow-up turn reopens the same record as a new run. Its open tool call, last message, usage and waiting-on-user flag come from its own thread's frames. A parent turn ending settles nothing. * fix(native-chat): close a Codex child's tool call by its item id alone A completion frame need not restate the tool it ran, so reading the tool name before closing left the call open and the record naming a finished tool. * test(native-chat): pin the Codex child-work evidence and every hop to the host's records Child turn start/end/follow-up, open tool call, last message, usage, waiting, the persistent command a child owns and its monitoring display, a primary turn end settling nothing, and session end. End to end through the real adapter: evidence after the journal and the legacy republish, and the parent state the records imply equals today's at every frame of a scripted session. Through the production runtime: a Codex session's child work reaches the status sink under its own address, and a provider exit ends it there. * test(native-chat): a Codex child's new run never inherits the last run's open call * test(native-chat): a Codex session with no child-work sink holds no evidence * test(native-chat): deliver a Codex child's announcement twice, as Codex does, before counting edges * refactor(native-chat): hand the Codex producer's pending edge over directly * fix(native-chat): name every Codex turn state in the outcome map; type the runtime test's fake opener * fix(native-chat): a Codex child's turn ends on the error that ends it, or on its thread closing Codex can end a child's turn with no turn/completed: an error it will not retry is that turn's own end (the verdict the transcript already settles the same turn on), and a closed thread ran its last turn. The executions, the one owner of child turn state, now end the turn on both, so the strip drops the child and its record settles (failed, or unknown for a close) together, instead of reading working for the life of the session. A systemError status is not an ending: Codex raises it for errors that leave the turn running. A child fact whose frame names no turn now belongs to the turn the child is running, instead of counting for every run. * test(native-chat): a Codex child's turn ending by fatal error or thread close settles strip and record together * test(native-chat): the Codex parity script reads a waiting child through the shared fold's waiting arm * test(native-chat): a Codex child row's journal attempt is its record's generation The journal numbers a Codex child's runs by the turns it observed on the child's thread; the host record numbers them by the runs its evidence opened. Both are keyed by the child's own turn id, so they must agree run for run, including when Codex reports the child's first turn before the spawn that announces it. * test(native-chat): a Codex session's end settles its live children and keeps the ended ones The host no longer erases a session's children when its provider goes away: a child still running settles with an outcome nobody reported, and a child that had already ended keeps what it said. The producer tests now expect exactly that, from the close path and from an unexpected exit. * fix(native-chat): a Codex subagent's shell is its open tool until the process exits Codex runs every agent shell through unified exec, so every subagent shell arrives with the source the persistent-command tracker keys on. The producer skipped those items, so a working subagent never named its shell, and an approved command (started on the approval path, completed from unified exec) stayed its open tool until the turn ended. The tracker still records the process separately, so a command that outlives the turn reads as monitoring. * fix(native-chat): a Codex shell becomes a subagent's own work only once it outlives its turn Codex runs every agent shell through unified exec and never says when one is left running, so the producer turned every shell, even a millisecond `rg`, into a command record the moment it started. Each settled into the session's pool of 32 settled records, so a busy turn evicted a finished subagent's record (its outcome row would vanish) and listed dozens of finished shells beside it. A command now becomes a record at the first turn boundary of the thread that launched it while its process still runs: until then it is the agent's open call. A shell that exits within its turn never becomes a record. * refactor(native-chat): child records keep every settled child and can be removed outright Settled child records now stay until the host drops the session's row; the 32-record trim is gone. A producer can say work stopped with nothing to report, and its record (and the handles it answered to) goes instead of settling. Evidence stays host-internal: the producer and the store share one process. * fix(native-chat): a Codex command is live work from its start until its process stops The command tracker is now the one owner of a Codex command's lifetime. It admits every command whatever `source` Codex tags it with (the approval path starts one as `agent`), and ends it when its process exits, when its thread closes (Codex stops the processes first, so no exit ever arrives), or when the session ends. The producer mirrors that one-to-one: a live record from the start, removed when the command stops, never settled. This removes the turn-boundary rule: a command that was only recorded at its turn's end left the parent reading done for one publish when the main agent's turn ended with a shell still running. The parity script now checks the parent at every journal write, not only at frame end. * fix(native-chat): a Codex command whose approval its turn abandoned never ran Codex starts an approval's command item before it asks, and when the turn ends with the question unanswered (the user stops at the approval), it drops the question and never completes the item. The command tracker admitted that start as a running process, so the strip kept a phantom command row and the session row read working until the session ended. The prompt registry, which owns which approvals are still unanswered, reports the command approvals a turn ended without; the tracker ends those commands with the frame that ended the turn. An answered approval keeps its command. * test(native-chat): start the Codex child-work runtime test without the removed hold Main no longer has host.hold: creating the session starts its child, and nothing a viewer does keeps it running. The test attaches and asserts the one child that attach started, then drives it as before.
This commit is contained in:
@@ -1,15 +1,20 @@
|
||||
import type { AgentSessionBackgroundTask } from '../../shared/agent-session-wire'
|
||||
import type { CodexBackgroundTaskEvent } from './codex-background-task-frames'
|
||||
import { codexCommandOutlivesTurn } from './codex-command-lifecycle'
|
||||
import { readRecord, readString } from './codex-item-field-readers'
|
||||
import { readCodexThreadItem } from './codex-structured-item-translation'
|
||||
import { MAX_CODEX_ITEM_STREAM_METADATA_BYTES } from './codex-item-stream-retention'
|
||||
import type { CodexAbandonedCommand } from './codex-prompt-registry'
|
||||
|
||||
const MAX_SETTLED_COMMANDS = 128
|
||||
const MAX_DESCRIPTION_CHARS = 512
|
||||
|
||||
type Command = { threadId: string; task: AgentSessionBackgroundTask; bytes: number }
|
||||
|
||||
/** A command process starting, or ending: it exited, its thread closed, or the session ended. */
|
||||
export type CodexBackgroundCommandChange =
|
||||
| { type: 'started'; threadId: string; task: AgentSessionBackgroundTask }
|
||||
| { type: 'ended'; threadId: string; taskId: string }
|
||||
|
||||
/** The label's reserved share of the description. Reserved, not merely capped:
|
||||
* a label free to spend the whole budget clips away the command it qualifies,
|
||||
* leaving a command row naming an agent and no command — the failure this
|
||||
@@ -65,28 +70,17 @@ export class CodexBackgroundCommandTracker {
|
||||
)
|
||||
}
|
||||
|
||||
observe(event: CodexBackgroundTaskEvent): void {
|
||||
observe(event: CodexBackgroundTaskEvent): CodexBackgroundCommandChange | null {
|
||||
const parsed = this.parse(event)
|
||||
if (!parsed || this.settled.has(parsed.key)) {
|
||||
return
|
||||
return null
|
||||
}
|
||||
const { key, command, completed } = parsed
|
||||
const existing = this.commands.get(key)
|
||||
if (completed) {
|
||||
if (existing) {
|
||||
this.liveBytes -= existing.bytes
|
||||
this.commands.delete(key)
|
||||
}
|
||||
const bytes = Buffer.byteLength(key, 'utf8') + 256
|
||||
if (this.liveBytes + bytes <= this.maxMetadataBytes) {
|
||||
this.settled.set(key, bytes)
|
||||
this.settledBytes += bytes
|
||||
}
|
||||
this.trimSettled()
|
||||
return
|
||||
return this.end(key)
|
||||
}
|
||||
if (existing) {
|
||||
return
|
||||
if (this.commands.has(key)) {
|
||||
return null
|
||||
}
|
||||
if (this.liveBytes + command.bytes > this.maxMetadataBytes) {
|
||||
throw new Error('Codex command metadata was not admitted before observation')
|
||||
@@ -94,6 +88,19 @@ export class CodexBackgroundCommandTracker {
|
||||
this.commands.set(key, command)
|
||||
this.liveBytes += command.bytes
|
||||
this.trimSettled()
|
||||
return { type: 'started', threadId: command.threadId, task: command.task }
|
||||
}
|
||||
|
||||
/** The thread closed: Codex stops its processes first, so none of them can report an exit. */
|
||||
endThread(threadId: string): CodexBackgroundCommandChange[] {
|
||||
return [...this.commands]
|
||||
.filter(([, command]) => command.threadId === threadId)
|
||||
.flatMap(([key]) => this.end(key) ?? [])
|
||||
}
|
||||
|
||||
/** Its approval went unanswered until its turn ended, so its process never started. */
|
||||
endUnapproved(command: CodexAbandonedCommand): CodexBackgroundCommandChange | null {
|
||||
return this.end(JSON.stringify([command.threadId, command.itemId]))
|
||||
}
|
||||
|
||||
tasks(
|
||||
@@ -113,11 +120,45 @@ export class CodexBackgroundCommandTracker {
|
||||
})
|
||||
}
|
||||
|
||||
clear(): void {
|
||||
/** The live commands one thread launched, as the strip would publish them. */
|
||||
threadTasks(threadId: string): AgentSessionBackgroundTask[] {
|
||||
return [...this.commands.values()]
|
||||
.filter((command) => command.threadId === threadId)
|
||||
.map((command) => command.task)
|
||||
}
|
||||
|
||||
/** The session ended, and every command with it. */
|
||||
clear(): CodexBackgroundCommandChange[] {
|
||||
const ended = [...this.commands.values()].map(
|
||||
({ threadId, task }): CodexBackgroundCommandChange => ({
|
||||
type: 'ended',
|
||||
threadId,
|
||||
taskId: task.id
|
||||
})
|
||||
)
|
||||
this.commands.clear()
|
||||
this.settled.clear()
|
||||
this.liveBytes = 0
|
||||
this.settledBytes = 0
|
||||
return ended
|
||||
}
|
||||
|
||||
/** Retires the key so a replayed frame cannot start the command again. */
|
||||
private end(key: string): CodexBackgroundCommandChange | null {
|
||||
const existing = this.commands.get(key)
|
||||
if (existing) {
|
||||
this.liveBytes -= existing.bytes
|
||||
this.commands.delete(key)
|
||||
}
|
||||
const bytes = Buffer.byteLength(key, 'utf8') + 256
|
||||
if (this.liveBytes + bytes <= this.maxMetadataBytes) {
|
||||
this.settled.set(key, bytes)
|
||||
this.settledBytes += bytes
|
||||
}
|
||||
this.trimSettled()
|
||||
return existing
|
||||
? { type: 'ended', threadId: existing.threadId, taskId: existing.task.id }
|
||||
: null
|
||||
}
|
||||
|
||||
private trimSettled(): void {
|
||||
@@ -140,8 +181,10 @@ export class CodexBackgroundCommandTracker {
|
||||
if (event.method !== 'item/started' && event.method !== 'item/completed') {
|
||||
return null
|
||||
}
|
||||
// Any command may outlive its turn; `source` says only how Codex launched it. A stdin write
|
||||
// starts no process: it reaches one already tracked.
|
||||
const item = readCodexThreadItem(readRecord(event.params).item)
|
||||
if (!item || !codexCommandOutlivesTurn(item)) {
|
||||
if (item?.type !== 'commandExecution' || item.source === 'unifiedExecInteraction') {
|
||||
return null
|
||||
}
|
||||
const key = JSON.stringify([event.threadId, item.id])
|
||||
|
||||
@@ -7,6 +7,7 @@ import {
|
||||
import { codexChildTurnState } from './codex-subagent-executions'
|
||||
import { readRecord } from './codex-item-field-readers'
|
||||
import { readCodexThreadItem } from './codex-structured-item-translation'
|
||||
import { readCodexProviderVerdict } from './codex-structured-journal-provider-verdicts'
|
||||
import { readCodexTurnId } from './codex-structured-thread-facts'
|
||||
|
||||
export type CodexBackgroundTaskFrame =
|
||||
@@ -15,6 +16,8 @@ export type CodexBackgroundTaskFrame =
|
||||
agentThreadId: string
|
||||
label: string | null
|
||||
parentTurnId: string | null | undefined
|
||||
/** The reporting thread, for a `started` activity: the agent that spawned the child. */
|
||||
spawnerThreadId: string | undefined
|
||||
}
|
||||
| {
|
||||
kind: 'turn'
|
||||
@@ -22,6 +25,15 @@ export type CodexBackgroundTaskFrame =
|
||||
turnId: string
|
||||
state: NativeChatSubagentState
|
||||
}
|
||||
| {
|
||||
/** A child turn that ended with no `turn/completed`. No `turnId`: the one it is running. */
|
||||
kind: 'turn-ended'
|
||||
threadId: string
|
||||
turnId: string | null
|
||||
state: CodexChildTurnEnding
|
||||
}
|
||||
|
||||
type CodexChildTurnEnding = Extract<NativeChatSubagentState, 'failed' | 'unverifiable'>
|
||||
|
||||
export type CodexBackgroundTaskEvent = {
|
||||
method: string
|
||||
@@ -29,10 +41,34 @@ export type CodexBackgroundTaskEvent = {
|
||||
params: unknown
|
||||
}
|
||||
|
||||
/**
|
||||
* The two ways Codex ends a child's turn without `turn/completed`. An `error` it will not retry is
|
||||
* that turn's own end: the verdict the transcript settles the same turn on. A closed thread ran
|
||||
* its last turn, and Codex never said how it went. A `systemError` status is neither: Codex raises
|
||||
* it for errors that leave the turn running too (a refused steer), and a turn one ends also
|
||||
* carries the `error`.
|
||||
*/
|
||||
function readCodexChildTurnEnding(
|
||||
event: CodexBackgroundTaskEvent
|
||||
): CodexBackgroundTaskFrame | null {
|
||||
if (readCodexProviderVerdict(event.method, event.params) === 'turn-failed') {
|
||||
const turnId = readCodexTurnId(event.params)
|
||||
return { kind: 'turn-ended', threadId: event.threadId, turnId, state: 'failed' }
|
||||
}
|
||||
return event.method === 'thread/closed'
|
||||
? { kind: 'turn-ended', threadId: event.threadId, turnId: null, state: 'unverifiable' }
|
||||
: null
|
||||
}
|
||||
|
||||
export function readCodexBackgroundTaskFrame(
|
||||
event: CodexBackgroundTaskEvent,
|
||||
primaryThreadId: string
|
||||
): CodexBackgroundTaskFrame | null {
|
||||
// The session's own turn ends through the journal's turn boundaries, never here.
|
||||
const ending = event.threadId === primaryThreadId ? null : readCodexChildTurnEnding(event)
|
||||
if (ending) {
|
||||
return ending
|
||||
}
|
||||
if (event.method === 'turn/started' || event.method === 'turn/completed') {
|
||||
const turnId = readCodexTurnId(event.params)
|
||||
if (turnId === null) {
|
||||
@@ -67,6 +103,8 @@ export function readCodexBackgroundTaskFrame(
|
||||
parentTurnId:
|
||||
activity.kind === 'started' || activity.kind === 'interacted'
|
||||
? readCodexTurnId(event.params)
|
||||
: undefined
|
||||
: undefined,
|
||||
// Only `started` names the spawner: other kinds ride whichever agent acted.
|
||||
spawnerThreadId: activity.kind === 'started' ? event.threadId : undefined
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,25 +2,50 @@ import type {
|
||||
AgentSessionBackgroundTask,
|
||||
AgentSessionBackgroundTaskState
|
||||
} from '../../shared/agent-session-wire'
|
||||
import type { AgentChildWorkEvidence } from '../../shared/agent-status-child-work-evidence'
|
||||
import {
|
||||
readCodexBackgroundTaskFrame,
|
||||
type CodexBackgroundTaskEvent
|
||||
} from './codex-background-task-frames'
|
||||
import { CodexSubagentExecutions } from './codex-subagent-executions'
|
||||
import { CodexBackgroundCommandTracker } from './codex-background-command-tracker'
|
||||
import { CodexChildWorkEvidence } from './codex-child-work-evidence'
|
||||
import type { CodexAbandonedCommand } from './codex-prompt-registry'
|
||||
import type { CodexStructuredSessionAdapterDeps } from './codex-structured-session-state'
|
||||
import { boundSubagentField } from './codex-subagent-group-body'
|
||||
|
||||
/** Where a session's child-work evidence goes, and the host clock that stamps it. */
|
||||
export type CodexChildWorkSink = {
|
||||
deliver: (evidence: AgentChildWorkEvidence[]) => void
|
||||
now: () => number
|
||||
}
|
||||
|
||||
export function codexChildWorkSink(
|
||||
sessionId: string,
|
||||
deps: Pick<CodexStructuredSessionAdapterDeps, 'onChildWorkEvidence' | 'now'>
|
||||
): CodexChildWorkSink {
|
||||
return {
|
||||
deliver: (evidence) => deps.onChildWorkEvidence?.(sessionId, evidence),
|
||||
now: () => deps.now?.() ?? Date.now()
|
||||
}
|
||||
}
|
||||
|
||||
/** Projects the same child execution facts the durable roster consumes. */
|
||||
export class CodexBackgroundTaskTracker {
|
||||
private publishedFingerprint = '[]'
|
||||
private publishedState: AgentSessionBackgroundTaskState | null = null
|
||||
private readonly commands: CodexBackgroundCommandTracker
|
||||
private readonly childWork: CodexChildWorkEvidence
|
||||
|
||||
constructor(
|
||||
private readonly primaryThreadId: string,
|
||||
private readonly executions = new CodexSubagentExecutions()
|
||||
private readonly executions = new CodexSubagentExecutions(),
|
||||
private readonly childWorkSink?: CodexChildWorkSink
|
||||
) {
|
||||
this.commands = new CodexBackgroundCommandTracker(primaryThreadId)
|
||||
this.childWork = new CodexChildWorkEvidence(primaryThreadId, executions, (threadId) =>
|
||||
this.commands.threadTasks(threadId)
|
||||
)
|
||||
}
|
||||
|
||||
get state(): AgentSessionBackgroundTaskState | null {
|
||||
@@ -32,20 +57,38 @@ export class CodexBackgroundTaskTracker {
|
||||
return this.commands.canObserve(event)
|
||||
}
|
||||
|
||||
observe(event: CodexBackgroundTaskEvent): boolean {
|
||||
/** `unapproved`: commands whose approval the journal dropped with this frame's turn ending. */
|
||||
observe(
|
||||
event: CodexBackgroundTaskEvent,
|
||||
unapproved: readonly CodexAbandonedCommand[] = []
|
||||
): boolean {
|
||||
const itemEvent = event.method === 'item/started' || event.method === 'item/completed'
|
||||
if (itemEvent) {
|
||||
this.commands.observe(event)
|
||||
}
|
||||
const command = itemEvent ? this.commands.observe(event) : null
|
||||
const commands = [
|
||||
...unapproved.flatMap((abandoned) => this.commands.endUnapproved(abandoned) ?? []),
|
||||
...(event.method === 'thread/closed'
|
||||
? this.commands.endThread(event.threadId)
|
||||
: command
|
||||
? [command]
|
||||
: [])
|
||||
]
|
||||
const frame = readCodexBackgroundTaskFrame(event, this.primaryThreadId)
|
||||
if (!frame) {
|
||||
return itemEvent ? this.refresh() : false
|
||||
}
|
||||
if (frame.kind === 'subagent') {
|
||||
this.executions.register(frame.agentThreadId, frame.label, frame.parentTurnId)
|
||||
} else if (frame.threadId !== this.primaryThreadId) {
|
||||
if (frame?.kind === 'subagent') {
|
||||
this.executions.register(
|
||||
frame.agentThreadId,
|
||||
frame.label,
|
||||
frame.parentTurnId,
|
||||
frame.spawnerThreadId
|
||||
)
|
||||
} else if (frame?.kind === 'turn-ended') {
|
||||
this.executions.endTurn(frame.threadId, frame.turnId, frame.state)
|
||||
} else if (frame && frame.threadId !== this.primaryThreadId) {
|
||||
this.executions.observeTurn(frame.threadId, frame.turnId, frame.state)
|
||||
}
|
||||
this.childWork.observe(event, frame, commands)
|
||||
if (!frame) {
|
||||
return itemEvent || commands.length > 0 ? this.refresh() : false
|
||||
}
|
||||
// A primary-turn frame only prompts a republish: turn end reveals children,
|
||||
// it never settles them. Codex `spawn_agent` children keep reporting well
|
||||
// past their parent turn, so nothing here may sweep the roster.
|
||||
@@ -54,10 +97,25 @@ export class CodexBackgroundTaskTracker {
|
||||
|
||||
clear(): boolean {
|
||||
this.executions.clear()
|
||||
this.commands.clear()
|
||||
this.childWork.clear(this.commands.clear())
|
||||
return this.refresh()
|
||||
}
|
||||
|
||||
/** Everything the frames observed since the last drain said about the session's child work. */
|
||||
drainChildWorkEvidence(observedAt: number): AgentChildWorkEvidence[] {
|
||||
return this.childWork.drain(observedAt)
|
||||
}
|
||||
|
||||
/** Hand the pending evidence to the host. Callers run this after the journal wrote the frame
|
||||
* and the parent's own row republished, so a child record never lands ahead of either. */
|
||||
publishChildWork(): void {
|
||||
// Drained even with no sink, so undelivered evidence never accumulates.
|
||||
const evidence = this.drainChildWorkEvidence(this.childWorkSink?.now() ?? Date.now())
|
||||
if (evidence.length > 0) {
|
||||
this.childWorkSink?.deliver(evidence)
|
||||
}
|
||||
}
|
||||
|
||||
private tasks(): AgentSessionBackgroundTask[] {
|
||||
const children = this.executions.workingChildren()
|
||||
const agents: AgentSessionBackgroundTask[] = children.map((child, index) => ({
|
||||
|
||||
@@ -0,0 +1,623 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { createAgentChildWorkAdmission } from '../../shared/agent-status-child-work-admission'
|
||||
import type { AgentChildWorkRecord } from '../../shared/agent-status-child-work'
|
||||
import type { AgentChildWorkEvidence } from '../../shared/agent-status-child-work-evidence'
|
||||
import { reconcileAgentChildWorkEvidence } from '../../shared/agent-status-child-work-reconciliation'
|
||||
import {
|
||||
agentChildWorkOwnedLiveness,
|
||||
deriveAgentChildDisplayState,
|
||||
projectAgentChildWorkViews
|
||||
} from '../../shared/agent-status-child-work-view'
|
||||
import { createAgentStatusStore } from '../../shared/agent-status-store'
|
||||
import { makeStructuredAgentStatusSubject } from '../../shared/agent-status-subject'
|
||||
import type { CodexBackgroundTaskEvent } from './codex-background-task-frames'
|
||||
import { CodexBackgroundTaskTracker } from './codex-background-task-tracker'
|
||||
|
||||
const PRIMARY = 'thread-parent'
|
||||
const PARENT_TURN = 'turn-parent'
|
||||
const CHILD = 'thread-child'
|
||||
const parent = makeStructuredAgentStatusSubject(
|
||||
{
|
||||
executionHostId: 'local',
|
||||
wslDistro: null,
|
||||
workspaceId: 'workspace-1',
|
||||
workspaceKind: 'folder'
|
||||
},
|
||||
'session-1'
|
||||
)
|
||||
|
||||
function turn(
|
||||
method: 'turn/started' | 'turn/completed',
|
||||
threadId: string,
|
||||
turnId: string,
|
||||
status = 'completed'
|
||||
): CodexBackgroundTaskEvent {
|
||||
return { method, threadId, params: { threadId, turn: { id: turnId, status } } }
|
||||
}
|
||||
|
||||
function spawned(
|
||||
child = CHILD,
|
||||
reporter = PRIMARY,
|
||||
name = 'audit_build'
|
||||
): CodexBackgroundTaskEvent {
|
||||
return {
|
||||
method: 'item/started',
|
||||
threadId: reporter,
|
||||
params: {
|
||||
threadId: reporter,
|
||||
turnId: PARENT_TURN,
|
||||
item: {
|
||||
type: 'subAgentActivity',
|
||||
id: `activity-${child}`,
|
||||
kind: 'started',
|
||||
agentThreadId: child,
|
||||
agentPath: `/root/${name}`
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function item(
|
||||
method: 'item/started' | 'item/completed',
|
||||
threadId: string,
|
||||
turnId: string,
|
||||
fields: Record<string, unknown>
|
||||
): CodexBackgroundTaskEvent {
|
||||
return { method, threadId, params: { threadId, turnId, item: fields } }
|
||||
}
|
||||
|
||||
function shell(id: string, command: string, status = 'inProgress', source = 'agent') {
|
||||
return { type: 'commandExecution', id, command, source, status }
|
||||
}
|
||||
|
||||
function harness() {
|
||||
const tracker = new CodexBackgroundTaskTracker(PRIMARY)
|
||||
const store = createAgentStatusStore({ epoch: 'epoch-1', mode: 'authority' })
|
||||
expect(store.applyMutation({ parent: { subject: parent } })).not.toBeNull()
|
||||
let minted = 0
|
||||
const admission = createAgentChildWorkAdmission(store, {
|
||||
mintChildWorkId: () => `child-${++minted}`
|
||||
})
|
||||
let clock = 1_000
|
||||
const log: AgentChildWorkEvidence[][] = []
|
||||
const send = (...events: CodexBackgroundTaskEvent[]): void => {
|
||||
for (const event of events) {
|
||||
tracker.observe(event)
|
||||
clock += 10
|
||||
const evidence = tracker.drainChildWorkEvidence(clock)
|
||||
log.push(evidence)
|
||||
reconcileAgentChildWorkEvidence({ store, admission, parent, provider: 'codex', evidence })
|
||||
}
|
||||
}
|
||||
const records = (): AgentChildWorkRecord[] => store.getChildren(parent)
|
||||
const byKind = (kind: AgentChildWorkRecord['kind']) =>
|
||||
records().filter((record) => record.kind === kind)
|
||||
const display = (childWorkId: string) => {
|
||||
const children = records()
|
||||
const views = projectAgentChildWorkViews(
|
||||
children,
|
||||
children.flatMap((child) => store.getAliasesForChild(child.childWorkId))
|
||||
)
|
||||
const view = views.find((candidate) => candidate.id === childWorkId)
|
||||
return view && deriveAgentChildDisplayState(view, agentChildWorkOwnedLiveness(views, view.id))
|
||||
}
|
||||
return { tracker, store, send, records, byKind, display, log }
|
||||
}
|
||||
|
||||
/** A child spawned by the parent turn and running its first turn. */
|
||||
function runningChild() {
|
||||
const run = harness()
|
||||
run.send(turn('turn/started', PRIMARY, PARENT_TURN), spawned(), turn('turn/started', CHILD, 'c1'))
|
||||
return run
|
||||
}
|
||||
|
||||
describe('Codex child-work evidence', () => {
|
||||
it('records a spawned child by its thread, with its turn as the run', () => {
|
||||
const { records, store, log, send } = runningChild()
|
||||
// Codex delivers the announcement a second time, on `item/completed`.
|
||||
send({ ...spawned(), method: 'item/completed' })
|
||||
expect(records()).toEqual([
|
||||
expect.objectContaining({
|
||||
kind: 'agent',
|
||||
membership: 'live',
|
||||
state: 'working',
|
||||
residency: 'background',
|
||||
description: 'audit_build',
|
||||
invocation: { invocationId: 'c1', generation: 1 },
|
||||
stoppable: false
|
||||
})
|
||||
])
|
||||
const aliases = store.getAliasesForChild(records()[0]!.childWorkId)
|
||||
expect(aliases.map(({ aliasKind, alias }) => [aliasKind, alias])).toEqual([
|
||||
['thread_id', CHILD],
|
||||
['turn_id', 'c1']
|
||||
])
|
||||
// The host hears the child once.
|
||||
expect(log.flat().filter((edge) => edge.type === 'live')).toHaveLength(1)
|
||||
})
|
||||
|
||||
it('makes no record for a child whose turn began before its announcement, until it lands', () => {
|
||||
const { send, records } = harness()
|
||||
send(turn('turn/started', PRIMARY, PARENT_TURN), turn('turn/started', CHILD, 'c1'))
|
||||
expect(records()).toEqual([])
|
||||
send(spawned())
|
||||
expect(records()).toEqual([expect.objectContaining({ membership: 'live', state: 'working' })])
|
||||
})
|
||||
|
||||
it.each([
|
||||
['completed', 'succeeded'],
|
||||
['interrupted', 'cancelled'],
|
||||
['failed', 'failed'],
|
||||
['somethingNew', 'unknown']
|
||||
])('settles a child whose own turn ended %s as %s', (status, outcome) => {
|
||||
const { send, tracker, records } = runningChild()
|
||||
send(turn('turn/completed', CHILD, 'c1', status))
|
||||
expect(records()).toEqual([
|
||||
expect.objectContaining({ membership: 'settled', state: 'done', outcome })
|
||||
])
|
||||
// Today's strip drops the child the moment its turn ends; only the record keeps its ending.
|
||||
expect(tracker.state).toBeNull()
|
||||
})
|
||||
|
||||
const childError = (
|
||||
turnId: string | undefined,
|
||||
willRetry: boolean
|
||||
): CodexBackgroundTaskEvent => ({
|
||||
method: 'error',
|
||||
threadId: CHILD,
|
||||
params: {
|
||||
threadId: CHILD,
|
||||
...(turnId ? { turnId } : {}),
|
||||
willRetry,
|
||||
error: { message: 'boom' }
|
||||
}
|
||||
})
|
||||
const childClosed: CodexBackgroundTaskEvent = {
|
||||
method: 'thread/closed',
|
||||
threadId: CHILD,
|
||||
params: { threadId: CHILD }
|
||||
}
|
||||
|
||||
it.each([
|
||||
['an error naming its turn that Codex will not retry', childError('c1', false), 'failed'],
|
||||
['an error naming no turn that Codex will not retry', childError(undefined, false), 'failed'],
|
||||
['its thread closing', childClosed, 'unknown']
|
||||
])(
|
||||
'settles a working child whose turn ended with no turn/completed, by %s, in the strip and the record together',
|
||||
(_label, ending, outcome) => {
|
||||
const { send, tracker, records } = runningChild()
|
||||
expect(tracker.state?.tasks).toHaveLength(1)
|
||||
send(ending)
|
||||
expect(records()).toEqual([
|
||||
expect.objectContaining({ membership: 'settled', state: 'done', outcome })
|
||||
])
|
||||
expect(tracker.state).toBeNull()
|
||||
// The first ending a turn gets stands.
|
||||
send(turn('turn/completed', CHILD, 'c1', 'completed'))
|
||||
expect(records()).toEqual([expect.objectContaining({ outcome })])
|
||||
}
|
||||
)
|
||||
|
||||
it('keeps a child working through a retried error and a systemError status: its turn runs on', () => {
|
||||
const { send, tracker, records } = runningChild()
|
||||
send(childError('c1', true), {
|
||||
method: 'thread/status/changed',
|
||||
threadId: CHILD,
|
||||
params: { threadId: CHILD, status: { type: 'systemError' } }
|
||||
})
|
||||
expect(records()).toEqual([expect.objectContaining({ membership: 'live', state: 'working' })])
|
||||
expect(tracker.state?.tasks).toHaveLength(1)
|
||||
// A fatal error naming a turn the child already finished ends nothing.
|
||||
send(turn('turn/completed', CHILD, 'c1'), childError('c1', false))
|
||||
expect(records()).toEqual([expect.objectContaining({ outcome: 'succeeded' })])
|
||||
})
|
||||
|
||||
it('never settles a child on its PARENT turn ending: children outlive the turn', () => {
|
||||
const { send, records } = runningChild()
|
||||
send(turn('turn/completed', PRIMARY, PARENT_TURN))
|
||||
expect(records()).toEqual([expect.objectContaining({ membership: 'live', state: 'working' })])
|
||||
})
|
||||
|
||||
it('reopens the same record for a follow-up turn on a finished child, as a new run', () => {
|
||||
const { send, records } = runningChild()
|
||||
send(turn('turn/completed', CHILD, 'c1'))
|
||||
const [finished] = records()
|
||||
send(turn('turn/started', CHILD, 'c2'))
|
||||
expect(records()).toEqual([
|
||||
expect.objectContaining({
|
||||
childWorkId: finished!.childWorkId,
|
||||
membership: 'live',
|
||||
state: 'working',
|
||||
invocation: { invocationId: 'c2', generation: 2 },
|
||||
previousInvocations: [
|
||||
expect.objectContaining({
|
||||
fence: { invocationId: 'c1', generation: 1 },
|
||||
outcome: 'succeeded'
|
||||
})
|
||||
]
|
||||
})
|
||||
])
|
||||
// A late ending of the first run neither ends nor restarts the second.
|
||||
send(turn('turn/completed', CHILD, 'c1', 'failed'))
|
||||
expect(records()).toEqual([
|
||||
expect.objectContaining({
|
||||
membership: 'live',
|
||||
invocation: { invocationId: 'c2', generation: 2 }
|
||||
})
|
||||
])
|
||||
send(turn('turn/completed', CHILD, 'c2', 'interrupted'))
|
||||
expect(records()).toEqual([
|
||||
expect.objectContaining({
|
||||
childWorkId: finished!.childWorkId,
|
||||
membership: 'settled',
|
||||
outcome: 'cancelled'
|
||||
})
|
||||
])
|
||||
})
|
||||
|
||||
it('says which tool the child has open, the way a CLI row names a Codex shell', () => {
|
||||
const { send, byKind } = runningChild()
|
||||
send(item('item/started', CHILD, 'c1', shell('cmd-1', 'npm test')))
|
||||
expect(byKind('agent')[0]?.operation).toEqual({
|
||||
toolName: 'Bash',
|
||||
input: 'npm test',
|
||||
basis: 'open',
|
||||
observedAt: 1_040
|
||||
})
|
||||
send(
|
||||
item('item/started', CHILD, 'c1', {
|
||||
type: 'mcpToolCall',
|
||||
id: 'mcp-1',
|
||||
server: 'github',
|
||||
tool: 'search_issues',
|
||||
arguments: { query: 'flaky' },
|
||||
status: 'inProgress'
|
||||
})
|
||||
)
|
||||
expect(byKind('agent')[0]?.operation).toMatchObject({
|
||||
toolName: 'mcp__github__search_issues',
|
||||
input: 'flaky'
|
||||
})
|
||||
// The newer call ends first: the child is still running the older one, since it opened.
|
||||
send(
|
||||
item('item/completed', CHILD, 'c1', { type: 'mcpToolCall', id: 'mcp-1', status: 'completed' })
|
||||
)
|
||||
expect(byKind('agent')[0]?.operation).toEqual({
|
||||
toolName: 'Bash',
|
||||
input: 'npm test',
|
||||
basis: 'open',
|
||||
observedAt: 1_040
|
||||
})
|
||||
send(item('item/completed', CHILD, 'c1', shell('cmd-1', 'npm test', 'completed')))
|
||||
expect(byKind('agent')[0]?.operation).toBeUndefined()
|
||||
})
|
||||
|
||||
it('names a unified-exec shell as the open call until its process exits', () => {
|
||||
const { send, byKind } = runningChild()
|
||||
// Codex runs every agent shell through unified exec, not only the ones that outlive a turn.
|
||||
send(
|
||||
item(
|
||||
'item/started',
|
||||
CHILD,
|
||||
'c1',
|
||||
shell('exec-1', 'npm test', 'inProgress', 'unifiedExecStartup')
|
||||
)
|
||||
)
|
||||
expect(byKind('agent')[0]?.operation).toMatchObject({ toolName: 'Bash', input: 'npm test' })
|
||||
send(
|
||||
item(
|
||||
'item/completed',
|
||||
CHILD,
|
||||
'c1',
|
||||
shell('exec-1', 'npm test', 'completed', 'unifiedExecStartup')
|
||||
)
|
||||
)
|
||||
expect(byKind('agent')[0]?.operation).toBeUndefined()
|
||||
// An approved command starts on the approval path and completes from unified exec.
|
||||
send(item('item/started', CHILD, 'c1', shell('exec-2', 'touch ~/marker')))
|
||||
expect(byKind('agent')[0]?.operation).toMatchObject({
|
||||
toolName: 'Bash',
|
||||
input: 'touch ~/marker'
|
||||
})
|
||||
send(
|
||||
item(
|
||||
'item/completed',
|
||||
CHILD,
|
||||
'c1',
|
||||
shell('exec-2', 'touch ~/marker', 'completed', 'unifiedExecStartup')
|
||||
)
|
||||
)
|
||||
expect(byKind('agent')[0]?.operation).toBeUndefined()
|
||||
})
|
||||
|
||||
it("never carries a run's open call into the next run when its ending was lost", () => {
|
||||
const { send, byKind } = runningChild()
|
||||
send(item('item/started', CHILD, 'c1', shell('cmd-1', 'npm test')))
|
||||
send(turn('turn/started', CHILD, 'c2'))
|
||||
expect(byKind('agent')[0]).toMatchObject({ invocation: { invocationId: 'c2', generation: 2 } })
|
||||
expect(byKind('agent')[0]?.operation).toBeUndefined()
|
||||
})
|
||||
|
||||
it('keeps what the child said last, and its usage, through to how it ended', () => {
|
||||
const { send, byKind } = runningChild()
|
||||
send(
|
||||
item('item/completed', CHILD, 'c1', {
|
||||
type: 'agentMessage',
|
||||
id: 'msg-1',
|
||||
text: 'Two tests\nflake on CI'
|
||||
}),
|
||||
{
|
||||
method: 'thread/tokenUsage/updated',
|
||||
threadId: CHILD,
|
||||
params: { threadId: CHILD, tokenUsage: { total: { totalTokens: 4_200 } } }
|
||||
}
|
||||
)
|
||||
expect(byKind('agent')[0]).toMatchObject({
|
||||
lastMessage: 'Two tests flake on CI',
|
||||
totalTokens: 4_200
|
||||
})
|
||||
send(turn('turn/completed', CHILD, 'c1'))
|
||||
expect(byKind('agent')[0]).toMatchObject({
|
||||
membership: 'settled',
|
||||
outcome: 'succeeded',
|
||||
lastMessage: 'Two tests flake on CI',
|
||||
totalTokens: 4_200
|
||||
})
|
||||
// A new run has said nothing yet.
|
||||
send(turn('turn/started', CHILD, 'c2'))
|
||||
expect(byKind('agent')[0]).not.toHaveProperty('lastMessage')
|
||||
})
|
||||
|
||||
it('files a message whose frame names no turn under the run that said it, never the next', () => {
|
||||
const { send, byKind } = runningChild()
|
||||
send({
|
||||
method: 'item/completed',
|
||||
threadId: CHILD,
|
||||
params: { threadId: CHILD, item: { type: 'agentMessage', id: 'msg-1', text: 'Done' } }
|
||||
})
|
||||
expect(byKind('agent')[0]?.lastMessage).toBe('Done')
|
||||
send(turn('turn/completed', CHILD, 'c1'), turn('turn/started', CHILD, 'c2'))
|
||||
send(turn('turn/completed', CHILD, 'c2'))
|
||||
expect(byKind('agent')[0]).toMatchObject({ outcome: 'succeeded' })
|
||||
expect(byKind('agent')[0]).not.toHaveProperty('lastMessage')
|
||||
})
|
||||
|
||||
it('reads a child waiting on the user from its own thread status', () => {
|
||||
const { send, byKind } = runningChild()
|
||||
const status = (status: unknown): CodexBackgroundTaskEvent => ({
|
||||
method: 'thread/status/changed',
|
||||
threadId: CHILD,
|
||||
params: { threadId: CHILD, status }
|
||||
})
|
||||
send(status({ type: 'active', activeFlags: ['waitingOnApproval'] }))
|
||||
expect(byKind('agent')[0]?.state).toBe('waiting')
|
||||
send(status({ type: 'active', activeFlags: [] }))
|
||||
expect(byKind('agent')[0]?.state).toBe('working')
|
||||
send(status({ type: 'active', activeFlags: ['waitingOnUserInput'] }))
|
||||
send(turn('turn/completed', CHILD, 'c1'))
|
||||
send(turn('turn/started', CHILD, 'c2'))
|
||||
// The wait ended with the turn that asked.
|
||||
expect(byKind('agent')[0]?.state).toBe('working')
|
||||
})
|
||||
|
||||
it('records a command from its start until its process exits, then removes it', () => {
|
||||
const { send, byKind, display } = runningChild()
|
||||
send(
|
||||
item(
|
||||
'item/started',
|
||||
CHILD,
|
||||
'c1',
|
||||
shell('exec-1', 'npm run dev', 'inProgress', 'unifiedExecStartup')
|
||||
)
|
||||
)
|
||||
const [agent] = byKind('agent')
|
||||
// While the child's turn runs, the command is also the tool it has open.
|
||||
expect(agent?.operation).toMatchObject({ toolName: 'Bash', input: 'npm run dev' })
|
||||
expect(byKind('command')).toEqual([
|
||||
expect.objectContaining({
|
||||
membership: 'live',
|
||||
description: 'npm run dev',
|
||||
residency: 'background',
|
||||
parentChildWorkId: agent!.childWorkId,
|
||||
firstObservedAt: 1_040
|
||||
})
|
||||
])
|
||||
send(turn('turn/completed', CHILD, 'c1'))
|
||||
expect(byKind('agent')[0]).toMatchObject({ membership: 'settled', outcome: 'succeeded' })
|
||||
expect(byKind('command')).toEqual([expect.objectContaining({ membership: 'live' })])
|
||||
expect(display(agent!.childWorkId)).toBe('monitoring')
|
||||
send(
|
||||
item('item/completed', CHILD, 'c1', {
|
||||
...shell('exec-1', 'npm run dev', 'completed', 'unifiedExecStartup'),
|
||||
exitCode: 1
|
||||
})
|
||||
)
|
||||
expect(byKind('command')).toEqual([])
|
||||
expect(display(agent!.childWorkId)).toBe('done')
|
||||
})
|
||||
|
||||
it('records an approved command while it runs, whatever source Codex starts it with', () => {
|
||||
const { send, tracker, byKind, display } = runningChild()
|
||||
// The approval path starts the item as `agent`; unified exec reports its exit.
|
||||
send(item('item/started', CHILD, 'c1', shell('exec-2', 'npm run dev')))
|
||||
const [agent] = byKind('agent')
|
||||
expect(byKind('command')).toEqual([
|
||||
expect.objectContaining({ membership: 'live', parentChildWorkId: agent!.childWorkId })
|
||||
])
|
||||
send(turn('turn/completed', CHILD, 'c1'))
|
||||
expect(display(agent!.childWorkId)).toBe('monitoring')
|
||||
expect(tracker.state?.tasks).toEqual([expect.objectContaining({ kind: 'command' })])
|
||||
send(
|
||||
item(
|
||||
'item/completed',
|
||||
CHILD,
|
||||
'c1',
|
||||
shell('exec-2', 'npm run dev', 'completed', 'unifiedExecStartup')
|
||||
)
|
||||
)
|
||||
expect(byKind('command')).toEqual([])
|
||||
expect(display(agent!.childWorkId)).toBe('done')
|
||||
expect(tracker.state).toBeNull()
|
||||
})
|
||||
|
||||
it("removes a closed thread's running commands: Codex stops them and never reports their exit", () => {
|
||||
const { send, tracker, byKind, display } = runningChild()
|
||||
send(
|
||||
item(
|
||||
'item/started',
|
||||
CHILD,
|
||||
'c1',
|
||||
shell('exec-1', 'npm run dev', 'inProgress', 'unifiedExecStartup')
|
||||
),
|
||||
turn('turn/completed', CHILD, 'c1'),
|
||||
turn('turn/completed', PRIMARY, PARENT_TURN)
|
||||
)
|
||||
const [agent] = byKind('agent')
|
||||
expect(display(agent!.childWorkId)).toBe('monitoring')
|
||||
send({ method: 'thread/closed', threadId: CHILD, params: { threadId: CHILD } })
|
||||
expect(byKind('command')).toEqual([])
|
||||
expect(display(agent!.childWorkId)).toBe('done')
|
||||
expect(tracker.state).toBeNull()
|
||||
// The exit Codex could not deliver starts nothing if it ever arrives.
|
||||
send(
|
||||
item(
|
||||
'item/completed',
|
||||
CHILD,
|
||||
'c1',
|
||||
shell('exec-1', 'npm run dev', 'completed', 'unifiedExecStartup')
|
||||
)
|
||||
)
|
||||
expect(byKind('command')).toEqual([])
|
||||
})
|
||||
|
||||
it("names the owner of a command launched before the host held its child's record", () => {
|
||||
const { send, byKind } = harness()
|
||||
send(
|
||||
turn('turn/started', PRIMARY, PARENT_TURN),
|
||||
turn('turn/started', CHILD, 'c1'),
|
||||
item(
|
||||
'item/started',
|
||||
CHILD,
|
||||
'c1',
|
||||
shell('exec-1', 'tail -f log', 'inProgress', 'unifiedExecStartup')
|
||||
)
|
||||
)
|
||||
expect(byKind('command')).toEqual([
|
||||
expect.not.objectContaining({ parentChildWorkId: expect.anything() })
|
||||
])
|
||||
send(spawned())
|
||||
expect(byKind('command')).toEqual([
|
||||
expect.objectContaining({
|
||||
description: 'tail -f log',
|
||||
parentChildWorkId: byKind('agent')[0]?.childWorkId
|
||||
})
|
||||
])
|
||||
})
|
||||
|
||||
it("records the session's own command with no owner from its start", () => {
|
||||
const { send, byKind } = harness()
|
||||
send(
|
||||
turn('turn/started', PRIMARY, PARENT_TURN),
|
||||
item(
|
||||
'item/started',
|
||||
PRIMARY,
|
||||
PARENT_TURN,
|
||||
shell('exec-9', 'sleep 90', 'inProgress', 'unifiedExecStartup')
|
||||
)
|
||||
)
|
||||
expect(byKind('command')).toEqual([
|
||||
expect.objectContaining({ membership: 'live', description: 'sleep 90' })
|
||||
])
|
||||
expect(byKind('command')[0]).not.toHaveProperty('parentChildWorkId')
|
||||
send(turn('turn/completed', PRIMARY, PARENT_TURN))
|
||||
expect(byKind('command')).toEqual([expect.objectContaining({ membership: 'live' })])
|
||||
})
|
||||
|
||||
it('leaves no record behind a finished shell, so none displaces a finished child', () => {
|
||||
const { send, byKind, records, log } = runningChild()
|
||||
send(turn('turn/completed', CHILD, 'c1'))
|
||||
const before = log.length
|
||||
// Codex runs every shell, however short, the way it runs one left running.
|
||||
for (let index = 0; index < 40; index += 1) {
|
||||
const id = `exec-${index}`
|
||||
send(
|
||||
item('item/started', PRIMARY, PARENT_TURN, {
|
||||
...shell(id, 'rg foo', 'inProgress', 'unifiedExecStartup'),
|
||||
durationMs: 0
|
||||
})
|
||||
)
|
||||
expect(byKind('command')).toEqual([expect.objectContaining({ description: 'rg foo' })])
|
||||
send(
|
||||
item('item/completed', PRIMARY, PARENT_TURN, {
|
||||
...shell(id, 'rg foo', 'completed', 'unifiedExecStartup'),
|
||||
exitCode: 0
|
||||
})
|
||||
)
|
||||
expect(byKind('command')).toEqual([])
|
||||
}
|
||||
send(turn('turn/completed', PRIMARY, PARENT_TURN))
|
||||
expect(records()).toEqual([
|
||||
expect.objectContaining({ kind: 'agent', membership: 'settled', outcome: 'succeeded' })
|
||||
])
|
||||
// One edge when each shell starts and one when it exits, as the strip republishes today.
|
||||
expect(
|
||||
log
|
||||
.slice(before)
|
||||
.flat()
|
||||
.map((edge) => edge.type)
|
||||
).toEqual(Array.from({ length: 40 }, () => ['live', 'removed']).flat())
|
||||
})
|
||||
|
||||
it('names the child that spawned a nested child as its owner', () => {
|
||||
const { send, byKind } = runningChild()
|
||||
send(
|
||||
spawned('thread-grandchild', CHILD, 'lint'),
|
||||
turn('turn/started', 'thread-grandchild', 'g1')
|
||||
)
|
||||
const nested = byKind('agent').find((record) => record.description === 'lint')
|
||||
const owner = byKind('agent').find((record) => record.description === 'audit_build')
|
||||
expect(nested?.parentChildWorkId).toBe(owner?.childWorkId)
|
||||
})
|
||||
|
||||
it('holds no evidence for a session with nowhere to deliver it', () => {
|
||||
const tracker = new CodexBackgroundTaskTracker(PRIMARY)
|
||||
tracker.observe(spawned())
|
||||
tracker.observe(turn('turn/started', CHILD, 'c1'))
|
||||
tracker.publishChildWork()
|
||||
expect(tracker.drainChildWorkEvidence(1)).toEqual([])
|
||||
})
|
||||
|
||||
it('settles every live child with no reported outcome, and removes every command, when the provider session ends', () => {
|
||||
const { send, tracker, records, store } = runningChild()
|
||||
send(
|
||||
item(
|
||||
'item/started',
|
||||
PRIMARY,
|
||||
PARENT_TURN,
|
||||
shell('exec-1', 'npm run dev', 'inProgress', 'unifiedExecStartup')
|
||||
),
|
||||
turn('turn/completed', PRIMARY, PARENT_TURN)
|
||||
)
|
||||
expect(records()).toHaveLength(2)
|
||||
tracker.clear()
|
||||
const evidence = tracker.drainChildWorkEvidence(9_000)
|
||||
expect(evidence).toEqual([
|
||||
{
|
||||
type: 'removed',
|
||||
observedAt: 9_000,
|
||||
handle: { idKind: 'task_id', id: 'codex-command:primary:exec-1' }
|
||||
},
|
||||
{ type: 'session-ended', observedAt: 9_000 }
|
||||
])
|
||||
reconcileAgentChildWorkEvidence({
|
||||
store,
|
||||
admission: createAgentChildWorkAdmission(store, { mintChildWorkId: () => 'unused' }),
|
||||
parent,
|
||||
provider: 'codex',
|
||||
evidence
|
||||
})
|
||||
expect(records()).toEqual([
|
||||
expect.objectContaining({ kind: 'agent', membership: 'settled', outcome: 'unknown' })
|
||||
])
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,312 @@
|
||||
// Codex child threads, and the commands they leave running, decoded into child-work evidence for
|
||||
// the host's records.
|
||||
//
|
||||
// The background-task tracker already follows which child exists, which turn it runs and how that
|
||||
// turn ended (the executions), and which command process is still running (the command tracker).
|
||||
// A command is a record from its process start until it stops, and then its record goes: the
|
||||
// tracker says when, and this module only mirrors it. This module keeps what only the records
|
||||
// read — the tool a child has open, what it said last, its usage, whether it waits on the user —
|
||||
// and after each frame re-derives the whole observation of the child that frame was about. Edges
|
||||
// are stamped with the host clock when drained, after the journal handled the frame, so the host
|
||||
// never holds a record ahead of the frame's own rows. A parent turn ending is never evidence here:
|
||||
// Codex children outlive the turn that spawned them, so only a child's own turn, or the session,
|
||||
// ends it.
|
||||
|
||||
import type { AgentSessionBackgroundTask } from '../../shared/agent-session-wire'
|
||||
import type {
|
||||
AgentChildWorkEvidence,
|
||||
AgentChildWorkLiveObservation
|
||||
} from '../../shared/agent-status-child-work-evidence'
|
||||
import type { CodexBackgroundCommandChange } from './codex-background-command-tracker'
|
||||
import type {
|
||||
CodexBackgroundTaskEvent,
|
||||
CodexBackgroundTaskFrame
|
||||
} from './codex-background-task-frames'
|
||||
import {
|
||||
codexChildMessageText,
|
||||
codexChildToolCall,
|
||||
codexChildTurnOutcome,
|
||||
codexToolCallEnded,
|
||||
type CodexChildToolCall
|
||||
} from './codex-child-work-translation'
|
||||
import { readRecord } from './codex-item-field-readers'
|
||||
import { readCodexThreadItem } from './codex-structured-item-translation'
|
||||
import { codexThreadWaitsOnUser, readCodexTurnId } from './codex-structured-thread-facts'
|
||||
import { CODEX_TOKEN_USAGE_METHOD, readCodexThreadTokenTotal } from './codex-subagent-activity'
|
||||
import type { CodexExecutionChild, CodexSubagentExecutions } from './codex-subagent-executions'
|
||||
|
||||
/** The executions' own child bound. */
|
||||
const MAX_CHILD_FACTS = 128
|
||||
const MAX_OPEN_CALLS_PER_CHILD = 16
|
||||
const CHILD_FRAME_METHODS: ReadonlySet<string> = new Set([
|
||||
'item/started',
|
||||
'item/completed',
|
||||
'thread/status/changed',
|
||||
CODEX_TOKEN_USAGE_METHOD
|
||||
])
|
||||
|
||||
/** A tool call a child has started and not finished. `openedAt` is the host clock of the first
|
||||
* drain that carried it, so a later edge keeps the time the call opened. */
|
||||
type OpenCall = CodexChildToolCall & { turnId: string | null; openedAt?: number }
|
||||
|
||||
type ChildFacts = {
|
||||
openCalls: Map<string, OpenCall>
|
||||
lastMessage?: { turnId: string | null; text: string }
|
||||
totalTokens?: number
|
||||
waiting: boolean
|
||||
/** The last observation handed to the host, so an unchanged re-derivation sends nothing. */
|
||||
published?: string
|
||||
}
|
||||
|
||||
export type CodexPendingChildWork = (observedAt: number) => AgentChildWorkEvidence
|
||||
|
||||
/** Evidence from a run only counts for that run: a fact recorded under another turn is stale. */
|
||||
function ofTurn<T extends { turnId: string | null }>(fact: T | undefined, turnId: string) {
|
||||
return fact?.turnId === turnId ? fact : undefined
|
||||
}
|
||||
|
||||
function commandLive(task: AgentSessionBackgroundTask, ownerId: string | null) {
|
||||
return (observedAt: number): AgentChildWorkEvidence => ({
|
||||
type: 'live',
|
||||
observedAt,
|
||||
child: {
|
||||
handle: { idKind: 'task_id', id: task.id },
|
||||
kind: 'command',
|
||||
// Its own process, not a turn's: no turn ending may settle it.
|
||||
residency: 'background',
|
||||
state: 'working',
|
||||
...(task.description ? { description: task.description } : {}),
|
||||
...(ownerId !== null ? { ownerId } : {}),
|
||||
stoppable: false
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
export class CodexChildWorkEvidence {
|
||||
private readonly facts = new Map<string, ChildFacts>()
|
||||
private pending: CodexPendingChildWork[] = []
|
||||
|
||||
constructor(
|
||||
private readonly primaryThreadId: string,
|
||||
private readonly executions: CodexSubagentExecutions,
|
||||
private readonly liveCommands: (threadId: string) => readonly AgentSessionBackgroundTask[]
|
||||
) {}
|
||||
|
||||
/** After the tracker applied the frame: which command processes it saw start or stop, and the
|
||||
* child the frame is about. */
|
||||
observe(
|
||||
event: CodexBackgroundTaskEvent,
|
||||
frame: CodexBackgroundTaskFrame | null,
|
||||
commands: readonly CodexBackgroundCommandChange[]
|
||||
): void {
|
||||
this.queueCommands(commands)
|
||||
const threadId = this.childThread(event, frame)
|
||||
if (threadId === null) {
|
||||
return
|
||||
}
|
||||
const facts = this.factsFor(threadId)
|
||||
if (facts && event.threadId === threadId) {
|
||||
this.record(facts, event)
|
||||
}
|
||||
this.queueChild(threadId)
|
||||
}
|
||||
|
||||
/** The provider session is gone, with the commands it ended: no child it still ran can report
|
||||
* its own ending. */
|
||||
clear(commands: readonly CodexBackgroundCommandChange[]): void {
|
||||
this.facts.clear()
|
||||
this.queueCommands(commands)
|
||||
this.pending.push((observedAt) => ({ type: 'session-ended', observedAt }))
|
||||
}
|
||||
|
||||
drain(observedAt: number): AgentChildWorkEvidence[] {
|
||||
const pending = this.pending
|
||||
this.pending = []
|
||||
return pending.map((edge) => edge(observedAt))
|
||||
}
|
||||
|
||||
private childThread(
|
||||
event: CodexBackgroundTaskEvent,
|
||||
frame: CodexBackgroundTaskFrame | null
|
||||
): string | null {
|
||||
const threadId =
|
||||
frame?.kind === 'subagent'
|
||||
? frame.agentThreadId
|
||||
: frame || CHILD_FRAME_METHODS.has(event.method)
|
||||
? event.threadId
|
||||
: null
|
||||
return threadId === this.primaryThreadId ? null : threadId
|
||||
}
|
||||
|
||||
/** A command belongs to the child thread that launched it; the session's own agent is no owner.
|
||||
* A stopped command leaves no record: it has nothing left to report. */
|
||||
private queueCommands(commands: readonly CodexBackgroundCommandChange[]): void {
|
||||
for (const command of commands) {
|
||||
if (command.type === 'started') {
|
||||
const ownerId = command.threadId === this.primaryThreadId ? null : command.threadId
|
||||
this.pending.push(commandLive(command.task, ownerId))
|
||||
continue
|
||||
}
|
||||
const { taskId } = command
|
||||
this.pending.push((observedAt) => ({
|
||||
type: 'removed',
|
||||
observedAt,
|
||||
handle: { idKind: 'task_id', id: taskId }
|
||||
}))
|
||||
}
|
||||
}
|
||||
|
||||
private record(facts: ChildFacts, event: CodexBackgroundTaskEvent): void {
|
||||
if (event.method === CODEX_TOKEN_USAGE_METHOD) {
|
||||
facts.totalTokens = readCodexThreadTokenTotal(event.params)?.totalTokens ?? facts.totalTokens
|
||||
return
|
||||
}
|
||||
if (event.method === 'thread/status/changed') {
|
||||
facts.waiting = codexThreadWaitsOnUser(event.params)
|
||||
return
|
||||
}
|
||||
// Every Codex agent shell is unified exec: it is the open call until its process exits.
|
||||
const item = readCodexThreadItem(readRecord(event.params).item)
|
||||
if (!item) {
|
||||
return
|
||||
}
|
||||
// A frame that names no turn belongs to the one the child is running.
|
||||
const turnId =
|
||||
readCodexTurnId(event.params) ??
|
||||
this.executions.find(event.threadId)?.execution?.turnId ??
|
||||
null
|
||||
const text = event.method === 'item/completed' ? codexChildMessageText(item) : undefined
|
||||
if (text) {
|
||||
facts.lastMessage = { turnId, text }
|
||||
}
|
||||
// An end closes the call by id alone: its closing frame need not restate what it ran.
|
||||
if (codexToolCallEnded(event.method, item)) {
|
||||
facts.openCalls.delete(item.id)
|
||||
return
|
||||
}
|
||||
const call = codexChildToolCall(item)
|
||||
if (call && !facts.openCalls.has(item.id)) {
|
||||
facts.openCalls.set(item.id, { ...call, turnId })
|
||||
for (const stale of [...facts.openCalls.keys()].slice(0, -MAX_OPEN_CALLS_PER_CHILD)) {
|
||||
facts.openCalls.delete(stale)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Re-derive the child's observation and hand it on when it changed. A child the provider has
|
||||
* not announced, or that never ran a turn, is no record. */
|
||||
private queueChild(threadId: string): void {
|
||||
const child = this.executions.find(threadId)
|
||||
const facts = this.facts.get(threadId)
|
||||
if (!child?.registered || !child.execution || !facts) {
|
||||
return
|
||||
}
|
||||
const { turnId, state } = child.execution
|
||||
if (state !== 'working') {
|
||||
// The turn is over, and so is every call it had open.
|
||||
facts.openCalls.clear()
|
||||
facts.waiting = false
|
||||
const lastMessage = ofTurn(facts.lastMessage, turnId)?.text
|
||||
const { totalTokens } = facts
|
||||
const outcome = codexChildTurnOutcome(state)
|
||||
this.publish(facts, JSON.stringify(['ended', turnId, state]), (observedAt) => ({
|
||||
type: 'ended',
|
||||
observedAt,
|
||||
handle: { idKind: 'thread_id', id: threadId, runId: turnId },
|
||||
outcome,
|
||||
...(lastMessage ? { lastMessage } : {}),
|
||||
...(totalTokens !== undefined ? { totalTokens } : {})
|
||||
}))
|
||||
return
|
||||
}
|
||||
for (const [itemId, call] of facts.openCalls) {
|
||||
if (!ofTurn(call, turnId)) {
|
||||
facts.openCalls.delete(itemId)
|
||||
}
|
||||
}
|
||||
const observation = this.liveAgent(threadId, child, facts, turnId)
|
||||
const openCall = [...facts.openCalls].at(-1)
|
||||
const announced = facts.published !== undefined
|
||||
this.publish(facts, JSON.stringify(['live', observation, openCall?.[0]]), (observedAt) => {
|
||||
if (!openCall) {
|
||||
return { type: 'live', observedAt, child: { ...observation, operation: null } }
|
||||
}
|
||||
const [, call] = openCall
|
||||
call.openedAt ??= observedAt
|
||||
const operation = {
|
||||
toolName: call.toolName,
|
||||
...(call.input ? { input: call.input } : {}),
|
||||
basis: 'open' as const,
|
||||
observedAt: call.openedAt
|
||||
}
|
||||
return { type: 'live', observedAt, child: { ...observation, operation } }
|
||||
})
|
||||
if (!announced) {
|
||||
this.requeueOwnedBy(threadId)
|
||||
}
|
||||
}
|
||||
|
||||
private liveAgent(
|
||||
threadId: string,
|
||||
child: Readonly<CodexExecutionChild>,
|
||||
facts: ChildFacts,
|
||||
turnId: string
|
||||
): AgentChildWorkLiveObservation {
|
||||
const lastMessage = ofTurn(facts.lastMessage, turnId)?.text
|
||||
const spawner = child.spawnerThreadId
|
||||
return {
|
||||
handle: { idKind: 'thread_id', id: threadId, runId: turnId },
|
||||
kind: 'agent',
|
||||
// A spawned child may outlive the turn that spawned it.
|
||||
residency: 'background',
|
||||
state: facts.waiting ? 'waiting' : 'working',
|
||||
// The agent path's last segment is the child's only label; today's row shows it there.
|
||||
...(child.label ? { description: child.label } : {}),
|
||||
...(facts.totalTokens !== undefined ? { totalTokens: facts.totalTokens } : {}),
|
||||
...(lastMessage ? { lastMessage } : {}),
|
||||
...(spawner && spawner !== this.primaryThreadId ? { ownerId: spawner } : {}),
|
||||
stoppable: false
|
||||
}
|
||||
}
|
||||
|
||||
private publish(facts: ChildFacts, fingerprint: string, edge: CodexPendingChildWork): void {
|
||||
if (facts.published !== fingerprint) {
|
||||
facts.published = fingerprint
|
||||
this.pending.push(edge)
|
||||
}
|
||||
}
|
||||
|
||||
/** Work a child launched before the host held its record was admitted with no owner; now that
|
||||
* the owner is recorded, say again whose it is. */
|
||||
private requeueOwnedBy(threadId: string): void {
|
||||
for (const task of this.liveCommands(threadId)) {
|
||||
this.pending.push(commandLive(task, threadId))
|
||||
}
|
||||
for (const spawned of this.executions.workingChildren()) {
|
||||
const facts = this.facts.get(spawned.agentThreadId)
|
||||
if (spawned.spawnerThreadId === threadId && facts?.published !== undefined) {
|
||||
facts.published = undefined
|
||||
this.queueChild(spawned.agentThreadId)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private factsFor(threadId: string): ChildFacts | undefined {
|
||||
const existing = this.facts.get(threadId)
|
||||
if (existing) {
|
||||
return existing
|
||||
}
|
||||
if (this.facts.size >= MAX_CHILD_FACTS) {
|
||||
const idle = [...this.facts.keys()].find(
|
||||
(id) => this.executions.find(id)?.execution?.state !== 'working'
|
||||
)
|
||||
if (idle === undefined) {
|
||||
return undefined
|
||||
}
|
||||
this.facts.delete(idle)
|
||||
}
|
||||
const facts: ChildFacts = { openCalls: new Map(), waiting: false }
|
||||
this.facts.set(threadId, facts)
|
||||
return facts
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,88 @@
|
||||
// Codex items and statuses, read in the child-work vocabulary.
|
||||
//
|
||||
// A tool is named the way Codex names it to its own hooks (`Bash`, `apply_patch`,
|
||||
// `mcp__server__tool`), so a structured Codex child running a shell reads exactly as a Codex
|
||||
// CLI agent running one does.
|
||||
|
||||
import type { AgentChildWorkOutcome } from '../../shared/agent-status-child-work'
|
||||
import type { NativeChatSubagentState } from '../../shared/native-chat-types'
|
||||
import {
|
||||
deriveFallbackToolInputPreview,
|
||||
deriveToolInputPreview
|
||||
} from '../../shared/agent-hook-listener/tool-input-preview'
|
||||
import { readRecord, readString, readTextContent } from './codex-item-field-readers'
|
||||
import type { CodexThreadItem } from './codex-structured-item-translation'
|
||||
|
||||
/** Raw provider text kept for a record; admission folds it to its own one-line bound. */
|
||||
export const CODEX_CHILD_WORK_TEXT_MAX_CHARS = 2_048
|
||||
|
||||
export type CodexChildToolCall = { toolName: string; input?: string }
|
||||
|
||||
function bounded(text: string | null | undefined): string | undefined {
|
||||
return text ? text.slice(0, CODEX_CHILD_WORK_TEXT_MAX_CHARS) : undefined
|
||||
}
|
||||
|
||||
function withInput(toolName: string, input: string | undefined): CodexChildToolCall {
|
||||
const preview = bounded(input)
|
||||
return preview ? { toolName, input: preview } : { toolName }
|
||||
}
|
||||
|
||||
function firstChangePath(changes: unknown): string | undefined {
|
||||
const [first] = Array.isArray(changes) ? changes : []
|
||||
return readString(readRecord(first), 'path') ?? undefined
|
||||
}
|
||||
|
||||
/** The tool a thread item runs, or null for an item that is not a tool call (a message, a
|
||||
* thought, a plan). */
|
||||
export function codexChildToolCall(item: CodexThreadItem): CodexChildToolCall | null {
|
||||
switch (item.type) {
|
||||
case 'commandExecution':
|
||||
return withInput('Bash', deriveToolInputPreview('Bash', { command: item.command }))
|
||||
case 'fileChange':
|
||||
return withInput('apply_patch', firstChangePath(item.changes))
|
||||
case 'mcpToolCall': {
|
||||
const server = readString(item, 'server')
|
||||
const tool = readString(item, 'tool')
|
||||
if (!tool) {
|
||||
return null
|
||||
}
|
||||
return withInput(
|
||||
server ? `mcp__${server}__${tool}` : tool,
|
||||
deriveFallbackToolInputPreview(item.arguments)
|
||||
)
|
||||
}
|
||||
case 'webSearch':
|
||||
return withInput('web_search', readString(item, 'query') ?? undefined)
|
||||
default:
|
||||
return null
|
||||
}
|
||||
}
|
||||
|
||||
/** Whether an item frame says the call is over, whatever frame carried it. */
|
||||
export function codexToolCallEnded(method: string, item: CodexThreadItem): boolean {
|
||||
const status = readString(item, 'status')
|
||||
return method === 'item/completed' || (status !== null && status !== 'inProgress')
|
||||
}
|
||||
|
||||
/** What a child said: an assistant message's text. */
|
||||
export function codexChildMessageText(item: CodexThreadItem): string | undefined {
|
||||
return item.type === 'agentMessage'
|
||||
? bounded(readString(item, 'text') ?? readTextContent(item, 'content'))
|
||||
: undefined
|
||||
}
|
||||
|
||||
/** A child turn's ending. Codex states three; anything else is an ending nobody classified. */
|
||||
export function codexChildTurnOutcome(state: NativeChatSubagentState): AgentChildWorkOutcome {
|
||||
switch (state) {
|
||||
case 'completed':
|
||||
return 'succeeded'
|
||||
case 'failed':
|
||||
return 'failed'
|
||||
case 'stopped':
|
||||
return 'cancelled'
|
||||
case 'unverifiable':
|
||||
case 'working':
|
||||
case 'idle':
|
||||
return 'unknown'
|
||||
}
|
||||
}
|
||||
@@ -108,3 +108,22 @@ describe('Codex prompt claim lifetime', () => {
|
||||
expect(prompt.deref()).toBeUndefined()
|
||||
})
|
||||
})
|
||||
|
||||
describe('Codex abandoned command approvals', () => {
|
||||
it("reports only a command's own approval that its turn ended unanswered, once", () => {
|
||||
const registry = new CodexPromptRegistry()
|
||||
const ask = (id: number, method: string, params: Record<string, string>) =>
|
||||
registry.register({ id, method, params: { threadId: 'thread', turnId: 'turn', ...params } })
|
||||
ask(1, 'item/commandExecution/requestApproval', { itemId: 'unanswered' })
|
||||
const answered = ask(2, 'item/commandExecution/requestApproval', { itemId: 'answered' })
|
||||
ask(3, 'item/commandExecution/requestApproval', { itemId: 'parent', approvalId: 'sub' })
|
||||
ask(4, 'item/fileChange/requestApproval', { itemId: 'patch' })
|
||||
if (!answered) {
|
||||
throw new Error('Fixture prompt was refused')
|
||||
}
|
||||
registry.forget(answered)
|
||||
registry.clearTurn('thread', 'turn')
|
||||
expect(registry.takeAbandonedCommands()).toEqual([{ threadId: 'thread', itemId: 'unanswered' }])
|
||||
expect(registry.takeAbandonedCommands()).toEqual([])
|
||||
})
|
||||
})
|
||||
|
||||
@@ -31,6 +31,9 @@ export type CodexPendingPrompt = {
|
||||
answers: Map<string, string>
|
||||
}
|
||||
|
||||
export type CodexAbandonedCommand = { threadId: string; itemId: string }
|
||||
const NO_ABANDONED_COMMANDS: readonly CodexAbandonedCommand[] = []
|
||||
|
||||
export type CodexPromptClaim = {
|
||||
readonly itemId: string
|
||||
readonly prompt: CodexPendingPrompt
|
||||
@@ -54,6 +57,7 @@ export class CodexPromptRegistry {
|
||||
private readonly journalItemIds = new Map<string, string>()
|
||||
private readonly boundPrompts = new Map<string, CodexPendingPrompt>()
|
||||
private readonly claims = new Map<CodexPendingPrompt, CodexPromptClaim>()
|
||||
private abandonedCommands: CodexAbandonedCommand[] = []
|
||||
|
||||
get sizes(): { prompts: number; journalBindings: number } {
|
||||
return { prompts: this.byAddress.size, journalBindings: this.journalItemIds.size }
|
||||
@@ -232,14 +236,33 @@ export class CodexPromptRegistry {
|
||||
)
|
||||
for (const prompt of prompts) {
|
||||
this.forget(prompt)
|
||||
// Codex abandons a turn's unanswered prompts: a command still awaiting approval never ran.
|
||||
// An `approvalId` asks for a subcommand, not the item's own command.
|
||||
if (
|
||||
prompt.method === CODEX_COMMAND_APPROVAL_METHOD &&
|
||||
prompt.promptKey === prompt.codexItemId
|
||||
) {
|
||||
this.abandonedCommands.push({ threadId: prompt.threadId, itemId: prompt.codexItemId })
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** The commands whose approval a turn ended without, since the last call. */
|
||||
takeAbandonedCommands(): readonly CodexAbandonedCommand[] {
|
||||
if (this.abandonedCommands.length === 0) {
|
||||
return NO_ABANDONED_COMMANDS
|
||||
}
|
||||
const taken = this.abandonedCommands
|
||||
this.abandonedCommands = []
|
||||
return taken
|
||||
}
|
||||
|
||||
clear(): void {
|
||||
this.byAddress.clear()
|
||||
this.journalItemIds.clear()
|
||||
this.boundPrompts.clear()
|
||||
this.claims.clear()
|
||||
this.abandonedCommands = []
|
||||
}
|
||||
|
||||
private address(threadId: string, promptKey: string): string {
|
||||
|
||||
@@ -0,0 +1,525 @@
|
||||
// A Codex session's frames, through the real adapter, into the host's child records: the order the
|
||||
// host receives them in, and whether the parent row the records imply is today's row.
|
||||
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import {
|
||||
foldAgentLeadStatus,
|
||||
type AgentLeadStatusResolution
|
||||
} from '../../shared/agent-lead-status-fold'
|
||||
import { createAgentChildWorkAdmission } from '../../shared/agent-status-child-work-admission'
|
||||
import type { AgentChildWorkRecord } from '../../shared/agent-status-child-work'
|
||||
import { agentChildWorkLiveness } from '../../shared/agent-status-child-work-liveness'
|
||||
import { reconcileAgentChildWorkEvidence } from '../../shared/agent-status-child-work-reconciliation'
|
||||
import {
|
||||
agentChildWorkOwnedLiveness,
|
||||
deriveAgentChildDisplayState,
|
||||
projectAgentChildWorkViews
|
||||
} from '../../shared/agent-status-child-work-view'
|
||||
import { createAgentStatusStore } from '../../shared/agent-status-store'
|
||||
import { agentJournalLinkageFields } from '../../shared/agent-session-journal-producer'
|
||||
import type {
|
||||
AgentJournalItemBody,
|
||||
AgentJournalProducerLinkage
|
||||
} from '../../shared/agent-session-journal-types'
|
||||
import { makeStructuredAgentStatusSubject } from '../../shared/agent-status-subject'
|
||||
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
|
||||
import { fakeCodex, identityFor, THREAD_ID } from './codex-structured-session-adapter-fixture'
|
||||
import { CodexStructuredSessionAdapter } from './codex-structured-session-adapter'
|
||||
|
||||
const parent = makeStructuredAgentStatusSubject(
|
||||
{
|
||||
executionHostId: 'local',
|
||||
wslDistro: null,
|
||||
workspaceId: 'workspace-1',
|
||||
workspaceKind: 'folder'
|
||||
},
|
||||
'session-1'
|
||||
)
|
||||
const REVIEWER = 'thread-reviewer'
|
||||
const TESTER = 'thread-tester'
|
||||
const LINTER = 'thread-linter'
|
||||
|
||||
type Frame = { method: string; params: Record<string, unknown> }
|
||||
type Delivery = { kind: 'journal' | 'legacy' | 'evidence'; detail: string }
|
||||
type Liveness = ReturnType<typeof agentChildWorkLiveness>
|
||||
/** What the records and today's strip each said when the journal wrote or published a row. */
|
||||
type JournalMoment = { recorded: Liveness; legacy: Liveness }
|
||||
|
||||
const turn = (
|
||||
method: 'turn/started' | 'turn/completed',
|
||||
threadId: string,
|
||||
id: string,
|
||||
status = 'completed'
|
||||
): Frame => ({
|
||||
method,
|
||||
params: { threadId, turn: { id, status } }
|
||||
})
|
||||
const spawned = (child: string, name: string, parentTurn: string): Frame => ({
|
||||
method: 'item/started',
|
||||
params: {
|
||||
threadId: THREAD_ID,
|
||||
turnId: parentTurn,
|
||||
item: {
|
||||
type: 'subAgentActivity',
|
||||
id: `spawn-${child}`,
|
||||
kind: 'started',
|
||||
agentThreadId: child,
|
||||
agentPath: `/root/${name}`
|
||||
}
|
||||
}
|
||||
})
|
||||
const item = (
|
||||
method: 'item/started' | 'item/completed',
|
||||
threadId: string,
|
||||
turnId: string,
|
||||
fields: Record<string, unknown>
|
||||
): Frame => ({
|
||||
method,
|
||||
params: { threadId, turnId, item: fields }
|
||||
})
|
||||
const status = (threadId: string, activeFlags: string[]): Frame => ({
|
||||
method: 'thread/status/changed',
|
||||
params: { threadId, status: { type: 'active', activeFlags } }
|
||||
})
|
||||
|
||||
async function producer() {
|
||||
const codex = fakeCodex()
|
||||
const store = createAgentStatusStore({ epoch: 'epoch-1', mode: 'authority' })
|
||||
expect(store.applyMutation({ parent: { subject: parent } })).not.toBeNull()
|
||||
let minted = 0
|
||||
const admission = createAgentChildWorkAdmission(store, {
|
||||
mintChildWorkId: () => `child-${++minted}`
|
||||
})
|
||||
const deliveries: Delivery[] = []
|
||||
const moments: JournalMoment[] = []
|
||||
const records = (): AgentChildWorkRecord[] => store.getChildren(parent)
|
||||
const recordedLiveness = () =>
|
||||
agentChildWorkLiveness(records().filter((record) => record.membership === 'live'))
|
||||
// Each journal write publishes the parent's row, so the records must imply its state right then.
|
||||
const moment = () =>
|
||||
moments.push({
|
||||
recorded: recordedLiveness(),
|
||||
legacy: agentChildWorkLiveness(adapter.backgroundTaskState('session-1')?.tasks)
|
||||
})
|
||||
const adapter = new CodexStructuredSessionAdapter({
|
||||
resolveLaunch: async () => ({
|
||||
command: 'codex',
|
||||
args: ['app-server'],
|
||||
cwd: '/work/repo',
|
||||
codexHome: null,
|
||||
resumeThreadId: null
|
||||
}),
|
||||
openConnection: codex.openConnection,
|
||||
readProcessStartTime: async () => 1_700_000_000_000,
|
||||
now: () => 1_700_000_000_500,
|
||||
onBackgroundTasksChanged: (_sessionId, state) =>
|
||||
deliveries.push({ kind: 'legacy', detail: String(state?.tasks?.length ?? 0) }),
|
||||
onChildWorkEvidence: (sessionId, evidence) => {
|
||||
expect(sessionId).toBe('session-1')
|
||||
deliveries.push({ kind: 'evidence', detail: evidence.map((edge) => edge.type).join(',') })
|
||||
reconcileAgentChildWorkEvidence({ store, admission, parent, provider: 'codex', evidence })
|
||||
}
|
||||
})
|
||||
const rows: { body: AgentJournalItemBody; linkage: AgentJournalProducerLinkage }[] = []
|
||||
const journal: StructuredAgentSessionEventSink = {
|
||||
appendItem: (identity, body, options) => {
|
||||
deliveries.push({ kind: 'journal', detail: JSON.stringify(identity) })
|
||||
rows.push({ body, linkage: agentJournalLinkageFields(options) })
|
||||
moment()
|
||||
},
|
||||
appendTombstone: () => {},
|
||||
publish: moment
|
||||
}
|
||||
await adapter.acquire({
|
||||
identity: identityFor('session-1'),
|
||||
fence: 7,
|
||||
spawnToken: 'spawn-9',
|
||||
events: journal
|
||||
})
|
||||
const send = (frame: Frame): Delivery[] => {
|
||||
const from = deliveries.length
|
||||
codex.connections[0]!.handlers.onNotification?.(frame.method, frame.params)
|
||||
return deliveries.slice(from)
|
||||
}
|
||||
/** The journal moments one frame produced. */
|
||||
const momentsOf = (frame: Frame): JournalMoment[] => {
|
||||
const from = moments.length
|
||||
send(frame)
|
||||
return moments.slice(from)
|
||||
}
|
||||
const byDescription = (description: string) =>
|
||||
records().find((record) => record.description === description)
|
||||
const display = (description: string) => {
|
||||
const children = records()
|
||||
const views = projectAgentChildWorkViews(
|
||||
children,
|
||||
children.flatMap((child) => store.getAliasesForChild(child.childWorkId))
|
||||
)
|
||||
const view = views.find((candidate) => candidate.description === description)
|
||||
return view && deriveAgentChildDisplayState(view, agentChildWorkOwnedLiveness(views, view.id))
|
||||
}
|
||||
/** The producer stamp on the newest journal row that carries this text. */
|
||||
const stampOf = (text: string) =>
|
||||
rows.findLast((row) => JSON.stringify(row.body).includes(text))?.linkage
|
||||
return {
|
||||
adapter,
|
||||
codex,
|
||||
send,
|
||||
momentsOf,
|
||||
records,
|
||||
recordedLiveness,
|
||||
byDescription,
|
||||
display,
|
||||
stampOf
|
||||
}
|
||||
}
|
||||
|
||||
const fold = (
|
||||
leadState: 'working' | 'done',
|
||||
childWorkLiveness: Liveness
|
||||
): AgentLeadStatusResolution => foldAgentLeadStatus({ leadState, childWorkLiveness })
|
||||
/** The strip has no word for a child waiting on a human: to it, that child is working. */
|
||||
const asStrip = (liveness: Liveness): Liveness => (liveness === 'waiting' ? 'working' : liveness)
|
||||
const shellFrame = (
|
||||
method: 'item/started' | 'item/completed',
|
||||
threadId: string,
|
||||
turnId: string,
|
||||
id: string,
|
||||
command: string,
|
||||
source = 'unifiedExecStartup'
|
||||
): Frame =>
|
||||
item(method, threadId, turnId, {
|
||||
type: 'commandExecution',
|
||||
id,
|
||||
command,
|
||||
source,
|
||||
status: method === 'item/started' ? 'inProgress' : 'completed',
|
||||
...(method === 'item/completed' ? { exitCode: 0 } : {})
|
||||
})
|
||||
|
||||
describe('Codex structured child-work producer', () => {
|
||||
it('delivers evidence only after the journal wrote the frame and the legacy row republished', async () => {
|
||||
const { send, records } = await producer()
|
||||
send(turn('turn/started', THREAD_ID, 'p1'))
|
||||
send(turn('turn/started', REVIEWER, 'r1'))
|
||||
const deliveries = send(spawned(REVIEWER, 'review', 'p1'))
|
||||
const kinds = deliveries.map((delivery) => delivery.kind)
|
||||
// The frame's own rows, then the parent's republished row, and only then its children.
|
||||
expect(kinds.filter((kind) => kind === 'journal').length).toBeGreaterThan(0)
|
||||
expect(kinds.slice(kinds.indexOf('legacy'))).toEqual(['legacy', 'evidence'])
|
||||
expect(records()).toEqual([
|
||||
expect.objectContaining({ description: 'review', membership: 'live' })
|
||||
])
|
||||
})
|
||||
|
||||
it('records the parent state today reads, at every journal write, while adding outcome and activity', async () => {
|
||||
const { adapter, momentsOf, records, recordedLiveness, byDescription, display } =
|
||||
await producer()
|
||||
const steps: {
|
||||
frame: Frame
|
||||
lead: 'working' | 'done'
|
||||
childWaits?: true
|
||||
check?: () => void
|
||||
}[] = [
|
||||
{ frame: turn('turn/started', THREAD_ID, 'p1'), lead: 'working' },
|
||||
// Codex reports the child's turn before its announcement.
|
||||
{
|
||||
frame: turn('turn/started', REVIEWER, 'r1'),
|
||||
lead: 'working',
|
||||
check: () => expect(records()).toEqual([])
|
||||
},
|
||||
{ frame: spawned(REVIEWER, 'review', 'p1'), lead: 'working' },
|
||||
// An approved command: Codex starts it on the approval path and reports its exit from
|
||||
// unified exec.
|
||||
{
|
||||
frame: shellFrame('item/started', REVIEWER, 'r1', 'cmd-1', 'npm test', 'agent'),
|
||||
lead: 'working',
|
||||
check: () => {
|
||||
expect(byDescription('review')?.operation).toMatchObject({
|
||||
toolName: 'Bash',
|
||||
input: 'npm test',
|
||||
basis: 'open'
|
||||
})
|
||||
expect(byDescription('npm test')).toMatchObject({
|
||||
membership: 'live',
|
||||
parentChildWorkId: byDescription('review')?.childWorkId
|
||||
})
|
||||
}
|
||||
},
|
||||
// The child starts a dev server it will leave running past its own turn.
|
||||
{
|
||||
frame: shellFrame('item/started', REVIEWER, 'r1', 'exec-1', 'npm run dev'),
|
||||
lead: 'working',
|
||||
check: () =>
|
||||
expect(byDescription('npm run dev')).toMatchObject({
|
||||
membership: 'live',
|
||||
parentChildWorkId: byDescription('review')?.childWorkId
|
||||
})
|
||||
},
|
||||
{
|
||||
frame: shellFrame('item/completed', REVIEWER, 'r1', 'cmd-1', 'npm test'),
|
||||
lead: 'working',
|
||||
check: () => {
|
||||
// A finished command leaves nothing behind.
|
||||
expect(byDescription('npm test')).toBeUndefined()
|
||||
// The dev server is still the child's open call while its turn runs.
|
||||
expect(byDescription('review')?.operation).toMatchObject({
|
||||
toolName: 'Bash',
|
||||
input: 'npm run dev'
|
||||
})
|
||||
}
|
||||
},
|
||||
{
|
||||
frame: item('item/completed', REVIEWER, 'r1', {
|
||||
type: 'agentMessage',
|
||||
id: 'msg-1',
|
||||
text: 'Dev server is up'
|
||||
}),
|
||||
lead: 'working'
|
||||
},
|
||||
// The parent's turn ends first; its child keeps running.
|
||||
{
|
||||
frame: turn('turn/completed', THREAD_ID, 'p1'),
|
||||
lead: 'done',
|
||||
check: () =>
|
||||
expect(byDescription('review')).toMatchObject({ membership: 'live', state: 'working' })
|
||||
},
|
||||
// The legacy task list carries no child state, so only the records can say a child waits,
|
||||
// and the shared fold ranks that wait above the parent's own state.
|
||||
{
|
||||
frame: status(REVIEWER, ['waitingOnApproval']),
|
||||
lead: 'done',
|
||||
childWaits: true,
|
||||
check: () => {
|
||||
expect(byDescription('review')?.state).toBe('waiting')
|
||||
expect(agentChildWorkLiveness(adapter.backgroundTaskState('session-1')?.tasks)).toBe(
|
||||
'working'
|
||||
)
|
||||
}
|
||||
},
|
||||
{ frame: status(REVIEWER, []), lead: 'done' },
|
||||
{
|
||||
frame: turn('turn/completed', REVIEWER, 'r1'),
|
||||
lead: 'done',
|
||||
check: () => {
|
||||
expect(byDescription('review')).toMatchObject({
|
||||
membership: 'settled',
|
||||
outcome: 'succeeded',
|
||||
lastMessage: 'Dev server is up'
|
||||
})
|
||||
expect(byDescription('npm run dev')).toMatchObject({
|
||||
membership: 'live',
|
||||
parentChildWorkId: byDescription('review')?.childWorkId
|
||||
})
|
||||
// Finished, but a shell it launched still runs: the CLI parent rule reads monitoring.
|
||||
expect(display('review')).toBe('monitoring')
|
||||
}
|
||||
},
|
||||
{ frame: turn('turn/started', THREAD_ID, 'p2'), lead: 'working' },
|
||||
// The parent asks the finished child a follow-up: the same record, a new run.
|
||||
{
|
||||
frame: turn('turn/started', REVIEWER, 'r2'),
|
||||
lead: 'working',
|
||||
check: () =>
|
||||
expect(byDescription('review')).toMatchObject({
|
||||
childWorkId: 'child-1',
|
||||
membership: 'live',
|
||||
invocation: { invocationId: 'r2', generation: 2 },
|
||||
previousInvocations: [expect.objectContaining({ outcome: 'succeeded' })]
|
||||
})
|
||||
},
|
||||
{ frame: spawned(TESTER, 'test', 'p2'), lead: 'working' },
|
||||
{ frame: turn('turn/started', TESTER, 't1'), lead: 'working' },
|
||||
{ frame: spawned(LINTER, 'lint', 'p2'), lead: 'working' },
|
||||
{ frame: turn('turn/started', LINTER, 'l1'), lead: 'working' },
|
||||
// Codex ends this child's turn with an error it will not retry, and no turn/completed.
|
||||
{
|
||||
frame: {
|
||||
method: 'error',
|
||||
params: { threadId: LINTER, turnId: 'l1', willRetry: false, error: { message: 'boom' } }
|
||||
},
|
||||
lead: 'working',
|
||||
check: () =>
|
||||
expect(byDescription('lint')).toMatchObject({ membership: 'settled', outcome: 'failed' })
|
||||
},
|
||||
{
|
||||
frame: turn('turn/completed', REVIEWER, 'r2', 'interrupted'),
|
||||
lead: 'working',
|
||||
check: () =>
|
||||
expect(byDescription('review')).toMatchObject({
|
||||
membership: 'settled',
|
||||
outcome: 'cancelled'
|
||||
})
|
||||
},
|
||||
{ frame: turn('turn/completed', THREAD_ID, 'p2'), lead: 'done' },
|
||||
{
|
||||
frame: turn('turn/completed', TESTER, 't1', 'failed'),
|
||||
lead: 'done',
|
||||
check: () =>
|
||||
expect(byDescription('test')).toMatchObject({ membership: 'settled', outcome: 'failed' })
|
||||
},
|
||||
{
|
||||
frame: shellFrame('item/completed', REVIEWER, 'r1', 'exec-1', 'npm run dev'),
|
||||
lead: 'done',
|
||||
check: () => {
|
||||
expect(byDescription('npm run dev')).toBeUndefined()
|
||||
expect(display('review')).toBe('interrupted')
|
||||
}
|
||||
}
|
||||
]
|
||||
let journalMoments = 0
|
||||
for (const [index, step] of steps.entries()) {
|
||||
const frameMoments = momentsOf(step.frame)
|
||||
journalMoments += frameMoments.length
|
||||
for (const [at, { recorded, legacy }] of frameMoments.entries()) {
|
||||
expect({ index, at, parent: fold(step.lead, asStrip(recorded)) }).toEqual({
|
||||
index,
|
||||
at,
|
||||
parent: fold(step.lead, legacy)
|
||||
})
|
||||
}
|
||||
const legacy = agentChildWorkLiveness(adapter.backgroundTaskState('session-1')?.tasks)
|
||||
const recorded = recordedLiveness()
|
||||
const expected = step.childWaits ? 'waiting' : legacy
|
||||
expect({ index, parent: fold(step.lead, recorded) }).toEqual({
|
||||
index,
|
||||
parent: fold(step.lead, expected)
|
||||
})
|
||||
expect({ index, liveness: recorded }).toEqual({ index, liveness: expected })
|
||||
step.check?.()
|
||||
}
|
||||
expect(journalMoments).toBeGreaterThan(steps.length)
|
||||
const settled = records()
|
||||
await adapter.closeSession('session-1')
|
||||
// Every child had already ended; closing the session changes none of what they said.
|
||||
expect(records()).toEqual(settled)
|
||||
expect(
|
||||
records().map(({ description, membership, outcome }) => ({
|
||||
description,
|
||||
membership,
|
||||
outcome
|
||||
}))
|
||||
).toEqual([
|
||||
{ description: 'review', membership: 'settled', outcome: 'cancelled' },
|
||||
{ description: 'test', membership: 'settled', outcome: 'failed' },
|
||||
{ description: 'lint', membership: 'settled', outcome: 'failed' }
|
||||
])
|
||||
expect(adapter.backgroundTaskState('session-1')).toBeUndefined()
|
||||
})
|
||||
|
||||
it("never reads done while the main agent's own shell runs past its turn", async () => {
|
||||
const { momentsOf, send, recordedLiveness } = await producer()
|
||||
send(turn('turn/started', THREAD_ID, 'p1'))
|
||||
send(shellFrame('item/started', THREAD_ID, 'p1', 'exec-dev', 'npm run dev'))
|
||||
const monitoring = { stateName: 'working', workingMode: 'monitoring' }
|
||||
// The turn ends with the dev server running: straight to monitoring, never done in between.
|
||||
const turnEnd = momentsOf(turn('turn/completed', THREAD_ID, 'p1'))
|
||||
expect(turnEnd.length).toBeGreaterThan(0)
|
||||
for (const { recorded } of turnEnd) {
|
||||
expect(fold('done', recorded)).toEqual(monitoring)
|
||||
}
|
||||
expect(fold('done', recordedLiveness())).toEqual(monitoring)
|
||||
momentsOf(shellFrame('item/completed', THREAD_ID, 'p1', 'exec-dev', 'npm run dev'))
|
||||
expect(fold('done', recordedLiveness())).toEqual({ stateName: 'done' })
|
||||
})
|
||||
|
||||
it('ends an approval left unanswered when its turn ends: Codex never ran the command', async () => {
|
||||
const { adapter, codex, send, records, recordedLiveness, display, byDescription } =
|
||||
await producer()
|
||||
const approve = (threadId: string, turnId: string, itemId: string, command: string) => {
|
||||
// The approval path starts the item before it asks, and drops the question at turn end.
|
||||
send(shellFrame('item/started', threadId, turnId, itemId, command, 'agent'))
|
||||
codex.connections[0]!.handlers.onServerRequest?.({
|
||||
id: `approval-${itemId}`,
|
||||
method: 'item/commandExecution/requestApproval',
|
||||
params: { itemId, threadId, turnId }
|
||||
})
|
||||
}
|
||||
send(turn('turn/started', THREAD_ID, 'p1'))
|
||||
send(turn('turn/started', REVIEWER, 'r1'))
|
||||
send(spawned(REVIEWER, 'review', 'p1'))
|
||||
approve(REVIEWER, 'r1', 'call-child', 'npm run e2e')
|
||||
approve(THREAD_ID, 'p1', 'call-main', 'npm run dev')
|
||||
expect(records().filter((record) => record.kind === 'command')).toHaveLength(2)
|
||||
// The user stops the child, then the main agent, each at its approval.
|
||||
send(turn('turn/completed', REVIEWER, 'r1', 'interrupted'))
|
||||
expect(byDescription('npm run e2e')).toBeUndefined()
|
||||
expect(display('review')).toBe('interrupted')
|
||||
send(turn('turn/completed', THREAD_ID, 'p1', 'interrupted'))
|
||||
expect(byDescription('npm run dev')).toBeUndefined()
|
||||
expect(adapter.backgroundTaskState('session-1')).toBeNull()
|
||||
expect(fold('done', recordedLiveness())).toEqual({ stateName: 'done' })
|
||||
})
|
||||
|
||||
it('keeps an answered approval running past its turn', async () => {
|
||||
const { adapter, codex, send, byDescription } = await producer()
|
||||
send(turn('turn/started', THREAD_ID, 'p1'))
|
||||
send(shellFrame('item/started', THREAD_ID, 'p1', 'call-1', 'npm run dev', 'agent'))
|
||||
codex.connections[0]!.handlers.onServerRequest?.({
|
||||
id: 'approval-1',
|
||||
method: 'item/commandExecution/requestApproval',
|
||||
params: { itemId: 'call-1', threadId: THREAD_ID, turnId: 'p1' }
|
||||
})
|
||||
await adapter.answerPrompt({
|
||||
sessionId: 'session-1',
|
||||
itemId: 'call-1',
|
||||
kind: 'approval',
|
||||
response: { kind: 'option', optionId: 'accept' },
|
||||
fence: 7,
|
||||
commit: async () => {}
|
||||
})
|
||||
send(turn('turn/completed', THREAD_ID, 'p1'))
|
||||
expect(byDescription('npm run dev')).toMatchObject({ membership: 'live' })
|
||||
expect(adapter.backgroundTaskState('session-1')?.tasks).toEqual([
|
||||
expect.objectContaining({ kind: 'command', description: 'npm run dev' })
|
||||
])
|
||||
})
|
||||
|
||||
it("numbers a child's runs as the journal does: a row's attempt is its record's generation", async () => {
|
||||
const { send, byDescription, stampOf } = await producer()
|
||||
const says = (turnId: string, text: string) =>
|
||||
item('item/completed', REVIEWER, turnId, { type: 'agentMessage', id: `msg-${text}`, text })
|
||||
// The journal stamps a child row with its run only once it is past the first.
|
||||
const runs = (text: string) => {
|
||||
const stamp = stampOf(text)
|
||||
return {
|
||||
agentId: stamp?.agentId,
|
||||
attempt: stamp ? (stamp.attempt ?? 1) : undefined,
|
||||
generation: byDescription('review')?.invocation.generation
|
||||
}
|
||||
}
|
||||
send(turn('turn/started', THREAD_ID, 'p1'))
|
||||
// Codex reports the child's first turn before the spawn that announces it.
|
||||
send(turn('turn/started', REVIEWER, 'r1'))
|
||||
send(spawned(REVIEWER, 'review', 'p1'))
|
||||
send(says('r1', 'run 1'))
|
||||
expect(runs('run 1')).toEqual({ agentId: REVIEWER, attempt: 1, generation: 1 })
|
||||
send(turn('turn/completed', REVIEWER, 'r1'))
|
||||
// Each follow-up the parent sends is the child's next run, on both sides.
|
||||
for (const run of [2, 3]) {
|
||||
send(turn('turn/started', REVIEWER, `r${run}`))
|
||||
send(says(`r${run}`, `run ${run}`))
|
||||
expect(runs(`run ${run}`)).toEqual({ agentId: REVIEWER, attempt: run, generation: run })
|
||||
send(turn('turn/completed', REVIEWER, `r${run}`))
|
||||
}
|
||||
})
|
||||
|
||||
it('settles a live child with no reported outcome when the provider exits unexpectedly', async () => {
|
||||
const { codex, send, records } = await producer()
|
||||
send(turn('turn/started', THREAD_ID, 'p1'))
|
||||
send(spawned(REVIEWER, 'review', 'p1'))
|
||||
send(turn('turn/started', REVIEWER, 'r1'))
|
||||
expect(records()).toEqual([
|
||||
expect.objectContaining({ description: 'review', membership: 'live', state: 'working' })
|
||||
])
|
||||
codex.connections[0]!.handlers.onExit?.(new Error('provider exited'))
|
||||
expect(records()).toEqual([
|
||||
expect.objectContaining({
|
||||
description: 'review',
|
||||
membership: 'settled',
|
||||
state: 'done',
|
||||
outcome: 'unknown'
|
||||
})
|
||||
])
|
||||
})
|
||||
})
|
||||
@@ -8,7 +8,7 @@ import {
|
||||
closeFailedCodexAcquisition,
|
||||
stopSupersededCodexAcquisition
|
||||
} from './codex-structured-acquisition-lifecycle'
|
||||
import { CodexBackgroundTaskTracker } from './codex-background-task-tracker'
|
||||
import { CodexBackgroundTaskTracker, codexChildWorkSink } from './codex-background-task-tracker'
|
||||
import { CodexSubagentExecutions } from './codex-subagent-executions'
|
||||
import { createCodexDispatchEchoes } from './codex-structured-dispatch-echo'
|
||||
import { createCodexJournalTranslator } from './codex-structured-journal-translation'
|
||||
@@ -44,6 +44,8 @@ import type { CodexStructuredTurnCancellation } from './codex-structured-turn-ca
|
||||
import type { CodexStructuredNotificationRetry } from './codex-structured-notification-retry'
|
||||
import type { deliverCodexServerRequest } from './codex-structured-provider-events'
|
||||
|
||||
const TURN_BOUNDARIES: ReadonlySet<string> = new Set(['turn/started', 'turn/completed'])
|
||||
|
||||
export async function acquireCodexStructuredSession(input: {
|
||||
input: StructuredAgentSessionAcquireInput
|
||||
deps: CodexStructuredSessionAdapterDeps
|
||||
@@ -138,7 +140,7 @@ export async function acquireCodexStructuredSession(input: {
|
||||
{
|
||||
onNotification: (method, params) => {
|
||||
// Stamped at receipt, ahead of any pre-publication buffering or retry.
|
||||
const observedAt = isCodexTurnBoundary(method) ? (deps.now?.() ?? Date.now()) : undefined
|
||||
const observedAt = TURN_BOUNDARIES.has(method) ? (deps.now?.() ?? Date.now()) : undefined
|
||||
const dispatchSequenceAtReceipt =
|
||||
method === 'turn/started' ? dispatchEchoes.latestSequence() : undefined
|
||||
input.deliver(
|
||||
@@ -240,6 +242,8 @@ export async function acquireCodexStructuredSession(input: {
|
||||
throw new Error(`codex app-server for session ${sessionId} exited while being acquired`)
|
||||
}
|
||||
acquisitions.deleteIfCurrent(sessionId, attempt)
|
||||
// Where this session's child work goes: the host's records, after each frame is journaled.
|
||||
const sink = codexChildWorkSink(sessionId, deps)
|
||||
const session: CodexSession = {
|
||||
connection,
|
||||
...codexSessionLifecycle(acquireInput.fence, acquired.acquisitionGeneration as string),
|
||||
@@ -254,7 +258,7 @@ export async function acquireCodexStructuredSession(input: {
|
||||
...(catalogAccess ? { catalogAccess } : {}),
|
||||
dispatchEchoes,
|
||||
translator,
|
||||
backgroundTasks: new CodexBackgroundTaskTracker(opened.threadId, subagentExecutions),
|
||||
backgroundTasks: new CodexBackgroundTaskTracker(opened.threadId, subagentExecutions, sink),
|
||||
forceCloseUnexpected: (reason) =>
|
||||
input.forceCloseUnexpected(
|
||||
sessionId,
|
||||
@@ -299,7 +303,3 @@ export async function acquireCodexStructuredSession(input: {
|
||||
attempt.finish()
|
||||
}
|
||||
}
|
||||
|
||||
function isCodexTurnBoundary(method: string): boolean {
|
||||
return method === 'turn/started' || method === 'turn/completed'
|
||||
}
|
||||
|
||||
@@ -161,9 +161,11 @@ export class CodexStructuredSessionAdapter implements StructuredAgentSessionAdap
|
||||
)
|
||||
// After the admission check, so a refused frame is observed by the strip
|
||||
// only on the retry that also reaches the journal.
|
||||
if (session.backgroundTasks.observe(event)) {
|
||||
if (session.backgroundTasks.observe(event, session.prompts.takeAbandonedCommands())) {
|
||||
this.deps.onBackgroundTasksChanged?.(event.sessionId, session.backgroundTasks.state)
|
||||
}
|
||||
// After the journal and the parent's republished row, never ahead of either.
|
||||
session.backgroundTasks.publishChildWork()
|
||||
}
|
||||
if (event.type === 'ended') {
|
||||
this.compactions.ended(event.sessionId)
|
||||
|
||||
@@ -56,6 +56,8 @@ export function handleCodexSessionExit(input: {
|
||||
session.dispatchEchoes.clear()
|
||||
session.backgroundTasks.clear()
|
||||
input.onBackgroundTasksChanged?.(input.sessionId, null)
|
||||
// Every close path funnels here, so the session's children end with it on each one.
|
||||
session.backgroundTasks.publishChildWork()
|
||||
session.unbindReadingControl?.()
|
||||
input.onEvent?.(event)
|
||||
session.prompts.clear()
|
||||
|
||||
@@ -12,6 +12,7 @@ import type {
|
||||
import { CodexAcquisitionWindow } from './codex-structured-acquisition-window'
|
||||
import type { CodexDispatchEchoes } from './codex-structured-dispatch-echo'
|
||||
import type { AgentSessionBackgroundTaskState } from '../../shared/agent-session-wire'
|
||||
import type { AgentChildWorkEvidence } from '../../shared/agent-status-child-work-evidence'
|
||||
import type { CodexBackgroundTaskTracker } from './codex-background-task-tracker'
|
||||
import type { CodexJournalTranslator } from './codex-structured-journal-translation'
|
||||
import type { CodexTurnProcessSnapshot } from './codex-structured-turn-processes'
|
||||
@@ -78,6 +79,8 @@ export type CodexStructuredSessionAdapterDeps = {
|
||||
sessionId: string,
|
||||
state: AgentSessionBackgroundTaskState | null
|
||||
) => void
|
||||
/** What the session's child work did, delivered after the journal handled the frame. */
|
||||
onChildWorkEvidence?: (sessionId: string, evidence: AgentChildWorkEvidence[]) => void
|
||||
/** A send admitted earlier: its identity once Codex echoes it, or its rejection when the turn
|
||||
* Codex answered it into ended without taking it. */
|
||||
onDispatchSettledLate?: (
|
||||
|
||||
@@ -75,3 +75,13 @@ export function codexThreadStoppedRunning(payload: unknown): boolean {
|
||||
const type = record(record(payload)?.status)?.type
|
||||
return type === 'idle' || type === 'systemError'
|
||||
}
|
||||
|
||||
/** An `active` thread flags each request it has open on the user (an approval, a question). */
|
||||
export function codexThreadWaitsOnUser(payload: unknown): boolean {
|
||||
const status = record(record(payload)?.status)
|
||||
const flags = status?.type === 'active' ? status.activeFlags : null
|
||||
return (
|
||||
Array.isArray(flags) &&
|
||||
flags.some((flag) => flag === 'waitingOnApproval' || flag === 'waitingOnUserInput')
|
||||
)
|
||||
}
|
||||
|
||||
@@ -94,6 +94,20 @@ export class CodexSubagentExecutions {
|
||||
return { child, execution }
|
||||
}
|
||||
|
||||
/** A child turn that ended with no `turn/completed`. With no turn named, the one the child is
|
||||
* running ended; a child running none has nothing to end. The first ending a turn gets stands. */
|
||||
endTurn(
|
||||
agentThreadId: string,
|
||||
turnId: string | null,
|
||||
state: Exclude<NativeChatSubagentState, 'working'>
|
||||
): void {
|
||||
const current = this.children.get(agentThreadId)?.execution
|
||||
const ended = turnId ?? (current?.state === 'working' ? current.turnId : null)
|
||||
if (current && ended !== null) {
|
||||
this.observeTurn(agentThreadId, ended, state)
|
||||
}
|
||||
}
|
||||
|
||||
/** Survives the child's turn, so a row outliving that turn can still name it. */
|
||||
label(agentThreadId: string): string | null {
|
||||
return this.children.get(agentThreadId)?.label ?? null
|
||||
@@ -109,6 +123,11 @@ export class CodexSubagentExecutions {
|
||||
return this.children.get(agentThreadId)?.turnOrdinals.get(turnId) ?? null
|
||||
}
|
||||
|
||||
/** The child as last observed, without creating one. */
|
||||
find(agentThreadId: string): Readonly<CodexExecutionChild> | undefined {
|
||||
return this.children.get(agentThreadId)
|
||||
}
|
||||
|
||||
workingChildren(): CodexExecutionChild[] {
|
||||
return [...this.children.values()].filter(
|
||||
(child) => child.registered && child.execution?.state === 'working'
|
||||
|
||||
@@ -247,6 +247,8 @@ async function install(deps: StructuredAgentSessionRuntimeDeps): Promise<Install
|
||||
modelCatalog: agentModelCatalogStore,
|
||||
onBackgroundTasksChanged: (sessionId, state) =>
|
||||
host?.publishBackgroundTaskState(sessionId, state),
|
||||
onChildWorkEvidence: (sessionId, evidence) =>
|
||||
host?.publishChildWorkEvidence(sessionId, evidence),
|
||||
onDispatchSettledLate,
|
||||
onPrimaryThreadStoppedRunning: ({ sessionId }) => {
|
||||
void host
|
||||
|
||||
@@ -0,0 +1,133 @@
|
||||
// The production runtime hands a Codex session's child work to the status sink, under the address
|
||||
// the session's own row landed under, and ends it with the provider.
|
||||
|
||||
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 type {
|
||||
CodexAppServerConnection,
|
||||
CodexAppServerConnectionHandlers,
|
||||
openCodexAppServerConnection
|
||||
} from '../codex/codex-app-server-connection'
|
||||
import {
|
||||
HOST_TEST_SESSION as SESSION,
|
||||
hostTestAttachParams
|
||||
} from '../native-chat/agent-session-wire/structured-agent-session-host-test-data'
|
||||
import type { StructuredAgentSessionStatusSink } from '../native-chat/agent-session-wire/structured-agent-session-status-feed'
|
||||
import {
|
||||
ensureStructuredAgentSessionHost,
|
||||
stopStructuredAgentSessionRuntime
|
||||
} from './structured-agent-session-runtime'
|
||||
|
||||
const THREAD = 'thread-runtime-child-work'
|
||||
const CHILD = 'thread-runtime-reviewer'
|
||||
const ROUTES: Record<string, unknown> = {
|
||||
'thread/start': { thread: { id: THREAD } },
|
||||
'model/list': {
|
||||
data: [
|
||||
{
|
||||
model: 'gpt-test',
|
||||
displayName: 'GPT Test',
|
||||
hidden: false,
|
||||
supportedReasoningEfforts: [],
|
||||
defaultReasoningEffort: null,
|
||||
isDefault: true
|
||||
}
|
||||
],
|
||||
nextCursor: null
|
||||
}
|
||||
}
|
||||
|
||||
describe('structured Codex child work through the production runtime', () => {
|
||||
let root: string | null = null
|
||||
|
||||
afterEach(async () => {
|
||||
await stopStructuredAgentSessionRuntime()
|
||||
if (root) {
|
||||
await rm(root, { recursive: true, force: true })
|
||||
root = null
|
||||
}
|
||||
})
|
||||
|
||||
it("hands its subagents to the status sink under the session's own address", async () => {
|
||||
root = await mkdtemp(join(tmpdir(), 'orca-runtime-codex-child-work-'))
|
||||
const connections: CodexAppServerConnectionHandlers[] = []
|
||||
const openConnection: typeof openCodexAppServerConnection = async (_launch, handlers = {}) => {
|
||||
connections.push(handlers)
|
||||
const connection: CodexAppServerConnection = {
|
||||
pid: 4321,
|
||||
closed: false,
|
||||
request: async (method) => (method in ROUTES ? ROUTES[method] : {}),
|
||||
notify: () => {},
|
||||
respond: () => {},
|
||||
respondWithError: () => {},
|
||||
close: async () => true
|
||||
}
|
||||
return connection
|
||||
}
|
||||
const childWork: Parameters<
|
||||
NonNullable<StructuredAgentSessionStatusSink['publishChildWork']>
|
||||
>[] = []
|
||||
const host = await ensureStructuredAgentSessionHost({
|
||||
stateDirectory: root,
|
||||
hostId: 'local',
|
||||
claimKeyId: 'key-1',
|
||||
resolveWorkspacePath: async () => root!,
|
||||
resolveClaudeAuthPolicy: () => ({ stripAuthEnv: true }),
|
||||
resolveCodexCommand: () => 'codex',
|
||||
resolveEnvironment: async () => ({ PATH: process.env.PATH }),
|
||||
openCodexConnection: openConnection,
|
||||
readProcessStartTime: async () => 1_700_000_000_000,
|
||||
statusSink: {
|
||||
publish: () => {},
|
||||
forget: () => {},
|
||||
publishChildWork: (...args) => childWork.push(args)
|
||||
}
|
||||
})
|
||||
const attachParams = hostTestAttachParams(null, { providerHandle: undefined })
|
||||
attachParams.envelope.clientOperationId = `${Date.now()}-${'1'.padStart(32, '0')}`
|
||||
const attached = await host.attach({ callerKey: 'runtime-test' }, attachParams)
|
||||
expect(attached).toMatchObject({ ok: true })
|
||||
// Creating the session starts its child; nothing else has to keep it running.
|
||||
expect(connections).toHaveLength(1)
|
||||
const notify = (method: string, params: Record<string, unknown>) =>
|
||||
connections[0]?.onNotification?.(method, params)
|
||||
notify('turn/started', { threadId: THREAD, turn: { id: 'turn-1' } })
|
||||
notify('turn/started', { threadId: CHILD, turn: { id: 'child-turn-1' } })
|
||||
notify('item/started', {
|
||||
threadId: THREAD,
|
||||
turnId: 'turn-1',
|
||||
item: {
|
||||
type: 'subAgentActivity',
|
||||
id: 'spawn-1',
|
||||
kind: 'started',
|
||||
agentThreadId: CHILD,
|
||||
agentPath: '/root/review'
|
||||
}
|
||||
})
|
||||
const subject = expect.objectContaining({ kind: 'structured-session', sessionId: SESSION })
|
||||
expect(childWork).toEqual([
|
||||
[
|
||||
subject,
|
||||
[
|
||||
expect.objectContaining({
|
||||
type: 'live',
|
||||
child: expect.objectContaining({
|
||||
handle: { idKind: 'thread_id', id: CHILD, runId: 'child-turn-1' },
|
||||
description: 'review'
|
||||
})
|
||||
})
|
||||
],
|
||||
'codex'
|
||||
]
|
||||
])
|
||||
// The provider dies: its session's end is reported under the same address.
|
||||
connections[0]?.onExit?.(new Error('scripted provider exit'))
|
||||
expect(childWork.at(-1)).toEqual([
|
||||
subject,
|
||||
[expect.objectContaining({ type: 'session-ended' })],
|
||||
'codex'
|
||||
])
|
||||
})
|
||||
})
|
||||
@@ -16,11 +16,17 @@ const CHILD_ALIAS_KEY_PREFIX = 'agent-child-work-alias-v1:'
|
||||
const MAX_ALIAS_PART_LENGTH = 512
|
||||
|
||||
/**
|
||||
* `thread_id` names a child by its own provider thread (a Codex subagent). The hook lane registers
|
||||
* a Claude `agent_id` under `task_id` (it is the same registry id) and a Codex `agent_id` under
|
||||
* `thread_id`; no `agent_id` kind exists on purpose.
|
||||
* `thread_id` names a child by its own provider thread (a Codex subagent); `turn_id` names one run
|
||||
* of such a child, as `tool_use_id` names one run of a task. The hook lane registers a Claude
|
||||
* `agent_id` under `task_id` (it is the same registry id) and a Codex `agent_id` under `thread_id`;
|
||||
* no `agent_id` kind exists on purpose.
|
||||
*/
|
||||
export const AGENT_CHILD_WORK_ALIAS_KINDS = ['task_id', 'tool_use_id', 'thread_id'] as const
|
||||
export const AGENT_CHILD_WORK_ALIAS_KINDS = [
|
||||
'task_id',
|
||||
'tool_use_id',
|
||||
'thread_id',
|
||||
'turn_id'
|
||||
] as const
|
||||
export type AgentChildWorkAliasKind = (typeof AGENT_CHILD_WORK_ALIAS_KINDS)[number]
|
||||
const ALIAS_KIND_SET: ReadonlySet<string> = new Set(AGENT_CHILD_WORK_ALIAS_KINDS)
|
||||
|
||||
|
||||
@@ -20,7 +20,15 @@ import { agentStatusSubjectsEqual, type AgentStatusSubject } from './agent-statu
|
||||
* a handle is unique per parent and provider without being unique across producers. */
|
||||
export const STRUCTURED_CHILD_WORK_PRODUCER_ID = 'structured-session-child-work'
|
||||
const SEGMENT_ID = STRUCTURED_CHILD_WORK_PRODUCER_ID
|
||||
const RUN_ALIAS_KIND: AgentChildWorkAliasKind = 'tool_use_id'
|
||||
/** A run is named in the provider's own terms: a task runs under its spawn call, a thread under
|
||||
* its turn. The kinds stay apart so a turn id can never pass for a spawn call. */
|
||||
const RUN_ALIAS_KIND_BY_ID_KIND = {
|
||||
task_id: 'tool_use_id',
|
||||
thread_id: 'turn_id'
|
||||
} as const satisfies Record<AgentChildWorkEvidenceHandle['idKind'], AgentChildWorkAliasKind>
|
||||
const RUN_ALIAS_KINDS: ReadonlySet<AgentChildWorkAliasKind> = new Set(
|
||||
Object.values(RUN_ALIAS_KIND_BY_ID_KIND)
|
||||
)
|
||||
|
||||
export const STRUCTURED_CHILD_WORK_PROVENANCE: AgentChildWorkProvenance = {
|
||||
source: 'structured-session',
|
||||
@@ -70,7 +78,13 @@ export function agentChildWorkHandleAliases(
|
||||
return [
|
||||
{ segmentId: SEGMENT_ID, aliasKind: handle.idKind, alias: handle.id },
|
||||
...(handle.runId !== undefined && handle.runId !== handle.id
|
||||
? [{ segmentId: SEGMENT_ID, aliasKind: RUN_ALIAS_KIND, alias: handle.runId }]
|
||||
? [
|
||||
{
|
||||
segmentId: SEGMENT_ID,
|
||||
aliasKind: RUN_ALIAS_KIND_BY_ID_KIND[handle.idKind],
|
||||
alias: handle.runId
|
||||
}
|
||||
]
|
||||
: [])
|
||||
]
|
||||
}
|
||||
@@ -123,14 +137,14 @@ export function resolveAgentChildWorkHandle(
|
||||
: { child: child ?? null, ambiguous: false, highestGeneration }
|
||||
}
|
||||
|
||||
/** The owner a handle id names, by its stable id or by the run handle it spawned under. */
|
||||
/** The owner a handle id names, by its stable id or by the spawn call it runs under. */
|
||||
export function resolveAgentChildWorkOwner(
|
||||
scope: AgentChildWorkEvidenceScope,
|
||||
ownerId: string
|
||||
): string | undefined {
|
||||
const resolution = resolveAgentChildWorkHandle(
|
||||
scope,
|
||||
['task_id', 'thread_id', RUN_ALIAS_KIND],
|
||||
['task_id', 'thread_id', 'tool_use_id'],
|
||||
ownerId
|
||||
)
|
||||
return resolution?.child?.childWorkId
|
||||
@@ -164,7 +178,7 @@ export function currentAgentChildWorkAliases(
|
||||
return {
|
||||
...(stable ? { stable } : {}),
|
||||
stableId: stable?.id,
|
||||
runId: current.find((alias) => alias.aliasKind === RUN_ALIAS_KIND)?.alias,
|
||||
runId: current.find((alias) => RUN_ALIAS_KINDS.has(alias.aliasKind))?.alias,
|
||||
aliases: current.map((alias) => ({
|
||||
segmentId: alias.segmentId,
|
||||
aliasKind: alias.aliasKind,
|
||||
@@ -183,7 +197,7 @@ export function isPreviousAgentChildWorkRun(
|
||||
.getAliasesForChild(child.childWorkId)
|
||||
.some(
|
||||
(alias) =>
|
||||
alias.aliasKind === RUN_ALIAS_KIND &&
|
||||
RUN_ALIAS_KINDS.has(alias.aliasKind) &&
|
||||
alias.alias === runId &&
|
||||
!agentChildWorkFencesEqual(alias.fence, child.invocation)
|
||||
)
|
||||
|
||||
@@ -2,7 +2,8 @@
|
||||
//
|
||||
// A producer decodes provider frames into these edges and the host folds them into the one
|
||||
// record per child it owns. Edges carry facts, not records: which child is live, what it is
|
||||
// doing, how it ended. Only a child's own ending settles it, or the end of its session.
|
||||
// doing, how it ended, or that it is gone. Only a child's own ending settles it, or the end of its
|
||||
// session. Edges are host-internal: the producer and the store share one process.
|
||||
|
||||
import type { AgentChildWorkAliasKind } from './agent-status-child-work-alias'
|
||||
import type {
|
||||
@@ -13,9 +14,9 @@ import type {
|
||||
AgentChildWorkState
|
||||
} from './agent-status-child-work'
|
||||
|
||||
/** How the provider names one child. `id` is the stable handle today's wire already publishes
|
||||
* (a Claude task id); `runId` names the current run when the provider mints one per run (the
|
||||
* spawn call), and a different one is the provider starting the child again. */
|
||||
/** How the provider names one child. `id` is the stable handle: a task id, or the child's own
|
||||
* thread. `runId` names the current run when the provider mints one per run (a task's spawn
|
||||
* call, a thread's turn), and a different one is the provider starting the child again. */
|
||||
export type AgentChildWorkEvidenceHandle = {
|
||||
idKind: Extract<AgentChildWorkAliasKind, 'task_id' | 'thread_id'>
|
||||
id: string
|
||||
@@ -74,8 +75,17 @@ export type AgentChildWorkEndedEvidence = {
|
||||
* with an outcome nobody reported. Settled children stay; the parent's removal drops them. */
|
||||
export type AgentChildWorkSessionEndedEvidence = { type: 'session-ended'; observedAt: number }
|
||||
|
||||
/** Work that leaves nothing to report once it stops, such as a command whose process exited: its
|
||||
* record goes rather than settles. For work that owns no other record. */
|
||||
export type AgentChildWorkRemovedEvidence = {
|
||||
type: 'removed'
|
||||
observedAt: number
|
||||
handle: AgentChildWorkEvidenceHandle
|
||||
}
|
||||
|
||||
export type AgentChildWorkEvidence =
|
||||
| AgentChildWorkLiveEvidence
|
||||
| AgentChildWorkOperationEvidence
|
||||
| AgentChildWorkEndedEvidence
|
||||
| AgentChildWorkRemovedEvidence
|
||||
| AgentChildWorkSessionEndedEvidence
|
||||
|
||||
@@ -5,10 +5,7 @@ import type {
|
||||
AgentChildWorkEvidence,
|
||||
AgentChildWorkLiveObservation
|
||||
} from './agent-status-child-work-evidence'
|
||||
import {
|
||||
reconcileAgentChildWorkEvidence,
|
||||
STRUCTURED_CHILD_WORK_MAX_SETTLED
|
||||
} from './agent-status-child-work-reconciliation'
|
||||
import { reconcileAgentChildWorkEvidence } from './agent-status-child-work-reconciliation'
|
||||
import { STRUCTURED_CHILD_WORK_MAX_LIVE } from './agent-status-child-work-evidence-admission'
|
||||
import { createAgentStatusStore, type AgentStatusStore } from './agent-status-store'
|
||||
import { makeStructuredAgentStatusSubject } from './agent-status-subject'
|
||||
@@ -399,17 +396,9 @@ describe('structured child-work reconciliation', () => {
|
||||
expect(records(store)).toHaveLength(STRUCTURED_CHILD_WORK_MAX_LIVE)
|
||||
})
|
||||
|
||||
it('keeps a bounded settled history, never dropping a child that owns live work', () => {
|
||||
it('keeps every settled child until the parent row goes', () => {
|
||||
const { store, apply } = harness()
|
||||
apply(live(child('owner')))
|
||||
apply(live(child('shell', { kind: 'command', ownerId: 'owner' })))
|
||||
apply({
|
||||
type: 'ended',
|
||||
observedAt: 101,
|
||||
handle: { idKind: 'task_id', id: 'owner' },
|
||||
outcome: 'succeeded'
|
||||
})
|
||||
for (let index = 0; index < STRUCTURED_CHILD_WORK_MAX_SETTLED + 1; index += 1) {
|
||||
for (let index = 0; index < 100; index += 1) {
|
||||
apply(live(child(`done-${index}`), 200 + index), {
|
||||
type: 'ended',
|
||||
observedAt: 200 + index,
|
||||
@@ -418,9 +407,22 @@ describe('structured child-work reconciliation', () => {
|
||||
})
|
||||
}
|
||||
const settled = records(store).filter((record) => record.membership === 'settled')
|
||||
expect(settled).toHaveLength(STRUCTURED_CHILD_WORK_MAX_SETTLED)
|
||||
expect(settled.map((record) => record.description)).toContain('Task owner')
|
||||
expect(settled.map((record) => record.description)).not.toContain('Task done-0')
|
||||
expect(settled.map((record) => record.description)).not.toContain('Task done-1')
|
||||
expect(settled).toHaveLength(100)
|
||||
expect(settled.map((record) => record.description)).toContain('Task done-0')
|
||||
})
|
||||
|
||||
it('removes work that stopped with nothing to report, and only the record it names', () => {
|
||||
const { store, apply } = harness()
|
||||
apply(live(child('owner')), live(child('shell', { kind: 'command', ownerId: 'owner' })))
|
||||
expect(
|
||||
apply({ type: 'removed', observedAt: 200, handle: { idKind: 'task_id', id: 'shell' } })
|
||||
).toMatchObject({ removed: 1, settled: 0 })
|
||||
expect(only(store)).toMatchObject({ description: 'Task owner', membership: 'live' })
|
||||
expect(store.getAliasesForChild('child-2')).toEqual([])
|
||||
// A handle it no longer answers to removes nothing.
|
||||
expect(
|
||||
apply({ type: 'removed', observedAt: 201, handle: { idKind: 'task_id', id: 'shell' } })
|
||||
).toMatchObject({ removed: 0 })
|
||||
expect(records(store)).toHaveLength(1)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -1,14 +1,17 @@
|
||||
// Fold one structured session's child-work evidence into the host's records.
|
||||
//
|
||||
// The store holds the only current record per child; evidence patches it. A child settles on its
|
||||
// own ending, or `unknown` when its session ends while it is still live. It owns only the records
|
||||
// its own producer admitted, and never claims an outcome the evidence did not report.
|
||||
// own ending, or `unknown` when its session ends while it is still live; work with nothing to
|
||||
// report once it stops is removed instead. Settled children stay until the host drops the parent's
|
||||
// row. It owns only the records its own producer admitted, and never claims an outcome the
|
||||
// evidence did not report.
|
||||
|
||||
import type { AgentChildWorkAdmission } from './agent-status-child-work-admission'
|
||||
import type {
|
||||
AgentChildWorkEndedEvidence,
|
||||
AgentChildWorkEvidence,
|
||||
AgentChildWorkOperationEvidence
|
||||
AgentChildWorkOperationEvidence,
|
||||
AgentChildWorkRemovedEvidence
|
||||
} from './agent-status-child-work-evidence'
|
||||
import {
|
||||
applyAgentChildWorkLive,
|
||||
@@ -26,9 +29,6 @@ import {
|
||||
|
||||
export type { AgentChildWorkReconcileOutcome } from './agent-status-child-work-evidence-admission'
|
||||
|
||||
/** Settled children kept per session. The oldest go first, never one that owns live work. */
|
||||
export const STRUCTURED_CHILD_WORK_MAX_SETTLED = 32
|
||||
|
||||
export type AgentChildWorkReconcileInput = AgentChildWorkEvidenceScope & {
|
||||
admission: AgentChildWorkAdmission
|
||||
evidence: readonly AgentChildWorkEvidence[]
|
||||
@@ -98,32 +98,20 @@ function settleLive(ctx: ReconcileContext, observedAt: number): void {
|
||||
}
|
||||
}
|
||||
|
||||
function removeChildren(ctx: ReconcileContext, childWorkIds: string[]): void {
|
||||
if (childWorkIds.length > 0 && ctx.store.applyMutation({ removeChildren: childWorkIds })) {
|
||||
ctx.outcome.removed += childWorkIds.length
|
||||
}
|
||||
}
|
||||
|
||||
/** Oldest-settled first; a settled child that still owns live work stays so its work keeps an owner. */
|
||||
function trimSettled(ctx: ReconcileContext): void {
|
||||
const owned = ownedStructuredChildWork(ctx)
|
||||
const settled = owned.filter((record) => record.membership === 'settled')
|
||||
const excess = settled.length - STRUCTURED_CHILD_WORK_MAX_SETTLED
|
||||
if (excess <= 0) {
|
||||
/** The work is gone and has no ending to keep: its record, and the handles it answered to, go. */
|
||||
function applyRemoved(ctx: ReconcileContext, edge: AgentChildWorkRemovedEvidence): void {
|
||||
const resolution = resolveAgentChildWorkHandle(ctx, [edge.handle.idKind], edge.handle.id)
|
||||
if (resolution?.ambiguous) {
|
||||
ctx.outcome.rejected.push({ handleId: edge.handle.id, reason: 'ambiguous' })
|
||||
return
|
||||
}
|
||||
const owners = new Set(
|
||||
owned.flatMap((record) =>
|
||||
record.membership === 'live' && record.parentChildWorkId ? [record.parentChildWorkId] : []
|
||||
)
|
||||
)
|
||||
const removable = settled
|
||||
.filter((record) => !owners.has(record.childWorkId))
|
||||
.sort((a, b) => (a.settledAt ?? a.observedAt) - (b.settledAt ?? b.observedAt))
|
||||
removeChildren(
|
||||
ctx,
|
||||
removable.slice(0, excess).map((record) => record.childWorkId)
|
||||
)
|
||||
const existing = resolution?.child
|
||||
if (!existing || agentChildWorkRunVerdict(ctx, existing, edge.handle.runId) === 'previous') {
|
||||
return
|
||||
}
|
||||
if (ctx.store.applyMutation({ removeChildren: [existing.childWorkId] })) {
|
||||
ctx.outcome.removed += 1
|
||||
}
|
||||
}
|
||||
|
||||
/** Apply one batch of evidence. The parent must already be held: the store refuses a child whose
|
||||
@@ -145,13 +133,11 @@ export function reconcileAgentChildWorkEvidence(
|
||||
applyOperation(ctx, edge)
|
||||
} else if (edge.type === 'ended') {
|
||||
applyEnded(ctx, edge)
|
||||
} else if (edge.type === 'removed') {
|
||||
applyRemoved(ctx, edge)
|
||||
} else {
|
||||
settleLive(ctx, edge.observedAt)
|
||||
}
|
||||
}
|
||||
// Only a settle adds settled history; skipping the scan otherwise keeps progress edges cheap.
|
||||
if (ctx.outcome.settled > 0) {
|
||||
trimSettled(ctx)
|
||||
}
|
||||
return ctx.outcome
|
||||
}
|
||||
|
||||
@@ -60,7 +60,8 @@ export type AgentChildWorkViewAlias = Pick<
|
||||
const PROVIDER_ID_ALIAS_RANK: Record<AgentChildWorkAliasKind, number> = {
|
||||
task_id: 0,
|
||||
thread_id: 1,
|
||||
tool_use_id: 2
|
||||
tool_use_id: 2,
|
||||
turn_id: 3
|
||||
}
|
||||
const PROVIDER_ID_ALIAS_ORDER = [...AGENT_CHILD_WORK_ALIAS_KINDS].sort(
|
||||
(left, right) => PROVIDER_ID_ALIAS_RANK[left] - PROVIDER_ID_ALIAS_RANK[right]
|
||||
|
||||
Reference in New Issue
Block a user