mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 00:03:15 +00:00
fix(orchestration): preserve pane identity and exclude coordinator dispatches
This commit is contained in:
@@ -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}`])
|
||||
}
|
||||
)
|
||||
})
|
||||
|
||||
@@ -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<typeof SendParams>
|
||||
type SendReceipt = <T extends object>(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:<id>` 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
|
||||
|
||||
Reference in New Issue
Block a user