mirror of
https://github.com/stablyai/orca.git
synced 2026-10-03 16:02:11 +00:00
refactor(orchestration): canonicalize a cleared session through the one id hook and the party resolver
- canonicalOrcaSessionId now walks a session's /clear lineage to its root; the parallel session identity and lost-worker rule are deleted, so the caller resolver, recipient routing, reach and idle-edge mailboxes all resolve a session through resolveOrcaSessionParty. - A Dispatch row's assignee_orca_session_id goes through the same hook. - The terminal-view delivery lane is gone with the terminal handoff: no PTY is bound to a session, so a chat's mail is always a session turn. - Pins a send to a Run-less chat's session address after a restart, when the send must start the agent-session host before routing reads its record.
This commit is contained in:
@@ -10,9 +10,7 @@ import type { StructuredPointerTarget } from './orchestration/structured-mailbox
|
||||
import { releaseRestoredStructuredPointerClaims } from './orchestration/structured-pointer-claim-restore'
|
||||
import {
|
||||
handleLessCoordinatorSessionId,
|
||||
findConnectedPtyBoundToSession,
|
||||
structuredSessionAddressTarget,
|
||||
structuredSessionMailDestination,
|
||||
structuredSessionMailTarget,
|
||||
structuredSessionIdleEdgeMailboxes
|
||||
} from './orchestration/structured-session-mail-target'
|
||||
@@ -198,8 +196,8 @@ export class OrcaRuntimeWithGetPtyRecordForPaneKey extends OrcaRuntimeWithPruneM
|
||||
|
||||
/**
|
||||
* Every structured session's status change reaches here. At its idle edge, retry what is parked
|
||||
* on it and re-derive the mailboxes it owns, so mail it could not take earlier (evicted, closed,
|
||||
* in its terminal view) is pointed again. Workers and chats alike: this is not per-dispatch.
|
||||
* on it and re-derive the mailboxes it owns, so mail it could not take earlier (evicted, closed)
|
||||
* is pointed again. Workers and chats alike: this is not per-dispatch.
|
||||
*/
|
||||
onStructuredSessionStatusForMail(summary: {
|
||||
sessionId: string
|
||||
@@ -214,16 +212,6 @@ export class OrcaRuntimeWithGetPtyRecordForPaneKey extends OrcaRuntimeWithPruneM
|
||||
structuredSessionIdleEdgeMailboxes(summary.sessionId, openDb).forEach(deliver)
|
||||
}
|
||||
|
||||
/** The terminal of a session's terminal view, while a TUI owns it; the PTY lane types there. */
|
||||
getTerminalViewHandleForSession(sessionId: string): string | null {
|
||||
const destination = structuredSessionMailDestination(sessionId, this._orchestrationDb)
|
||||
const pty =
|
||||
destination?.view === 'terminal-view'
|
||||
? findConnectedPtyBoundToSession(this.ptysById.values(), destination.sessionId)
|
||||
: undefined
|
||||
return pty?.paneKey ? this.getTerminalHandleForPaneKey(pty.paneKey) : null
|
||||
}
|
||||
|
||||
/** Settlement drops anything parked for the session; nothing will ever redrive it again. */
|
||||
forgetStructuredSessionMail(sessionId: string): void {
|
||||
this.orchestrationStructuredMailboxPointerDelivery.forgetSession(sessionId)
|
||||
|
||||
@@ -2,7 +2,6 @@
|
||||
import { OrchestrationStructuredMailboxPointerDelivery } from './orchestration/structured-mailbox-pointer-delivery'
|
||||
import { createStructuredMailboxPointerHost } from './orchestration/structured-mailbox-pointer-host'
|
||||
import { isStructuredWorkerHandle } from './structured-worker-identity'
|
||||
import { agentSessionPtyWriteGate } from './agent-session-pty-write-gate'
|
||||
import { resolveStructuredWorkerAuthority } from './structured-worker-authority'
|
||||
import { OrcaRuntimeWithRuntimeId } from './orca-runtime-runtime-id'
|
||||
import { RuntimeTerminalAgentPresence } from './runtime-terminal-agent-presence'
|
||||
@@ -184,8 +183,7 @@ export class OrcaRuntimeWithStopRequestedPtyIds extends OrcaRuntimeWithRuntimeId
|
||||
onRoutedMessageTypes: (mailboxHandle, types) =>
|
||||
this.messageWaiters.notifyRouted(mailboxHandle, types),
|
||||
onForeignMailboxRouted: (mailboxHandle, messageType) =>
|
||||
this.notifyMessageArrived(mailboxHandle, messageType),
|
||||
getBoundSessionIdForPty: (ptyId) => agentSessionPtyWriteGate.boundSessionId(ptyId)
|
||||
this.notifyMessageArrived(mailboxHandle, messageType)
|
||||
})
|
||||
|
||||
protected readonly orchestrationMailboxDeliveryTarget = new OrchestrationMailboxDeliveryTarget({
|
||||
@@ -195,8 +193,7 @@ export class OrcaRuntimeWithStopRequestedPtyIds extends OrcaRuntimeWithRuntimeId
|
||||
isStructuredWorkerHandle: (handle) => isStructuredWorkerHandle(handle),
|
||||
canProbePtyLiveness: () => Boolean(this.ptyController?.probePtyLiveness),
|
||||
controllerKnowsPtyIsLive: (ptyId) => this.controllerKnowsPtyIsLive(ptyId),
|
||||
isLeafPtyProvenAbsent: (ptyId) => this.isLeafPtyProvenAbsent(ptyId),
|
||||
getTerminalViewHandleForSession: (sessionId) => this.getTerminalViewHandleForSession(sessionId)
|
||||
isLeafPtyProvenAbsent: (ptyId) => this.isLeafPtyProvenAbsent(ptyId)
|
||||
})
|
||||
|
||||
protected readonly orchestrationMailboxPointerDelivery = new OrchestrationMailboxPointerDelivery({
|
||||
|
||||
@@ -1,7 +1,38 @@
|
||||
import type { OrcaSessionId } from '../../../shared/orca-session-address'
|
||||
import { isOrcaSessionId, type OrcaSessionId } from '../../../shared/orca-session-address'
|
||||
import {
|
||||
clearedInto,
|
||||
readAgentSessionRecordStore,
|
||||
type AgentSessionRecordReader
|
||||
} from './structured-session-lineage'
|
||||
|
||||
/** The Orca session id orchestration addresses a session by; every session-to-party step calls this. */
|
||||
export function canonicalOrcaSessionId(orcaSessionId: OrcaSessionId): OrcaSessionId {
|
||||
// Later lineage canonicalization (a `/clear`ed session to its lineage root) plugs in here.
|
||||
return orcaSessionId
|
||||
/**
|
||||
* The Orca session id orchestration addresses a session by: the first session of its `/clear`
|
||||
* lineage, so a cleared chat keeps the address, Runs and mail it had. Every session-to-party step
|
||||
* calls this. Without a record store there is no lineage to read, and the id stands for itself.
|
||||
*/
|
||||
export function canonicalOrcaSessionId(
|
||||
orcaSessionId: OrcaSessionId,
|
||||
store: AgentSessionRecordReader | null = readAgentSessionRecordStore()
|
||||
): OrcaSessionId {
|
||||
if (!store) {
|
||||
return orcaSessionId
|
||||
}
|
||||
const clearedFrom = new Map<string, string>()
|
||||
for (const record of store.listRecords()) {
|
||||
const next = clearedInto(record)
|
||||
if (next) {
|
||||
clearedFrom.set(next, record.sessionId)
|
||||
}
|
||||
}
|
||||
// A clear chain is acyclic by construction; the visited set only bounds a corrupt store.
|
||||
let root: string = orcaSessionId
|
||||
const earlier = new Set([root])
|
||||
let prior = clearedFrom.get(root)
|
||||
while (prior && !earlier.has(prior)) {
|
||||
earlier.add(prior)
|
||||
root = prior
|
||||
prior = clearedFrom.get(root)
|
||||
}
|
||||
// Record ids are minted as Orca session ids; one that is not cannot name the conversation.
|
||||
return isOrcaSessionId(root) ? root : orcaSessionId
|
||||
}
|
||||
|
||||
@@ -9,7 +9,7 @@ import { recordedCreatorIdentity, type DispatchCreator } from '../dispatch-depth
|
||||
import type { OrchestrationDb } from '../orchestration-db'
|
||||
import { transitionLifecycleWithDb } from '../lifecycle-transition'
|
||||
import { taskNotFoundError, taskNotStartableError } from '../../task-dispatch-refusal'
|
||||
import { structuredWorkerOrcaSessionIdForIncarnation } from '../../../structured-worker-identity'
|
||||
import { dispatchAssigneeOrcaSessionId } from '../../dispatch-assignee-orca-session-id'
|
||||
|
||||
export function createDispatchContext(
|
||||
this: OrchestrationDb,
|
||||
@@ -66,7 +66,7 @@ export function createDispatchContext(
|
||||
launchTokenHash: launchTokenHash ?? null,
|
||||
assigneeHandle,
|
||||
assigneePaneKey: assigneePaneKey ?? null,
|
||||
assigneeOrcaSessionId: structuredWorkerOrcaSessionIdForIncarnation(processIncarnation),
|
||||
assigneeOrcaSessionId: dispatchAssigneeOrcaSessionId(processIncarnation),
|
||||
processIncarnation: processIncarnation ?? null,
|
||||
creatorDispatchId,
|
||||
...recordedCreatorIdentity(params.creator),
|
||||
|
||||
@@ -2,7 +2,7 @@ import { randomBytes } from 'node:crypto'
|
||||
import { OrchestrationError } from '../../orchestration-error'
|
||||
import { hashDispatchCapability } from '../dispatch-capability-hash'
|
||||
import type { OrchestrationDb } from '../orchestration-db'
|
||||
import { structuredWorkerOrcaSessionIdForIncarnation } from '../../../structured-worker-identity'
|
||||
import { dispatchAssigneeOrcaSessionId } from '../../dispatch-assignee-orca-session-id'
|
||||
|
||||
export function prepareStartingWorkerAuthority(
|
||||
this: OrchestrationDb,
|
||||
@@ -64,7 +64,7 @@ export function prepareStartingWorkerAuthority(
|
||||
.run(
|
||||
params.handle,
|
||||
params.paneKey,
|
||||
structuredWorkerOrcaSessionIdForIncarnation(params.processIncarnation),
|
||||
dispatchAssigneeOrcaSessionId(params.processIncarnation),
|
||||
params.processIncarnation,
|
||||
params.hostScope ?? null,
|
||||
hashDispatchCapability(capability),
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import type { WorkerDispatchRow } from '../../types'
|
||||
import type { OrchestrationDb } from '../orchestration-db'
|
||||
import { structuredWorkerOrcaSessionIdForIncarnation } from '../../../structured-worker-identity'
|
||||
import { dispatchAssigneeOrcaSessionId } from '../../dispatch-assignee-orca-session-id'
|
||||
|
||||
/**
|
||||
* A start that dies before `prepareStartingWorkerAuthority` never filled the Dispatch context in,
|
||||
@@ -29,7 +29,7 @@ export function recordFailedStartDispatchIdentity(
|
||||
.run(
|
||||
resource.terminal_handle,
|
||||
resource.pane_key,
|
||||
structuredWorkerOrcaSessionIdForIncarnation(resource.process_incarnation),
|
||||
dispatchAssigneeOrcaSessionId(resource.process_incarnation),
|
||||
resource.process_incarnation,
|
||||
resource.host_scope,
|
||||
worker.dispatch_id
|
||||
|
||||
@@ -0,0 +1,11 @@
|
||||
import type { OrcaSessionId } from '../../../shared/orca-session-address'
|
||||
import { structuredWorkerOrcaSessionIdForIncarnation } from '../structured-worker-identity'
|
||||
import { canonicalOrcaSessionId } from './canonical-orca-session-id'
|
||||
|
||||
/** The Orca session id a Dispatch row stores for the structured worker a process incarnation names. */
|
||||
export function dispatchAssigneeOrcaSessionId(
|
||||
processIncarnation: string | null | undefined
|
||||
): OrcaSessionId | null {
|
||||
const orcaSessionId = structuredWorkerOrcaSessionIdForIncarnation(processIncarnation)
|
||||
return orcaSessionId === null ? null : canonicalOrcaSessionId(orcaSessionId)
|
||||
}
|
||||
@@ -1,6 +1,5 @@
|
||||
import type { OrchestrationDb } from './db'
|
||||
import type { OrchestrationMailboxLeaf } from './mailbox-owner'
|
||||
import { handleLessCoordinatorSessionId } from './structured-session-mail-target'
|
||||
|
||||
type OrchestrationMailboxDeliveryTargetDependencies = {
|
||||
getDb: () => OrchestrationDb | null
|
||||
@@ -11,8 +10,6 @@ type OrchestrationMailboxDeliveryTargetDependencies = {
|
||||
canProbePtyLiveness: () => boolean
|
||||
controllerKnowsPtyIsLive: (ptyId: string) => boolean
|
||||
isLeafPtyProvenAbsent: (ptyId: string) => Promise<boolean>
|
||||
/** The terminal of a structured session's terminal view, while a TUI owns that session. */
|
||||
getTerminalViewHandleForSession?: (sessionId: string) => string | null
|
||||
}
|
||||
|
||||
export class OrchestrationMailboxDeliveryTarget {
|
||||
@@ -34,12 +31,8 @@ export class OrchestrationMailboxDeliveryTarget {
|
||||
const remote =
|
||||
dispatchId && !dispatch ? db?.getRemoteDispatchAttachment?.(dispatchId) : undefined
|
||||
const paneKey = dispatch?.assignee_pane_key ?? remote?.pane_key
|
||||
const run = runId ? db?.getRun(runId) : undefined
|
||||
const coordinatorSessionId = run ? handleLessCoordinatorSessionId(run) : null
|
||||
const ownerHandle = runId
|
||||
? coordinatorSessionId
|
||||
? this.deps.getTerminalViewHandleForSession?.(coordinatorSessionId)
|
||||
: run?.coordinator_handle
|
||||
? db?.getRun(runId)?.coordinator_handle
|
||||
: ((paneKey ? this.deps.getTerminalHandleForPaneKey(paneKey) : null) ??
|
||||
dispatch?.assignee_handle ??
|
||||
remote?.terminal_handle)
|
||||
|
||||
@@ -1,8 +1,5 @@
|
||||
import type { AgentStatus } from '../../../shared/agent-detection'
|
||||
import { isOrcaSessionId } from '../../../shared/orca-session-address'
|
||||
import type { OrchestrationDb } from './db'
|
||||
import type { OrchestrationCallerIdentity } from './orchestration-caller-identity'
|
||||
import { sessionOrchestrationIdentity } from './structured-session-mail-address'
|
||||
|
||||
export type OrchestrationMailboxLeaf = {
|
||||
tabId: string
|
||||
@@ -33,8 +30,6 @@ type OrchestrationMailboxOwnerDependencies = {
|
||||
getTerminalProcessIncarnation: (terminalHandle: string) => string | null
|
||||
onRoutedMessageTypes: (mailboxHandle: string, types: readonly string[]) => void
|
||||
onForeignMailboxRouted: (mailboxHandle: string, messageType: string) => void
|
||||
/** The structured session a PTY is the terminal view of, if any. */
|
||||
getBoundSessionIdForPty?: (ptyId: string) => string | null
|
||||
}
|
||||
|
||||
export class OrchestrationMailboxOwner {
|
||||
@@ -62,23 +57,17 @@ export class OrchestrationMailboxOwner {
|
||||
return null
|
||||
}
|
||||
const paneKey = `${leaf.tabId}:${leaf.leafId}`
|
||||
const session = this.sessionMailOwner(db, leaf)
|
||||
const sessionRun = session ? db.getCurrentRunForCoordinator?.(session) : null
|
||||
const run = sessionRun ?? db.getCurrentRunForPane?.(paneKey)
|
||||
const run = db.getCurrentRunForPane?.(paneKey)
|
||||
if (run) {
|
||||
const address = sessionRun && session ? session.address : terminalHandle
|
||||
return this.resolveRunMailbox(db, leaf, address, run.id, requestedMailbox, options)
|
||||
return this.resolveRunMailbox(db, leaf, terminalHandle, run.id, requestedMailbox, options)
|
||||
}
|
||||
|
||||
const sessionDispatch = session
|
||||
? db.getActiveDispatchForIdentity?.(session.address, session.paneKey ?? undefined)
|
||||
: null
|
||||
const dispatch = sessionDispatch ?? db.getActiveDispatchForIdentity?.(terminalHandle, paneKey)
|
||||
const dispatch = db.getActiveDispatchForIdentity?.(terminalHandle, paneKey)
|
||||
if (dispatch) {
|
||||
return this.resolveDispatchMailbox(
|
||||
db,
|
||||
leaf,
|
||||
sessionDispatch && session ? session.address : terminalHandle,
|
||||
terminalHandle,
|
||||
dispatch.id,
|
||||
dispatch.run_id,
|
||||
requestedMailbox,
|
||||
@@ -103,21 +92,6 @@ export class OrchestrationMailboxOwner {
|
||||
return !requestedMailbox || requestedMailbox === terminalHandle ? terminalHandle : null
|
||||
}
|
||||
|
||||
/**
|
||||
* A terminal view speaks for its session (the CLI there acts as the session), so its Run and
|
||||
* Dispatch are the session's. A PTY-born worker adopted into a session still owns its mail by
|
||||
* terminal, which is why callers fall back to the pane when the session owns nothing.
|
||||
*/
|
||||
private sessionMailOwner(
|
||||
db: OrchestrationDb,
|
||||
leaf: OrchestrationMailboxLeaf
|
||||
): OrchestrationCallerIdentity | null {
|
||||
const sessionId = leaf.ptyId ? this.deps.getBoundSessionIdForPty?.(leaf.ptyId) : null
|
||||
return sessionId && isOrcaSessionId(sessionId)
|
||||
? sessionOrchestrationIdentity(sessionId, db)
|
||||
: null
|
||||
}
|
||||
|
||||
routeForeignDirectMessages(leaf: OrchestrationMailboxLeaf): RoutedOrchestrationMailbox[] {
|
||||
const db = this.deps.getDb()
|
||||
if (!db) {
|
||||
@@ -138,13 +112,10 @@ export class OrchestrationMailboxOwner {
|
||||
if (!ownerRunId) {
|
||||
return []
|
||||
}
|
||||
const session = this.sessionMailOwner(db, leaf)
|
||||
const sessionOwnsRun =
|
||||
session !== null && db.getCurrentRunForCoordinator?.(session)?.id === ownerRunId
|
||||
const routed = db.routeForeignDirectMessagesToOwnedMailboxes?.(
|
||||
sessionOwnsRun ? session.address : terminalHandle,
|
||||
terminalHandle,
|
||||
ownerRunId,
|
||||
sessionOwnsRun ? (session.paneKey ?? undefined) : paneKey
|
||||
paneKey
|
||||
)
|
||||
if (routed?.hasMore) {
|
||||
this.scheduleDirectReconciliation(leaf)
|
||||
|
||||
@@ -0,0 +1,44 @@
|
||||
/**
|
||||
* A structured session's `/clear` lineage, read off the durable session records. `/clear` continues
|
||||
* a chat in a new session; the committed clear on the old record names the session that replaced it.
|
||||
* Derived from the records every time; nothing is rewritten at a clear.
|
||||
*/
|
||||
|
||||
import type { AgentSessionRecord } from '../../../shared/agent-session-record'
|
||||
import { getStructuredAgentSessionHost } from '../../native-chat/agent-session-wire/structured-agent-session-registry'
|
||||
|
||||
export type AgentSessionRecordReader = {
|
||||
getRecord: (sessionId: string) => AgentSessionRecord | null
|
||||
listRecords: () => AgentSessionRecord[]
|
||||
/** Absent on a store that predates tab visibility; every session then counts as open. */
|
||||
getVisibleSessionTabIndex?: () => { present: boolean; sessionIds: string[] }
|
||||
}
|
||||
|
||||
/** Null until the agent-session host is installed; callers that must see records ensure it first. */
|
||||
export function readAgentSessionRecordStore(): AgentSessionRecordReader | null {
|
||||
return getStructuredAgentSessionHost()?.deps.store ?? null
|
||||
}
|
||||
|
||||
/** The session a committed `/clear` continued this one in, if any. */
|
||||
export function clearedInto(record: AgentSessionRecord): string | null {
|
||||
const command = record.conversationCommand
|
||||
return command?.command === 'clear' && command.phase === 'committed'
|
||||
? (command.replacementSessionId ?? null)
|
||||
: null
|
||||
}
|
||||
|
||||
/** The session running the lineage now; null when the chain names a session with no record. */
|
||||
export function lineageLiveSession(
|
||||
store: AgentSessionRecordReader,
|
||||
sessionId: string
|
||||
): AgentSessionRecord | null {
|
||||
let live = store.getRecord(sessionId)
|
||||
const later = new Set([sessionId])
|
||||
let next = live ? clearedInto(live) : null
|
||||
while (live && next && !later.has(next)) {
|
||||
later.add(next)
|
||||
live = store.getRecord(next)
|
||||
next = live ? clearedInto(live) : null
|
||||
}
|
||||
return live
|
||||
}
|
||||
@@ -3,10 +3,8 @@
|
||||
* told is its public address. Recipient routing and pointer delivery both read these rules off the
|
||||
* durable session record, so the two can never disagree about which sessions mail can reach.
|
||||
*
|
||||
* The address names a conversation, not one session of it. `/clear` continues a chat in a new
|
||||
* session, and the conversation keeps the address of its first session (its lineage root): a Run it
|
||||
* coordinates, mail sent to it, and what it sends all stay under that one spelling, and any session
|
||||
* of the lineage names it. Derived from the records every time; nothing is rewritten at a clear.
|
||||
* The address names a conversation, not one session of it: any session of a `/clear` lineage names
|
||||
* the lineage root's address (`canonicalOrcaSessionId`), and mail reaches the lineage's live session.
|
||||
*
|
||||
* A released lease does not end a session. The host evicts a chat nobody is looking at 15s after
|
||||
* its last turn and hands its lease back, and mail must wake it again (resume on demand). For mail,
|
||||
@@ -14,77 +12,13 @@
|
||||
*/
|
||||
|
||||
import type { AgentSessionRecord } from '../../../shared/agent-session-record'
|
||||
import {
|
||||
formatOrcaSessionAddress,
|
||||
isOrcaSessionId,
|
||||
type OrcaSessionId
|
||||
} from '../../../shared/orca-session-address'
|
||||
import { getStructuredAgentSessionHost } from '../../native-chat/agent-session-wire/structured-agent-session-registry'
|
||||
import {
|
||||
isRecordedStructuredWorkerSession,
|
||||
resolveStructuredWorkerIdentityForSession
|
||||
} from '../structured-worker-authority'
|
||||
import { isOrcaSessionId, type OrcaSessionId } from '../../../shared/orca-session-address'
|
||||
import { ORCHESTRATION_SESSION_CALLER_ERROR_CODES as CODES } from '../../../shared/orchestration-session-caller-codes'
|
||||
import { structuredWorkerHostScope } from '../structured-worker-identity'
|
||||
import type { OrchestrationDb } from './db'
|
||||
import type { OrchestrationCallerIdentity } from './orchestration-caller-identity'
|
||||
|
||||
export type AgentSessionRecordReader = {
|
||||
getRecord: (sessionId: string) => AgentSessionRecord | null
|
||||
listRecords: () => AgentSessionRecord[]
|
||||
/** Absent on a store that predates tab visibility; every session then counts as open. */
|
||||
getVisibleSessionTabIndex?: () => { present: boolean; sessionIds: string[] }
|
||||
}
|
||||
|
||||
export function readAgentSessionRecordStore(): AgentSessionRecordReader | null {
|
||||
return getStructuredAgentSessionHost()?.deps.store ?? null
|
||||
}
|
||||
|
||||
/** The session a committed `/clear` continued this one in, if any. */
|
||||
function clearedInto(record: AgentSessionRecord): string | null {
|
||||
const command = record.conversationCommand
|
||||
return command?.command === 'clear' && command.phase === 'committed'
|
||||
? (command.replacementSessionId ?? null)
|
||||
: null
|
||||
}
|
||||
|
||||
export type StructuredSessionLineage = {
|
||||
/** The conversation's first session: its address for as long as the conversation lasts. */
|
||||
rootSessionId: string
|
||||
/** The session running the conversation now; null when the chain names a session with no record. */
|
||||
live: AgentSessionRecord | null
|
||||
}
|
||||
|
||||
/** Any session of a `/clear` lineage, resolved to the lineage's root and its live end. */
|
||||
export function structuredSessionLineage(
|
||||
store: AgentSessionRecordReader,
|
||||
sessionId: string
|
||||
): StructuredSessionLineage {
|
||||
const clearedFrom = new Map<string, string>()
|
||||
for (const record of store.listRecords()) {
|
||||
const next = clearedInto(record)
|
||||
if (next) {
|
||||
clearedFrom.set(next, record.sessionId)
|
||||
}
|
||||
}
|
||||
// A clear chain is acyclic by construction; the visited sets only bound a corrupt store.
|
||||
let rootSessionId = sessionId
|
||||
const earlier = new Set([sessionId])
|
||||
let prior = clearedFrom.get(rootSessionId)
|
||||
while (prior && !earlier.has(prior)) {
|
||||
earlier.add(prior)
|
||||
rootSessionId = prior
|
||||
prior = clearedFrom.get(rootSessionId)
|
||||
}
|
||||
let live = store.getRecord(sessionId)
|
||||
const later = new Set([sessionId])
|
||||
let next = live ? clearedInto(live) : null
|
||||
while (live && next && !later.has(next)) {
|
||||
later.add(next)
|
||||
live = store.getRecord(next)
|
||||
next = live ? clearedInto(live) : null
|
||||
}
|
||||
return { rootSessionId, live }
|
||||
}
|
||||
import { OrchestrationError } from './orchestration-error'
|
||||
import { resolveOrcaSessionParty, type OrchestrationSessionParty } from './orchestration-party'
|
||||
import { lineageLiveSession, type AgentSessionRecordReader } from './structured-session-lineage'
|
||||
|
||||
export type OrcaAgentSessionLookup =
|
||||
| { kind: 'found'; record: AgentSessionRecord }
|
||||
@@ -122,17 +56,14 @@ export function structuredSessionMailReach(
|
||||
record: AgentSessionRecord,
|
||||
db: OrchestrationDb | null | undefined
|
||||
): StructuredSessionMailReach {
|
||||
const live = structuredSessionLineage(store, record.sessionId).live
|
||||
const live = lineageLiveSession(store, record.sessionId)
|
||||
if (!live) {
|
||||
return { kind: 'ended', reason: 'continuation-missing' }
|
||||
}
|
||||
if (!structuredWorkerHostScope(live.location)) {
|
||||
return { kind: 'other-host' }
|
||||
}
|
||||
const identity = isOrcaSessionId(live.sessionId)
|
||||
? sessionOrchestrationIdentity(live.sessionId, db, store)
|
||||
: null
|
||||
if (db && identity && hasLostStructuredWorkerIdentity(identity, db)) {
|
||||
if (isOrcaSessionId(live.sessionId) && !addressableSessionParty(live.sessionId, db)) {
|
||||
// Why: it can no longer act (the caller resolver refuses it), so mail to it could never be read.
|
||||
return { kind: 'ended', reason: 'worker-identity-lost' }
|
||||
}
|
||||
@@ -145,58 +76,19 @@ export function structuredSessionMailReach(
|
||||
}
|
||||
|
||||
/**
|
||||
* Which view carries a pointer to the session right now. A terminal view owns the session while a
|
||||
* TUI holds its lease, and its PTY takes the pointer; otherwise the host sends a session turn,
|
||||
* resuming an evicted session for it.
|
||||
* The party a session resolves to, or null for a worker whose identity this host lost: the party
|
||||
* resolver refuses it in every role, so it can never read mail and none is owed to it.
|
||||
*/
|
||||
export function structuredSessionDeliveryView(
|
||||
record: AgentSessionRecord
|
||||
): 'session-turn' | 'terminal-view' {
|
||||
return record.lease.runtimeKind === 'tui' && record.lease.claimStatus !== 'released'
|
||||
? 'terminal-view'
|
||||
: 'session-turn'
|
||||
}
|
||||
|
||||
/**
|
||||
* Who a session is to orchestration: its conversation's Orca session id (the lineage root's), plus
|
||||
* the handle and pane a structured worker was minted. The caller resolver and mail delivery both take
|
||||
* it from here, so a session is matched the same way whether it is sending, checking, or being
|
||||
* delivered to. Without a record store there is no lineage to read, and the id stands for itself.
|
||||
*/
|
||||
export function sessionOrchestrationIdentity(
|
||||
export function addressableSessionParty(
|
||||
sessionId: OrcaSessionId,
|
||||
db: OrchestrationDb | null | undefined,
|
||||
store: AgentSessionRecordReader | null = readAgentSessionRecordStore()
|
||||
): OrchestrationCallerIdentity & { orcaSessionId: OrcaSessionId } {
|
||||
const lineage = store ? structuredSessionLineage(store, sessionId) : null
|
||||
const root = lineage?.rootSessionId ?? sessionId
|
||||
// Record ids are minted as Orca session ids; one that is not cannot name the conversation.
|
||||
const orcaSessionId = isOrcaSessionId(root) ? root : sessionId
|
||||
// A worker identity is minted for one session, so it is looked up on the one running now.
|
||||
const worker = resolveStructuredWorkerIdentityForSession(
|
||||
lineage?.live?.sessionId ?? sessionId,
|
||||
db
|
||||
)
|
||||
return {
|
||||
orcaSessionId,
|
||||
address: worker?.handle ?? formatOrcaSessionAddress(orcaSessionId),
|
||||
terminalHandle: worker?.handle ?? null,
|
||||
paneKey: worker?.paneKey ?? null
|
||||
db: OrchestrationDb | null | undefined
|
||||
): OrchestrationSessionParty | null {
|
||||
try {
|
||||
return resolveOrcaSessionParty(sessionId, db)
|
||||
} catch (error) {
|
||||
if (error instanceof OrchestrationError && error.code === CODES.notLive) {
|
||||
return null
|
||||
}
|
||||
throw error
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* A structured worker whose worker identity this host no longer has. It may not act handle-less (that
|
||||
* would split one worker into two identities), so mail to it could never be read either. Checked on
|
||||
* the conversation, whose root session is the one a Dispatch assigned.
|
||||
*/
|
||||
export function hasLostStructuredWorkerIdentity(
|
||||
identity: OrchestrationCallerIdentity,
|
||||
db: OrchestrationDb
|
||||
): boolean {
|
||||
return (
|
||||
identity.terminalHandle === null &&
|
||||
identity.orcaSessionId !== null &&
|
||||
isRecordedStructuredWorkerSession(identity.orcaSessionId, db)
|
||||
)
|
||||
}
|
||||
|
||||
@@ -25,12 +25,15 @@ vi.mock('../../native-chat/agent-session-wire/structured-agent-session-registry'
|
||||
const { OrcaRuntimeWithGetPtyRecordForPaneKey } =
|
||||
await import('../orca-runtime-get-pty-record-for-pane-key')
|
||||
const { OrchestrationDb } = await import('./db')
|
||||
const { agentSessionPtyWriteGate } = await import('../agent-session-pty-write-gate')
|
||||
const { sessionOrchestrationIdentity } = await import('./structured-session-mail-address')
|
||||
const { resolveOrcaSessionParty } = await import('./orchestration-party')
|
||||
const {
|
||||
mintStructuredWorkerHandle,
|
||||
mintStructuredWorkerPaneKey,
|
||||
structuredWorkerProcessIncarnation
|
||||
} = await import('../structured-worker-identity')
|
||||
|
||||
const CHAT = testOrcaSessionId('4a1f6c2e-8b3d-4e7a-9c15-0d2b6e8f1a37')
|
||||
const CHAT_ADDRESS = formatOrcaSessionAddress(CHAT)
|
||||
const TERMINAL_VIEW_PANE = 'tab_view:77777777-7777-4777-8777-777777777777'
|
||||
|
||||
/** The real methods through the real prototype chain; a re-declared copy would pin nothing. */
|
||||
class MailTargetProbe extends OrcaRuntimeWithGetPtyRecordForPaneKey {
|
||||
@@ -99,7 +102,6 @@ beforeEach(() => {
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
agentSessionPtyWriteGate.detachRecordLookup()
|
||||
db.close()
|
||||
})
|
||||
|
||||
@@ -118,12 +120,6 @@ describe('a Run whose coordinator is a chat (an Orca session id, no handle)', ()
|
||||
expect(probe().target(`run:${runId}`)).toEqual({ sessionId: CHAT, dispatchId: null })
|
||||
})
|
||||
|
||||
it('leaves the mailbox to the PTY lane while the chat is in its terminal view', () => {
|
||||
installStore(chatRecord({ runtimeKind: 'tui' }))
|
||||
const runId = chatCoordinatedRun()
|
||||
expect(probe().target(`run:${runId}`)).toBeNull()
|
||||
})
|
||||
|
||||
it('does not deliver to a chat that was closed, cleared into no known session, or runs on another host', () => {
|
||||
const runId = chatCoordinatedRun()
|
||||
installStore(chatRecord(), false)
|
||||
@@ -185,25 +181,6 @@ describe('a session addressed directly', () => {
|
||||
})
|
||||
})
|
||||
|
||||
describe('a chat in its terminal view', () => {
|
||||
it('is reached through the PTY bound to it', () => {
|
||||
installStore(chatRecord({ runtimeKind: 'tui' }))
|
||||
agentSessionPtyWriteGate.attachRecordLookup(() => null)
|
||||
agentSessionPtyWriteGate.bindPty('pty-view', CHAT)
|
||||
const runtime = probe({
|
||||
ptysById: new Map([
|
||||
['pty-view', { ptyId: 'pty-view', connected: true, paneKey: TERMINAL_VIEW_PANE }]
|
||||
]),
|
||||
getTerminalHandleForPaneKey: (paneKey: string) =>
|
||||
paneKey === TERMINAL_VIEW_PANE ? 'term_view' : null
|
||||
})
|
||||
expect(runtime.getTerminalViewHandleForSession(CHAT)).toBe('term_view')
|
||||
// Its native view does not answer at the same time.
|
||||
installStore(chatRecord())
|
||||
expect(runtime.getTerminalViewHandleForSession(CHAT)).toBeNull()
|
||||
})
|
||||
})
|
||||
|
||||
describe('the idle edge of a structured session', () => {
|
||||
it('re-derives and delivers the mailboxes the session owns, and nothing while it works', () => {
|
||||
installStore(chatRecord())
|
||||
@@ -382,10 +359,23 @@ describe('a coordinator chat continued by /clear', () => {
|
||||
objective: 'o',
|
||||
coordinatorHandle: null,
|
||||
coordinatorPaneKey: null,
|
||||
coordinatorOrcaSessionId: sessionOrchestrationIdentity(sessionId, db).orcaSessionId
|
||||
coordinatorOrcaSessionId: resolveOrcaSessionParty(sessionId, db).orcaSessionId
|
||||
}).id
|
||||
}
|
||||
|
||||
it('stores a Dispatch assignee by the lineage root of the session its incarnation names', () => {
|
||||
installLineage(SUCCESSOR)
|
||||
const dispatch = db.createDispatchContext({
|
||||
taskId: db.createTask({ runId: chatCoordinatedRun(), spec: 'work' }).id,
|
||||
assigneeHandle: mintStructuredWorkerHandle(),
|
||||
assigneePaneKey: mintStructuredWorkerPaneKey(SUCCESSOR),
|
||||
processIncarnation: structuredWorkerProcessIncarnation(SUCCESSOR),
|
||||
creator: { kind: 'system' },
|
||||
maxDepth: Number.MAX_SAFE_INTEGER
|
||||
})
|
||||
expect(db.getDispatchContextById(dispatch.id)?.assignee_orca_session_id).toBe(CHAT)
|
||||
})
|
||||
|
||||
it("never unbinds the successor's own Run", () => {
|
||||
// The strand this pins: the successor's own run-create, then its idle edge rebinding the
|
||||
// predecessor's Run through an exclusive bind, which unbound the Run the successor created.
|
||||
@@ -397,7 +387,7 @@ describe('a coordinator chat continued by /clear', () => {
|
||||
expect(idleEdge(SUCCESSOR)).toContain(`run:${ownRun}`)
|
||||
idleEdge(SUCCESSOR)
|
||||
|
||||
const successor = sessionOrchestrationIdentity(SUCCESSOR, db)
|
||||
const successor = resolveOrcaSessionParty(SUCCESSOR, db)
|
||||
expect(db.getRunRaw(ownRun)).toMatchObject({
|
||||
coordinator_orca_session_id: successor.orcaSessionId,
|
||||
consumer_generation: generation
|
||||
@@ -418,9 +408,7 @@ describe('a coordinator chat continued by /clear', () => {
|
||||
consumer_generation: generation
|
||||
})
|
||||
for (const member of [CHAT, MIDDLE, SUCCESSOR]) {
|
||||
expect(db.getCurrentRunForCoordinator(sessionOrchestrationIdentity(member, db))?.id).toBe(
|
||||
runId
|
||||
)
|
||||
expect(db.getCurrentRunForCoordinator(resolveOrcaSessionParty(member, db))?.id).toBe(runId)
|
||||
}
|
||||
expect(probe().target(`run:${runId}`)).toEqual({ sessionId: SUCCESSOR, dispatchId: null })
|
||||
})
|
||||
@@ -434,7 +422,7 @@ describe('a coordinator chat continued by /clear', () => {
|
||||
installLineage(MIDDLE, SUCCESSOR)
|
||||
idleEdge(SUCCESSOR)
|
||||
|
||||
const successor = sessionOrchestrationIdentity(SUCCESSOR, db)
|
||||
const successor = resolveOrcaSessionParty(SUCCESSOR, db)
|
||||
expect(db.getRunRaw(middleRun)).toMatchObject({
|
||||
coordinator_orca_session_id: successor.orcaSessionId,
|
||||
consumer_generation: generation
|
||||
@@ -445,7 +433,7 @@ describe('a coordinator chat continued by /clear', () => {
|
||||
it('reaches the live session through any spelling of the conversation, and stores one', () => {
|
||||
installLineage(MIDDLE, SUCCESSOR)
|
||||
for (const member of [CHAT, MIDDLE, SUCCESSOR]) {
|
||||
expect(sessionOrchestrationIdentity(member, db)).toMatchObject({
|
||||
expect(resolveOrcaSessionParty(member, db)).toMatchObject({
|
||||
orcaSessionId: CHAT,
|
||||
address: CHAT_ADDRESS
|
||||
})
|
||||
|
||||
@@ -1,11 +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: its native view takes the pointer as a session turn
|
||||
* (the structured lane), its terminal view as bytes typed into the PTY that owns it (the PTY lane).
|
||||
* Exactly one view answers for a session at a time, so the two lanes never both claim a mailbox.
|
||||
* addressed directly at `session:<id>`. The session is resolved here, never a pane, and takes the
|
||||
* pointer as a session turn.
|
||||
*/
|
||||
|
||||
import {
|
||||
@@ -14,32 +11,19 @@ import {
|
||||
parseOrcaSessionAddress,
|
||||
type OrcaSessionId
|
||||
} from '../../../shared/orca-session-address'
|
||||
import { agentSessionPtyWriteGate } from '../agent-session-pty-write-gate'
|
||||
import type { OrchestrationDb } from './db'
|
||||
import { currentRunCoordinatorOrcaSessionId } from './db/runs/run-coordinator-orca-session'
|
||||
import type { StructuredPointerTarget } from './structured-mailbox-pointer-delivery'
|
||||
import {
|
||||
readAgentSessionRecordStore,
|
||||
sessionOrchestrationIdentity,
|
||||
structuredSessionDeliveryView,
|
||||
structuredSessionMailReach,
|
||||
type AgentSessionRecordReader
|
||||
addressableSessionParty,
|
||||
structuredSessionMailReach
|
||||
} from './structured-session-mail-address'
|
||||
import {
|
||||
readAgentSessionRecordStore,
|
||||
type AgentSessionRecordReader
|
||||
} from './structured-session-lineage'
|
||||
import type { RunRow } from './types'
|
||||
|
||||
/** The live PTY the write gate binds to a session: its terminal view, when a TUI owns it. */
|
||||
export function findConnectedPtyBoundToSession<T extends { ptyId: string; connected: boolean }>(
|
||||
ptys: Iterable<T>,
|
||||
sessionId: string
|
||||
): T | undefined {
|
||||
for (const pty of ptys) {
|
||||
if (pty.connected && agentSessionPtyWriteGate.boundSessionId(pty.ptyId) === sessionId) {
|
||||
return pty
|
||||
}
|
||||
}
|
||||
return undefined
|
||||
}
|
||||
|
||||
/**
|
||||
* The session a Run's coordinator binding names when that binding has no handle. A structured
|
||||
* worker coordinates by its own handle and resolves through it, so only a handle-less binding names
|
||||
@@ -62,29 +46,18 @@ export function handleLessCoordinatorSessionId(
|
||||
}
|
||||
|
||||
/**
|
||||
* The session that takes a pointer for `sessionId`'s conversation now (its live session, whichever
|
||||
* session of the lineage was named) and through which view; null when mail cannot reach it here.
|
||||
* The structured-lane target for `sessionId`'s conversation: its live session, whichever session of
|
||||
* the lineage was named; null when mail cannot reach it here.
|
||||
*/
|
||||
export function structuredSessionMailDestination(
|
||||
export function structuredSessionMailTarget(
|
||||
sessionId: string,
|
||||
db: OrchestrationDb | null | undefined,
|
||||
store: AgentSessionRecordReader | null = readAgentSessionRecordStore()
|
||||
): { sessionId: string; view: 'session-turn' | 'terminal-view' } | null {
|
||||
): StructuredPointerTarget | null {
|
||||
const record = store?.getRecord(sessionId)
|
||||
const reach = store && record ? structuredSessionMailReach(store, record, db) : null
|
||||
return reach?.kind === 'reachable'
|
||||
? { sessionId: reach.session.sessionId, view: structuredSessionDeliveryView(reach.session) }
|
||||
: null
|
||||
}
|
||||
|
||||
/** The structured-lane target for a session whose native view takes the pointer. */
|
||||
export function structuredSessionMailTarget(
|
||||
sessionId: string,
|
||||
db: OrchestrationDb | null | undefined
|
||||
): StructuredPointerTarget | null {
|
||||
const destination = structuredSessionMailDestination(sessionId, db)
|
||||
return destination?.view === 'session-turn'
|
||||
? { sessionId: destination.sessionId, dispatchId: null }
|
||||
? { sessionId: reach.session.sessionId, dispatchId: null }
|
||||
: null
|
||||
}
|
||||
|
||||
@@ -106,16 +79,16 @@ export function structuredSessionAddressTarget(
|
||||
/**
|
||||
* Every mailbox a session reads for itself: the Runs it coordinates and its own direct mail.
|
||||
* Re-derived from the database on each idle edge rather than remembered, so mail that arrived
|
||||
* while the session could not take it (closed, evicted, in the other view) is found again.
|
||||
* while the session could not take it (closed, evicted) is found again.
|
||||
*/
|
||||
export function structuredSessionOwnedMailboxes(sessionId: string, db: OrchestrationDb): string[] {
|
||||
if (!isOrcaSessionId(sessionId)) {
|
||||
const party = isOrcaSessionId(sessionId) ? addressableSessionParty(sessionId, db) : null
|
||||
if (!party) {
|
||||
return []
|
||||
}
|
||||
const identity = sessionOrchestrationIdentity(sessionId, db)
|
||||
const mailboxes = db.runsBoundToCoordinator(identity).map((run) => `run:${run.id}`)
|
||||
if (db.getUnreadDirectMessageTypes(identity.address).length > 0) {
|
||||
mailboxes.push(identity.address)
|
||||
const mailboxes = db.runsBoundToCoordinator(party).map((run) => `run:${run.id}`)
|
||||
if (db.getUnreadDirectMessageTypes(party.address).length > 0) {
|
||||
mailboxes.push(party.address)
|
||||
}
|
||||
return mailboxes
|
||||
}
|
||||
|
||||
@@ -1,86 +0,0 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import { OrchestrationDb } from './db'
|
||||
import { OrchestrationMailboxDeliveryTarget } from './mailbox-delivery-target'
|
||||
import { OrchestrationMailboxOwner, type OrchestrationMailboxLeaf } from './mailbox-owner'
|
||||
import { testOrcaSessionId } from '../../../shared/orca-session-address-test-fixture'
|
||||
|
||||
const CHAT = testOrcaSessionId('4a1f6c2e-8b3d-4e7a-9c15-0d2b6e8f1a37')
|
||||
const VIEW_TAB = 'tab_view'
|
||||
const VIEW_LEAF = '77777777-7777-4777-8777-777777777777'
|
||||
|
||||
let db: OrchestrationDb
|
||||
|
||||
beforeEach(() => {
|
||||
db = new OrchestrationDb(':memory:')
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
db.close()
|
||||
})
|
||||
|
||||
const viewLeaf: OrchestrationMailboxLeaf = {
|
||||
tabId: VIEW_TAB,
|
||||
leafId: VIEW_LEAF,
|
||||
ptyId: 'pty-view',
|
||||
writable: true,
|
||||
lastAgentStatus: 'idle',
|
||||
lastAgentStatusObservedLive: true,
|
||||
lastOscTitle: null
|
||||
}
|
||||
|
||||
function owner(boundSessionId: string | null): OrchestrationMailboxOwner {
|
||||
return new OrchestrationMailboxOwner({
|
||||
getDb: () => db,
|
||||
getLeaf: () => viewLeaf,
|
||||
getLeafKey: (tabId, leafId) => `${tabId}::${leafId}`,
|
||||
getTerminalHandleForLeafKey: () => 'term_view',
|
||||
getTerminalProcessIncarnation: () => null,
|
||||
onRoutedMessageTypes: vi.fn(),
|
||||
onForeignMailboxRouted: vi.fn(),
|
||||
getBoundSessionIdForPty: (ptyId) => (ptyId === 'pty-view' ? boundSessionId : null)
|
||||
})
|
||||
}
|
||||
|
||||
function chatRun(): string {
|
||||
return db.createRun({
|
||||
objective: 'o',
|
||||
coordinatorHandle: null,
|
||||
coordinatorPaneKey: null,
|
||||
coordinatorOrcaSessionId: CHAT
|
||||
}).id
|
||||
}
|
||||
|
||||
describe("a chat's terminal view reads the chat's mail", () => {
|
||||
it('owns the Run the chat coordinates, which no pane binding names', () => {
|
||||
// The CLI in a terminal view acts as the session, so the Run is bound by Orca session id and the pane
|
||||
// path finds nothing: without this the idle edge of the view never pointed coordinator mail.
|
||||
const runId = chatRun()
|
||||
expect(owner(CHAT).resolve(viewLeaf)).toBe(`run:${runId}`)
|
||||
expect(owner(CHAT).resolve(viewLeaf, `run:${runId}`)).toBe(`run:${runId}`)
|
||||
expect(owner(null).resolve(viewLeaf)).toBe('term_view')
|
||||
})
|
||||
|
||||
it('falls back to the pane for a terminal whose session owns nothing', () => {
|
||||
const runId = db.createRun({
|
||||
objective: 'o',
|
||||
coordinatorHandle: 'term_view',
|
||||
coordinatorPaneKey: `${VIEW_TAB}:${VIEW_LEAF}`
|
||||
}).id
|
||||
expect(owner(CHAT).resolve(viewLeaf)).toBe(`run:${runId}`)
|
||||
})
|
||||
|
||||
it("resolves a chat-coordinated Run's mailbox to the terminal view's handle", () => {
|
||||
const runId = chatRun()
|
||||
const target = new OrchestrationMailboxDeliveryTarget({
|
||||
getDb: () => db,
|
||||
getTerminalHandleForPaneKey: () => null,
|
||||
hasTerminalHandle: (handle) => handle === 'term_view',
|
||||
isStructuredWorkerHandle: () => false,
|
||||
canProbePtyLiveness: () => false,
|
||||
controllerKnowsPtyIsLive: () => true,
|
||||
isLeafPtyProvenAbsent: async () => false,
|
||||
getTerminalViewHandleForSession: (sessionId) => (sessionId === CHAT ? 'term_view' : null)
|
||||
})
|
||||
expect(target.resolveTerminalHandle(`run:${runId}`)).toBe('term_view')
|
||||
})
|
||||
})
|
||||
@@ -56,6 +56,8 @@ function installHost(options: {
|
||||
hostRef.current = {
|
||||
deps: {
|
||||
store: {
|
||||
// No committed /clear: each session is its own lineage's root.
|
||||
listRecords: () => [],
|
||||
getRecord: (sessionId: string) =>
|
||||
({
|
||||
sessionId,
|
||||
|
||||
@@ -12,7 +12,6 @@
|
||||
*/
|
||||
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import { setStructuredAgentSessionHost } from '../../../native-chat/agent-session-wire/structured-agent-session-registry'
|
||||
import { OrcaRuntimeService } from '../../orca-runtime'
|
||||
import { OrchestrationDb } from '../../orchestration/db'
|
||||
import {
|
||||
@@ -29,6 +28,12 @@ const STRUCTURED_HANDLE = 'structworker_worker'
|
||||
const TERMINAL_HANDLE = 'term_worker'
|
||||
|
||||
const structuredPreambles: string[] = []
|
||||
// The session host the code under test reads; a structural fake, so no host type is claimed.
|
||||
const hostRef = vi.hoisted((): { current: unknown } => ({ current: null }))
|
||||
|
||||
vi.mock('../../../native-chat/agent-session-wire/structured-agent-session-registry', () => ({
|
||||
getStructuredAgentSessionHost: () => hostRef.current
|
||||
}))
|
||||
|
||||
vi.mock('./orchestration/worker/worker-topology', async (importOriginal) => ({
|
||||
...(await importOriginal<Record<string, unknown>>()),
|
||||
@@ -70,7 +75,7 @@ function installStructuredCoordinator(handle: string, sessionId: string): string
|
||||
worktreeId: WORKTREE,
|
||||
hostScope: { kind: 'local', hostId: 'local' }
|
||||
})
|
||||
setStructuredAgentSessionHost({
|
||||
hostRef.current = {
|
||||
hasSession: () => true,
|
||||
deps: {
|
||||
store: {
|
||||
@@ -83,10 +88,12 @@ function installStructuredCoordinator(handle: string, sessionId: string): string
|
||||
deathEvidence: null,
|
||||
runtimeFence: 1
|
||||
}
|
||||
})
|
||||
}),
|
||||
// No committed /clear: each session is its own lineage's root.
|
||||
listRecords: () => []
|
||||
}
|
||||
}
|
||||
} as never)
|
||||
}
|
||||
return paneKey
|
||||
}
|
||||
|
||||
@@ -150,7 +157,7 @@ describe('a worker cannot tell which mode it is running in', () => {
|
||||
|
||||
afterEach(() => {
|
||||
db.close()
|
||||
setStructuredAgentSessionHost(null)
|
||||
hostRef.current = null
|
||||
structuredWorkerIdentities.clear()
|
||||
vi.restoreAllMocks()
|
||||
})
|
||||
|
||||
@@ -5,7 +5,7 @@ import type { OrcaRuntimeService } from '../../../../orca-runtime'
|
||||
import { resolveOrchestrationParty } from '../../../../orchestration/orchestration-party'
|
||||
import { isEquivalentPaneKey } from '../../../../orchestration/db/pane-key-match'
|
||||
import { CURRENT_CONTRACT_VERSION } from '../../../../orchestration/db/contract-constants'
|
||||
import { readAgentSessionRecordStore } from '../../../../orchestration/structured-session-mail-address'
|
||||
import { readAgentSessionRecordStore } from '../../../../orchestration/structured-session-lineage'
|
||||
import {
|
||||
readSessionRecipient,
|
||||
refuseUndeliverableSessionRecipient,
|
||||
|
||||
@@ -21,9 +21,9 @@ import {
|
||||
import { ORCHESTRATION_SESSION_CALLER_ERROR_CODES as CODES } from '../../../../../../shared/orchestration-session-caller-codes'
|
||||
import {
|
||||
lookupOrcaAgentSession,
|
||||
structuredSessionMailReach,
|
||||
type AgentSessionRecordReader
|
||||
structuredSessionMailReach
|
||||
} from '../../../../orchestration/structured-session-mail-address'
|
||||
import type { AgentSessionRecordReader } from '../../../../orchestration/structured-session-lineage'
|
||||
import type { OrchestrationDb } from '../../../../orchestration/db'
|
||||
|
||||
/** `address` is the named session's own spelling; the mailbox mail lands in is its identity address. */
|
||||
|
||||
@@ -29,10 +29,8 @@ import {
|
||||
resolveOrcaSessionParty,
|
||||
resolveOrchestrationParty
|
||||
} from '../orchestration/orchestration-party'
|
||||
import {
|
||||
lookupOrcaAgentSession,
|
||||
type AgentSessionRecordReader
|
||||
} from '../orchestration/structured-session-mail-address'
|
||||
import { lookupOrcaAgentSession } from '../orchestration/structured-session-mail-address'
|
||||
import { readAgentSessionRecordStore } from '../orchestration/structured-session-lineage'
|
||||
import { structuredWorkerHostScope } from '../structured-worker-identity'
|
||||
import type { RpcRequest } from './core'
|
||||
|
||||
|
||||
@@ -7,6 +7,7 @@ import {
|
||||
} from '../structured-worker-identity'
|
||||
import {
|
||||
ADDRESS_X,
|
||||
ADDRESS_Y,
|
||||
createSessionCallerHarness,
|
||||
idOf,
|
||||
isRecord,
|
||||
@@ -16,6 +17,7 @@ import {
|
||||
SESSION_X,
|
||||
SESSION_Y,
|
||||
sessionRecord,
|
||||
WORKER_HANDLE,
|
||||
type SessionCallerHarness
|
||||
} from './orchestration-session-caller-test-fixture'
|
||||
|
||||
@@ -217,3 +219,114 @@ describe('a live structured worker addressed by its session id', () => {
|
||||
expect(await flaglessCheck()).toMatchObject({ messages: [{ subject: 'hello' }] })
|
||||
})
|
||||
})
|
||||
|
||||
describe('mail sent to a session address reaches the mailbox that session reads', () => {
|
||||
let h: SessionCallerHarness
|
||||
|
||||
beforeEach(() => {
|
||||
h = createSessionCallerHarness(hostRef)
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
h.close()
|
||||
vi.restoreAllMocks()
|
||||
})
|
||||
|
||||
async function as(sessionId: string | undefined, method: string, params: Row): Promise<Row> {
|
||||
return resultOf(await h.dispatch(orchestrationRequest(method, params, { sessionId })))
|
||||
}
|
||||
|
||||
function sendFromTerminal(to: string): Promise<Row> {
|
||||
return as(undefined, 'orchestration.send', { from: WORKER_HANDLE, to, subject: 'hello' })
|
||||
}
|
||||
|
||||
it("files it under the chat's current Run, as a terminal coordinator's pane does", async () => {
|
||||
await as(SESSION_X, 'orchestration.runCreate', { objective: 'first' })
|
||||
const current = idOf(
|
||||
(await as(SESSION_X, 'orchestration.runCreate', { objective: 'next' })).run
|
||||
)
|
||||
|
||||
const { message } = await sendFromTerminal(ADDRESS_X)
|
||||
|
||||
expect(message).toMatchObject({ to_handle: `run:${current}`, run_id: current })
|
||||
expect(await as(SESSION_X, 'orchestration.check', {})).toMatchObject({
|
||||
runId: current,
|
||||
messages: [{ subject: 'hello' }]
|
||||
})
|
||||
})
|
||||
|
||||
it('delivers it to a chat with no Run, which reads its direct mailbox', async () => {
|
||||
const { message } = await sendFromTerminal(ADDRESS_X)
|
||||
|
||||
expect(message).toMatchObject({ to_handle: ADDRESS_X })
|
||||
expect(await as(SESSION_X, 'orchestration.check', {})).toMatchObject({
|
||||
messages: [{ subject: 'hello' }]
|
||||
})
|
||||
})
|
||||
|
||||
it('reaches a chat with no Run after a restart, once the send has started the session host', async () => {
|
||||
// After an app restart the agent-session host starts lazily; routing reads the session record
|
||||
// synchronously, so the send must start the host before it resolves the recipient.
|
||||
const store = hostRef.current
|
||||
hostRef.current = null
|
||||
vi.mocked(h.runtime.ensureStructuredAgentSessionHost).mockImplementation(async () => {
|
||||
hostRef.current = store
|
||||
})
|
||||
|
||||
const { message } = await sendFromTerminal(ADDRESS_X)
|
||||
|
||||
expect(message).toMatchObject({ to_handle: ADDRESS_X })
|
||||
})
|
||||
|
||||
it('refuses an Orca session this host does not run, or has no record of', async () => {
|
||||
h.records.set(
|
||||
SESSION_X,
|
||||
sessionRecord(SESSION_X, { location: { executionHostId: 'ssh:devbox' } })
|
||||
)
|
||||
h.records.delete(SESSION_Y)
|
||||
|
||||
for (const [to, code] of [
|
||||
[ADDRESS_X, 'session_caller_host_boundary'],
|
||||
[ADDRESS_Y, 'session_caller_unknown']
|
||||
]) {
|
||||
const response = await h.dispatch(
|
||||
orchestrationRequest('orchestration.send', { from: WORKER_HANDLE, to, subject: 's' })
|
||||
)
|
||||
expect(response).toMatchObject({ ok: false, error: { code } })
|
||||
}
|
||||
})
|
||||
|
||||
it("routes a structured worker's session address to the Dispatch it is working", async () => {
|
||||
const handle = mintStructuredWorkerHandle()
|
||||
const paneKey = mintStructuredWorkerPaneKey(SESSION_Y)
|
||||
structuredWorkerIdentities.register({
|
||||
handle,
|
||||
sessionId: SESSION_Y,
|
||||
agent: 'claude',
|
||||
paneKey,
|
||||
processIncarnation: structuredWorkerProcessIncarnation(SESSION_Y),
|
||||
worktreeId: 'wt_1',
|
||||
hostScope: { kind: 'local', hostId: 'local' }
|
||||
})
|
||||
const runId = idOf((await as(SESSION_X, 'orchestration.runCreate', { objective: 'o' })).run)
|
||||
const dispatch = h.db.createDispatchContext({
|
||||
taskId: h.db.createTask({ runId, spec: 'work' }).id,
|
||||
assigneeHandle: handle,
|
||||
assigneePaneKey: paneKey,
|
||||
processIncarnation: structuredWorkerProcessIncarnation(SESSION_Y),
|
||||
creator: { kind: 'session', orcaSessionId: SESSION_X },
|
||||
maxDepth: Number.MAX_SAFE_INTEGER
|
||||
})
|
||||
|
||||
const { message } = await as(SESSION_X, 'orchestration.send', {
|
||||
to: ADDRESS_Y,
|
||||
subject: 'to the worker'
|
||||
})
|
||||
|
||||
expect(message).toMatchObject({ to_handle: `dispatch:${dispatch.id}`, run_id: runId })
|
||||
expect(await as(SESSION_Y, 'orchestration.check', { peek: true })).toMatchObject({
|
||||
dispatchId: dispatch.id,
|
||||
messages: [{ subject: 'to the worker' }]
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user