mirror of
https://github.com/stablyai/orca.git
synced 2026-10-07 08:02:21 +00:00
feat(native-chat): a Claude subagent waiting on a permission prompt reads as waiting (#22634)
* feat(native-chat): a Claude subagent waiting on a permission prompt reads as waiting
A subagent's permission request reaches the parent session's callback naming the
subagent that asked (agent_id) and the tool call it gates (tool_use_id). The
pending request is recorded with the asking agent. On every drain the child-work
producer re-derives which children a pending request blocks and hands that set
to the Claude child decoder, the one owner of each child's live edges: a blocked
child reads waiting on every live edge it reports, and a child that starts or
stops waiting is a live edge of its own. Answering, denying or cancelling the
request returns the child to its prior live state; nothing is stored beyond the
pending requests.
A live task_updated carrying an error now reaches the record as the child's last
message, without an ending or a new state.
The replay test drives a scrubbed capture of the real CLI (foreground allow,
deny, interrupt, background allow, and the main agent's own request) through the
real adapter into the host's child records.
* test(native-chat): a subagent's request names it before its tool call is read
* docs(agent-status): a subagent asking for approval waits in every lane; the parent row keeps the session's own attention
* test(native-chat): hand canUseTool the asking agent without widening the helper's cast
* test(native-chat): an interrupted Claude subagent settles cancelled, not failed
A captured interrupt shows the spawn call's error result ("The user doesn't want
to proceed…") arriving before the subagent's own `task_updated {status: killed}`.
The spawn result ends nothing (the child ends only on its own terminal frame), so
the child stays live until its `killed` status settles it cancelled. A genuine
failure, captured with the subagent on a model that does not exist, sends its
`failed` status before the error result and still ends failed. Both captures now
replay through the real adapter into the host's records.
* fix(native-chat): a Claude subagent's prompt makes the parent row wait, not block
A subagent's pending prompt made the whole session `attention`, which reads as
the main agent's own `blocked` and outranks the fold's waiting arm, so the
parent row read blocked where a CLI Claude parent reads waiting. The main
agent's state now reads only its own pending prompts.
- A Claude prompt row carries the linkage of the agent that raised it: the one
the permission request names, or the owner of the tool call it gates. The
same join decides which child reads waiting, so the two cannot disagree.
- The status summary projects the session's own status from root prompts only;
every other reader (delivery gates, teardown, restart) still asks whether
anyone is waiting on a human.
- An answer keeps the prompt row's linkage by the journal's own rule: a
revision that names no producer keeps the row's existing one.
- The child-tool queries gain the prompt's producer, so a prompt row and a
child record answer "which agent" from the same join.
* test(native-chat): say which ids the permission capture scrubs and which are its own
* fix(native-chat): the status clock dates attention by the session's own asks only
The session's status is now `attention` only for its own pending prompt, so the
clock's fallback to a subagent's ask could no longer be reached, and it read the
journal by a different rule than the status it dates. Both now read root prompts.
The journal also stamps a Codex subagent's prompt with its thread (#22532), so a
Codex child's approval is that child's wait in the Codex lane too. Two tests
written for the earlier rule are updated: a subagent's ask leaves a running
session `working` on its turn's clock, and a Codex child's answered approval
leaves the settled parent's Activity row done with nothing unread.
* fix(native-chat): a completion still says the user is asked when a subagent asks
The turn-completion feed marked a completion `awaitingUser` from the status
summary's `attention`. The status now means the session's own agent is waiting,
so a subagent's pending approval stopped reaching the completion. The projection
now also says whether anyone is waiting on the user, as the delivery gates,
teardown and restart ask it, and the completion reads that.
The waiting-subagent replay answers its prompt with the adapter's current
response shape.
* revert(native-chat): a live Claude task's error stays out of the child's last message
No capture shows a live task_updated carrying an error, and it is unrelated to
a subagent waiting on a permission request; it leaves this PR.
* test(native-chat): settle the Claude session's startup before replaying a subagent's request
A startup frame drained child work during the first await, so answering a
request freed the child even with the answer's own republish removed.
* fix(native-chat): a Claude subagent's prompt row names it as its other rows do
The prompt row stamped only the asking agent's id, so a nested subagent's
request lost the agent that spawned it, its spawn call and its run. It now takes
the linkage the asker's own rows take: the gated tool call's, when that names the
same agent, else the one resolved through the agent's spawn call. The provider's
agent id stays the asker's id.
* fix(native-chat): a subagent's request makes the parent row wait without a child record
The parent row learned that a subagent needed the user only from that
subagent's child record, so a request no record carried (a Codex child the host
never registered, a Claude task past the live cap) left the row working or done
while the approval card sat in the chat.
"Someone in this session must answer" is now one derived session fact. The
projection names two facts instead of a mode flag: the main agent's own status
(attention only for its own request) and structuredAgentSessionAwaitsUser (any
pending prompt). The status summary publishes the second as an optional
awaitsUser, and the shared fold reads it: the main agent's own ask is blocked,
otherwise awaitsUser or a waiting child record makes the row wait. Every caller
picks the fact it means: the completion edge's awaitingUser and the delivery
gate read awaitsUser; the quit snapshot folds the same two inputs the sidebar
does.
* fix(native-chat): a client that predates awaitsUser still reads a subagent's request as attention
A status summary's status is now the main agent's own, so a client built before
the split would read a subagent's request as working (or idle) and fold it with
code that has no awaitsUser input. Clients advertise
agent-session.status-awaits-user.v1; at agentSession.subscribeStatus the host
sends any client that does not the pre-split summary: attention whenever
awaitsUser is set, without the main agent's own tool line, verdict and clock.
The feed and every in-process reader keep the canonical summary. Transitional,
like the turn-item downgrade.
* test(native-chat): a Codex subagent's approval makes its settled parent's Activity row wait
The test pinned the parent row done while a Codex child asked, through a harness
that fed no child records, so it proved nothing about the ask. It now drives the
ask twice through the real host status store: with no child record (the
session's awaitsUser alone) and with the child's own record waiting from
thread/status/changed. Both read waiting with needsAttention while the ask is
open, then done with nothing unread.
* docs(agent-status): a subagent's request reaches the parent row through awaitsUser in every structured lane
The store reference said a Codex child's request still read as the main agent's
blocked and that only the Codex hook lane fed a waiting child. Both structured
lanes stamp the asking child and feed child records, and awaitsUser carries the
request when no record does. The liveness comment goes back to main's: a child's
blocked is a failed task on an older host's legacy rows.
* fix(native-chat): the restart dialog still headlines a subagent's pending approval
The quit snapshot now records the main agent's own state, so a subagent asking
while the main agent worked recorded `working` and the dialog said "Was
mid-reply" where it used to say "Waiting for your approval". The headline now
comes from the snapshot's pending prompt, whoever raised it, with the existing
copy; `state` stays the main agent's own.
* test(orchestration): a subagent's pending approval holds structured mail delivery
Scoping the delivery gate to the main agent's own request left every gate test
green; a subagent's request now has its own case.
* fix(native-chat): a subagent's request is dated by when it was raised, on every client
Since the summary's clock became the main agent's own, nothing dated a wait
that only a subagent's request held: a pre-split client was sent attention with
no clock, where the old host dated it by the subagent's prompt, and a new
client's waiting row fell back to the time it first saw it, so after a reload a
request the user had already read could read unread again.
The session fact is now when someone started being asked: awaitsUserSince, the
oldest pending prompt whoever raised it, and its presence is what awaitsUser
meant. A row waiting on someone else's request takes that as its clock; the
downgrade for a client without the capability dates its attention by it, which
is what the old host published. A cross-version test pinned to the last
pre-split release runs the same journals through that release's projection and
through this one plus the downgrade, and compares the whole summary. The Codex
end-to-end test also reads the host's own status row, and keeps a read ask read
through a later row and a reload.
* test(native-chat): the pre-split parity check compares only the fields the split owns
An additive summary field is safe for old clients, so comparing whole summaries
against the pinned release would redden on one. The wire comment now says how
the downgrade dates attention: the main agent's own oldest ask, else
awaitsUserSince.
* test(runtime): an aged host-held working summary states that nobody is asked
The test built its working summary by overriding the status of a published
approval summary, which still carried awaitsUserSince, so the row correctly
read waiting. It now drops the request as its scenario says.
* fix(native-chat): the chat's subagent block says waiting when the strip does
While a Claude subagent's request was open, the sidebar and the composer strip
read waiting but the subagent block in the chat history a few pixels above
still read "Kicked off 1 subagent working": it shows the journal's roster
state, and the journal records no wait.
The structured chat now hands its transcript the subagents the strip shows
waiting, read from the host's child records through the strip's own row model
and matched by the provider id the roster names each one by. A running entry
the host says is waiting reads waiting in the group row, its entry and its
section head, with the strip's word and the question colour; it reads the
journal's state again as soon as the host stops reporting the wait.
* fix(native-chat): a collapsed subagent group shows a wait beside a failed sibling
A failed sibling took the group row's one alert slot, so a group with a waiting,
a working and a failed child read "1 working +1 failed" and hid the wait; it
now reads "1 working +1 waiting +1 failed". The waiting set keeps its identity
while a child frame changes no wait, so the transcript's subagent rows do not
re-render on every frame, and the test of a wait ending now updates one mounted
row instead of remounting it.
* refactor(claude): one needs-input state on the parent; the asking subagent alone reads waiting
Drop the split of the main agent's own status from a session-wide "someone must
answer" fact: awaitsUserSince, the agent-session.status-awaits-user.v1
capability and its old-client downgrade, and every reader change that only
consumed them (fold, equality, ingest, delivery gate, turn-completion feed,
quit snapshot, resume headline, status clock, status bridge, attention
dispatch) go back to main. The parent row again reads one needs-input state
for a pending request whoever asked, dated as before.
Kept: a request's owner recorded once on its prompt row with full producer
linkage; the asking subagent's own record reads waiting, re-derived on every
update; the chat history's subagent block reads that same state; an answered
subagent request stays in its subagent's group.
A subagent now waits only on a request the user can still answer (its card
open, no answer underway), and the adapter frees it before the host records an
answer or dismissal. So a waiting child record always sits beside the pending
card, and main's fold never reads the parent as waiting on it: no window after
an answer, and no ~3 s wait after a card dismissed by Stop.
* fix(claude): a subagent waits only beside its committed card
A subagent's wait was pushed to the host as soon as its request arrived,
while the request's card row reached the journal at least a microtask later.
So every subagent request published the parent row as waiting before
blocked (the main agent's own fold reads a waiting child that way), and
Activity got an extra unread "waiting" event that main never shows.
The card is now the one record of an open request. The translator records
the asker on the card once (its row's linkage) and counts the card open only
after the sink confirms its rows landed, then publishes the wait; anything
that closes the card (an answer underway, a dismissal handed to the host,
Claude's own withdrawal, the session's end) frees the subagent first. So
every publish that shows a subagent waiting also shows its pending card, and
the parent reads one needs-input state, exactly as on main.
This retires the registry's view of pending requests (unclaimed(), the
asking-child join) and the translator's holdsOpen. The prompt row's linkage
takes one rule: the agent the provider names, else the gated call's owner.
The parent-row proof now runs through the real deferred sink, durable
journal and status feed, publishing as production does, and checks at every
publish that waiting subagents have pending cards and that the parent row
matches a host fed no waits.
* fix(native-chat): a closed sink's dropped writes never read as landed
The sink's written() resolved ok when the sink was closed with writes still
queued, so a subagent's prompt card could count as open with no row in the
journal. written() now reports a close that dropped writes admitted so far
as not landed; drained() and lifecycleBarrier() keep reading a closed sink as
settled.
Tests: a card never opens when its sink closes first; a card Claude withdraws
while the sink holds the cancelled row back closes at once; two subagents
asking at once, and the main agent asking beside a subagent, keep the parent
row as before with each waiting subagent beside its own card; a process that
dies mid-request leaves no subagent waiting.
* fix(claude): a withdrawn subagent request frees its child before its card closes
Main's sink now hands each write to the journal as it is submitted, and an idle journal commits it
and runs the publication at once. Claude's own withdrawal of a subagent's request therefore closed
the card and published the parent row before the child's wait was freed, so one publish showed the
subagent waiting beside no pending card (fg-interrupt replay). The child's wait now also requires the
request to still be open in the registry, and a withdrawal republishes child work before the
journal takes the close.
* refactor(claude): trim subagent request waiting to the common pattern and its essential tests
The chat history no longer marks a subagent block as waiting: the approval card itself carries the
request, and the asking subagent's row in the sidebar and composer strip reads waiting, as before.
NativeChatWaitingSubagentsProvider, native-chat-waiting-subagents.ts and their renderer changes go.
A subagent waits while its request is still open and unanswered in the prompt registry and its card
has landed in the journal. The registry check also covers a withdrawal under backpressure, so the
card list no longer filters pending cancellations itself.
Tests: one integration file replays the captured CLI frames through the real adapter, sink, journal
and status feed (renamed claude-subagent-permission-request.test.ts), with the asking subagent's
state timeline, attribution, nested linkage and a card write that waits for the journal. The
producer-harness waiting test, the redundant prompt-card cases, the harness reducer swap and three
unused captures (deny, interrupt, main agent, failed subagent) are removed.
* test(claude): pin the parent row's dating when a subagent asked first, and narrow the oracle's claim
This commit is contained in:
@@ -288,7 +288,10 @@ Every lane, Codex included, combines through the fold. A child waiting on a
|
||||
human is a fold input (`childWorkLiveness: 'waiting'`, derived from the child's
|
||||
own `waiting` state; a child's `blocked` means it failed and stays live work)
|
||||
and makes the row wait whatever the main agent is doing, unless the main agent
|
||||
is itself asking. Only the Codex hook lane feeds that input today. Known
|
||||
is itself asking. The Codex hook lane feeds it from its child transcripts, and
|
||||
the structured lanes from child records, which read `waiting` for a Codex child
|
||||
thread's approval or input flag and for a Claude subagent's open permission
|
||||
request. Known
|
||||
divergences, pinned by name in the parity table
|
||||
(`src/shared/main-agent-status-parity.test.ts`) where they are reachable, so a
|
||||
reader does not mistake them for drift:
|
||||
@@ -299,8 +302,13 @@ reader does not mistake them for drift:
|
||||
main agent event overwrites the slot, so the row stops reading `waiting`
|
||||
while the child is still asking, and a second asking child replaces the
|
||||
first.
|
||||
- The structured lane has no per-child wait: a child's pending prompt makes
|
||||
the session `attention`, which reads as the main agent's own `blocked`.
|
||||
- In the structured lane a child's pending prompt also makes the session
|
||||
`attention`, which reads as the main agent's own `blocked`: one needs-input
|
||||
state whoever asked. A Claude subagent reads `waiting` only while the
|
||||
journal holds its card pending: from after the card's row is written until
|
||||
just before anyone closes it, so every publish that shows the child waiting
|
||||
also shows the session's `attention`, and the row never reads `waiting` for
|
||||
a Claude subagent's request.
|
||||
- The Codex hook lane drops its roster on a root `Stop` when it tracks no
|
||||
child transcripts, so a still-running or still-asking child stops holding
|
||||
the row.
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -12,6 +12,7 @@ export function invokeCanUseTool(
|
||||
input?: Record<string, unknown>
|
||||
suggestions?: unknown[]
|
||||
signal?: AbortSignal
|
||||
agentID?: string
|
||||
} = {}
|
||||
): { promise: Promise<unknown>; settled: () => boolean } {
|
||||
const options = {
|
||||
@@ -20,9 +21,10 @@ export function invokeCanUseTool(
|
||||
signal: extra.signal ?? new AbortController().signal,
|
||||
...(extra.suggestions ? { suggestions: extra.suggestions } : {})
|
||||
} as unknown as Parameters<NonNullable<ClaudeStreamJsonConnectionHandlers['canUseTool']>>[2]
|
||||
const asked = extra.agentID ? { ...options, agentID: extra.agentID } : options
|
||||
let done = false
|
||||
const promise = Promise.resolve(
|
||||
connection.handlers.canUseTool?.(toolName, extra.input ?? {}, options)
|
||||
connection.handlers.canUseTool?.(toolName, extra.input ?? {}, asked)
|
||||
).finally(() => {
|
||||
done = true
|
||||
})
|
||||
|
||||
@@ -1,8 +1,13 @@
|
||||
// Which agent a tool call or a frame belongs to, answered from the journal's own linkage, so a
|
||||
// child's record and its open operation name the agent its rows name.
|
||||
// Which agent a tool call, a frame or a prompt belongs to, answered from the journal's own linkage,
|
||||
// so a child's record, its open operation and its prompt rows all name the agent its rows name.
|
||||
|
||||
import type { AgentJournalProducerLinkage } from '../../shared/agent-session-journal-types'
|
||||
import type { ClaudePendingPrompt } from './claude-prompt-registry'
|
||||
import type { ClaudeToolUse } from './claude-structured-item-translation'
|
||||
import type { ClaudeSubagentLinkageSource } from './claude-subagent-linkage'
|
||||
import type {
|
||||
ClaudeAgentLinkageSource,
|
||||
ClaudeSubagentLinkageSource
|
||||
} from './claude-subagent-linkage'
|
||||
import type { ClaudeToolOriginRegistry } from './claude-tool-origin-registry'
|
||||
|
||||
export type ClaudeChildToolQueries = {
|
||||
@@ -11,12 +16,17 @@ export type ClaudeChildToolQueries = {
|
||||
childToolOwner: (toolUseId: string) => string | null
|
||||
/** The child a frame's `parent_tool_use_id` names, and its newest call still awaiting a result. */
|
||||
childActivity: (parentToolUseId: string) => { agentId: string; openTool: ClaudeToolUse | null }
|
||||
/** The linkage a prompt row carries, as the asking agent's other rows carry it; none for the
|
||||
* session's own agent. */
|
||||
promptProducer: (
|
||||
prompt: Pick<ClaudePendingPrompt, 'agentId' | 'toolUseId'>
|
||||
) => AgentJournalProducerLinkage
|
||||
}
|
||||
|
||||
export function claudeChildToolQueries(deps: {
|
||||
tools: ReadonlyMap<string, ClaudeToolUse>
|
||||
toolOrigins: Pick<ClaudeToolOriginRegistry, 'childOwnerRef'>
|
||||
linkage: Pick<ClaudeSubagentLinkageSource, 'settledLinkageFor'>
|
||||
linkage: Pick<ClaudeSubagentLinkageSource, 'settledLinkageFor'> & ClaudeAgentLinkageSource
|
||||
}): ClaudeChildToolQueries {
|
||||
const childToolOwner = (toolUseId: string): string | null => {
|
||||
const ownerRef = deps.toolOrigins.childOwnerRef(toolUseId)
|
||||
@@ -35,6 +45,14 @@ export function claudeChildToolQueries(deps: {
|
||||
}
|
||||
const { agentId } = deps.linkage.settledLinkageFor(parentToolUseId).linkage
|
||||
return { agentId: agentId ?? parentToolUseId, openTool }
|
||||
},
|
||||
// The provider names the asker when it can; otherwise the gated call's owner is the asker.
|
||||
promptProducer: (prompt) => {
|
||||
if (prompt.agentId !== undefined) {
|
||||
return deps.linkage.linkageForAgent(prompt.agentId)
|
||||
}
|
||||
const ownerRef = deps.toolOrigins.childOwnerRef(prompt.toolUseId)
|
||||
return ownerRef === null ? {} : deps.linkage.settledLinkageFor(ownerRef).linkage
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -212,3 +212,51 @@ describe('Claude child-work decoder', () => {
|
||||
expect(decoder.drain(600)).toEqual([expect.not.objectContaining({ restart: true })])
|
||||
})
|
||||
})
|
||||
|
||||
describe('Claude children waiting on a permission request', () => {
|
||||
const states = (edges: ReturnType<ClaudeChildWorkDecoder['drain']>) =>
|
||||
edges.flatMap((edge) =>
|
||||
edge.type === 'live' ? [`${edge.child.handle.id} ${edge.child.state}`] : [edge.type]
|
||||
)
|
||||
|
||||
it('reports a live child that starts or stops waiting, and only then', () => {
|
||||
const decoder = decoderWith(backgroundAgent)
|
||||
decoder.observeWaiting(new Set(['agent-bg', 'agent-untracked']))
|
||||
expect(states(decoder.drain(500))).toEqual(['agent-bg waiting'])
|
||||
decoder.observeWaiting(new Set(['agent-bg', 'agent-untracked']))
|
||||
expect(decoder.drain(600)).toEqual([])
|
||||
decoder.observeWaiting(new Set())
|
||||
expect(states(decoder.drain(700))).toEqual(['agent-bg working'])
|
||||
})
|
||||
|
||||
it('reads every live report of a waiting child as waiting, a start included', () => {
|
||||
const decoder = new ClaudeChildWorkDecoder()
|
||||
// The request can name a child before its start is read.
|
||||
decoder.observeWaiting(new Set(['agent-bg']))
|
||||
decoder.observe(backgroundAgent)
|
||||
decoder.observe(foregroundAgent)
|
||||
decoder.observe(system('task_progress', { task_id: 'agent-bg', last_tool_name: 'Bash' }))
|
||||
expect(states(decoder.drain(500))).toEqual([
|
||||
'agent-bg waiting',
|
||||
'agent-fg working',
|
||||
'agent-bg waiting'
|
||||
])
|
||||
})
|
||||
|
||||
it('settles a waiting child on its own ending, and forgets every wait with the session', () => {
|
||||
const decoder = decoderWith(backgroundAgent)
|
||||
decoder.observeWaiting(new Set(['agent-bg']))
|
||||
decoder.observe(system('task_notification', { task_id: 'agent-bg', status: 'stopped' }))
|
||||
expect(states(decoder.drain(500))).toEqual(['agent-bg waiting', 'ended'])
|
||||
// Nothing is live to stop waiting.
|
||||
decoder.observeWaiting(new Set())
|
||||
expect(decoder.drain(600)).toEqual([])
|
||||
decoder.observe(system('task_started', { ...backgroundAgent, tool_use_id: 'toolu_bg_2' }))
|
||||
decoder.observeWaiting(new Set(['agent-bg']))
|
||||
decoder.clear()
|
||||
// The session's end frees every wait: nothing its drain carries reads waiting.
|
||||
const edges = states(decoder.drain(700))
|
||||
expect(edges.at(-1)).toBe('session-ended')
|
||||
expect(edges).not.toContain('agent-bg waiting')
|
||||
})
|
||||
})
|
||||
|
||||
@@ -5,7 +5,9 @@
|
||||
// or a spawn call returning is the parent's view of the child, not the child's, and the CLI sends
|
||||
// every child its own terminal frame, so none of them settles one. When Orca ends the session and
|
||||
// proves its tree gone, what is still live is stopped (`stopLive`); any other end leaves it for the
|
||||
// host to settle as unknown. Edges wait here until the frame is journaled, then take the host clock.
|
||||
// host to settle as unknown. A live child blocked on a permission request reads waiting; no task
|
||||
// frame says so, so the caller hands over which children a request blocks. Edges wait here until
|
||||
// the frame is journaled, then take the host clock.
|
||||
|
||||
import type {
|
||||
AgentChildWorkKind,
|
||||
@@ -55,7 +57,8 @@ function observation(
|
||||
id: string,
|
||||
task: DecodedClaudeTask,
|
||||
facts: ClaudeTaskFacts,
|
||||
observedAt: number
|
||||
observedAt: number,
|
||||
waiting: boolean
|
||||
): AgentChildWorkLiveObservation {
|
||||
const operation: AgentChildWorkOperation | undefined = facts.toolName
|
||||
? { toolName: facts.toolName, basis: 'reported', observedAt }
|
||||
@@ -68,7 +71,7 @@ function observation(
|
||||
},
|
||||
kind: task.kind,
|
||||
residency: task.backgrounded ? 'background' : 'foreground',
|
||||
state: task.running || task.kind !== 'monitor' ? 'working' : 'monitoring',
|
||||
state: waiting ? 'waiting' : task.running || task.kind !== 'monitor' ? 'working' : 'monitoring',
|
||||
// The published row names a task's type as both its name and its agent type.
|
||||
...(task.name ? { name: task.name, agentType: task.name } : {}),
|
||||
...(task.description ? { description: task.description } : {}),
|
||||
@@ -84,6 +87,8 @@ export class ClaudeChildWorkDecoder {
|
||||
private readonly live = new Map<string, DecodedClaudeTask>()
|
||||
/** Ended task ids, with the spawn call each ended under. */
|
||||
private readonly ended = new Map<string, string | undefined>()
|
||||
/** The children a pending permission request blocks, as the caller last derived them. */
|
||||
private waiting: ReadonlySet<string> = new Set()
|
||||
private pending: PendingEdge[] = []
|
||||
|
||||
observe(message: Record<string, unknown>): void {
|
||||
@@ -124,6 +129,19 @@ export class ClaudeChildWorkDecoder {
|
||||
}
|
||||
}
|
||||
|
||||
/** Which children a pending request blocks, re-derived by the caller before every drain. A
|
||||
* live child that starts or stops waiting is a live edge of its own. */
|
||||
observeWaiting(waiting: ReadonlySet<string>): void {
|
||||
const previous = this.waiting
|
||||
this.waiting = waiting
|
||||
for (const id of new Set([...previous, ...waiting])) {
|
||||
const task = this.live.get(id)
|
||||
if (task && previous.has(id) !== waiting.has(id)) {
|
||||
this.report(id, task, {})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Orca ended the session and proved its process tree gone: what still ran is stopped. The
|
||||
* ending is Orca's, not the child's, so a frame of the child's own still replaces it. */
|
||||
stopLive(): void {
|
||||
@@ -136,6 +154,7 @@ export class ClaudeChildWorkDecoder {
|
||||
clear(): void {
|
||||
this.live.clear()
|
||||
this.ended.clear()
|
||||
this.waiting = new Set()
|
||||
this.pending.push((observedAt) => ({ type: 'session-ended', observedAt }))
|
||||
}
|
||||
|
||||
@@ -237,10 +256,15 @@ export class ClaudeChildWorkDecoder {
|
||||
}
|
||||
this.ended.delete(id)
|
||||
this.live.set(id, task)
|
||||
this.report(id, task, facts, restart)
|
||||
}
|
||||
|
||||
/** Whether the child waits is read at drain, so every live edge a drain carries agrees. */
|
||||
private report(id: string, task: DecodedClaudeTask, facts: ClaudeTaskFacts, restart = false) {
|
||||
this.pending.push((observedAt) => ({
|
||||
type: 'live',
|
||||
observedAt,
|
||||
child: observation(id, task, facts, observedAt),
|
||||
child: observation(id, task, facts, observedAt, this.waiting.has(id)),
|
||||
...(restart ? { restart: true } : {})
|
||||
}))
|
||||
}
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
// Claude child work as the host records it: the outcome vocabulary, what a progress frame says,
|
||||
// the owner of each child through the journal's own linkage, and the tool a child has open.
|
||||
// the owner of each child through the journal's own linkage, which children a pending request
|
||||
// blocks, and the tool a child has open.
|
||||
// The task frames themselves are read by `claude-child-work-decoder`; everything here is drained
|
||||
// after the journal handled the frame, so the host never admits evidence ahead of its rows.
|
||||
|
||||
@@ -9,7 +10,7 @@ import type {
|
||||
AgentChildWorkLiveObservation
|
||||
} from '../../shared/agent-status-child-work-evidence'
|
||||
import { taskText, taskUsageTotalTokens } from './claude-background-task-frames'
|
||||
import type { ClaudeSession } from './claude-structured-session-state'
|
||||
import type { ClaudeSession, ClaudeStructuredSessionEvent } from './claude-structured-session-state'
|
||||
import { deriveToolInputPreview } from '../../shared/agent-hook-listener/tool-input-preview'
|
||||
import {
|
||||
claudeRecord,
|
||||
@@ -114,15 +115,43 @@ export function claudeChildOperation(
|
||||
]
|
||||
}
|
||||
|
||||
/** The subagents whose request is still open with no answer underway (the registry's fact) and
|
||||
* whose card has landed (the journal's): a waiting child always sits beside its pending card, and
|
||||
* its parent reads that card, never the child's wait. A request is claimed or forgotten in the
|
||||
* registry before its card closes, so however it ends its child is freed first. */
|
||||
function claudeWaitingChildIds(
|
||||
session: Pick<ClaudeSession, 'prompts' | 'translator'>
|
||||
): Set<string> {
|
||||
const waiting = new Set<string>()
|
||||
for (const card of session.translator?.journalPrompts.openCards() ?? []) {
|
||||
if (session.prompts.awaitsAnswer(card.promptKey)) {
|
||||
waiting.add(card.asker)
|
||||
}
|
||||
}
|
||||
return waiting
|
||||
}
|
||||
|
||||
/** Settles once the card a prompt event raised is written, for a sink that writes it later. */
|
||||
export function claudePromptCardWritten(
|
||||
session: Pick<ClaudeSession, 'translator'> | null | undefined,
|
||||
event: ClaudeStructuredSessionEvent
|
||||
): Promise<void> | undefined {
|
||||
return event.type === 'prompt'
|
||||
? session?.translator?.journalPrompts.whenWritten(event.prompt.promptKey)
|
||||
: undefined
|
||||
}
|
||||
|
||||
/** Everything one frame (or a close) said about the session's child work, owners named. */
|
||||
export function drainClaudeChildWork(
|
||||
session: Pick<ClaudeSession, 'childWork' | 'translator'> | null | undefined,
|
||||
session: Pick<ClaudeSession, 'childWork' | 'translator' | 'prompts'> | null | undefined,
|
||||
message: Record<string, unknown> | null,
|
||||
observedAt: number
|
||||
): AgentChildWorkEvidence[] {
|
||||
if (!session) {
|
||||
return []
|
||||
}
|
||||
// Re-derived from the open cards on every drain, so no wait outlives its card.
|
||||
session.childWork.observeWaiting(claudeWaitingChildIds(session))
|
||||
return [
|
||||
...withClaudeChildWorkOwners(
|
||||
session.childWork.drain(observedAt),
|
||||
|
||||
@@ -10,7 +10,10 @@ import type { ClaudeCommandStart } from './claude-command-turn'
|
||||
|
||||
export type ClaudeJournalTranslator = {
|
||||
handle: (event: ClaudeStructuredSessionEvent) => void
|
||||
journalPrompts: Pick<ClaudeJournalPrompts, 'resolve' | 'handOver' | 'cancel'>
|
||||
journalPrompts: Pick<
|
||||
ClaudeJournalPrompts,
|
||||
'resolve' | 'handOver' | 'cancel' | 'openCards' | 'whenWritten'
|
||||
>
|
||||
/** The open turn's provider id — the same id its journal row carries, and the one
|
||||
* a client's Stop names. Sole owner: no reader keeps a copy to disagree with. */
|
||||
readonly currentTurnId: string | null
|
||||
|
||||
@@ -27,6 +27,8 @@ export type ClaudePendingPrompt = ClaudePromptPresentation & {
|
||||
suggestions: PermissionUpdate[]
|
||||
questionIds: readonly string[]
|
||||
settle: ClaudePromptSettle
|
||||
/** The subagent the provider says asked; absent when the session's own agent did. */
|
||||
agentId?: string
|
||||
}
|
||||
|
||||
export type ClaudePromptRegistration = ClaudePromptPresentation & {
|
||||
@@ -36,6 +38,7 @@ export type ClaudePromptRegistration = ClaudePromptPresentation & {
|
||||
input: Record<string, unknown>
|
||||
suggestions: PermissionUpdate[]
|
||||
settle: ClaudePromptSettle
|
||||
agentId?: string
|
||||
}
|
||||
|
||||
export type ClaudePromptClaim = {
|
||||
@@ -74,6 +77,7 @@ export class ClaudePromptRegistry {
|
||||
const toolUseId = readClaudePromptString(registration.toolUseId)
|
||||
const toolName = readClaudePromptString(registration.toolName)
|
||||
const input = isClaudePromptRecord(registration.input) ? registration.input : null
|
||||
const agentId = readClaudePromptString(registration.agentId)
|
||||
if (!toolUseId || !toolName || !input) {
|
||||
return null
|
||||
}
|
||||
@@ -94,12 +98,19 @@ export class ClaudePromptRegistry {
|
||||
...(registration.matchedAskRule ? { matchedAskRule: registration.matchedAskRule } : {}),
|
||||
...(registration.subject ? { subject: registration.subject } : {}),
|
||||
questionIds: questions.map(questionId),
|
||||
settle: registration.settle
|
||||
settle: registration.settle,
|
||||
...(agentId ? { agentId } : {})
|
||||
}
|
||||
this.prompts.set(prompt.promptKey, prompt)
|
||||
return prompt
|
||||
}
|
||||
|
||||
/** The request is still open and nobody is answering it yet. */
|
||||
awaitsAnswer(promptKey: string): boolean {
|
||||
const prompt = this.prompts.get(promptKey)
|
||||
return prompt !== undefined && !this.claims.has(prompt)
|
||||
}
|
||||
|
||||
/** True only if the prompt was still pending; lets abort and answer settle once. */
|
||||
forgetIfPending(prompt: ClaudePendingPrompt): boolean {
|
||||
if (!this.prompts.has(prompt.promptKey)) {
|
||||
|
||||
@@ -240,7 +240,9 @@ describe('answerClaudePrompt', () => {
|
||||
journalPrompts: {
|
||||
resolve: resolvePrompt,
|
||||
handOver: () => () => {},
|
||||
cancel: () => ({ accepted: true })
|
||||
cancel: () => ({ accepted: true }),
|
||||
openCards: () => [][Symbol.iterator](),
|
||||
whenWritten: () => undefined
|
||||
},
|
||||
currentTurnId: null,
|
||||
commandTurnId: null,
|
||||
|
||||
@@ -51,7 +51,8 @@ export function buildClaudePermissionCallbacks(deps: ClaudePermissionCallbackDep
|
||||
toolUseId: options.toolUseID,
|
||||
input,
|
||||
suggestions: options.suggestions ?? [],
|
||||
settle
|
||||
settle,
|
||||
...(options.agentID ? { agentId: options.agentID } : {})
|
||||
})
|
||||
if (!prompt) {
|
||||
settle(denySafeResult(options.toolUseID))
|
||||
|
||||
@@ -0,0 +1,79 @@
|
||||
// A subagent's prompt card counts as open once its rows land, until anyone closes or takes it over.
|
||||
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import { AGENT_JOURNAL_THREAD_SCOPE } from '../../shared/agent-session-journal-types'
|
||||
import type { StructuredAgentSessionSinkBarrier } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
|
||||
import { ClaudeJournalPrompts } from './claude-structured-journal-prompts'
|
||||
import type { ClaudeStructuredSessionEvent } from './claude-structured-session-state'
|
||||
|
||||
const PROMPT: Extract<ClaudeStructuredSessionEvent, { type: 'prompt' }> = {
|
||||
type: 'prompt',
|
||||
sessionId: 'session-1',
|
||||
prompt: {
|
||||
requestId: 'req-1',
|
||||
promptKey: 'req-1',
|
||||
toolUseId: 'toolu-1',
|
||||
toolName: 'Bash',
|
||||
kind: 'approval',
|
||||
input: { command: 'touch f' },
|
||||
suggestions: [],
|
||||
questionIds: [],
|
||||
settle: vi.fn()
|
||||
}
|
||||
}
|
||||
|
||||
function cards(options: { accepted?: boolean; asker?: string } = {}) {
|
||||
let land: (barrier: StructuredAgentSessionSinkBarrier) => void = () => {}
|
||||
const prompts = new ClaudeJournalPrompts({
|
||||
sink: {
|
||||
appendItem: () => {},
|
||||
tryAppendItem: () =>
|
||||
options.accepted === false
|
||||
? { accepted: false, reason: 'backpressure' }
|
||||
: { accepted: true },
|
||||
appendTombstone: () => {},
|
||||
publish: () => {},
|
||||
written: () =>
|
||||
new Promise((resolve) => {
|
||||
land = resolve
|
||||
})
|
||||
},
|
||||
turnScope: () => AGENT_JOURNAL_THREAD_SCOPE,
|
||||
producerOf: () =>
|
||||
options.asker === undefined ? {} : { agentId: options.asker, producerKind: 'agent' }
|
||||
})
|
||||
prompts.handle(PROMPT)
|
||||
const open = () => [...prompts.openCards()]
|
||||
const landed = async (barrier: StructuredAgentSessionSinkBarrier = { ok: true }) => {
|
||||
const written = prompts.whenWritten('req-1')
|
||||
land(barrier)
|
||||
await written
|
||||
}
|
||||
return {
|
||||
prompts,
|
||||
open,
|
||||
landed,
|
||||
land: (barrier: StructuredAgentSessionSinkBarrier) => land(barrier)
|
||||
}
|
||||
}
|
||||
|
||||
describe("a subagent's prompt card", () => {
|
||||
it('never opens when the sink refused its rows or failed writing them', async () => {
|
||||
const refused = cards({ asker: 'agent-1', accepted: false })
|
||||
await refused.landed()
|
||||
expect(refused.open()).toEqual([])
|
||||
const failed = cards({ asker: 'agent-1' })
|
||||
await failed.landed({ ok: false, error: new Error('write failed') })
|
||||
expect(failed.open()).toEqual([])
|
||||
})
|
||||
|
||||
it('opens a card handed back before its rows landed once they land', async () => {
|
||||
const { prompts, open, land } = cards({ asker: 'agent-1' })
|
||||
const written = prompts.whenWritten('req-1')
|
||||
const handBack = prompts.handOver('req-1')
|
||||
handBack()
|
||||
land({ ok: true })
|
||||
await written
|
||||
expect(open()).toEqual([{ promptKey: 'req-1', asker: 'agent-1' }])
|
||||
})
|
||||
})
|
||||
@@ -1,6 +1,7 @@
|
||||
import type {
|
||||
AgentJournalApprovalItem,
|
||||
AgentJournalItemIdentity,
|
||||
AgentJournalProducerLinkage,
|
||||
AgentJournalQuestionItem,
|
||||
AgentJournalTurnScope
|
||||
} from '../../shared/agent-session-journal-types'
|
||||
@@ -29,8 +30,17 @@ type ClaudeJournalPrompt = {
|
||||
type ClaudeJournalPromptEntry = {
|
||||
items: ClaudeJournalPrompt[]
|
||||
cancellationPending: boolean
|
||||
/** The subagent that raised it, as its rows name it; absent for the session's own agent. */
|
||||
asker?: string
|
||||
/** Its rows have landed in the journal, so a reader of the journal sees the card. */
|
||||
written: boolean
|
||||
/** Settles once `written` is decided, for a sink that writes later. */
|
||||
landed?: Promise<void>
|
||||
}
|
||||
|
||||
/** A card the journal holds pending, by the subagent that raised it. */
|
||||
export type ClaudeOpenPromptCard = { promptKey: string; asker: string }
|
||||
|
||||
function cancelledPromptBody(
|
||||
body: AgentJournalApprovalItem | AgentJournalQuestionItem
|
||||
): AgentJournalApprovalItem | AgentJournalQuestionItem {
|
||||
@@ -63,30 +73,38 @@ export class ClaudeJournalPrompts {
|
||||
sessionId: string
|
||||
prompt: Extract<ClaudeStructuredSessionEvent, { type: 'prompt' }>['prompt']
|
||||
}) => ClaudeQuestionItem[]
|
||||
/** The agent that raised the prompt, as a row's producer linkage; empty for the session's own. */
|
||||
producerOf?: (
|
||||
prompt: Extract<ClaudeStructuredSessionEvent, { type: 'prompt' }>['prompt']
|
||||
) => AgentJournalProducerLinkage
|
||||
}
|
||||
) {}
|
||||
|
||||
/**
|
||||
* Prompt rows carry NO producer linkage, and cannot.
|
||||
*
|
||||
* A prompt is not a transcript frame: it reaches Orca through the SDK's
|
||||
* permission callback, whose options carry a request id and the tool awaiting
|
||||
* approval and no parent reference of any kind. So when a subagent asks, the
|
||||
* row cannot name it — unattributable at this site, not deliberately root.
|
||||
*
|
||||
* No reader is wrong because of it. A pending prompt projects the session as
|
||||
* `attention` whoever raised it, which is the truth: the USER has to answer.
|
||||
* A prompt row carries the linkage of the agent that raised it: the permission callback names the
|
||||
* subagent that asked, or the tool call it gates names one. The pending row still makes the session
|
||||
* `attention` whoever asked; the linkage files the card under that subagent.
|
||||
*/
|
||||
handle(event: Extract<ClaudeStructuredSessionEvent, { type: 'prompt' }>): void {
|
||||
const producer = this.deps.producerOf?.(event.prompt) ?? {}
|
||||
const items: ClaudeJournalPrompt[] = []
|
||||
const turnScope = this.deps.turnScope()
|
||||
let admitted = true
|
||||
const append = (identity: AgentJournalItemIdentity, body: ClaudeJournalPrompt['body']) => {
|
||||
const options = { ...producer, turnScope }
|
||||
if (this.deps.sink.tryAppendItem) {
|
||||
admitted &&= this.deps.sink.tryAppendItem(identity, body, options).accepted
|
||||
} else {
|
||||
this.deps.sink.appendItem(identity, body, options)
|
||||
}
|
||||
}
|
||||
if (event.prompt.kind === 'question') {
|
||||
for (const question of (this.deps.questionItems ?? claudeQuestionItems)({
|
||||
sessionId: event.sessionId,
|
||||
prompt: event.prompt
|
||||
})) {
|
||||
items.push({ ...question, turnScope })
|
||||
this.deps.sink.appendItem(question.identity, question.body, { turnScope })
|
||||
append(question.identity, question.body)
|
||||
this.deps.bindPromptItemId?.(agentJournalItemKey(question.identity), event.prompt.promptKey)
|
||||
}
|
||||
} else {
|
||||
@@ -96,12 +114,39 @@ export class ClaudeJournalPrompts {
|
||||
})
|
||||
const body = claudeApprovalItem(event.prompt)
|
||||
items.push({ identity, body, turnScope })
|
||||
this.deps.sink.appendItem(identity, body, { turnScope })
|
||||
append(identity, body)
|
||||
this.deps.bindPromptItemId?.(agentJournalItemKey(identity), event.prompt.promptKey)
|
||||
}
|
||||
this.deletePrompt(event.prompt.promptKey)
|
||||
this.items.set(event.prompt.promptKey, { items, cancellationPending: false })
|
||||
const entry: ClaudeJournalPromptEntry = {
|
||||
items,
|
||||
cancellationPending: false,
|
||||
...(producer.agentId ? { asker: producer.agentId } : {}),
|
||||
written: false
|
||||
}
|
||||
this.items.set(event.prompt.promptKey, entry)
|
||||
this.deps.sink.publish()
|
||||
if (admitted) {
|
||||
this.markWritten(entry)
|
||||
}
|
||||
}
|
||||
|
||||
/** A sink that cannot say when its writes land wrote them already. */
|
||||
private markWritten(entry: ClaudeJournalPromptEntry): void {
|
||||
const written = this.deps.sink.written?.()
|
||||
if (!written) {
|
||||
entry.written = true
|
||||
return
|
||||
}
|
||||
entry.landed = written.then((barrier) => {
|
||||
entry.written = barrier.ok
|
||||
})
|
||||
}
|
||||
|
||||
/** Settles once the card's rows have landed, or proved they never will; nothing for a card a
|
||||
* sink wrote at once, or one already closed. */
|
||||
whenWritten(promptKey: string): Promise<void> | undefined {
|
||||
return this.items.get(promptKey)?.landed
|
||||
}
|
||||
|
||||
private admitCancellation(promptKey: string): StructuredAgentSessionSinkAdmission {
|
||||
@@ -202,6 +247,15 @@ export class ClaudeJournalPrompts {
|
||||
this.deletePrompt(promptKey)
|
||||
}
|
||||
|
||||
/** Subagents' cards whose rows have landed and that nobody has closed or taken over yet. */
|
||||
*openCards(): IterableIterator<ClaudeOpenPromptCard> {
|
||||
for (const [promptKey, entry] of this.items) {
|
||||
if (entry.written && entry.asker !== undefined) {
|
||||
yield { promptKey, asker: entry.asker }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** The host records the card itself, so nothing here writes it any more. The returned undo hands
|
||||
* it back when that record fails, so Claude's own withdrawal can still close it. */
|
||||
handOver(promptKey: string): () => void {
|
||||
@@ -209,7 +263,9 @@ export class ClaudeJournalPrompts {
|
||||
this.deletePrompt(promptKey)
|
||||
return () => {
|
||||
if (entry && !this.items.has(promptKey)) {
|
||||
this.items.set(promptKey, { items: entry.items, cancellationPending: false })
|
||||
// The same entry, so a write still landing marks the card it hands back.
|
||||
entry.cancellationPending = false
|
||||
this.items.set(promptKey, entry)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -67,7 +67,11 @@ export function createClaudeJournalTranslator(
|
||||
const tools = new Map<string, ClaudeToolUse>()
|
||||
// Every row joins the root turn open when it is written, whoever produced it.
|
||||
const turnScope = () => turn.turnScope
|
||||
const prompts = new ClaudeJournalPrompts({ ...deps, turnScope })
|
||||
const prompts = new ClaudeJournalPrompts({
|
||||
...deps,
|
||||
turnScope,
|
||||
producerOf: (prompt) => childQueries.promptProducer(prompt)
|
||||
})
|
||||
const streamedBlocks = createClaudeStreamedBlockRegistry()
|
||||
const turn = new ClaudeOpenTurn({
|
||||
sink: deps.sink,
|
||||
|
||||
@@ -36,7 +36,7 @@ import {
|
||||
} from './claude-structured-session-exit-lifecycle'
|
||||
import type { AgentSessionBackgroundTaskState } from '../../shared/agent-session-wire'
|
||||
import { resolveClaudeProviderHistoryWindow } from './claude-structured-history-window'
|
||||
import { drainClaudeChildWork } from './claude-child-work-evidence'
|
||||
import { claudePromptCardWritten, drainClaudeChildWork } from './claude-child-work-evidence'
|
||||
import {
|
||||
answerClaudeStructuredPrompt,
|
||||
cancelClaudeStructuredTurn,
|
||||
@@ -152,6 +152,10 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
|
||||
} else if (event.type === 'message') {
|
||||
session?.childWork.observe(event.message)
|
||||
session?.backgroundTasks.observe(event.message, event.startsTurn === true)
|
||||
} else if (event.type === 'prompt-cancelled') {
|
||||
// A withdrawn request frees its child before its card closes: the journal may take that
|
||||
// write, and publish it, as it is submitted.
|
||||
this.publishChildWork(event.sessionId, session)
|
||||
}
|
||||
if (event.type === 'message' && session?.commands.observe(event.message)) {
|
||||
session.events?.publish()
|
||||
@@ -159,13 +163,15 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
|
||||
session?.translator?.handle(event)
|
||||
this.deps.onEvent?.(event)
|
||||
this.publishChildWork(event.sessionId, session, event.type === 'message' ? event.message : null)
|
||||
// A subagent's card holds it waiting only once its row is written: its wait goes out after.
|
||||
void claudePromptCardWritten(session, event)?.then(() => this.publishChildWork(event.sessionId))
|
||||
}
|
||||
|
||||
/** After the journal handled the frame, which republished the parent's own row: the host never
|
||||
* holds a child record ahead of the rows that frame wrote, and never before its parent. */
|
||||
private publishChildWork(
|
||||
sessionId: string,
|
||||
session: ClaudeSession | null | undefined,
|
||||
session: ClaudeSession | null | undefined = this.sessions.get(sessionId),
|
||||
message: Record<string, unknown> | null = null
|
||||
): void {
|
||||
const evidence = drainClaudeChildWork(session, message, this.deps.now?.() ?? Date.now())
|
||||
@@ -197,8 +203,23 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
|
||||
stopEndsSession = (): boolean => true
|
||||
awaitStoppedRequestEnd = claudeStoppedRequestEndWait(this.sessions)
|
||||
routePromptCancel = claudePromptCancelRoute
|
||||
dismissPrompt: StructuredAgentSessionAdapter['dismissPrompt'] = (request) =>
|
||||
dismissClaudeStructuredPrompt({ request, sessions: this.sessions })
|
||||
dismissPrompt: NonNullable<StructuredAgentSessionAdapter['dismissPrompt']> = (request) =>
|
||||
this.freeingAsker(request, (freeing) =>
|
||||
dismissClaudeStructuredPrompt({ request: freeing, sessions: this.sessions })
|
||||
)
|
||||
/** An answered or dismissed request frees the child it blocked before the host records the card,
|
||||
* so no row reads the child waiting beside a closed card; no provider frame says so first. */
|
||||
private freeingAsker = <R extends { sessionId: string; commit: () => Promise<void> }>(
|
||||
request: R,
|
||||
settle: (request: R) => Promise<void>
|
||||
): Promise<void> => {
|
||||
const free = () => this.publishChildWork(request.sessionId)
|
||||
const commit = async (): Promise<void> => {
|
||||
free()
|
||||
await request.commit()
|
||||
}
|
||||
return settle({ ...request, commit }).finally(free)
|
||||
}
|
||||
stopBackgroundTasks: NonNullable<StructuredAgentSessionAdapter['stopBackgroundTasks']> = async (
|
||||
input
|
||||
) => {
|
||||
@@ -240,7 +261,9 @@ export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAda
|
||||
return session ? claudeHoldsDispatch(session) : false
|
||||
}
|
||||
answerPrompt: StructuredAgentSessionAdapter['answerPrompt'] = (request) =>
|
||||
answerClaudeStructuredPrompt({ request, sessions: this.sessions })
|
||||
this.freeingAsker(request, (freeing) =>
|
||||
answerClaudeStructuredPrompt({ request: freeing, sessions: this.sessions })
|
||||
)
|
||||
setOption: StructuredAgentSessionAdapter['setOption'] = (input) =>
|
||||
setClaudeStructuredSessionOption(
|
||||
this.session(input.sessionId),
|
||||
|
||||
@@ -38,6 +38,12 @@ export type ClaudeSubagentLinkageSource = {
|
||||
) => Exclude<ClaudeSubagentLinkageVerdict, { kind: 'pending' }>
|
||||
}
|
||||
|
||||
/** The linkage of an agent the provider names by its own task id rather than by a frame's
|
||||
* `parent_tool_use_id`, as a permission request does. */
|
||||
export type ClaudeAgentLinkageSource = {
|
||||
linkageForAgent: (agentId: string) => AgentJournalProducerLinkage
|
||||
}
|
||||
|
||||
/** What the roster knows about one child, reduced to what attribution needs. */
|
||||
export type ClaudeSubagentLinkageEntry = { attempt: number }
|
||||
|
||||
@@ -54,6 +60,8 @@ export type ClaudeSubagentLinkageDeps = {
|
||||
* `parent_tool_use_id` is one of those ids, and this is the only route from
|
||||
* it to the agent that actually spawned the grandchild. */
|
||||
childOwnerRefOf?: (toolUseId: string) => string | null
|
||||
/** The spawn call the roster tracked a child under, by its canonical id. */
|
||||
spawnRefOf?: (canonicalId: string) => string | null
|
||||
}
|
||||
|
||||
/** How far a sidechain is followed when naming a row's parent. Depth beyond
|
||||
@@ -66,7 +74,9 @@ const MAX_PARENT_RESOLUTION_DEPTH = 8
|
||||
* parent, which is what an unrecorded owner also means. */
|
||||
type ParentAgentVerdict = { kind: 'known'; agentId?: string } | { kind: 'pending' }
|
||||
|
||||
export class ClaudeSubagentLinkage implements ClaudeSubagentLinkageSource {
|
||||
export class ClaudeSubagentLinkage
|
||||
implements ClaudeSubagentLinkageSource, ClaudeAgentLinkageSource
|
||||
{
|
||||
constructor(private readonly deps: ClaudeSubagentLinkageDeps) {}
|
||||
|
||||
linkageFor = (parentToolUseId: string): ClaudeSubagentLinkageVerdict =>
|
||||
@@ -85,6 +95,22 @@ export class ClaudeSubagentLinkage implements ClaudeSubagentLinkageSource {
|
||||
: verdict
|
||||
}
|
||||
|
||||
/** Resolved through the agent's spawn call, so it states what that agent's own rows state; the
|
||||
* provider's id stays the agent's id. Before the roster tracks the agent, its id alone. */
|
||||
linkageForAgent = (agentId: string): AgentJournalProducerLinkage => {
|
||||
const spawnRef = this.deps.spawnRefOf?.(agentId) ?? null
|
||||
if (spawnRef === null) {
|
||||
return { agentId, producerKind: 'agent' }
|
||||
}
|
||||
const { linkage } = this.settledLinkageFor(spawnRef)
|
||||
const { parentAgentId, ...rest } = linkage
|
||||
return {
|
||||
...rest,
|
||||
agentId,
|
||||
...(parentAgentId === undefined || parentAgentId === agentId ? {} : { parentAgentId })
|
||||
}
|
||||
}
|
||||
|
||||
private resolve(
|
||||
parentToolUseId: string,
|
||||
settled: boolean,
|
||||
|
||||
@@ -0,0 +1,651 @@
|
||||
// A Claude subagent's permission request, replayed from a capture of the real CLI through the real
|
||||
// adapter, deferred sink, durable journal and status feed, published as production publishes it:
|
||||
// on the sink's own publish, on a journal commit (a microtask later), and after child work. At every
|
||||
// status publish the subagents read waiting must each have a pending card in that same journal, and
|
||||
// the parent row must match a second host fed the same evidence with no subagent ever waiting: a
|
||||
// subagent's wait never reaches its parent's row. Both hosts read the same journal, so what the
|
||||
// prompt rows' linkage changes on the parent (its dating) is pinned where it shows.
|
||||
|
||||
import { readFileSync } from 'node:fs'
|
||||
import { mkdtemp, rm } from 'node:fs/promises'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { isDeepStrictEqual } from 'node:util'
|
||||
import { afterEach, beforeEach, describe, expect, it } from 'vitest'
|
||||
import { parseAgentJournalItemKey } from '../../shared/agent-session-journal-item-key'
|
||||
import type {
|
||||
AgentJournalProducerLinkage,
|
||||
AgentJournalRenderItem,
|
||||
AgentJournalResolution
|
||||
} from '../../shared/agent-session-journal-types'
|
||||
import type { AgentChildWorkEvidence } from '../../shared/agent-status-child-work-evidence'
|
||||
import type { AgentStatusStructuredSessionSubject } from '../../shared/agent-status-subject'
|
||||
import { AgentHookServer } from '../agent-hooks/server'
|
||||
import type { AgentSessionJournal } from '../native-chat/agent-session-journal/journal-store'
|
||||
import { createTrackedJournalOpener } from '../native-chat/agent-session-journal/journal-host-database-test-support'
|
||||
import { settleStructuredAgentSessionDeadGeneration } from '../native-chat/agent-session-wire/structured-agent-session-dead-generation-settlement'
|
||||
import { createDeferredStructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
|
||||
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
|
||||
import { createStructuredAgentSessionLogger } from '../native-chat/agent-session-wire/structured-agent-session-logger'
|
||||
import { StructuredAgentSessionStatusFeed } from '../native-chat/agent-session-wire/structured-agent-session-status-feed'
|
||||
import { indexedStatusFeedSession } from '../native-chat/agent-session-wire/structured-agent-session-status-feed-test-session'
|
||||
import { invokeCanUseTool } from './claude-can-use-tool-test-support'
|
||||
import { system, toolUse } from './claude-child-work-producer-harness.test-fixture'
|
||||
import { ClaudeStructuredSessionAdapter } from './claude-structured-session-adapter'
|
||||
import type { ClaudeStructuredSessionEvent } from './claude-structured-session-state'
|
||||
import {
|
||||
fakeClaude,
|
||||
identityFor,
|
||||
PROVIDER_SESSION_ID
|
||||
} from './claude-structured-session-test-support'
|
||||
|
||||
const SESSION = 'session-1'
|
||||
const FENCE = 7
|
||||
|
||||
type Captured = { from: 'cli' | 'orca'; frame: Record<string, unknown> }
|
||||
|
||||
function isRecord(value: unknown): value is Record<string, unknown> {
|
||||
return typeof value === 'object' && value !== null && !Array.isArray(value)
|
||||
}
|
||||
|
||||
function text(value: unknown): string {
|
||||
return typeof value === 'string' ? value : ''
|
||||
}
|
||||
|
||||
function captured(name: string): Captured[] {
|
||||
const path = join(__dirname, '__fixtures__', 'claude-subagent-permission-frames.json')
|
||||
const parsed: unknown = JSON.parse(readFileSync(path, 'utf8'))
|
||||
const events = isRecord(parsed) && isRecord(parsed.scenarios) ? parsed.scenarios[name] : null
|
||||
if (!Array.isArray(events)) {
|
||||
throw new Error(`no captured scenario ${name}`)
|
||||
}
|
||||
return events.flatMap((event) =>
|
||||
isRecord(event) && (event.from === 'cli' || event.from === 'orca') && isRecord(event.frame)
|
||||
? [{ from: event.from, frame: event.frame }]
|
||||
: []
|
||||
)
|
||||
}
|
||||
|
||||
/** The subagent that asks in a capture. */
|
||||
function askerOf(name: string): string {
|
||||
return text(captured(name).find((event) => event.frame.subtype === 'task_started')?.frame.task_id)
|
||||
}
|
||||
|
||||
const ASKER = askerOf('fg-allow')
|
||||
|
||||
/** The same evidence from a producer that never reads a subagent waiting. */
|
||||
function withoutWaits(evidence: AgentChildWorkEvidence[]): AgentChildWorkEvidence[] {
|
||||
return evidence.map((edge) =>
|
||||
edge.type === 'live' && edge.child.state === 'waiting'
|
||||
? { ...edge, child: { ...edge.child, state: 'working' } }
|
||||
: edge
|
||||
)
|
||||
}
|
||||
|
||||
function parentRow(server: AgentHookServer) {
|
||||
const row = server.getStatusSnapshot()[0]
|
||||
return (
|
||||
row && {
|
||||
state: row.state,
|
||||
workingMode: row.workingMode,
|
||||
mainAgent: row.mainAgent,
|
||||
stateStartedAt: row.stateStartedAt
|
||||
}
|
||||
)
|
||||
}
|
||||
|
||||
function pendingCards(items: readonly AgentJournalRenderItem[]): AgentJournalRenderItem[] {
|
||||
return items.filter(
|
||||
(item) =>
|
||||
(item.body.kind === 'approval' || item.body.kind === 'question') &&
|
||||
item.body.resolution.state === 'pending'
|
||||
)
|
||||
}
|
||||
|
||||
const journals = createTrackedJournalOpener()
|
||||
let root: string
|
||||
/** What the host does with the adapter's events, where a test needs it. */
|
||||
const hooks: { onEvent?: (event: ClaudeStructuredSessionEvent) => void } = {}
|
||||
|
||||
beforeEach(async () => {
|
||||
root = await mkdtemp(join(tmpdir(), 'orca-claude-subagent-request-'))
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
hooks.onEvent = undefined
|
||||
await journals.closeAll()
|
||||
await rm(root, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
/** One status publish: what the parent row read beside the main-equivalent host, and whether a
|
||||
* subagent read waiting without a pending card of its own in the journal it was projected from. */
|
||||
type Publish = { waiting: string[]; askers: string[]; row: unknown; unwaited: unknown }
|
||||
|
||||
async function pipeline() {
|
||||
let clock = 1_700_000_000_000
|
||||
const now = () => (clock += 1)
|
||||
const journal: AgentSessionJournal = await journals.open({
|
||||
identity: identityFor(SESSION),
|
||||
now,
|
||||
stateDirectory: join(root, SESSION)
|
||||
})
|
||||
const server = new AgentHookServer()
|
||||
const unwaited = new AgentHookServer()
|
||||
const publishes: Publish[] = []
|
||||
let published: AgentStatusStructuredSessionSubject | undefined
|
||||
/** The child's row as the host shows it, which every surface reads. */
|
||||
const viewOf = (providerId: string) =>
|
||||
published &&
|
||||
server.getStructuredChildWorkViews(published).find((view) => view.providerId === providerId)
|
||||
const record = (subject: AgentStatusStructuredSessionSubject): void => {
|
||||
published = subject
|
||||
const cards = pendingCards(journal.snapshot().items)
|
||||
publishes.push({
|
||||
waiting: server
|
||||
.getStructuredChildWorkViews(subject)
|
||||
.flatMap((view) =>
|
||||
view.state === 'waiting' && view.membership === 'live' ? [view.providerId ?? '?'] : []
|
||||
),
|
||||
askers: cards.flatMap((card) => (card.agentId ? [card.agentId] : [])),
|
||||
row: parentRow(server),
|
||||
unwaited: parentRow(unwaited)
|
||||
})
|
||||
}
|
||||
const feed = new StructuredAgentSessionStatusFeed({
|
||||
logger: createStructuredAgentSessionLogger(),
|
||||
sessions: new Map([
|
||||
[
|
||||
SESSION,
|
||||
indexedStatusFeedSession({
|
||||
journal,
|
||||
child: { generation: 'spawn-9', fence: FENCE, phase: 'ready' },
|
||||
provider: 'claude'
|
||||
})
|
||||
]
|
||||
]),
|
||||
getRecord: () => null,
|
||||
now,
|
||||
statusSink: () => ({
|
||||
publish: (summary, subject) => {
|
||||
server.ingestStructuredStatus(summary, subject)
|
||||
unwaited.ingestStructuredStatus(summary, subject)
|
||||
record(subject)
|
||||
},
|
||||
forget: (subject) => {
|
||||
server.dropStructuredStatus(subject)
|
||||
unwaited.dropStructuredStatus(subject)
|
||||
},
|
||||
publishChildWork: (subject, evidence, provider) => {
|
||||
server.ingestStructuredChildWork(subject, evidence, provider)
|
||||
unwaited.ingestStructuredChildWork(subject, withoutWaits(evidence), provider)
|
||||
},
|
||||
readChildWork: (subject) => server.getStructuredChildWorkViews(subject)
|
||||
})
|
||||
})
|
||||
// As the host publishes: the sink's own publish, and every journal commit a microtask later.
|
||||
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging(SESSION))
|
||||
const publishStatus = () => feed.publish(SESSION, journal)
|
||||
const target = { journal, fence: FENCE, publish: publishStatus }
|
||||
deferred.bind(target)
|
||||
let queued = false
|
||||
journal.observeCommits(() => {
|
||||
if (!queued) {
|
||||
queued = true
|
||||
queueMicrotask(() => {
|
||||
queued = false
|
||||
feed.publish(SESSION, journal)
|
||||
})
|
||||
}
|
||||
})
|
||||
const claude = fakeClaude()
|
||||
const adapter = new ClaudeStructuredSessionAdapter({
|
||||
resolveLaunch: async () => ({
|
||||
pathToClaudeCodeExecutable: 'claude',
|
||||
options: {},
|
||||
cwd: '/work/repo',
|
||||
claudeConfigDir: '/accounts/claude',
|
||||
providerSessionId: PROVIDER_SESSION_ID,
|
||||
resumeLeafUuid: null,
|
||||
resumesTranscript: false,
|
||||
continuesChain: false
|
||||
}),
|
||||
openConnection: claude.openConnection,
|
||||
readProcessStartTime: async () => 1_700_000_000_000,
|
||||
now,
|
||||
persistHandle: async () => {},
|
||||
onChildWorkEvidence: (sessionId, evidence) => feed.publishChildWork(sessionId, evidence),
|
||||
onEvent: (event) => hooks.onEvent?.(event)
|
||||
})
|
||||
await adapter.acquire({
|
||||
identity: identityFor(SESSION),
|
||||
fence: FENCE,
|
||||
spawnToken: 'spawn-9',
|
||||
events: deferred.sink
|
||||
})
|
||||
await adapter.awaitStarted(SESSION)
|
||||
const settle = async (): Promise<void> => {
|
||||
for (let round = 0; round < 3; round += 1) {
|
||||
expect(await deferred.drained()).toEqual({ ok: true })
|
||||
await new Promise((resolve) => setTimeout(resolve, 0))
|
||||
}
|
||||
}
|
||||
/** What the host's own commit writes for the card: its resolution, at the live turn. */
|
||||
const hostRecords =
|
||||
(resolution: Pick<AgentJournalResolution, 'state' | 'selectedOptionId'>, itemId = cardId()) =>
|
||||
async (): Promise<void> => {
|
||||
await deferred.drained()
|
||||
const card = pendingCards(journal.snapshot().items).find((item) => item.itemId === itemId)
|
||||
const identity = card && parseAgentJournalItemKey(card.itemId)
|
||||
if (card?.body.kind !== 'approval' || !identity) {
|
||||
throw new Error('no pending card to record')
|
||||
}
|
||||
await journal.appendItem(
|
||||
identity,
|
||||
{ ...card.body, resolution: { ...resolution, resolvedBy: 'client-1', resolvedAt: now() } },
|
||||
{ fence: FENCE, turnScope: journal.liveTurnScope() }
|
||||
)
|
||||
}
|
||||
const cardId = (): string => {
|
||||
const card = pendingCards(journal.snapshot().items)[0]
|
||||
if (!card) {
|
||||
throw new Error('no pending card')
|
||||
}
|
||||
return card.itemId
|
||||
}
|
||||
const connection = claude.connections[0]!
|
||||
const aborts = new Map<string, AbortController>()
|
||||
/** Feeds a captured frame as the CLI or the SDK hands it over, then lets every write land. The
|
||||
* SDK hands `canUseTool` no agent id when `withoutAgentId`, as when the CLI names none. */
|
||||
const step = async (
|
||||
{ from, frame }: Captured,
|
||||
options: { withoutAgentId?: boolean } = {}
|
||||
): Promise<void> => {
|
||||
const request = isRecord(frame.request) ? frame.request : null
|
||||
if (frame.type === 'control_request' && request?.subtype === 'can_use_tool') {
|
||||
const controller = new AbortController()
|
||||
aborts.set(text(frame.request_id), controller)
|
||||
const agentId = options.withoutAgentId ? '' : text(request.agent_id)
|
||||
invokeCanUseTool(
|
||||
connection,
|
||||
text(request.tool_name),
|
||||
text(frame.request_id),
|
||||
text(request.tool_use_id),
|
||||
{
|
||||
input: isRecord(request.input) ? request.input : {},
|
||||
signal: controller.signal,
|
||||
...(agentId ? { agentID: agentId } : {})
|
||||
}
|
||||
)
|
||||
} else if (from === 'cli') {
|
||||
connection.handlers.onMessage?.({ ...frame, session_id: PROVIDER_SESSION_ID })
|
||||
}
|
||||
await settle()
|
||||
}
|
||||
/** Replays a capture up to and including its permission request. */
|
||||
const ask = async (name: string, options: { withoutAgentId?: boolean } = {}): Promise<void> => {
|
||||
const events = captured(name)
|
||||
const at = events.findIndex((event) => event.frame.type === 'control_request')
|
||||
for (const event of events.slice(0, at + 1)) {
|
||||
await step(event, options)
|
||||
}
|
||||
}
|
||||
/** Publishes that broke the invariant or told the parent rows apart. */
|
||||
const violations = () =>
|
||||
publishes.filter(
|
||||
(entry) =>
|
||||
entry.waiting.some((id) => !entry.askers.includes(id)) ||
|
||||
!isDeepStrictEqual(entry.row, entry.unwaited)
|
||||
)
|
||||
/** Raises a request the capture does not hold, as the SDK would. */
|
||||
const raise = (requestId: string, toolUseId: string, agentId?: string): void => {
|
||||
const controller = new AbortController()
|
||||
aborts.set(requestId, controller)
|
||||
invokeCanUseTool(connection, 'Bash', requestId, toolUseId, {
|
||||
input: { command: `echo ${requestId}` },
|
||||
signal: controller.signal,
|
||||
...(agentId ? { agentID: agentId } : {})
|
||||
})
|
||||
}
|
||||
return {
|
||||
adapter,
|
||||
journal,
|
||||
deferred,
|
||||
target,
|
||||
publishStatus,
|
||||
connection,
|
||||
now,
|
||||
raise,
|
||||
publishes,
|
||||
violations,
|
||||
viewOf,
|
||||
settle,
|
||||
step,
|
||||
ask,
|
||||
cardId,
|
||||
hostRecords,
|
||||
aborts
|
||||
}
|
||||
}
|
||||
|
||||
type Pipeline = Awaited<ReturnType<typeof pipeline>>
|
||||
|
||||
function answer(
|
||||
run: Pipeline,
|
||||
optionId: 'allow' | 'deny',
|
||||
commit = run.hostRecords({ state: 'resolved', selectedOptionId: optionId }),
|
||||
itemId = run.cardId()
|
||||
) {
|
||||
return run.adapter.answerPrompt({
|
||||
sessionId: SESSION,
|
||||
itemId,
|
||||
kind: 'approval',
|
||||
response: { kind: 'option', optionId },
|
||||
fence: FENCE,
|
||||
commit
|
||||
})
|
||||
}
|
||||
|
||||
/** The user closes the card without answering; the Stop that follows ends the request. */
|
||||
function dismiss(run: Pipeline) {
|
||||
return run.adapter.dismissPrompt({
|
||||
sessionId: SESSION,
|
||||
itemId: run.cardId(),
|
||||
fence: FENCE,
|
||||
answer: false,
|
||||
commit: run.hostRecords({ state: 'cancelled', selectedOptionId: null })
|
||||
})
|
||||
}
|
||||
|
||||
/** Replays a whole capture, answering through the app's own path where Orca answered, and returns
|
||||
* the asking subagent's row after each event, labelled. */
|
||||
async function replay(name: string): Promise<Pipeline & { timeline: string[] }> {
|
||||
const run = await pipeline()
|
||||
const timeline: string[] = []
|
||||
for (const event of captured(name)) {
|
||||
const response = isRecord(event.frame.response) ? event.frame.response : null
|
||||
const request = isRecord(event.frame.request) ? event.frame.request : null
|
||||
let label = text(event.frame.subtype) || text(event.frame.type)
|
||||
if (event.from === 'orca' && event.frame.type === 'control_response' && response) {
|
||||
const behavior = isRecord(response.response) ? text(response.response.behavior) : ''
|
||||
await answer(run, behavior === 'deny' ? 'deny' : 'allow')
|
||||
await run.settle()
|
||||
label = behavior
|
||||
} else {
|
||||
await run.step(event)
|
||||
label = request?.subtype === 'can_use_tool' ? 'can_use_tool' : label
|
||||
}
|
||||
const view = run.viewOf(askerOf(name))
|
||||
timeline.push(`${label} -> ${view ? `${view.membership} ${view.state}` : 'none'}`)
|
||||
}
|
||||
return { ...run, timeline }
|
||||
}
|
||||
|
||||
/** Root spawns agent-a, which spawns agent-n; `asker` asks for agent-n's Bash call, which is read
|
||||
* after the request unless `toolCallFirst`. */
|
||||
async function nested(asker: string, toolCallFirst = false): Promise<Pipeline> {
|
||||
const run = await pipeline()
|
||||
const spawn = (id: string, taskId: string, parentRef: string | null) => [
|
||||
toolUse(id, 'Agent', { description: taskId, prompt: 'go' }, parentRef),
|
||||
system('task_started', {
|
||||
task_id: taskId,
|
||||
tool_use_id: id,
|
||||
description: taskId,
|
||||
task_type: 'local_agent'
|
||||
})
|
||||
]
|
||||
const bash = toolUse('toolu_bash_n', 'Bash', { command: 'touch n' }, 'toolu_n')
|
||||
for (const frame of [
|
||||
...spawn('toolu_a', 'agent-a', null),
|
||||
...spawn('toolu_n', 'agent-n', 'toolu_a'),
|
||||
...(toolCallFirst ? [bash] : [])
|
||||
]) {
|
||||
await run.step({ from: 'cli', frame })
|
||||
}
|
||||
run.raise('req-n', 'toolu_bash_n', asker)
|
||||
await run.step({
|
||||
from: 'cli',
|
||||
frame: toolUse('toolu_read_n', 'Read', { file_path: 'n' }, 'toolu_n')
|
||||
})
|
||||
return run
|
||||
}
|
||||
|
||||
function linkageOf(item: AgentJournalProducerLinkage | undefined) {
|
||||
return {
|
||||
agentId: item?.agentId,
|
||||
parentAgentId: item?.parentAgentId,
|
||||
providerParentRef: item?.providerParentRef,
|
||||
producerKind: item?.producerKind,
|
||||
attempt: item?.attempt
|
||||
}
|
||||
}
|
||||
|
||||
describe("a Claude subagent's permission request", () => {
|
||||
it.each([
|
||||
[
|
||||
'fg-allow',
|
||||
[
|
||||
'session_state_changed -> live working',
|
||||
'can_use_tool -> live waiting',
|
||||
'allow -> live working'
|
||||
]
|
||||
],
|
||||
// The parent's own turn ends while its background subagent is still asking.
|
||||
[
|
||||
'bg-allow',
|
||||
[
|
||||
'session_state_changed -> live working',
|
||||
'can_use_tool -> live waiting',
|
||||
'assistant -> live waiting',
|
||||
'success -> live waiting',
|
||||
'allow -> live working'
|
||||
]
|
||||
]
|
||||
])(
|
||||
'waits from its request until answered, beside its card, the parent row as before (%s)',
|
||||
async (name, around) => {
|
||||
const run = await replay(name)
|
||||
const at = run.timeline.indexOf('can_use_tool -> live waiting')
|
||||
expect(run.timeline.slice(at - 1, at - 1 + around.length)).toEqual(around)
|
||||
expect(run.violations()).toEqual([])
|
||||
}
|
||||
)
|
||||
|
||||
it('reads blocked as soon as the request arrives, dated by its card', async () => {
|
||||
const run = await pipeline()
|
||||
await run.ask('fg-allow')
|
||||
const card = pendingCards(run.journal.snapshot().items)[0]
|
||||
expect(run.publishes.at(-1)).toMatchObject({
|
||||
waiting: [card?.agentId],
|
||||
row: {
|
||||
state: 'blocked',
|
||||
stateStartedAt: card?.observedAt,
|
||||
mainAgent: { state: 'blocked', stateStartedAt: card?.observedAt }
|
||||
}
|
||||
})
|
||||
expect(run.viewOf(ASKER)?.operation).toMatchObject({ toolName: 'Bash' })
|
||||
expect(run.violations()).toEqual([])
|
||||
})
|
||||
|
||||
it('waits only once its card is written, though the write waits for the journal', async () => {
|
||||
const run = await pipeline()
|
||||
const events = captured('fg-allow')
|
||||
const at = events.findIndex((event) => event.frame.type === 'control_request')
|
||||
for (const event of events.slice(0, at)) {
|
||||
await run.step(event)
|
||||
}
|
||||
// No journal is bound: the card's row waits in the sink while the host publishes.
|
||||
run.deferred.unbind()
|
||||
run.raise('req-late', 'toolu_late', ASKER)
|
||||
await new Promise((resolve) => setTimeout(resolve, 0))
|
||||
run.publishStatus()
|
||||
expect(run.publishes.at(-1)?.waiting).toEqual([])
|
||||
run.deferred.bind(run.target)
|
||||
await run.settle()
|
||||
expect(run.publishes.at(-1)?.waiting).toEqual([ASKER])
|
||||
expect(run.violations()).toEqual([])
|
||||
})
|
||||
|
||||
it.each(['allowed', 'denied', 'dismissed and left to the Stop', 'withdrawn by Claude'] as const)(
|
||||
'frees the subagent when its request is %s',
|
||||
async (how) => {
|
||||
const run = await pipeline()
|
||||
await run.ask('fg-allow')
|
||||
if (how === 'allowed' || how === 'denied') {
|
||||
await answer(run, how === 'allowed' ? 'allow' : 'deny')
|
||||
} else if (how === 'dismissed and left to the Stop') {
|
||||
await dismiss(run)
|
||||
}
|
||||
await run.settle()
|
||||
// Claude withdraws the request itself; after a dismissal that closes nothing more.
|
||||
if (how !== 'allowed' && how !== 'denied') {
|
||||
for (const controller of run.aborts.values()) {
|
||||
controller.abort()
|
||||
}
|
||||
await run.settle()
|
||||
}
|
||||
expect(run.publishes.at(-1)?.waiting).toEqual([])
|
||||
expect(run.violations()).toEqual([])
|
||||
}
|
||||
)
|
||||
|
||||
it('waits again, beside its card, when the host fails to record the answer', async () => {
|
||||
const run = await pipeline()
|
||||
await run.ask('fg-allow')
|
||||
await expect(
|
||||
answer(run, 'allow', async () => {
|
||||
throw new Error('journal write failed')
|
||||
})
|
||||
).rejects.toThrow('journal write failed')
|
||||
await run.settle()
|
||||
expect(run.publishes.at(-1)?.waiting).toEqual([ASKER])
|
||||
expect(run.violations()).toEqual([])
|
||||
})
|
||||
|
||||
it('keeps an answered request in the rows of the subagent that asked', async () => {
|
||||
const run = await replay('fg-allow')
|
||||
const card = run.journal.snapshot().items.find((item) => item.body.kind === 'approval')
|
||||
expect(card?.agentId).toBe(ASKER)
|
||||
expect(card?.body).toMatchObject({
|
||||
resolution: { state: 'resolved', selectedOptionId: 'allow' }
|
||||
})
|
||||
})
|
||||
|
||||
it('names the subagent through the tool call it gates when the CLI does not', async () => {
|
||||
const run = await pipeline()
|
||||
await run.ask('fg-allow', { withoutAgentId: true })
|
||||
expect(pendingCards(run.journal.snapshot().items)[0]?.agentId).toBe(ASKER)
|
||||
expect(run.publishes.at(-1)?.waiting).toEqual([ASKER])
|
||||
expect(run.violations()).toEqual([])
|
||||
})
|
||||
|
||||
it("gives a nested subagent's request the linkage its own rows carry", async () => {
|
||||
const run = await nested('agent-n')
|
||||
const items = run.journal.snapshot().items
|
||||
const card = items.find((item) => item.body.kind === 'approval')
|
||||
const sibling = items.find(
|
||||
(item) => item.body.kind === 'tool-call' && item.agentId === 'agent-n'
|
||||
)
|
||||
expect(linkageOf(card)).toEqual(linkageOf(sibling))
|
||||
expect(card).toMatchObject({ agentId: 'agent-n', parentAgentId: 'agent-a' })
|
||||
})
|
||||
|
||||
it("files the request under the agent the CLI names when the gated call is another's", async () => {
|
||||
const run = await nested('agent-a', true)
|
||||
const items = run.journal.snapshot().items
|
||||
const card = items.find((item) => item.body.kind === 'approval')
|
||||
const askerRow = items.find(
|
||||
(item) => item.body.kind === 'tool-call' && item.agentId === 'agent-a'
|
||||
)
|
||||
expect(linkageOf(card)).toEqual(linkageOf(askerRow))
|
||||
})
|
||||
|
||||
it('waits each asking subagent beside its own card, and only those', async () => {
|
||||
const run = await pipeline()
|
||||
await run.ask('fg-allow')
|
||||
const started = captured('fg-allow').find((event) => event.frame.subtype === 'task_started')
|
||||
await run.step({
|
||||
from: 'cli',
|
||||
frame: { ...started?.frame, task_id: 'agent-two', uuid: 'u-two', tool_use_id: 'toolu_two' }
|
||||
})
|
||||
run.raise('req-two', 'toolu_two_bash', 'agent-two')
|
||||
await run.settle()
|
||||
expect([...(run.publishes.at(-1)?.waiting ?? [])].sort()).toEqual([ASKER, 'agent-two'].sort())
|
||||
const cardOf = (agentId: string) =>
|
||||
pendingCards(run.journal.snapshot().items).find((card) => card.agentId === agentId)?.itemId
|
||||
const first = cardOf(ASKER)
|
||||
await answer(
|
||||
run,
|
||||
'allow',
|
||||
run.hostRecords({ state: 'resolved', selectedOptionId: 'allow' }, first),
|
||||
first
|
||||
)
|
||||
await run.settle()
|
||||
expect(run.publishes.at(-1)?.waiting).toEqual(['agent-two'])
|
||||
const second = cardOf('agent-two')
|
||||
await answer(
|
||||
run,
|
||||
'allow',
|
||||
run.hostRecords({ state: 'resolved', selectedOptionId: 'allow' }, second),
|
||||
second
|
||||
)
|
||||
await run.settle()
|
||||
expect(run.publishes.at(-1)?.waiting).toEqual([])
|
||||
expect(run.violations()).toEqual([])
|
||||
})
|
||||
|
||||
it("stays blocked on the session's own request after its subagent's is answered", async () => {
|
||||
const run = await pipeline()
|
||||
await run.ask('fg-allow')
|
||||
run.raise('req-main', 'toolu_main_bash')
|
||||
await run.settle()
|
||||
const cards = pendingCards(run.journal.snapshot().items)
|
||||
expect(cards.map((card) => card.agentId ?? null).sort()).toEqual([ASKER, null].sort())
|
||||
// Dated by the session's own ask, though its subagent asked first: main's rule, as for Codex.
|
||||
const own = cards.find((card) => !card.agentId)?.observedAt
|
||||
expect(own).toBeGreaterThan(
|
||||
cards.find((card) => card.agentId === ASKER)?.observedAt ?? Infinity
|
||||
)
|
||||
expect(run.publishes.at(-1)).toMatchObject({
|
||||
waiting: [ASKER],
|
||||
row: { state: 'blocked', stateStartedAt: own, mainAgent: { stateStartedAt: own } }
|
||||
})
|
||||
const subagents = cards.find((card) => card.agentId === ASKER)?.itemId
|
||||
await answer(
|
||||
run,
|
||||
'allow',
|
||||
run.hostRecords({ state: 'resolved', selectedOptionId: 'allow' }, subagents),
|
||||
subagents
|
||||
)
|
||||
await run.settle()
|
||||
expect(run.publishes.at(-1)).toMatchObject({ waiting: [], row: { state: 'blocked' } })
|
||||
expect(run.violations()).toEqual([])
|
||||
})
|
||||
|
||||
it('stops the subagent waiting when the process dies mid-request', async () => {
|
||||
const run = await pipeline()
|
||||
await run.ask('fg-allow')
|
||||
expect(run.publishes.at(-1)?.waiting).toEqual([ASKER])
|
||||
let settled: Promise<unknown> | undefined
|
||||
hooks.onEvent = (event) => {
|
||||
if (event.type === 'ended') {
|
||||
settled = settleStructuredAgentSessionDeadGeneration({
|
||||
journal: run.journal,
|
||||
sessionId: SESSION,
|
||||
fence: FENCE,
|
||||
settlementId: 'provider-exit:test',
|
||||
verdict: { state: 'interrupted', completedAt: run.now() },
|
||||
pendingSubmissionReason: 'provider_exited_before_acknowledgement'
|
||||
})
|
||||
}
|
||||
}
|
||||
run.connection.handlers.onExit?.(new Error('claude crashed'))
|
||||
for (let attempt = 0; attempt < 20 && !settled; attempt += 1) {
|
||||
await new Promise((resolve) => setTimeout(resolve, 5))
|
||||
}
|
||||
await settled
|
||||
await run.settle()
|
||||
expect(pendingCards(run.journal.snapshot().items)).toEqual([])
|
||||
expect(run.publishes.at(-1)?.waiting).toEqual([])
|
||||
expect(run.violations()).toEqual([])
|
||||
})
|
||||
})
|
||||
@@ -27,7 +27,11 @@ import { writeClaudeSubagentGroupRow } from './claude-subagent-group-row'
|
||||
import { ClaudeSubagentIds } from './claude-subagent-id-aliases'
|
||||
import type { ClaudeJournaledRosterSource } from './claude-subagent-journaled-roster'
|
||||
import { ClaudeSubagentRosterGroups } from './claude-subagent-roster-groups'
|
||||
import { ClaudeSubagentLinkage, type ClaudeSubagentLinkageSource } from './claude-subagent-linkage'
|
||||
import {
|
||||
ClaudeSubagentLinkage,
|
||||
type ClaudeAgentLinkageSource,
|
||||
type ClaudeSubagentLinkageSource
|
||||
} from './claude-subagent-linkage'
|
||||
import { readClaudeSubagentTaskFrame } from './claude-subagent-task-frames'
|
||||
import {
|
||||
applyClaudeSubagentInvocation,
|
||||
@@ -72,7 +76,7 @@ export class ClaudeSubagentRoster {
|
||||
private readonly groups: ClaudeSubagentRosterGroups
|
||||
private readonly ids: ClaudeSubagentIds
|
||||
/** Who produced a row, for every write site journaling this session. */
|
||||
readonly linkage: ClaudeSubagentLinkageSource
|
||||
readonly linkage: ClaudeSubagentLinkageSource & ClaudeAgentLinkageSource
|
||||
/** Set by ANY `task_started`, including one the subagent filter rejects. Once
|
||||
* this CLI has proven it declares its tasks, child traffic for an id it never
|
||||
* announced is a nested tool or a grandchild, not a subagent. */
|
||||
@@ -91,7 +95,8 @@ export class ClaudeSubagentRoster {
|
||||
ids: this.ids,
|
||||
trackedFor: (canonicalId) => this.groups.locate(canonicalId)?.tracked ?? null,
|
||||
isForwardedParentTool: deps.isForwardedParentTool,
|
||||
childOwnerRefOf: deps.childOwnerRefOf
|
||||
childOwnerRefOf: deps.childOwnerRefOf,
|
||||
spawnRefOf: (canonicalId) => this.groups.locate(canonicalId)?.tracked.toolUseId ?? null
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -96,7 +96,9 @@ function sessionHoldingTurn(turnId: string | null): ReturnType<typeof sessionFor
|
||||
journalPrompts: {
|
||||
resolve: vi.fn(),
|
||||
handOver: () => () => {},
|
||||
cancel: () => ({ accepted: true })
|
||||
cancel: () => ({ accepted: true }),
|
||||
openCards: () => [][Symbol.iterator](),
|
||||
whenWritten: () => undefined
|
||||
},
|
||||
currentTurnId: turnId,
|
||||
commandTurnId: null,
|
||||
|
||||
@@ -51,6 +51,8 @@ export class StructuredAgentSessionSinkQueue {
|
||||
private backpressured = false
|
||||
private acceptedSequence = 0
|
||||
private settledSequence = 0
|
||||
/** The newest operation handed to the journal; a close drops the buffered rest unwritten. */
|
||||
private handedOverSequence = 0
|
||||
/** Settles once every operation handed over so far has, in handover order. */
|
||||
private handedOverSettled: Promise<void> = Promise.resolve()
|
||||
private readonly buffered: Admitted[] = []
|
||||
@@ -127,6 +129,15 @@ export class StructuredAgentSessionSinkQueue {
|
||||
return new Promise((resolve) => this.waiters.push({ through, resolve }))
|
||||
}
|
||||
|
||||
/** Like `barrier`, but a close that dropped writes admitted so far reads as not landed. */
|
||||
written = async (): Promise<StructuredAgentSessionSinkBarrier> => {
|
||||
const through = this.acceptedSequence
|
||||
const settled = await this.barrier()
|
||||
return settled.ok && this.handedOverSequence < through
|
||||
? { ok: false, error: new Error('the sink closed before its writes landed') }
|
||||
: settled
|
||||
}
|
||||
|
||||
submit(
|
||||
operation: Omit<StructuredAgentSessionSinkOperation, 'sequence'>,
|
||||
options: StructuredAgentSessionAppendOptions = {}
|
||||
@@ -193,6 +204,7 @@ export class StructuredAgentSessionSinkQueue {
|
||||
|
||||
private handOver(operation: Admitted, bound: StructuredAgentSessionEventTarget): void {
|
||||
const earlier = this.handedOverSettled
|
||||
this.handedOverSequence = operation.sequence
|
||||
const key = operation.publicationKey
|
||||
let outcome: Promise<unknown>
|
||||
if (key === undefined) {
|
||||
|
||||
@@ -243,6 +243,23 @@ describe('deferred structured agent-session event sink', () => {
|
||||
expect(log).toEqual([])
|
||||
})
|
||||
|
||||
it('says its writes landed only when they ran, though a close still settles drained()', async () => {
|
||||
const log: Recorded[] = []
|
||||
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
|
||||
deferred.bind(target(1, log))
|
||||
deferred.sink.appendItem(identity(0), BODY, { turnScope: AGENT_JOURNAL_THREAD_SCOPE })
|
||||
expect(await deferred.sink.written?.()).toEqual({ ok: true })
|
||||
|
||||
deferred.unbind()
|
||||
deferred.sink.appendItem(identity(1), BODY, { turnScope: AGENT_JOURNAL_THREAD_SCOPE })
|
||||
const written = deferred.sink.written?.()
|
||||
deferred.close()
|
||||
|
||||
expect(await written).toMatchObject({ ok: false })
|
||||
expect(await deferred.drained()).toEqual({ ok: true })
|
||||
expect(log).toEqual([{ call: 'appendItem', fence: 1, ordinal: 0 }])
|
||||
})
|
||||
|
||||
it('lands writes already handed to the journal when closed; only never-bound ones drop', async () => {
|
||||
const owed = await owingTarget(3)
|
||||
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
|
||||
@@ -253,6 +270,7 @@ describe('deferred structured agent-session event sink', () => {
|
||||
deferred.sink.appendItem(identity(2), BODY, { turnScope: AGENT_JOURNAL_THREAD_SCOPE })
|
||||
// Handed over, waiting behind the copy: none has landed when the sink closes.
|
||||
expect(owed.landed()).toEqual([])
|
||||
const written = deferred.sink.written?.()
|
||||
deferred.close()
|
||||
expect(
|
||||
deferred.sink.tryAppendItem?.(identity(3), BODY, { turnScope: AGENT_JOURNAL_THREAD_SCOPE })
|
||||
@@ -260,6 +278,7 @@ describe('deferred structured agent-session event sink', () => {
|
||||
await expect(deferred.drained()).resolves.toEqual({ ok: true })
|
||||
|
||||
expect(owed.landed()).toEqual([0, 1, 2])
|
||||
await expect(written).resolves.toEqual({ ok: true })
|
||||
expect(owed.journal.cursor().sequence).toBe(owed.history + 3)
|
||||
await owed.dispose()
|
||||
})
|
||||
|
||||
@@ -95,6 +95,9 @@ export type StructuredAgentSessionEventSink = {
|
||||
options?: StructuredAgentSessionAppendOptions
|
||||
): StructuredAgentSessionSinkAdmission
|
||||
publish(options?: StructuredAgentSessionPublishOptions): void
|
||||
/** Resolves `ok` once every write admitted so far has landed in the journal; not `ok` when one
|
||||
* failed or a close dropped it unwritten. */
|
||||
written?(): Promise<StructuredAgentSessionSinkBarrier>
|
||||
setActivity?(activity: AgentSessionTurnActivity | null): void
|
||||
tryAppendItem?(
|
||||
identity: AgentJournalItemIdentity,
|
||||
@@ -321,6 +324,7 @@ export function createDeferredStructuredAgentSessionEventSink(deps: {
|
||||
publish: (options = {}) => {
|
||||
publish(options)
|
||||
},
|
||||
written: queue.written,
|
||||
setActivity: (activity) => {
|
||||
queue.submit({
|
||||
bytes: Buffer.byteLength(JSON.stringify(activity), 'utf8') + 64,
|
||||
|
||||
@@ -231,11 +231,12 @@ const STORIES: Story[] = [
|
||||
],
|
||||
expect: { state: 'waiting', mainAgent: { state: 'done' } }
|
||||
},
|
||||
// KNOWN DIVERGENCE: no structured producer reports a waiting task; a child's pending prompt
|
||||
// is a session-level `attention`, so this lane blames the main agent for the child's request.
|
||||
// KNOWN DIVERGENCE: a child's pending prompt is the session's `attention`, so this lane reads
|
||||
// it as the main agent's own `blocked`, one needs-input state whoever asked, even beside the
|
||||
// child's own waiting record.
|
||||
structured: {
|
||||
status: 'attention',
|
||||
backgroundTasks: [AGENT_TASK],
|
||||
backgroundTasks: [{ ...AGENT_TASK, state: 'waiting' }],
|
||||
expect: { state: 'blocked', mainAgent: { state: 'blocked' } }
|
||||
},
|
||||
codex: {
|
||||
|
||||
Reference in New Issue
Block a user