mirror of
https://github.com/stablyai/orca.git
synced 2026-09-28 08:02:43 +00:00
resolve run method callers through one principal resolver
This commit is contained in:
@@ -0,0 +1,257 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { AgentSessionRecord } from '../../../../../shared/agent-session-record'
|
||||
import {
|
||||
agentSessionLeaseFixture,
|
||||
agentSessionRecordFixture
|
||||
} from '../../../../../shared/agent-session-record.test-fixture'
|
||||
import { OrchestrationDb } from '../../../orchestration/db'
|
||||
import { OrcaRuntimeService } from '../../../orca-runtime'
|
||||
import {
|
||||
mintStructuredWorkerPaneKey,
|
||||
structuredWorkerIdentities,
|
||||
structuredWorkerProcessIncarnation
|
||||
} from '../../../structured-worker-identity'
|
||||
import { readStructuredAgentSessionRecord } from '../../../structured-worker-authority'
|
||||
import type * as StructuredWorkerAuthority from '../../../structured-worker-authority'
|
||||
import { resolveCallerPrincipal } from './caller-principal'
|
||||
|
||||
vi.mock('../../../structured-worker-authority', async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof StructuredWorkerAuthority>()
|
||||
return { ...actual, readStructuredAgentSessionRecord: vi.fn() }
|
||||
})
|
||||
|
||||
const SESSION_ID = 'session-alpha-1'
|
||||
const STRUCTURED_HANDLE = 'structworker_11111111-2222-4333-8444-555555555555'
|
||||
const STRUCTURED_PANE_KEY = mintStructuredWorkerPaneKey(SESSION_ID)
|
||||
const COORDINATOR_PANE_KEY = 'tab_coord:aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa'
|
||||
|
||||
function nativeRecord(
|
||||
leaseOverrides: Parameters<typeof agentSessionLeaseFixture>[0] = {}
|
||||
): AgentSessionRecord {
|
||||
return agentSessionRecordFixture(
|
||||
agentSessionLeaseFixture({ runtimeKind: 'native', ...leaseOverrides })
|
||||
)
|
||||
}
|
||||
|
||||
describe('resolveCallerPrincipal', () => {
|
||||
let db: OrchestrationDb
|
||||
let runtime: OrcaRuntimeService
|
||||
|
||||
beforeEach(() => {
|
||||
db = new OrchestrationDb(':memory:')
|
||||
runtime = new OrcaRuntimeService()
|
||||
runtime.setOrchestrationDb(db)
|
||||
vi.spyOn(runtime, 'getTerminalPaneKey').mockImplementation((handle) =>
|
||||
handle === 'term_coord' ? COORDINATOR_PANE_KEY : null
|
||||
)
|
||||
structuredWorkerIdentities.clear()
|
||||
structuredWorkerIdentities.register({
|
||||
handle: STRUCTURED_HANDLE,
|
||||
sessionId: SESSION_ID,
|
||||
agent: 'claude',
|
||||
paneKey: STRUCTURED_PANE_KEY,
|
||||
processIncarnation: structuredWorkerProcessIncarnation(SESSION_ID),
|
||||
worktreeId: 'worktree-1',
|
||||
hostScope: { kind: 'local', hostId: 'local' }
|
||||
})
|
||||
vi.mocked(readStructuredAgentSessionRecord).mockReturnValue(nativeRecord())
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
structuredWorkerIdentities.clear()
|
||||
db.close()
|
||||
vi.restoreAllMocks()
|
||||
})
|
||||
|
||||
it('resolves a terminal handle to a pane principal', () => {
|
||||
const caller = resolveCallerPrincipal(runtime, {
|
||||
from: 'term_coord',
|
||||
requireBindableCaller: true
|
||||
})
|
||||
expect(caller.principal).toEqual({ kind: 'pane', paneKey: COORDINATOR_PANE_KEY })
|
||||
expect(caller.principalId).toBe(`pane:${COORDINATOR_PANE_KEY}`)
|
||||
expect(caller.waiterHandles).toEqual(['term_coord'])
|
||||
expect(caller.binding).toEqual({
|
||||
principalId: `pane:${COORDINATOR_PANE_KEY}`,
|
||||
terminalHandle: 'term_coord',
|
||||
paneKey: COORDINATOR_PANE_KEY
|
||||
})
|
||||
expect(caller.attested).toBe(false)
|
||||
})
|
||||
|
||||
it('resolves a structured worker handle to a session principal', () => {
|
||||
const caller = resolveCallerPrincipal(runtime, {
|
||||
from: STRUCTURED_HANDLE,
|
||||
requireBindableCaller: true
|
||||
})
|
||||
expect(caller.principal).toEqual({ kind: 'session', sessionId: SESSION_ID })
|
||||
expect(caller.principalId).toBe(`session:${SESSION_ID}`)
|
||||
// The binding still carries the handle and the minted pane key so dual-write and mail
|
||||
// routing stay byte-identical with today.
|
||||
expect(caller.binding).toEqual({
|
||||
principalId: `session:${SESSION_ID}`,
|
||||
terminalHandle: STRUCTURED_HANDLE,
|
||||
paneKey: STRUCTURED_PANE_KEY
|
||||
})
|
||||
expect(caller.attested).toBe(false)
|
||||
expect(caller.ownerGeneration).toBe('7')
|
||||
expect(caller.workspaceId).toBe('workspace-1')
|
||||
expect(caller.hostScope).toEqual({ kind: 'local', hostId: 'local' })
|
||||
expect(caller.waiterHandles).toEqual([STRUCTURED_HANDLE])
|
||||
})
|
||||
|
||||
it('returns one identical shape for both kinds', () => {
|
||||
const pane = resolveCallerPrincipal(runtime, {
|
||||
from: 'term_coord',
|
||||
requireBindableCaller: true
|
||||
})
|
||||
const session = resolveCallerPrincipal(runtime, {
|
||||
from: STRUCTURED_HANDLE,
|
||||
requireBindableCaller: true
|
||||
})
|
||||
expect(Object.keys(session).sort()).toEqual(Object.keys(pane).sort())
|
||||
})
|
||||
|
||||
it('refuses a declared session id with no bearer', () => {
|
||||
// Session ids are public (tab ids embed them); accepting one would let any RPC caller
|
||||
// impersonate any native session.
|
||||
expect(() =>
|
||||
resolveCallerPrincipal(runtime, { agentSessionId: SESSION_ID, requireBindableCaller: true })
|
||||
).toThrowError(expect.objectContaining({ code: 'consumer_fenced' }))
|
||||
expect(() =>
|
||||
resolveCallerPrincipal(runtime, {
|
||||
agentSessionId: SESSION_ID,
|
||||
runtimeFence: 7,
|
||||
requireBindableCaller: true
|
||||
})
|
||||
).toThrowError(expect.objectContaining({ code: 'consumer_fenced' }))
|
||||
})
|
||||
|
||||
it('refuses a declared session id alongside a pane bearer', () => {
|
||||
expect(() =>
|
||||
resolveCallerPrincipal(runtime, {
|
||||
from: 'term_coord',
|
||||
agentSessionId: SESSION_ID,
|
||||
requireBindableCaller: true
|
||||
})
|
||||
).toThrowError(expect.objectContaining({ code: 'consumer_fenced' }))
|
||||
})
|
||||
|
||||
it('treats declared session fields as corroboration for a structured bearer', () => {
|
||||
// Stale declared fence: the caller holds yesterday's identity.
|
||||
expect(() =>
|
||||
resolveCallerPrincipal(runtime, {
|
||||
from: STRUCTURED_HANDLE,
|
||||
runtimeFence: 6,
|
||||
requireBindableCaller: true
|
||||
})
|
||||
).toThrowError(expect.objectContaining({ code: 'consumer_fenced' }))
|
||||
// Declared session id mismatching the bearer's session.
|
||||
expect(() =>
|
||||
resolveCallerPrincipal(runtime, {
|
||||
from: STRUCTURED_HANDLE,
|
||||
agentSessionId: 'session-other',
|
||||
requireBindableCaller: true
|
||||
})
|
||||
).toThrowError(expect.objectContaining({ code: 'consumer_fenced' }))
|
||||
// Matching corroboration passes.
|
||||
expect(
|
||||
resolveCallerPrincipal(runtime, {
|
||||
from: STRUCTURED_HANDLE,
|
||||
agentSessionId: SESSION_ID,
|
||||
runtimeFence: 7,
|
||||
requireBindableCaller: true
|
||||
}).principalId
|
||||
).toBe(`session:${SESSION_ID}`)
|
||||
})
|
||||
|
||||
it('accepts fence 0 against a fence-0 lease', () => {
|
||||
// The lease validator allows fence 0; a .positive() schema would lock out fresh sessions.
|
||||
vi.mocked(readStructuredAgentSessionRecord).mockReturnValue(nativeRecord({ runtimeFence: 0 }))
|
||||
const caller = resolveCallerPrincipal(runtime, {
|
||||
from: STRUCTURED_HANDLE,
|
||||
runtimeFence: 0,
|
||||
requireBindableCaller: true
|
||||
})
|
||||
expect(caller.ownerGeneration).toBe('0')
|
||||
})
|
||||
|
||||
it.each([
|
||||
['claimStatus reserved', nativeRecord({ claimStatus: 'reserved' })],
|
||||
['claimStatus conflicted', nativeRecord({ claimStatus: 'conflicted' })],
|
||||
['claimStatus released', nativeRecord({ claimStatus: 'released' })],
|
||||
['handoff in progress', nativeRecord({ handoffStage: 'preparing' })],
|
||||
['unreconciled lease', nativeRecord({ unreconciled: true })],
|
||||
['terminal-owned lease', nativeRecord({ runtimeKind: 'tui' })],
|
||||
[
|
||||
'non-local host',
|
||||
{
|
||||
...nativeRecord(),
|
||||
location: { ...nativeRecord().location, executionHostId: 'ssh:remote' as const }
|
||||
}
|
||||
],
|
||||
['missing record', null]
|
||||
])('refuses a structured bearer whose lease is not current: %s', (_label, record) => {
|
||||
vi.mocked(readStructuredAgentSessionRecord).mockReturnValue(record)
|
||||
expect(() =>
|
||||
resolveCallerPrincipal(runtime, { from: STRUCTURED_HANDLE, requireBindableCaller: true })
|
||||
).toThrowError(expect.objectContaining({ code: 'consumer_fenced' }))
|
||||
})
|
||||
|
||||
it('requires some credential', () => {
|
||||
expect(() => resolveCallerPrincipal(runtime, { requireBindableCaller: true })).toThrowError(
|
||||
expect.objectContaining({ code: 'invalid_argument' })
|
||||
)
|
||||
})
|
||||
|
||||
it('still delegates evidence attestation, immediately or deferred', () => {
|
||||
vi.spyOn(runtime, 'verifyOrchestrationCompatibilityCaller').mockReturnValue({
|
||||
hostScope: { kind: 'local', hostId: 'local' },
|
||||
paneKey: COORDINATOR_PANE_KEY,
|
||||
terminalHandle: 'term_other',
|
||||
processIncarnation: 'incarnation-1',
|
||||
launchTokenHash: 'hash-1'
|
||||
})
|
||||
const evidence = { terminalHandle: 'term_other', paneKey: 'p', launchToken: 't' }
|
||||
expect(() =>
|
||||
resolveCallerPrincipal(runtime, {
|
||||
from: 'term_coord',
|
||||
evidence,
|
||||
requireBindableCaller: true
|
||||
})
|
||||
).toThrowError(expect.objectContaining({ code: 'consumer_fenced' }))
|
||||
const deferred = resolveCallerPrincipal(runtime, {
|
||||
from: 'term_coord',
|
||||
evidence,
|
||||
requireBindableCaller: true,
|
||||
deferEvidenceAssertion: true
|
||||
})
|
||||
expect(() => deferred.attestDeclaredCaller()).toThrowError(
|
||||
expect.objectContaining({ code: 'consumer_fenced' })
|
||||
)
|
||||
const sessionDeferred = resolveCallerPrincipal(runtime, {
|
||||
from: STRUCTURED_HANDLE,
|
||||
evidence,
|
||||
requireBindableCaller: true,
|
||||
deferEvidenceAssertion: true
|
||||
})
|
||||
expect(() => sessionDeferred.attestDeclaredCaller()).toThrowError(
|
||||
expect.objectContaining({ code: 'consumer_fenced' })
|
||||
)
|
||||
})
|
||||
|
||||
it('resolves null for a bearer with no stable pane unless bindability is required', () => {
|
||||
expect(resolveCallerPrincipal(runtime, { from: 'term_stale' })).toBeNull()
|
||||
expect(() =>
|
||||
resolveCallerPrincipal(runtime, { from: 'term_stale', requireBindableCaller: true })
|
||||
).toThrowError(expect.objectContaining({ code: 'stable_pane_required' }))
|
||||
})
|
||||
|
||||
it('prefers a declared pane key over the live pane fallback', () => {
|
||||
const caller = resolveCallerPrincipal(runtime, {
|
||||
from: 'term_stale',
|
||||
paneKey: 'tab_declared:bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb'
|
||||
})
|
||||
expect(caller?.principalId).toBe('pane:tab_declared:bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb')
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,226 @@
|
||||
/**
|
||||
* ONE caller resolver for the orchestration RPC boundary: every credential form in, one resolved
|
||||
* shape out. Methods consume the shape whole and never inspect the principal's kind — per-method
|
||||
* `if (params.agentSessionId)` branches are the defect class this module exists to prevent.
|
||||
*
|
||||
* Credential tiers: identity requires possession of an UNGUESSABLE BEARER — a `term_` handle, a
|
||||
* `structworker_` handle, or an attested pane key with a random leaf. A declared `agentSessionId`
|
||||
* or `runtimeFence` NEVER authenticates on its own: session ids are public (tab ids embed them in
|
||||
* plain text) and fences are small integers, so both are corroboration only. Deriving a session
|
||||
* principal from a pane key is legal exactly once — PR 1's one-time server-side backfill and
|
||||
* dual-write over rows the host itself wrote — and never at request time.
|
||||
*/
|
||||
import { z } from 'zod'
|
||||
import type { OrchestrationCompatibilityEvidence } from '../../../../../shared/orchestration-compatibility-evidence'
|
||||
import {
|
||||
formatOrchestrationPrincipal,
|
||||
type OrchestrationPrincipal
|
||||
} from '../../../../../shared/orchestration-principal'
|
||||
import { OrchestrationError } from '../../../orchestration/orchestration-error'
|
||||
import type { RunCoordinatorBinding } from '../../../orchestration/types'
|
||||
import type { WorkerTerminalHostScope } from '../../../orchestration/worker-terminal-process-liveness'
|
||||
import type {
|
||||
OrcaRuntimeService,
|
||||
OrchestrationCompatibilityCallerAuthority
|
||||
} from '../../../orca-runtime'
|
||||
import {
|
||||
readStructuredAgentSessionRecord,
|
||||
resolveStructuredWorkerIdentity
|
||||
} from '../../../structured-worker-authority'
|
||||
import {
|
||||
isStructuredWorkerHandle,
|
||||
structuredWorkerHostScope,
|
||||
structuredWorkerRecordIsCurrent
|
||||
} from '../../../structured-worker-identity'
|
||||
import { OptionalString } from '../../schemas'
|
||||
import { assertCallerHandleMatchesEvidence, resolveOrchestrationCaller } from './runs/run-scope'
|
||||
|
||||
/**
|
||||
* Shared wire fragment for methods that accept orchestration caller credentials. `runtimeFence`
|
||||
* uses `.min(0)` deliberately: the lease validator allows fence 0, so `.positive()` would lock
|
||||
* out freshly leased sessions.
|
||||
*/
|
||||
export const orchestrationCallerParamFields = {
|
||||
from: OptionalString,
|
||||
agentSessionId: OptionalString,
|
||||
runtimeFence: z.number().int().min(0).optional()
|
||||
}
|
||||
|
||||
export type OrchestrationCallerCredentials = {
|
||||
/** Terminal handle or structured-worker handle. Accepted forever; never required. */
|
||||
from?: string
|
||||
/** Caller-declared pane key (resolveRunScope's callerPaneKey passthrough). */
|
||||
paneKey?: string
|
||||
/**
|
||||
* Corroboration ONLY, never a credential: session ids are guessable (embedded in tab ids)
|
||||
* and fences are small integers. Session-kind resolution requires an unguessable bearer
|
||||
* (a structworker handle in `from` today; a host-baked session bearer in later PRs).
|
||||
*/
|
||||
agentSessionId?: string
|
||||
runtimeFence?: number
|
||||
evidence?: OrchestrationCompatibilityEvidence
|
||||
callerAuthority?: OrchestrationCompatibilityCallerAuthority
|
||||
/** Same contract as resolveOrchestrationCaller's requireStablePane. */
|
||||
requireBindableCaller?: boolean
|
||||
/** Same contract & warning as evidenceAssertedByCaller; pairs with attestDeclaredCaller(). */
|
||||
deferEvidenceAssertion?: boolean
|
||||
}
|
||||
|
||||
export type ResolvedOrchestrationCaller = Readonly<{
|
||||
principal: OrchestrationPrincipal
|
||||
/** Canonical `pane:<paneKey>` / `session:<sessionId>`. */
|
||||
principalId: string
|
||||
/** pane: processIncarnation (null when unproven); session: String(lease.runtimeFence). Opaque. */
|
||||
ownerGeneration: string | null
|
||||
hostScope: WorkerTerminalHostScope | null
|
||||
workspaceId: string | null
|
||||
/**
|
||||
* Live proven caller: pane = callerAuthority matches handle+pane. Session-kind is always false:
|
||||
* hook attestation can never mint authority for a structured handle (no PTY, no launch token),
|
||||
* so takeover-legacy must keep failing for that form.
|
||||
*/
|
||||
attested: boolean
|
||||
/** Opaque carrier for DB writers & mail routing; methods pass it through whole. */
|
||||
binding: RunCoordinatorBinding
|
||||
/** Mailbox waiter keys to cancel on rebind; methods iterate, never inspect. */
|
||||
waiterHandles: readonly string[]
|
||||
/** Deferred half of the attestation when deferEvidenceAssertion was set; no-op otherwise. */
|
||||
attestDeclaredCaller(): void
|
||||
}>
|
||||
|
||||
// Overloaded like resolveOrchestrationCaller: with requireBindableCaller:true the caller is
|
||||
// non-null or the resolver throws; otherwise a pane caller with no stable pane resolves to null.
|
||||
export function resolveCallerPrincipal(
|
||||
runtime: OrcaRuntimeService,
|
||||
credentials: OrchestrationCallerCredentials & { requireBindableCaller: true }
|
||||
): ResolvedOrchestrationCaller
|
||||
export function resolveCallerPrincipal(
|
||||
runtime: OrcaRuntimeService,
|
||||
credentials: OrchestrationCallerCredentials
|
||||
): ResolvedOrchestrationCaller | null
|
||||
export function resolveCallerPrincipal(
|
||||
runtime: OrcaRuntimeService,
|
||||
credentials: OrchestrationCallerCredentials
|
||||
): ResolvedOrchestrationCaller | null {
|
||||
if (credentials.from && isStructuredWorkerHandle(credentials.from)) {
|
||||
return resolveSessionPrincipal(runtime, credentials, credentials.from)
|
||||
}
|
||||
if (credentials.from) {
|
||||
return resolvePanePrincipal(runtime, credentials, credentials.from)
|
||||
}
|
||||
if (credentials.agentSessionId) {
|
||||
// Deliberate: session ids are public, so accepting a declared one would let any RPC caller
|
||||
// impersonate any native session.
|
||||
throw new OrchestrationError(
|
||||
'consumer_fenced',
|
||||
"A declared session id is not a credential. Call with the session's issued handle."
|
||||
)
|
||||
}
|
||||
throw new OrchestrationError(
|
||||
'invalid_argument',
|
||||
'Missing coordinator identity: pass --from or an agent session credential.'
|
||||
)
|
||||
}
|
||||
|
||||
/** Immediate unless deferred; the deferred half preserves run-use's resolve→takeover→attest order. */
|
||||
function evidenceAttestation(
|
||||
runtime: OrcaRuntimeService,
|
||||
from: string,
|
||||
credentials: OrchestrationCallerCredentials
|
||||
): () => void {
|
||||
if (!credentials.deferEvidenceAssertion) {
|
||||
assertCallerHandleMatchesEvidence(runtime, from, credentials.evidence)
|
||||
return () => {}
|
||||
}
|
||||
return () => assertCallerHandleMatchesEvidence(runtime, from, credentials.evidence)
|
||||
}
|
||||
|
||||
function resolveSessionPrincipal(
|
||||
runtime: OrcaRuntimeService,
|
||||
credentials: OrchestrationCallerCredentials,
|
||||
from: string
|
||||
): ResolvedOrchestrationCaller {
|
||||
const attestDeclaredCaller = evidenceAttestation(runtime, from, credentials)
|
||||
const identity = resolveStructuredWorkerIdentity(from, runtime.getOrchestrationDb())
|
||||
const record = identity ? readStructuredAgentSessionRecord(identity.sessionId) : null
|
||||
if (
|
||||
!identity ||
|
||||
!record ||
|
||||
!structuredWorkerRecordIsCurrent(record) ||
|
||||
record.lease.claimStatus !== 'live' ||
|
||||
record.lease.unreconciled ||
|
||||
record.lease.handoffStage !== null ||
|
||||
// Declared values corroborate the bearer: a mismatch means yesterday's identity.
|
||||
(credentials.agentSessionId !== undefined &&
|
||||
credentials.agentSessionId !== identity.sessionId) ||
|
||||
(credentials.runtimeFence !== undefined &&
|
||||
credentials.runtimeFence !== record.lease.runtimeFence)
|
||||
) {
|
||||
throw new OrchestrationError('consumer_fenced', 'The native session lease is not current.')
|
||||
}
|
||||
const principal: OrchestrationPrincipal = { kind: 'session', sessionId: identity.sessionId }
|
||||
const principalId = formatOrchestrationPrincipal(principal)
|
||||
return {
|
||||
principal,
|
||||
principalId,
|
||||
ownerGeneration: String(record.lease.runtimeFence),
|
||||
hostScope: structuredWorkerHostScope(record.location),
|
||||
workspaceId: record.location.workspaceId,
|
||||
attested: false,
|
||||
// Handle + minted pane key keep dual-write and mail routing byte-identical with today.
|
||||
binding: { principalId, terminalHandle: from, paneKey: identity.paneKey },
|
||||
waiterHandles: [from],
|
||||
attestDeclaredCaller
|
||||
}
|
||||
}
|
||||
|
||||
function resolvePanePrincipal(
|
||||
runtime: OrcaRuntimeService,
|
||||
credentials: OrchestrationCallerCredentials,
|
||||
from: string
|
||||
): ResolvedOrchestrationCaller | null {
|
||||
if (credentials.agentSessionId !== undefined) {
|
||||
// A pane bearer may not claim to be a session.
|
||||
throw new OrchestrationError(
|
||||
'consumer_fenced',
|
||||
'A pane caller cannot declare an agent session identity.'
|
||||
)
|
||||
}
|
||||
const resolved = resolveOrchestrationCaller(runtime, {
|
||||
callerTerminalHandle: from,
|
||||
callerEvidence: credentials.evidence,
|
||||
callerAuthority: credentials.callerAuthority,
|
||||
evidenceAssertedByCaller: credentials.deferEvidenceAssertion,
|
||||
requireStablePane:
|
||||
credentials.requireBindableCaller === true && credentials.paneKey === undefined
|
||||
})
|
||||
const callerAuthority = credentials.callerAuthority
|
||||
// Attested authority wins; a declared pane key slots ahead of the live-pane fallback only
|
||||
// (resolveRunScope's callerPaneKey contract).
|
||||
const paneKey =
|
||||
callerAuthority?.terminalHandle === from ? resolved : (credentials.paneKey ?? resolved)
|
||||
if (!paneKey) {
|
||||
return null
|
||||
}
|
||||
const principal: OrchestrationPrincipal = { kind: 'pane', paneKey }
|
||||
const principalId = formatOrchestrationPrincipal(principal)
|
||||
let authority: ReturnType<OrcaRuntimeService['getOrchestrationDispatchAuthority']> = null
|
||||
try {
|
||||
authority = runtime.getOrchestrationDispatchAuthority(from)
|
||||
} catch {
|
||||
authority = null
|
||||
}
|
||||
return {
|
||||
principal,
|
||||
principalId,
|
||||
ownerGeneration: authority?.processIncarnation ?? null,
|
||||
hostScope: authority?.hostScope ?? null,
|
||||
workspaceId: authority?.worktreeId ?? null,
|
||||
attested: callerAuthority?.terminalHandle === from && callerAuthority?.paneKey === paneKey,
|
||||
binding: { principalId, terminalHandle: from, paneKey },
|
||||
waiterHandles: [from],
|
||||
attestDeclaredCaller: credentials.deferEvidenceAssertion
|
||||
? () => assertCallerHandleMatchesEvidence(runtime, from, credentials.evidence)
|
||||
: () => {}
|
||||
}
|
||||
}
|
||||
@@ -1,6 +1,7 @@
|
||||
import type { OrchestrationCompatibilityEvidence } from '../../../../../../shared/orchestration-compatibility-evidence'
|
||||
import { orchestrationSkillRecoveryData } from '../../../../../../shared/orchestration-rpc-contract'
|
||||
import { OrchestrationError } from '../../../../orchestration/orchestration-error'
|
||||
import { resolveCallerPrincipal } from '../caller-principal'
|
||||
import type { RunRow } from '../../../../orchestration/types'
|
||||
import type {
|
||||
OrcaRuntimeService,
|
||||
@@ -11,6 +12,8 @@ export type RunScopeParams = {
|
||||
runId?: string
|
||||
callerTerminalHandle?: string
|
||||
callerPaneKey?: string
|
||||
callerAgentSessionId?: string
|
||||
callerRuntimeFence?: number
|
||||
requireCurrentConsumer: boolean
|
||||
legacyCoordinatorRunId?: string
|
||||
// Why: the caller's declared handle is a user param; this is the attested one to check it against.
|
||||
@@ -97,18 +100,24 @@ export function resolveRunScope(runtime: OrcaRuntimeService, params: RunScopePar
|
||||
orchestrationSkillRecoveryData()
|
||||
)
|
||||
}
|
||||
assertCallerHandleMatchesEvidence(runtime, params.callerTerminalHandle, params.callerEvidence)
|
||||
// Why: attestation must stay ahead of the legacy early-return; the bindability throw follows it.
|
||||
const caller = resolveCallerPrincipal(runtime, {
|
||||
from: params.callerTerminalHandle,
|
||||
paneKey: params.callerPaneKey,
|
||||
agentSessionId: params.callerAgentSessionId,
|
||||
runtimeFence: params.callerRuntimeFence,
|
||||
evidence: params.callerEvidence
|
||||
})
|
||||
if (explicit && params.legacyCoordinatorRunId === explicit.id) {
|
||||
return explicit
|
||||
}
|
||||
const paneKey = params.callerPaneKey ?? runtime.getTerminalPaneKey(params.callerTerminalHandle)
|
||||
if (!paneKey) {
|
||||
if (!caller) {
|
||||
throw new OrchestrationError(
|
||||
'stable_pane_required',
|
||||
'The coordinator terminal has no stable pane identity.'
|
||||
)
|
||||
}
|
||||
const current = db.getCurrentRunForPane(paneKey)
|
||||
const current = db.getCurrentRunForPrincipal(caller.principalId)
|
||||
if (!current) {
|
||||
if (explicit) {
|
||||
throw new OrchestrationError(
|
||||
|
||||
@@ -0,0 +1,189 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { RpcContext } from '../../../core'
|
||||
import type { AgentSessionRecord } from '../../../../../../shared/agent-session-record'
|
||||
import {
|
||||
agentSessionLeaseFixture,
|
||||
agentSessionRecordFixture
|
||||
} from '../../../../../../shared/agent-session-record.test-fixture'
|
||||
import type { OrchestrationDb } from '../../../../orchestration/db'
|
||||
import type { OrcaRuntimeService } from '../../../../orca-runtime'
|
||||
import {
|
||||
mintStructuredWorkerPaneKey,
|
||||
structuredWorkerIdentities,
|
||||
structuredWorkerProcessIncarnation
|
||||
} from '../../../../structured-worker-identity'
|
||||
import { readStructuredAgentSessionRecord } from '../../../../structured-worker-authority'
|
||||
import type * as StructuredWorkerAuthority from '../../../../structured-worker-authority'
|
||||
import { createOrchestrationRpcHarness } from '../rpc-test-harness'
|
||||
|
||||
vi.mock('../../../../structured-worker-authority', async (importOriginal) => {
|
||||
const actual = await importOriginal<typeof StructuredWorkerAuthority>()
|
||||
return { ...actual, readStructuredAgentSessionRecord: vi.fn() }
|
||||
})
|
||||
|
||||
const SESSION_ID = 'session-alpha-1'
|
||||
const STRUCTURED_HANDLE = 'structworker_11111111-2222-4333-8444-555555555555'
|
||||
const STRUCTURED_PANE_KEY = mintStructuredWorkerPaneKey(SESSION_ID)
|
||||
|
||||
function nativeRecord(
|
||||
leaseOverrides: Parameters<typeof agentSessionLeaseFixture>[0] = {}
|
||||
): AgentSessionRecord {
|
||||
return agentSessionRecordFixture(
|
||||
agentSessionLeaseFixture({ runtimeKind: 'native', ...leaseOverrides })
|
||||
)
|
||||
}
|
||||
|
||||
describe('run methods called with a structured session bearer', () => {
|
||||
const h = createOrchestrationRpcHarness()
|
||||
let db: OrchestrationDb
|
||||
let runtime: OrcaRuntimeService
|
||||
let ctx: RpcContext
|
||||
|
||||
beforeEach(() => {
|
||||
;({ db, runtime, ctx } = h.setup(false))
|
||||
structuredWorkerIdentities.clear()
|
||||
structuredWorkerIdentities.register({
|
||||
handle: STRUCTURED_HANDLE,
|
||||
sessionId: SESSION_ID,
|
||||
agent: 'claude',
|
||||
paneKey: STRUCTURED_PANE_KEY,
|
||||
processIncarnation: structuredWorkerProcessIncarnation(SESSION_ID),
|
||||
worktreeId: 'worktree-1',
|
||||
hostScope: { kind: 'local', hostId: 'local' }
|
||||
})
|
||||
vi.mocked(readStructuredAgentSessionRecord).mockReturnValue(nativeRecord())
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
structuredWorkerIdentities.clear()
|
||||
h.cleanup()
|
||||
})
|
||||
|
||||
async function call(name: string, params: Record<string, unknown>) {
|
||||
return h.call(name, params, ctx)
|
||||
}
|
||||
|
||||
function rawRun(id: string) {
|
||||
return db.getRunRaw(id)!
|
||||
}
|
||||
|
||||
it('creates, rebinds, and reads a Run through the session principal with a receipt identical to the pane path', async () => {
|
||||
const paneCreated = (await call('orchestration.runCreate', {
|
||||
objective: 'Pane-coordinated',
|
||||
from: 'term_coord'
|
||||
})) as { run: Record<string, unknown> & { id: string } }
|
||||
const created = (await call('orchestration.runCreate', {
|
||||
objective: 'Session-coordinated',
|
||||
from: STRUCTURED_HANDLE,
|
||||
runtimeFence: 7
|
||||
})) as { run: Record<string, unknown> & { id: string } }
|
||||
|
||||
expect(Object.keys(created.run).sort()).toEqual(Object.keys(paneCreated.run).sort())
|
||||
expect(rawRun(created.run.id).coordinator_principal).toBe(`session:${SESSION_ID}`)
|
||||
|
||||
const current = (await call('orchestration.runCurrent', { from: STRUCTURED_HANDLE })) as {
|
||||
run: { id: string } | null
|
||||
}
|
||||
expect(current.run?.id).toBe(created.run.id)
|
||||
|
||||
const rebound = (await call('orchestration.runUse', {
|
||||
id: paneCreated.run.id,
|
||||
from: STRUCTURED_HANDLE
|
||||
})) as { run: { id: string } }
|
||||
expect(rebound.run.id).toBe(paneCreated.run.id)
|
||||
expect(rawRun(paneCreated.run.id).coordinator_principal).toBe(`session:${SESSION_ID}`)
|
||||
// The rebind released the session's previous Run.
|
||||
expect(rawRun(created.run.id).coordinator_principal).toBeNull()
|
||||
})
|
||||
|
||||
it('dual-writes the handle and minted pane key so the mail path still resolves the Run', async () => {
|
||||
const created = (await call('orchestration.runCreate', {
|
||||
objective: 'Dual-write',
|
||||
from: STRUCTURED_HANDLE
|
||||
})) as { run: { id: string } }
|
||||
const raw = rawRun(created.run.id)
|
||||
expect(raw.coordinator_principal).toBe(`session:${SESSION_ID}`)
|
||||
expect(raw.coordinator_handle).toBe(STRUCTURED_HANDLE)
|
||||
expect(raw.coordinator_pane_key).toBe(STRUCTURED_PANE_KEY)
|
||||
expect(db.getCurrentRunForPane(STRUCTURED_PANE_KEY)?.id).toBe(created.run.id)
|
||||
})
|
||||
|
||||
it('cancels the prior Run waiters and fences its delivery when the session creates a new Run', async () => {
|
||||
const prior = (await call('orchestration.runCreate', {
|
||||
objective: 'Prior run',
|
||||
from: STRUCTURED_HANDLE
|
||||
})) as { run: { id: string; consumer_generation: number } }
|
||||
db.insertMessage({
|
||||
from: 'term_worker',
|
||||
to: `run:${prior.run.id}`,
|
||||
subject: 'pending',
|
||||
runId: prior.run.id
|
||||
})
|
||||
const delivery = db.getOrCreateRunDelivery({
|
||||
runId: prior.run.id,
|
||||
consumerGeneration: prior.run.consumer_generation
|
||||
})!
|
||||
const cancelled = vi.spyOn(runtime, 'cancelMessageWaiters')
|
||||
|
||||
await call('orchestration.runCreate', { objective: 'Next run', from: STRUCTURED_HANDLE })
|
||||
|
||||
// The #19648 defect: its session branch skipped the prior-run lookup and both cancels.
|
||||
expect(cancelled.mock.calls.map((c) => c[0])).toEqual(
|
||||
expect.arrayContaining([STRUCTURED_HANDLE, `run:${prior.run.id}`])
|
||||
)
|
||||
const deliveryStatus = db.db
|
||||
.prepare('SELECT status FROM deliveries WHERE id = ?')
|
||||
.get(delivery.delivery.id) as { status: string }
|
||||
expect(deliveryStatus.status).toBe('fenced')
|
||||
expect(rawRun(prior.run.id).coordinator_principal).toBeNull()
|
||||
})
|
||||
|
||||
it('refuses takeover-legacy from a session caller', async () => {
|
||||
await expect(
|
||||
call('orchestration.runUse', { id: 'run-x', from: STRUCTURED_HANDLE, takeoverLegacy: true })
|
||||
).rejects.toMatchObject({ code: 'legacy_read_only' })
|
||||
})
|
||||
|
||||
it('refuses a declared session id with no bearer and writes no row', async () => {
|
||||
await expect(
|
||||
call('orchestration.runCreate', {
|
||||
objective: 'Impersonation attempt',
|
||||
agentSessionId: SESSION_ID,
|
||||
runtimeFence: 7
|
||||
})
|
||||
).rejects.toMatchObject({ code: 'consumer_fenced' })
|
||||
expect(db.listRuns().runs.filter((run) => run.legacy === 0)).toHaveLength(0)
|
||||
})
|
||||
|
||||
it.each([
|
||||
['a stale declared fence', { runtimeFence: 6 }, nativeRecord()],
|
||||
['a reserved lease', {}, nativeRecord({ claimStatus: 'reserved' })],
|
||||
['a lease mid-handoff', {}, nativeRecord({ handoffStage: 'preparing' })]
|
||||
])('refuses all three methods for %s and writes no row', async (_label, extra, record) => {
|
||||
vi.mocked(readStructuredAgentSessionRecord).mockReturnValue(record)
|
||||
await expect(
|
||||
call('orchestration.runCreate', {
|
||||
objective: 'Refused',
|
||||
from: STRUCTURED_HANDLE,
|
||||
...extra
|
||||
})
|
||||
).rejects.toMatchObject({ code: 'consumer_fenced' })
|
||||
await expect(
|
||||
call('orchestration.runUse', { id: 'run-x', from: STRUCTURED_HANDLE, ...extra })
|
||||
).rejects.toMatchObject({ code: 'consumer_fenced' })
|
||||
await expect(
|
||||
call('orchestration.runCurrent', { from: STRUCTURED_HANDLE, ...extra })
|
||||
).rejects.toMatchObject({ code: 'consumer_fenced' })
|
||||
expect(db.listRuns().runs.filter((run) => run.legacy === 0)).toHaveLength(0)
|
||||
})
|
||||
|
||||
it('accepts a corroborating fence of 0 against a fence-0 lease', async () => {
|
||||
vi.mocked(readStructuredAgentSessionRecord).mockReturnValue(nativeRecord({ runtimeFence: 0 }))
|
||||
const created = (await call('orchestration.runCreate', {
|
||||
objective: 'Fresh lease',
|
||||
from: STRUCTURED_HANDLE,
|
||||
runtimeFence: 0
|
||||
})) as { run: { id: string } }
|
||||
expect(rawRun(created.run.id).coordinator_principal).toBe(`session:${SESSION_ID}`)
|
||||
})
|
||||
})
|
||||
@@ -3,21 +3,21 @@ import { defineMethod, type RpcMethod } from '../../../core'
|
||||
import { OptionalBoolean, OptionalString, requiredString } from '../../../schemas'
|
||||
import { ORCHESTRATION_RUN_PAGE_LIMIT } from '../../../../../../shared/orchestration-run-pagination'
|
||||
import { OrchestrationError } from '../../../../orchestration/orchestration-error'
|
||||
import { assertCallerHandleMatchesEvidence, resolveOrchestrationCaller } from './run-scope'
|
||||
import { orchestrationCallerParamFields, resolveCallerPrincipal } from '../caller-principal'
|
||||
import { exposeRun } from './run-receipt'
|
||||
|
||||
const RunCreateParams = z.object({
|
||||
objective: requiredString('Missing --objective'),
|
||||
from: requiredString('Missing coordinator terminal')
|
||||
...orchestrationCallerParamFields
|
||||
})
|
||||
|
||||
const RunUseParams = z.object({
|
||||
id: requiredString('Missing --id'),
|
||||
from: requiredString('Missing coordinator terminal'),
|
||||
takeoverLegacy: OptionalBoolean
|
||||
takeoverLegacy: OptionalBoolean,
|
||||
...orchestrationCallerParamFields
|
||||
})
|
||||
|
||||
const RunCurrentParams = z.object({ from: requiredString('Missing coordinator terminal') })
|
||||
const RunCurrentParams = z.object({ ...orchestrationCallerParamFields })
|
||||
const RunListParams = z.object({
|
||||
limit: z.number().int().min(1).max(ORCHESTRATION_RUN_PAGE_LIMIT).optional(),
|
||||
cursor: z.string().min(1).optional()
|
||||
@@ -29,19 +29,19 @@ export const ORCHESTRATION_RUN_METHODS: RpcMethod[] = [
|
||||
name: 'orchestration.runCreate',
|
||||
params: RunCreateParams,
|
||||
handler: (params, { orchestrationCompatibilityEvidence, runtime }) => {
|
||||
const paneKey = resolveOrchestrationCaller(runtime, {
|
||||
callerTerminalHandle: params.from,
|
||||
callerEvidence: orchestrationCompatibilityEvidence,
|
||||
requireStablePane: true
|
||||
const caller = resolveCallerPrincipal(runtime, {
|
||||
from: params.from,
|
||||
agentSessionId: params.agentSessionId,
|
||||
runtimeFence: params.runtimeFence,
|
||||
evidence: orchestrationCompatibilityEvidence,
|
||||
requireBindableCaller: true
|
||||
})
|
||||
const db = runtime.getOrchestrationDb()
|
||||
const priorRun = db.getCurrentRunForPane(paneKey)
|
||||
const run = db.createRun({
|
||||
objective: params.objective,
|
||||
coordinatorHandle: params.from,
|
||||
coordinatorPaneKey: paneKey
|
||||
})
|
||||
runtime.cancelMessageWaiters(params.from)
|
||||
const priorRun = db.getCurrentRunForPrincipal(caller.principalId)
|
||||
const run = db.createRun({ objective: params.objective, coordinator: caller.binding })
|
||||
for (const waiterHandle of caller.waiterHandles) {
|
||||
runtime.cancelMessageWaiters(waiterHandle)
|
||||
}
|
||||
if (priorRun) {
|
||||
runtime.cancelMessageWaiters(`run:${priorRun.id}`)
|
||||
}
|
||||
@@ -60,30 +60,28 @@ export const ORCHESTRATION_RUN_METHODS: RpcMethod[] = [
|
||||
orchestrationCompatibilityCallerAuthority: callerAuthority
|
||||
}
|
||||
) => {
|
||||
const paneKey = resolveOrchestrationCaller(runtime, {
|
||||
callerTerminalHandle: params.from,
|
||||
callerEvidence: orchestrationCompatibilityEvidence,
|
||||
const caller = resolveCallerPrincipal(runtime, {
|
||||
from: params.from,
|
||||
agentSessionId: params.agentSessionId,
|
||||
runtimeFence: params.runtimeFence,
|
||||
evidence: orchestrationCompatibilityEvidence,
|
||||
callerAuthority,
|
||||
requireStablePane: true,
|
||||
evidenceAssertedByCaller: true
|
||||
requireBindableCaller: true,
|
||||
deferEvidenceAssertion: true
|
||||
})
|
||||
if (
|
||||
params.takeoverLegacy &&
|
||||
(callerAuthority?.terminalHandle !== params.from || callerAuthority.paneKey !== paneKey)
|
||||
) {
|
||||
if (params.takeoverLegacy && !caller.attested) {
|
||||
throw new OrchestrationError(
|
||||
'legacy_read_only',
|
||||
'Legacy takeover must be invoked by the live coordinator agent terminal it will bind. No effects were applied.',
|
||||
{ effectsApplied: false }
|
||||
)
|
||||
}
|
||||
assertCallerHandleMatchesEvidence(runtime, params.from, orchestrationCompatibilityEvidence)
|
||||
caller.attestDeclaredCaller()
|
||||
const db = runtime.getOrchestrationDb()
|
||||
const priorRun = db.getCurrentRunForPane(paneKey)
|
||||
const priorRun = db.getCurrentRunForPrincipal(caller.principalId)
|
||||
const run = db.bindRun({
|
||||
runId: params.id,
|
||||
coordinatorHandle: params.from,
|
||||
coordinatorPaneKey: paneKey,
|
||||
coordinator: caller.binding,
|
||||
takeoverLegacy: params.takeoverLegacy,
|
||||
legacyCoordinatorAuthority
|
||||
})
|
||||
@@ -93,7 +91,9 @@ export const ORCHESTRATION_RUN_METHODS: RpcMethod[] = [
|
||||
`Run ${params.id} was not found or is inspect-only.`
|
||||
)
|
||||
}
|
||||
runtime.cancelMessageWaiters(params.from)
|
||||
for (const waiterHandle of caller.waiterHandles) {
|
||||
runtime.cancelMessageWaiters(waiterHandle)
|
||||
}
|
||||
runtime.cancelMessageWaiters(`run:${params.id}`)
|
||||
if (priorRun && priorRun.id !== params.id) {
|
||||
runtime.cancelMessageWaiters(`run:${priorRun.id}`)
|
||||
@@ -105,12 +105,14 @@ export const ORCHESTRATION_RUN_METHODS: RpcMethod[] = [
|
||||
name: 'orchestration.runCurrent',
|
||||
params: RunCurrentParams,
|
||||
handler: (params, { orchestrationCompatibilityEvidence, runtime }) => {
|
||||
const paneKey = resolveOrchestrationCaller(runtime, {
|
||||
callerTerminalHandle: params.from,
|
||||
callerEvidence: orchestrationCompatibilityEvidence,
|
||||
requireStablePane: true
|
||||
const caller = resolveCallerPrincipal(runtime, {
|
||||
from: params.from,
|
||||
agentSessionId: params.agentSessionId,
|
||||
runtimeFence: params.runtimeFence,
|
||||
evidence: orchestrationCompatibilityEvidence,
|
||||
requireBindableCaller: true
|
||||
})
|
||||
const run = runtime.getOrchestrationDb().getCurrentRunForPane(paneKey)
|
||||
const run = runtime.getOrchestrationDb().getCurrentRunForPrincipal(caller.principalId)
|
||||
return { run: run ? exposeRun(run) : null }
|
||||
}
|
||||
}),
|
||||
|
||||
Reference in New Issue
Block a user