From 6b737e9d4d3269bb06ef7dc7d8d7dbbe8c9fc0c2 Mon Sep 17 00:00:00 2001 From: Jinwoo-H Date: Thu, 10 Sep 2026 14:37:46 -0400 Subject: [PATCH] fix(orchestration): preserve pane identity and exclude coordinator dispatches --- .../messaging/send-group.test.ts | 63 +++++++++++++++++++ .../orchestration/messaging/send-group.ts | 21 ++++++- 2 files changed, 81 insertions(+), 3 deletions(-) diff --git a/src/main/runtime/rpc/methods/orchestration/messaging/send-group.test.ts b/src/main/runtime/rpc/methods/orchestration/messaging/send-group.test.ts index 6884738201d..b076be06593 100644 --- a/src/main/runtime/rpc/methods/orchestration/messaging/send-group.test.ts +++ b/src/main/runtime/rpc/methods/orchestration/messaging/send-group.test.ts @@ -565,4 +565,67 @@ describe('orchestration.send group addresses', () => { message: `No recipients resolved for group address: ${to}` }) }) + it.each(['@all', '@idle', '@codex'])( + 'excludes an owning coordinator with a self-Dispatch from %s', + async (to) => { + setupWithTerminals( + [ + makeSummary('term_coord', { agentIdentity: 'codex' }), + makeSummary('term_a', { agentIdentity: 'codex' }), + makeSummary('term_b', { agentIdentity: 'codex' }) + ], + { term_coord: 'idle', term_a: 'idle', term_b: 'idle' } + ) + createRootDispatch( + db, + db.createTask({ spec: 'coordinator context', runId: activeRunId }).id, + 'term_coord', + coordinatorPaneKey + ) + dispatchWorker('term_a') + const sibling = dispatchWorker('term_b') + const result = (await call('orchestration.send', { + from: 'term_a', + to, + subject: 'siblings only' + })) as GroupReceipt + expect(result.messages.map((m) => m.to_handle)).toEqual([`dispatch:${sibling}`]) + expect(db.getUnreadMessages(`run:${activeRunId}`)).toHaveLength(0) + } + ) + + it.each(['term_snapshot', 'term_original'])( + 'preserves pane identity across discovery when the recorded handle is %s', + async (recordedHandle) => { + const pane = 'tab_worker:11111111-1111-4111-8111-111111111111' + const snapshot = makeSummary('term_snapshot', { + tabId: 'tab_worker', + leafId: '11111111-1111-4111-8111-111111111111', + agentIdentity: 'codex' + }) + setupWithTerminals([makeSummary('term_coord'), snapshot]) + const dispatch = createRootDispatch( + db, + db.createTask({ spec: 'worker', runId: activeRunId }).id, + recordedHandle, + pane + ) + vi.spyOn(runtime, 'getTerminalHandleForPaneKey').mockImplementation((key) => + key === pane ? 'term_snapshot' : null + ) + vi.mocked(runtime.listTerminals).mockImplementation(async () => { + // The captured identity still belongs to this pane after its handle is reissued. + vi.mocked(runtime.getTerminalHandleForPaneKey).mockImplementation((key) => + key === pane ? 'term_new' : null + ) + return { terminals: [snapshot], totalCount: 1, truncated: false } + }) + const result = (await call('orchestration.send', { + from: 'term_coord', + to: '@codex', + subject: 'codex guidance' + })) as GroupReceipt + expect(result.messages.map((m) => m.to_handle)).toEqual([`dispatch:${dispatch.id}`]) + } + ) }) diff --git a/src/main/runtime/rpc/methods/orchestration/messaging/send-group.ts b/src/main/runtime/rpc/methods/orchestration/messaging/send-group.ts index 412d24fb964..d5ee366a40f 100644 --- a/src/main/runtime/rpc/methods/orchestration/messaging/send-group.ts +++ b/src/main/runtime/rpc/methods/orchestration/messaging/send-group.ts @@ -2,6 +2,7 @@ import type { MessagePriority, MessageType, OrchestrationDb } from '../../../../ import type { OrcaRuntimeService } from '../../../../orca-runtime' import { OrchestrationError } from '../../../../orchestration/orchestration-error' import { resolveGroupAddress } from '../../../../orchestration/groups' +import { isEquivalentPaneKey } from '../../../../orchestration/db/pane-key-match' import { resolveBareOrchestrationRecipient } from './recipient-routing' import { listAddressableStructuredWorkers, @@ -17,6 +18,8 @@ import type { z } from 'zod' type SendParamsInput = z.infer type SendReceipt = (receipt: T) => T & { warnings?: SendRecipientWarning[] } +type GroupAgentSnapshot = OrchestrationAddressableAgent & { tabId?: string; leafId?: string } + /** Run candidates already identify a durable mailbox. */ type GroupCandidate = OrchestrationAddressableAgent & { mailbox?: { to: string; runId: string } } @@ -25,7 +28,7 @@ function listRunGroupCandidates(args: { runtime: OrcaRuntimeService senderRunId: string groupAddress: string - agents: readonly OrchestrationAddressableAgent[] + agents: readonly GroupAgentSnapshot[] warnings: SendRecipientWarning[] }): GroupCandidate[] { const { db, runtime, senderRunId, groupAddress, agents, warnings } = args @@ -59,7 +62,19 @@ function listRunGroupCandidates(args: { to // Nested coordinators consume their child Run mailbox, not their parent Dispatch mailbox. const coordinated = paneKey ? db.getCurrentRunForPane(paneKey) : undefined - const agentIdentity = identityByHandle.get(handle) + if (coordinated?.id === row.runId) { + return [] + } + // Discovery can precede a handle remint; the pane still owns the captured identity. + const agentIdentity = + identityByHandle.get(handle) ?? + agents.find( + (agent) => + paneKey && + agent.tabId && + agent.leafId && + isEquivalentPaneKey(`${agent.tabId}:${agent.leafId}`, paneKey) + )?.agentIdentity return [ { handle, @@ -125,7 +140,7 @@ export async function sendGroupMessage(args: { // `@worktree:` names one workspace explicitly; every other group means the sender's Run. const worktreeGroup = groupAddress.toLowerCase().startsWith('@worktree:') let audienceRunId = worktreeGroup ? undefined : resolveAudienceRunId() - let agents: OrchestrationAddressableAgent[] = [] + let agents: GroupAgentSnapshot[] = [] if (worktreeGroup || !['@all', '@idle'].includes(groupAddress.toLowerCase())) { const { terminals } = await runtime.listTerminals(undefined, undefined, { includeVisualLayouts: false