From 28ff7f54f2060cd79e8a6b932265fa1d0b1abd53 Mon Sep 17 00:00:00 2001 From: Merge Sim Date: Sun, 6 Sep 2026 02:31:01 -0700 Subject: [PATCH] fix(orchestration): reach structured workers through group addresses MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `orca orchestration send --to @all` — and `@idle`, `@claude`, `@codex`, `@worktree:` — silently skipped every structured worker. Recipients came from `listTerminals`, which enumerates leaves and PTYs, and a structured session is on neither. The exclusion happened BEFORE per-recipient resolution, so the `SendRecipientWarning` machinery never ran: the caller got exit 0 and a receipt naming the workers that did resolve, and a broadcast "stop work" or "base moved" reached the PTY workers and nobody else. With every worker structured it degraded to `terminal_not_found`, which reads as "the group was empty". Fixed at the group-resolution site rather than inside `listTerminals`. That result is published to paired mobile and remote clients and to consumers that assume a summary carries a `ptyId` or is writable, so widening it is its own change under `docs/reference/remote-wire-compatibility.md`. Group addressing reads exactly three fields off a recipient, and `RuntimeTerminalSummary` already satisfies them structurally, so the resolver widens to that smaller shape and nothing here invents a `worktreePath` or a `branch`. Candidates are liveness- gated on the same observation the rest of the structured surface uses — mail addressed to a settled worker would be stored for a lane that will never deliver it — and once a worker IS a candidate, the existing per-recipient warnings cover it, so an unresolvable one is reported rather than dropped. `@idle` needed more than enumeration: `getAgentStatusForHandle` reaches a PTY probe that throws for a handle with no pane, so a structured worker would have been enumerated and then silently dropped from the one group address that selects on status. It now answers from the session's journal — and off the FULL reduced timeline, never a bounded tail. Settlement tombstones the running turn's lifecycle item rather than rewriting it, so on any page-sized read a long tool-calling turn looks identical to an idle session; `@idle` would then broadcast into a running turn, which Codex answers with `turn already running` and Claude queues behind. An unreadable session answers null, never idle. `terminal list` and `worktree ps` still omit structured workers; that is the wire-visible half and is deliberately not in this change. --- ...e-prune-mobile-session-tab-group-layout.ts | 8 + src/main/runtime/orchestration/groups.ts | 9 +- .../structured-mailbox-pointer-host.ts | 47 +++-- ...structured-worker-group-addressing.test.ts | 166 ++++++++++++++++++ .../structured-worker-group-addressing.ts | 64 +++++++ .../rpc/methods/orchestration-send-group.ts | 7 +- .../runtime/structured-worker-identity.ts | 5 + 7 files changed, 286 insertions(+), 20 deletions(-) create mode 100644 src/main/runtime/orchestration/structured-worker-group-addressing.test.ts create mode 100644 src/main/runtime/orchestration/structured-worker-group-addressing.ts diff --git a/src/main/runtime/orca-runtime-prune-mobile-session-tab-group-layout.ts b/src/main/runtime/orca-runtime-prune-mobile-session-tab-group-layout.ts index feb36204ff0..3322a8ad85e 100644 --- a/src/main/runtime/orca-runtime-prune-mobile-session-tab-group-layout.ts +++ b/src/main/runtime/orca-runtime-prune-mobile-session-tab-group-layout.ts @@ -24,6 +24,8 @@ import { FIRST_PANE_ID } from '../../shared/pane-key' import { isTerminalLeafId, makePaneKey, parsePaneKey } from '../../shared/stable-pane-id' import type { SleepingAgentLaunchConfig } from '../../shared/agent-session-resume' import { copySleepingAgentLaunchConfig } from './runtime-agent-launch-resolution' +import { resolveStructuredWorkerAuthority } from './structured-worker-authority' +import { structuredWorkerAgentStatus } from './orchestration/structured-worker-group-addressing' export class OrcaRuntimeWithPruneMobileSessionTabGroupLayout extends OrcaRuntimeWithScheduleMobileSessionTabsChanged { protected pruneMobileSessionTabGroupLayout( @@ -189,6 +191,12 @@ export class OrcaRuntimeWithPruneMobileSessionTabGroupLayout extends OrcaRuntime // Why: group address resolution (Section 4.5) queries per-handle status and must not throw on stale handles; return null on any error. getAgentStatusForHandle(handle: string): string | null { + // A structured worker has no pane and no title, so every PTY probe below answers null and + // `@idle` would enumerate it and then silently drop it. Its status is the journal's. + const structured = resolveStructuredWorkerAuthority(handle, this._orchestrationDb) + if (structured) { + return structuredWorkerAgentStatus(structured.identity.sessionId) + } try { const ptyId = this.getTerminalAgentStatusPtyId(handle) return this.getTerminalAgentStatusSnapshot(handle, ptyId).titleStatus diff --git a/src/main/runtime/orchestration/groups.ts b/src/main/runtime/orchestration/groups.ts index 64c0900bfb9..a4e548f06d0 100644 --- a/src/main/runtime/orchestration/groups.ts +++ b/src/main/runtime/orchestration/groups.ts @@ -1,5 +1,5 @@ -import type { RuntimeTerminalSummary } from '../../../shared/runtime-types' import type { TuiAgent } from '../../../shared/tui-agent' +import type { OrchestrationAddressableAgent } from './structured-worker-group-addressing' // Why: group addresses enable broadcast messaging to logical groups of agents. // Resolution is done at send-time: one message record per recipient, same thread_id, @@ -51,14 +51,17 @@ const GROUP_AGENT_IDS: Record = { * delivering is visible and recoverable — the sender sees no recipients; delivering to the wrong * agent is neither. */ -function terminalIsAgent(terminal: RuntimeTerminalSummary, agentName: AgentNameGroup): boolean { +function terminalIsAgent( + terminal: OrchestrationAddressableAgent, + agentName: AgentNameGroup +): boolean { return terminal.agentIdentity === GROUP_AGENT_IDS[agentName] } export function resolveGroupAddress( to: string, senderHandle: string, - terminals: RuntimeTerminalSummary[], + terminals: readonly OrchestrationAddressableAgent[], getAgentStatus: (handle: string) => string | null ): string[] { if (!isGroupAddress(to)) { diff --git a/src/main/runtime/orchestration/structured-mailbox-pointer-host.ts b/src/main/runtime/orchestration/structured-mailbox-pointer-host.ts index cb401c7f1f8..4b80b8df5d7 100644 --- a/src/main/runtime/orchestration/structured-mailbox-pointer-host.ts +++ b/src/main/runtime/orchestration/structured-mailbox-pointer-host.ts @@ -9,7 +9,10 @@ import { AGENT_SESSION_NOT_ATTACHED } from '../../native-chat/agent-session-wire/structured-agent-session-mutation-admission' import { getStructuredAgentSessionHost } from '../../native-chat/agent-session-wire/structured-agent-session-registry' import type { StructuredMailboxPointerHost } from './structured-mailbox-pointer-delivery' -import { structuredSessionGateFacts } from './structured-session-pointer-delivery' +import { + structuredSessionGateFacts, + type StructuredSessionGateFacts +} from './structured-session-pointer-delivery' /** Per-dispatch so one worker's nudges cannot exhaust the shared runtime operation-ledger budget. */ export function structuredPointerCallerKey(dispatchId: string): string { @@ -27,24 +30,36 @@ export function structuredSessionPointerCallerKey(sessionId: string): string { return `trusted-local:orchestration:session:${sessionId}` } +/** + * The idle gate for a structured session, read off its FULL reduced timeline. + * + * Never a bounded page. Settlement tombstones the running turn's lifecycle item rather than + * rewriting it to `completed`, so on any tail window an idle session and a busy one whose + * lifecycle item scrolled off look identical — and idle-with-history is the normal steady state of + * a working agent. Shared so the pointer lane and group addressing cannot disagree about it. + */ +export function readStructuredSessionGateFacts( + sessionId: string +): StructuredSessionGateFacts | null { + const host = getStructuredAgentSessionHost() + if (!host) { + return null + } + try { + return structuredSessionGateFacts(host.journalSnapshot(sessionId).items) + } catch (error) { + // Not attached is a retain reason, not a failure; anything else is still unreadable. + if ((error as Error)?.message !== AGENT_SESSION_NOT_ATTACHED.code) { + console.warn('[orchestration] structured journal unreadable', sessionId, error) + } + return null + } +} + export function createStructuredMailboxPointerHost(): StructuredMailboxPointerHost { return { readGateFacts(sessionId) { - const host = getStructuredAgentSessionHost() - if (!host) { - return null - } - try { - // The full reduced timeline, never a page: settlement tombstones the running turn's - // lifecycle item, so a bounded tail cannot tell an idle worker from a busy one. - return structuredSessionGateFacts(host.journalSnapshot(sessionId).items) - } catch (error) { - // Not attached is a retain reason, not a failure; anything else is still unreadable. - if ((error as Error)?.message !== AGENT_SESSION_NOT_ATTACHED.code) { - console.warn('[orchestration] structured journal unreadable', sessionId, error) - } - return null - } + return readStructuredSessionGateFacts(sessionId) }, currentFence(sessionId) { diff --git a/src/main/runtime/orchestration/structured-worker-group-addressing.test.ts b/src/main/runtime/orchestration/structured-worker-group-addressing.test.ts new file mode 100644 index 00000000000..7eff64bc61f --- /dev/null +++ b/src/main/runtime/orchestration/structured-worker-group-addressing.test.ts @@ -0,0 +1,166 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest' +import type { AgentJournalRenderItem } from '../../../shared/agent-session-journal-types' +import type { AgentSessionRecord } from '../../../shared/agent-session-record' + +const hostRef: { current: unknown } = { current: null } + +vi.mock('../../native-chat/agent-session-wire/structured-agent-session-registry', () => ({ + getStructuredAgentSessionHost: () => hostRef.current +})) + +const { listAddressableStructuredWorkers, structuredWorkerAgentStatus } = + await import('./structured-worker-group-addressing') +const { resolveGroupAddress } = await import('./groups') +const { + mintStructuredWorkerHandle, + mintStructuredWorkerPaneKey, + structuredWorkerIdentities, + structuredWorkerProcessIncarnation +} = await import('../structured-worker-identity') + +const SESSION_ID = 'a1b2c3d4-e5f6-4a7b-8c9d-0e1f2a3b4c5d' + +function idleTurn(): AgentJournalRenderItem { + return { + itemId: 'lifecycle-1', + body: { kind: 'status', text: 'done', turnLifecycle: { turnId: 't1', state: 'completed' } } + } as unknown as AgentJournalRenderItem +} + +function runningTurn(): AgentJournalRenderItem { + return { + itemId: 'lifecycle-1', + body: { kind: 'status', text: 'working', turnLifecycle: { turnId: 't1', state: 'running' } } + } as unknown as AgentJournalRenderItem +} + +function transcript(count: number): AgentJournalRenderItem[] { + return Array.from( + { length: count }, + (_unused, index) => + ({ + itemId: `tool-${index}`, + body: { kind: 'tool-call', name: 'Bash', input: {}, state: 'completed' } + }) as unknown as AgentJournalRenderItem + ) +} + +function installHost(options: { + items?: AgentJournalRenderItem[] + lease?: { runtimeKind: string; claimStatus: string } + hasSession?: boolean +}): void { + const lease = options.lease ?? { runtimeKind: 'native', claimStatus: 'live' } + hostRef.current = { + deps: { + store: { + getRecord: (sessionId: string) => + ({ + sessionId, + provider: 'codex', + location: { executionHostId: 'local', wslDistro: null }, + lease: { ...lease, runtimeFence: 1, deathEvidence: null } + }) as unknown as AgentSessionRecord + } + }, + hasSession: () => options.hasSession ?? true, + journalSnapshot: () => ({ items: options.items ?? [idleTurn()] }) + } +} + +function registerWorker(worktreeId = 'wt_1'): string { + const handle = mintStructuredWorkerHandle() + structuredWorkerIdentities.register({ + handle, + sessionId: SESSION_ID, + agent: 'codex', + paneKey: mintStructuredWorkerPaneKey(SESSION_ID), + processIncarnation: structuredWorkerProcessIncarnation(SESSION_ID), + worktreeId, + hostScope: { kind: 'local', hostId: 'local' } + }) + return handle +} + +const PTY_TERMINAL = { handle: 'term_a', worktreeId: 'wt_1', agentIdentity: 'claude' as const } + +describe('group addressing and structured workers', () => { + beforeEach(() => { + structuredWorkerIdentities.clear() + hostRef.current = null + }) + + it('enumerates a live structured worker as a candidate', () => { + const handle = registerWorker() + installHost({}) + expect(listAddressableStructuredWorkers()).toEqual([ + { handle, worktreeId: 'wt_1', agentIdentity: 'codex' } + ]) + }) + + it('leaves out a worker whose session is not proven live', () => { + // Addressing a settled worker would store mail no lane will ever deliver. + registerWorker() + installHost({ lease: { runtimeKind: 'native', claimStatus: 'live' }, hasSession: false }) + expect(listAddressableStructuredWorkers()).toEqual([]) + }) + + it('reaches a structured worker through @all', () => { + // The defect this pins: recipients came only from `listTerminals`, which enumerates leaves and + // PTYs, so a structured worker was excluded BEFORE per-recipient resolution — the warning + // machinery never ran and the sender got exit 0 with a receipt naming only who did resolve. + const handle = registerWorker() + installHost({}) + const recipients = [PTY_TERMINAL, ...listAddressableStructuredWorkers()] + expect(resolveGroupAddress('@all', 'term_sender', recipients, () => 'idle')).toContain(handle) + }) + + it('reaches a structured worker through @worktree: and @codex, but not @claude', () => { + const handle = registerWorker('wt_2') + installHost({}) + const recipients = [PTY_TERMINAL, ...listAddressableStructuredWorkers()] + expect(resolveGroupAddress('@worktree:wt_2', 'term_sender', recipients, () => 'idle')).toEqual([ + handle + ]) + expect(resolveGroupAddress('@codex', 'term_sender', recipients, () => 'idle')).toEqual([handle]) + expect(resolveGroupAddress('@claude', 'term_sender', recipients, () => 'idle')).toEqual([ + 'term_a' + ]) + }) + + it('reads @idle status off the FULL timeline, never a bounded tail', () => { + // The same trap that already cost this branch once: settlement tombstones the lifecycle item + // rather than rewriting it, so a long tool-calling turn pushes it arbitrarily far from the + // tail and any page-sized read reports a BUSY worker as idle — then `@idle` broadcasts into a + // running turn, which Codex refuses outright and Claude queues behind. + registerWorker() + installHost({ items: [runningTurn(), ...transcript(500)] }) + expect(structuredWorkerAgentStatus(SESSION_ID)).toBe('working') + }) + + it('answers idle only when no turn is running and no human is awaited', () => { + registerWorker() + installHost({ items: [idleTurn()] }) + expect(structuredWorkerAgentStatus(SESSION_ID)).toBe('idle') + installHost({ + items: [ + { + itemId: 'q1', + body: { + kind: 'question', + question: 'which?', + options: [], + resolution: { state: 'pending' } + } + } as unknown as AgentJournalRenderItem + ] + }) + expect(structuredWorkerAgentStatus(SESSION_ID)).toBe('attention') + }) + + it('answers null rather than idle when the session cannot be read', () => { + // Unknown must never read as idle, or `@idle` wakes a worker mid-turn. + hostRef.current = null + expect(structuredWorkerAgentStatus(SESSION_ID)).toBeNull() + }) +}) diff --git a/src/main/runtime/orchestration/structured-worker-group-addressing.ts b/src/main/runtime/orchestration/structured-worker-group-addressing.ts new file mode 100644 index 00000000000..8abad118de7 --- /dev/null +++ b/src/main/runtime/orchestration/structured-worker-group-addressing.ts @@ -0,0 +1,64 @@ +/** + * Structured workers as group-address recipients. + * + * `@all` and its siblings resolve recipients from `listTerminals`, which enumerates leaves and + * PTYs — so a structured worker was never a candidate. Worse, the exclusion happened BEFORE + * per-recipient resolution, so the `SendRecipientWarning` machinery never ran and the caller got + * exit 0 plus a receipt naming only the workers that did resolve. A broadcast "stop work" reached + * the PTY workers and silently missed the structured ones. + * + * Deliberately NOT solved by teaching `listTerminals` about structured sessions: that result is + * published to paired mobile and remote clients and to every consumer that assumes a summary has a + * `ptyId` or is writable, so it is its own change under + * `docs/reference/remote-wire-compatibility.md`. Group addressing needs three fields, and + * `RuntimeTerminalSummary` already satisfies them structurally — so the group resolver widens to + * the smaller shape instead, and nothing here has to invent a `worktreePath` or a `branch`. + */ + +import type { TuiAgent } from '../../../shared/tui-agent' +import { observeStructuredWorker, structuredWorkerAgent } from '../structured-worker-authority' +import { structuredWorkerIdentities } from '../structured-worker-identity' +import { readStructuredSessionGateFacts } from './structured-mailbox-pointer-host' + +/** The only facts group addressing reads off a recipient. */ +export type OrchestrationAddressableAgent = { + handle: string + worktreeId: string + /** Absent means "unknown", and `@claude`/`@codex` fail closed on it, exactly as for a pane. */ + agentIdentity?: TuiAgent +} + +/** + * Live structured workers of this runtime, as group-address candidates. + * + * Liveness-gated on the same observation the rest of the structured surface uses: a settled or + * handed-off worker is not a recipient, and addressing one would store mail no lane will deliver. + */ +export function listAddressableStructuredWorkers(): OrchestrationAddressableAgent[] { + return structuredWorkerIdentities + .list() + .filter((identity) => observeStructuredWorker(identity).status === 'live') + .map((identity) => ({ + handle: identity.handle, + worktreeId: identity.worktreeId, + agentIdentity: structuredWorkerAgent(identity) as TuiAgent + })) +} + +/** + * A structured worker's agent status, in the vocabulary `@idle` already matches on. + * + * Null when the session cannot be read: unknown must not read as idle, or a broadcast to `@idle` + * would wake a worker mid-turn — which Codex answers with `turn already running` and Claude queues + * behind the running turn. + */ +export function structuredWorkerAgentStatus(sessionId: string): string | null { + const facts = readStructuredSessionGateFacts(sessionId) + if (!facts) { + return null + } + if (facts.awaitingHuman) { + return 'attention' + } + return facts.turnRunning ? 'working' : 'idle' +} diff --git a/src/main/runtime/rpc/methods/orchestration-send-group.ts b/src/main/runtime/rpc/methods/orchestration-send-group.ts index aa3c8d47788..1d30c1cce34 100644 --- a/src/main/runtime/rpc/methods/orchestration-send-group.ts +++ b/src/main/runtime/rpc/methods/orchestration-send-group.ts @@ -3,6 +3,7 @@ import type { OrcaRuntimeService } from '../../orca-runtime' import { OrchestrationError } from '../../orchestration/orchestration-error' import { resolveGroupAddress } from '../../orchestration/groups' import { resolveBareOrchestrationRecipient } from './orchestration-recipient-routing' +import { listAddressableStructuredWorkers } from '../../orchestration/structured-worker-group-addressing' import { legacyWorkerDeliveryContract } from './orchestration-routing' import type { SendRecipientWarning } from './orchestration-recipient-routing' import type { SendParams } from './orchestration-schemas' @@ -42,7 +43,11 @@ export async function sendGroupMessage(args: { const { terminals } = await runtime.listTerminals(undefined, undefined, { includeVisualLayouts: false }) - const handles = resolveGroupAddress(groupAddress, from, terminals, (handle: string) => + // Structured workers are on no PTY surface, so `listTerminals` cannot see them and a broadcast + // silently missed every one. Composed here rather than inside `listTerminals`, whose result is + // published to paired clients and to consumers that assume a summary is writable. + const recipients = [...terminals, ...listAddressableStructuredWorkers()] + const handles = resolveGroupAddress(groupAddress, from, recipients, (handle: string) => runtime.getAgentStatusForHandle(handle) ) if (handles.length === 0) { diff --git a/src/main/runtime/structured-worker-identity.ts b/src/main/runtime/structured-worker-identity.ts index a82134f6dd6..161ae55dd5d 100644 --- a/src/main/runtime/structured-worker-identity.ts +++ b/src/main/runtime/structured-worker-identity.ts @@ -141,6 +141,11 @@ export class StructuredWorkerIdentityRegistry { return this.bySessionId.get(sessionId) ?? null } + /** Every worker this process knows about; callers apply their own liveness gate. */ + list(): StructuredWorkerIdentity[] { + return [...this.byHandle.values()] + } + forget(handle: string): void { const identity = this.byHandle.get(handle) if (!identity) {