mirror of
https://github.com/stablyai/orca.git
synced 2026-10-02 16:02:15 +00:00
fix(orchestration): deliver a cleared structured worker's mail to its live successor
A terminal keeps its handle across /clear; a structured worker's successor now does the same. Mail at the worker's handle, its dispatch mailbox, and a Run it coordinates resolved to the session minted for the worker, which /clear replaced. Each now walks the /clear lineage forward to the live session, which the caller resolver already treats as the worker.
This commit is contained in:
@@ -12,7 +12,8 @@ import {
|
||||
handleLessCoordinatorSessionId,
|
||||
structuredSessionAddressTarget,
|
||||
structuredSessionMailTarget,
|
||||
structuredSessionIdleEdgeMailboxes
|
||||
structuredSessionIdleEdgeMailboxes,
|
||||
structuredWorkerMailSessionId
|
||||
} from './orchestration/structured-session-mail-target'
|
||||
import {
|
||||
resolveTerminalIdentityFromProbes,
|
||||
@@ -242,8 +243,8 @@ export class OrcaRuntimeWithGetPtyRecordForPaneKey extends OrcaRuntimeWithPruneM
|
||||
if (!assignee) {
|
||||
return null
|
||||
}
|
||||
const identity = resolveStructuredWorkerAuthority(assignee, this._orchestrationDb)?.identity
|
||||
return identity ? { sessionId: identity.sessionId, dispatchId } : null
|
||||
const sessionId = this.liveStructuredWorkerSessionId(assignee)
|
||||
return sessionId ? { sessionId, dispatchId } : null
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -266,8 +267,8 @@ export class OrcaRuntimeWithGetPtyRecordForPaneKey extends OrcaRuntimeWithPruneM
|
||||
if (!coordinator) {
|
||||
return null
|
||||
}
|
||||
const identity = resolveStructuredWorkerAuthority(coordinator, this._orchestrationDb)?.identity
|
||||
return identity ? { sessionId: identity.sessionId, dispatchId: null } : null
|
||||
const workerSessionId = this.liveStructuredWorkerSessionId(coordinator)
|
||||
return workerSessionId ? { sessionId: workerSessionId, dispatchId: null } : null
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -289,11 +290,19 @@ export class OrcaRuntimeWithGetPtyRecordForPaneKey extends OrcaRuntimeWithPruneM
|
||||
// Answers null for anything that is not a live structured worker of THIS runtime, so `run:`
|
||||
// and PTY handles fall through to the PTY lane exactly as before.
|
||||
const identity = resolveStructuredWorkerAuthority(handle, db)?.identity
|
||||
if (!identity) {
|
||||
const sessionId = identity ? structuredWorkerMailSessionId(identity.sessionId) : null
|
||||
if (!identity || !sessionId) {
|
||||
return null
|
||||
}
|
||||
const dispatchId = db?.findActiveDispatchForAssignee?.(handle, identity.paneKey)?.id ?? null
|
||||
return { sessionId: identity.sessionId, dispatchId }
|
||||
return { sessionId, dispatchId }
|
||||
}
|
||||
|
||||
/** The live session behind a structured worker handle of this runtime; see
|
||||
* `structuredWorkerMailSessionId`. */
|
||||
private liveStructuredWorkerSessionId(handle: string): string | null {
|
||||
const identity = resolveStructuredWorkerAuthority(handle, this._orchestrationDb)?.identity
|
||||
return identity ? structuredWorkerMailSessionId(identity.sessionId) : null
|
||||
}
|
||||
|
||||
protected scheduleRestoredMessageRepoints(): void {
|
||||
|
||||
@@ -29,6 +29,7 @@ const { resolveOrcaSessionParty } = await import('./orchestration-party')
|
||||
const {
|
||||
mintStructuredWorkerHandle,
|
||||
mintStructuredWorkerPaneKey,
|
||||
structuredWorkerIdentities,
|
||||
structuredWorkerProcessIncarnation
|
||||
} = await import('../structured-worker-identity')
|
||||
|
||||
@@ -452,3 +453,79 @@ describe('a coordinator chat continued by /clear', () => {
|
||||
expect(db.getMessageById(direct.id)).toMatchObject({ to_handle: CHAT_ADDRESS, read: 0 })
|
||||
})
|
||||
})
|
||||
|
||||
describe('a structured worker continued by /clear', () => {
|
||||
const WORKER = testOrcaSessionId('9c2e4a61-3f7b-4d8e-b105-6a2d8e4f1c93')
|
||||
const WORKER_SUCCESSOR = testOrcaSessionId('clear-a1b2c3d4e5f60718293a4b5c6d7e8f9012345678')
|
||||
|
||||
afterEach(() => {
|
||||
structuredWorkerIdentities.clear()
|
||||
})
|
||||
|
||||
/** A worker minted for WORKER, whose conversation `/clear` continued in WORKER_SUCCESSOR. */
|
||||
function clearedWorker(): { handle: string; paneKey: string } {
|
||||
const store = installStore(null)
|
||||
const minted = agentSessionRecordFixture(
|
||||
agentSessionLeaseFixture({ sessionId: WORKER, runtimeKind: 'native' })
|
||||
)
|
||||
store.records.set(WORKER, {
|
||||
...minted,
|
||||
conversationCommand: {
|
||||
command: 'clear',
|
||||
state: 'completed',
|
||||
replacementSessionId: WORKER_SUCCESSOR,
|
||||
operationId: 'op-worker',
|
||||
callerKey: 'caller',
|
||||
phase: 'committed'
|
||||
}
|
||||
})
|
||||
store.records.set(
|
||||
WORKER_SUCCESSOR,
|
||||
agentSessionRecordFixture(
|
||||
agentSessionLeaseFixture({ sessionId: WORKER_SUCCESSOR, runtimeKind: 'native' })
|
||||
)
|
||||
)
|
||||
const identity = structuredWorkerIdentities.register({
|
||||
handle: mintStructuredWorkerHandle(),
|
||||
sessionId: WORKER,
|
||||
agent: 'codex',
|
||||
paneKey: mintStructuredWorkerPaneKey(WORKER),
|
||||
processIncarnation: structuredWorkerProcessIncarnation(WORKER),
|
||||
worktreeId: 'wt_1',
|
||||
hostScope: { kind: 'local', hostId: 'local' }
|
||||
})
|
||||
return { handle: identity.handle, paneKey: identity.paneKey }
|
||||
}
|
||||
|
||||
it("delivers mail at the worker's handle to the live successor, as a terminal keeps its handle", () => {
|
||||
// The strand this pins: the handle resolved to the session minted for it, which `/clear`
|
||||
// replaced, so the worker's own mail was pointed at a session that no longer runs its turns.
|
||||
const { handle } = clearedWorker()
|
||||
expect(probe().target(handle)).toEqual({ sessionId: WORKER_SUCCESSOR, dispatchId: null })
|
||||
})
|
||||
|
||||
it("delivers the worker's dispatch mailbox and a Run it coordinates to the live successor", () => {
|
||||
const { handle, paneKey } = clearedWorker()
|
||||
const dispatch = db.createDispatchContext({
|
||||
taskId: db.createTask({ runId: chatCoordinatedRun(), spec: 'work' }).id,
|
||||
assigneeHandle: handle,
|
||||
assigneePaneKey: paneKey,
|
||||
processIncarnation: structuredWorkerProcessIncarnation(WORKER),
|
||||
creator: { kind: 'system' },
|
||||
maxDepth: Number.MAX_SAFE_INTEGER
|
||||
})
|
||||
expect(probe().target(`dispatch:${dispatch.id}`)).toEqual({
|
||||
sessionId: WORKER_SUCCESSOR,
|
||||
dispatchId: dispatch.id
|
||||
})
|
||||
const workerRun = db.createRun({
|
||||
objective: 'o',
|
||||
coordinatorHandle: handle,
|
||||
coordinatorPaneKey: paneKey
|
||||
}).id
|
||||
expect(probe().target(`run:${workerRun}`)).toEqual({
|
||||
sessionId: WORKER_SUCCESSOR,
|
||||
dispatchId: null
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
/**
|
||||
* Where a mailbox owned by a structured session is delivered, for sessions that are not structured
|
||||
* workers: a chat that coordinates a Run (`run:<id>` with no coordinator handle) and a session
|
||||
* addressed directly at `session:<id>`. The session is resolved here, never a pane, and takes the
|
||||
* pointer as a session turn.
|
||||
* Where a mailbox owned by a structured session is delivered: a chat that coordinates a Run
|
||||
* (`run:<id>` with no coordinator handle), a session addressed directly at `session:<id>`, and the
|
||||
* live session behind a structured worker's handle. The session is resolved here, never a pane, and
|
||||
* takes the pointer as a session turn.
|
||||
*/
|
||||
|
||||
import {
|
||||
@@ -13,12 +13,14 @@ import {
|
||||
} from '../../../shared/orca-session-address'
|
||||
import type { OrchestrationDb } from './db'
|
||||
import { currentRunCoordinatorOrcaSessionId } from './db/runs/run-coordinator-orca-session'
|
||||
import { structuredWorkerHostScope } from '../structured-worker-identity'
|
||||
import type { StructuredPointerTarget } from './structured-mailbox-pointer-delivery'
|
||||
import {
|
||||
addressableSessionParty,
|
||||
structuredSessionMailReach
|
||||
} from './structured-session-mail-address'
|
||||
import {
|
||||
lineageLiveSession,
|
||||
readAgentSessionRecordStore,
|
||||
type AgentSessionRecordReader
|
||||
} from './structured-session-lineage'
|
||||
@@ -61,6 +63,18 @@ export function structuredSessionMailTarget(
|
||||
: null
|
||||
}
|
||||
|
||||
/**
|
||||
* The session a structured worker's mail reaches: the one minted for it, or that session's live
|
||||
* `/clear` successor, which carries on as the worker the way a terminal keeps its handle.
|
||||
*/
|
||||
export function structuredWorkerMailSessionId(
|
||||
mintedSessionId: string,
|
||||
store: AgentSessionRecordReader | null = readAgentSessionRecordStore()
|
||||
): string | null {
|
||||
const live = store ? lineageLiveSession(store, mintedSessionId) : null
|
||||
return live && structuredWorkerHostScope(live.location) ? live.sessionId : null
|
||||
}
|
||||
|
||||
/**
|
||||
* The target of a `session:<id>` mailbox; `undefined` when the handle is not a session address at
|
||||
* all, so other address forms keep their own resolution.
|
||||
|
||||
Reference in New Issue
Block a user