mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 00:03:15 +00:00
fix(orchestration): reach structured workers through group addresses
`orca orchestration send --to @all` — and `@idle`, `@claude`, `@codex`, `@worktree:<id>` — 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.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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<AgentNameGroup, TuiAgent> = {
|
||||
* 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)) {
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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()
|
||||
})
|
||||
})
|
||||
@@ -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'
|
||||
}
|
||||
@@ -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) {
|
||||
|
||||
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user