mirror of
https://github.com/stablyai/orca.git
synced 2026-10-09 08:02:35 +00:00
feat(agent-launch): host-assigned caller identity and a launch record written when the surface exists
A host restart while a launch waited for its agent to be ready left the launch
recorded as unknown, so a retry was refused ("couldn't confirm") although the
agent was running. The record is now written twice: when the surface exists,
with the launch as it stands and a prompt result that says only what creating
the surface settled, and again when the prompt result is known. The first
write is fired without waiting; the store commits in order.
Every request now carries one caller identity stamped by the dispatcher from
what the transport proved: the runtime socket is the local CLI, the desktop's
IPC channel is the desktop, a paired socket is its device. Params never set
it. The desktop is therefore admitted to replay-safe launches under its own
namespace instead of being refused, and existing CLI and device namespaces keep
their strings so saved rows still replay.
Admission opens only the record store, through its own slot; the chat host is
built on that same store, and stopping the runtime closes the slot whether or
not a host was built.
This commit is contained in:
@@ -11,7 +11,7 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import { executeAgentLaunch, type AgentLaunchExecution } from './agent-launch-executor'
|
||||
import { AgentLaunchStructuredSessionRefusedError } from './agent-launch-surface-factories'
|
||||
import type { AgentLaunchIntent } from '../../shared/agent-launch-intent'
|
||||
import type { AgentLaunchIntent, AgentLaunchResult } from '../../shared/agent-launch-intent'
|
||||
import { FLOATING_TERMINAL_WORKTREE_ID } from '../../shared/constants'
|
||||
|
||||
const STRUCTURED_PREFERENCE = {
|
||||
@@ -29,6 +29,7 @@ function harness(options: {
|
||||
terminalPromptDelivered?: boolean
|
||||
/** Whether the surface reports that its typed line took the offered prompt. */
|
||||
lineCarriesPrompt?: boolean
|
||||
onSurfacePublished?: AgentLaunchExecution['onSurfacePublished']
|
||||
}) {
|
||||
const calls: string[] = []
|
||||
const carried = (startupPrompt: string | undefined) =>
|
||||
@@ -96,7 +97,8 @@ function harness(options: {
|
||||
deliverStructuredPrompt,
|
||||
deliverTerminalPrompt
|
||||
},
|
||||
workspaces: { createWorktree }
|
||||
workspaces: { createWorktree },
|
||||
...(options.onSurfacePublished ? { onSurfacePublished: options.onSurfacePublished } : {})
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -593,3 +595,68 @@ describe('caller-supplied launch inputs', () => {
|
||||
expect(result.warning).toBeUndefined()
|
||||
})
|
||||
})
|
||||
|
||||
describe('the surface is published as the launch stands, before its prompt is delivered', () => {
|
||||
const PROMPTED_EXISTING: AgentLaunchIntent = {
|
||||
agent: 'claude',
|
||||
target: { kind: 'existing', worktree: 'wt-7' },
|
||||
prompt: { text: 'fix the build', delivery: 'submit' }
|
||||
}
|
||||
|
||||
function publishing(options: Parameters<typeof harness>[0]) {
|
||||
const published: AgentLaunchResult[] = []
|
||||
const launch = harness({
|
||||
...options,
|
||||
onSurfacePublished: (surface) => {
|
||||
launch.calls.push('published')
|
||||
published.push(surface)
|
||||
}
|
||||
})
|
||||
return { launch, published }
|
||||
}
|
||||
|
||||
it('records a prompt still owed as not delivered, then delivers it', async () => {
|
||||
const { launch, published } = publishing({ settings: {}, lineCarriesPrompt: false })
|
||||
|
||||
const result = await launch.run(PROMPTED_EXISTING)
|
||||
|
||||
expect(launch.calls).toEqual(['createTerminalAgent', 'published', 'deliverTerminalPrompt'])
|
||||
expect(published).toEqual([
|
||||
{
|
||||
outcome: { kind: 'terminal', handle: 'term_1' },
|
||||
worktreeId: 'wt-7',
|
||||
receipt: result.receipt,
|
||||
prompt: { delivery: 'submit', outcome: 'not-delivered' }
|
||||
}
|
||||
])
|
||||
expect(result.prompt).toEqual({ delivery: 'submit', outcome: 'handed-to-terminal' })
|
||||
})
|
||||
|
||||
it('records a prompt the launch command carried as already handed over', async () => {
|
||||
const { launch, published } = publishing({ settings: {}, lineCarriesPrompt: true })
|
||||
|
||||
const result = await launch.run(PROMPTED_EXISTING)
|
||||
|
||||
expect(published[0]?.prompt).toEqual({ delivery: 'submit', outcome: 'handed-to-terminal' })
|
||||
expect(published[0]).toEqual(result)
|
||||
expect(launch.deliverTerminalPrompt).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('records a chat before its first message is committed', async () => {
|
||||
const { launch, published } = publishing({})
|
||||
|
||||
const result = await launch.run(PROMPTED_EXISTING)
|
||||
|
||||
expect(launch.calls).toEqual([
|
||||
'createSupport',
|
||||
'createStructuredSession',
|
||||
'published',
|
||||
'deliverStructuredPrompt'
|
||||
])
|
||||
expect(published[0]).toEqual({
|
||||
...result,
|
||||
prompt: { delivery: 'submit', outcome: 'not-delivered' }
|
||||
})
|
||||
expect(result.prompt).toEqual({ delivery: 'submit', outcome: 'journaled', messageId: 'msg-1' })
|
||||
})
|
||||
})
|
||||
|
||||
@@ -38,6 +38,7 @@ import {
|
||||
HANDED_TO_TERMINAL,
|
||||
launchCommandPrompt,
|
||||
promptReceipt,
|
||||
settledAtCreation,
|
||||
settleLaunchPromptDisposal
|
||||
} from './agent-launch-prompt-delivery'
|
||||
import type { TuiAgent } from '../../shared/tui-agent'
|
||||
@@ -74,7 +75,12 @@ export type AgentLaunchExecution = {
|
||||
onSurfacePublished?: (surface: AgentLaunchPublishedSurface) => void
|
||||
}
|
||||
|
||||
export type AgentLaunchPublishedSurface = Pick<AgentLaunchResult, 'outcome' | 'worktreeId'>
|
||||
/**
|
||||
* The launch as it stands once its surface exists: a complete result whose prompt receipt says only
|
||||
* what creation itself settled — carried on the launch command, or not (yet) delivered. Complete so
|
||||
* a host that dies during the delivery still leaves a truthful answer behind.
|
||||
*/
|
||||
export type AgentLaunchPublishedSurface = AgentLaunchResult
|
||||
|
||||
export async function executeAgentLaunch(
|
||||
execution: AgentLaunchExecution
|
||||
@@ -101,11 +107,12 @@ export async function executeAgentLaunch(
|
||||
if (intent.reuseTerminal) {
|
||||
const reused = published(execution, {
|
||||
outcome: { kind: 'terminal', handle: intent.reuseTerminal.handle },
|
||||
worktreeId: existingWorktreeId(intent.target)
|
||||
worktreeId: existingWorktreeId(intent.target),
|
||||
receipt: preflight,
|
||||
...promptReceipt(intent, settledAtCreation({}))
|
||||
})
|
||||
return {
|
||||
...reused,
|
||||
receipt: preflight,
|
||||
...promptReceipt(
|
||||
intent,
|
||||
await deliverTerminalLaunchPrompt(execution, intent.reuseTerminal.handle, {
|
||||
@@ -124,12 +131,13 @@ export async function executeAgentLaunch(
|
||||
handle: placed.startupTerminalHandle,
|
||||
...(placed.startupTerminalPaneKey ? { paneKey: placed.startupTerminalPaneKey } : {})
|
||||
},
|
||||
worktreeId: placed.worktreeId
|
||||
worktreeId: placed.worktreeId,
|
||||
receipt: preflight,
|
||||
...(placed.warning ? { warning: placed.warning } : {}),
|
||||
...promptReceipt(intent, settledAtCreation(placed))
|
||||
})
|
||||
return {
|
||||
...startup,
|
||||
receipt: preflight,
|
||||
...(placed.warning ? { warning: placed.warning } : {}),
|
||||
...promptReceipt(
|
||||
intent,
|
||||
placed.promptRodeLaunchCommand
|
||||
@@ -179,11 +187,15 @@ export async function executeAgentLaunch(
|
||||
// not start while looking at it. Telling those apart needs `createManagedWorktree` to stop
|
||||
// multiplexing "couldn't copy untracked files" and "startup terminal failed" into one string.
|
||||
const warning = combineLaunchWarnings(placed.warning, created.warning)
|
||||
const surface = published(execution, { outcome: created.outcome, worktreeId: placed.worktreeId })
|
||||
return {
|
||||
...surface,
|
||||
const surface = published(execution, {
|
||||
outcome: created.outcome,
|
||||
worktreeId: placed.worktreeId,
|
||||
receipt: settled,
|
||||
...(warning ? { warning } : {}),
|
||||
...promptReceipt(intent, settledAtCreation(created))
|
||||
})
|
||||
return {
|
||||
...surface,
|
||||
...promptReceipt(intent, await settleLaunchPromptDisposal(execution, created))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -29,6 +29,17 @@ import type { AgentLaunchStructuredSurface } from './agent-launch-surface-factor
|
||||
export const HANDED_TO_TERMINAL: AgentLaunchPromptDisposal = { outcome: 'handed-to-terminal' }
|
||||
const NOT_DELIVERED: AgentLaunchPromptDisposal = { outcome: 'not-delivered' }
|
||||
|
||||
/**
|
||||
* What the receipt may say before any delivery runs: only what creating the surface already
|
||||
* settled. A launch command that carried the text has handed it over; anything else is not (yet)
|
||||
* delivered, which is also the truthful answer when the host dies before delivering it.
|
||||
*/
|
||||
export function settledAtCreation(created: {
|
||||
promptRodeLaunchCommand?: boolean
|
||||
}): AgentLaunchPromptDisposal {
|
||||
return created.promptRodeLaunchCommand ? HANDED_TO_TERMINAL : NOT_DELIVERED
|
||||
}
|
||||
|
||||
/** Each surface delivers its own way, so the disposal is decided where the surface is known. */
|
||||
export async function settleLaunchPromptDisposal(
|
||||
execution: AgentLaunchExecution,
|
||||
|
||||
@@ -8,10 +8,17 @@
|
||||
import { beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { RuntimeCapability } from '../../shared/protocol-version'
|
||||
import { AGENT_LAUNCH_RUNTIME_CAPABILITY } from '../../shared/agent-launch-runtime-capability'
|
||||
import {
|
||||
resolveRpcCallerIdentity,
|
||||
rpcCallerOperationKey,
|
||||
type RpcCallerIdentity
|
||||
} from '../runtime/rpc/rpc-caller-identity'
|
||||
|
||||
type AdvertisedClient = {
|
||||
clientKind?: 'mobile' | 'runtime'
|
||||
clientCapabilities?: readonly RuntimeCapability[]
|
||||
caller?: RpcCallerIdentity
|
||||
pairedDeviceId?: string
|
||||
}
|
||||
|
||||
const { handlers, advertised } = vi.hoisted(
|
||||
@@ -57,6 +64,7 @@ vi.mock('../runtime/rpc/dispatcher', () => ({
|
||||
|
||||
const { registerRuntimeHandlers } = await import('./runtime')
|
||||
const { supportsAgentLaunch } = await import('../runtime/rpc/methods/agent-launch')
|
||||
const { agentLaunchOperationCallerKey } = await import('../runtime/rpc/methods/agent-launch-replay')
|
||||
|
||||
function rendererEvent() {
|
||||
const mainFrame = {}
|
||||
@@ -126,4 +134,19 @@ describe('desktop renderer reaching agent.launch on its own main process', () =>
|
||||
)
|
||||
expect(supportsAgentLaunch(streaming)).toBe(true)
|
||||
})
|
||||
|
||||
it('names the desktop as its caller on both paths, so its launches replay under one identity', () => {
|
||||
invoke('runtime:call', { method: 'status.get' })
|
||||
invoke('runtime:subscribe', { subscriptionId: 'sub-1', method: 'session.tabs.watch' })
|
||||
|
||||
for (const client of [advertised.unary, advertised.streaming].map(onlyAdvertisedClient)) {
|
||||
const caller = resolveRpcCallerIdentity(client)
|
||||
expect(caller).toEqual({ kind: 'desktop' })
|
||||
expect(agentLaunchOperationCallerKey({ caller })).toBe('trusted-local:desktop')
|
||||
}
|
||||
// Negative control: without the transport's word the same renderer has no identity at all.
|
||||
const { caller: _named, ...unnamed } = onlyAdvertisedClient(advertised.unary)
|
||||
expect(resolveRpcCallerIdentity(unnamed)).toBeUndefined()
|
||||
expect(rpcCallerOperationKey({ kind: 'local-cli' })).toBe('trusted-local:runtime')
|
||||
})
|
||||
})
|
||||
|
||||
@@ -12,6 +12,7 @@ import type { ClientHostedBrowserRowsEvent } from '../../shared/client-hosted-br
|
||||
import { TERMINAL_FIT_RESTORE_DEADLINE_MS } from '../../shared/terminal-fit-restore-deadline'
|
||||
import { DESKTOP_RENDERER_RUNTIME_CLIENT_CAPABILITIES } from './desktop-renderer-runtime-capabilities'
|
||||
import { RpcDispatcher } from '../runtime/rpc/dispatcher'
|
||||
import { DESKTOP_RPC_CALLER } from '../runtime/rpc/rpc-caller-identity'
|
||||
import { ALL_RPC_METHODS } from '../runtime/rpc/methods'
|
||||
import { DesktopRuntimeSenderLifecycle } from './desktop-runtime-sender-lifecycle'
|
||||
|
||||
@@ -78,6 +79,7 @@ export function registerRuntimeHandlers(runtime: OrcaRuntimeService): void {
|
||||
},
|
||||
{
|
||||
clientId: 'desktop-renderer',
|
||||
caller: DESKTOP_RPC_CALLER,
|
||||
clientKind: 'runtime',
|
||||
connectionId: desktopSenders.connectionIdFor(event.sender),
|
||||
clientCapabilities: DESKTOP_RENDERER_RUNTIME_CLIENT_CAPABILITIES
|
||||
@@ -120,6 +122,7 @@ export function registerRuntimeHandlers(runtime: OrcaRuntimeService): void {
|
||||
{
|
||||
signal: controller.signal,
|
||||
clientId: 'desktop-renderer',
|
||||
caller: DESKTOP_RPC_CALLER,
|
||||
clientKind: 'runtime',
|
||||
connectionId,
|
||||
clientCapabilities: DESKTOP_RENDERER_RUNTIME_CLIENT_CAPABILITIES
|
||||
|
||||
@@ -0,0 +1,99 @@
|
||||
/**
|
||||
* The record store opened on its own, for launch admission, and the chat host built on it later.
|
||||
* The store is a single writer, so the property that matters is that there is only ever one.
|
||||
*/
|
||||
|
||||
import { mkdtemp, rm, writeFile } from 'node:fs/promises'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import { JournalHostDatabase } from '../native-chat/agent-session-journal/journal-host-database'
|
||||
import { createStructuredAgentSessionLogger } from '../native-chat/agent-session-wire/structured-agent-session-logger'
|
||||
import { openAgentSessionRecordStoreOnce } from './agent-session-record-store-slot'
|
||||
import {
|
||||
ensureStructuredAgentSessionHost,
|
||||
stopStructuredAgentSessionRuntime
|
||||
} from './structured-agent-session-runtime'
|
||||
|
||||
const HOST_ID = 'local'
|
||||
let stateDirectory: string
|
||||
|
||||
function location(directory = stateDirectory) {
|
||||
return {
|
||||
stateDirectory: directory,
|
||||
hostId: HOST_ID,
|
||||
logger: createStructuredAgentSessionLogger()
|
||||
}
|
||||
}
|
||||
|
||||
beforeEach(async () => {
|
||||
stateDirectory = await mkdtemp(join(tmpdir(), 'orca-record-store-slot-'))
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
await stopStructuredAgentSessionRuntime()
|
||||
vi.restoreAllMocks()
|
||||
await rm(stateDirectory, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
describe('the record store slot', () => {
|
||||
it('hands every caller the same store', async () => {
|
||||
const [first, second] = await Promise.all([
|
||||
openAgentSessionRecordStoreOnce(location()),
|
||||
openAgentSessionRecordStoreOnce(location())
|
||||
])
|
||||
|
||||
expect(second.store).toBe(first.store)
|
||||
})
|
||||
|
||||
it('builds a chat host installed after admission on the store admission opened', async () => {
|
||||
const { store } = await openAgentSessionRecordStoreOnce(location())
|
||||
const close = vi.spyOn(JournalHostDatabase.prototype, 'close')
|
||||
|
||||
const host = await ensureStructuredAgentSessionHost({
|
||||
...location(),
|
||||
claimKeyId: 'key-1',
|
||||
resolveWorkspacePath: async () => stateDirectory,
|
||||
resolveClaudeAuthPolicy: () => ({ stripAuthEnv: true }),
|
||||
resolveEnvironment: async () => ({})
|
||||
})
|
||||
expect(host.deps.store).toBe(store)
|
||||
|
||||
await stopStructuredAgentSessionRuntime()
|
||||
// The host's teardown closes the one connection; stop does not close it a second time.
|
||||
expect(close).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('closes a store no chat host was built on when the runtime stops', async () => {
|
||||
const { store } = await openAgentSessionRecordStoreOnce(location())
|
||||
const close = vi.spyOn(JournalHostDatabase.prototype, 'close')
|
||||
|
||||
await stopStructuredAgentSessionRuntime()
|
||||
|
||||
expect(close).toHaveBeenCalledOnce()
|
||||
const reopened = await openAgentSessionRecordStoreOnce(location())
|
||||
expect(reopened.store).not.toBe(store)
|
||||
})
|
||||
|
||||
it('refuses a second profile rather than opening a second store', async () => {
|
||||
await openAgentSessionRecordStoreOnce(location())
|
||||
const elsewhere = await mkdtemp(join(tmpdir(), 'orca-record-store-slot-other-'))
|
||||
try {
|
||||
await expect(openAgentSessionRecordStoreOnce(location(elsewhere))).rejects.toThrow(
|
||||
'agent_session_record_store_location_changed'
|
||||
)
|
||||
} finally {
|
||||
await rm(elsewhere, { recursive: true, force: true })
|
||||
}
|
||||
})
|
||||
|
||||
it('lets the next caller retry after an open fails', async () => {
|
||||
const notADirectory = join(stateDirectory, 'profile-file')
|
||||
await writeFile(notADirectory, 'not a profile')
|
||||
|
||||
await expect(openAgentSessionRecordStoreOnce(location(notADirectory))).rejects.toThrow()
|
||||
await expect(openAgentSessionRecordStoreOnce(location())).resolves.toMatchObject({
|
||||
store: expect.anything()
|
||||
})
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,73 @@
|
||||
// The process's one durable agent-session record store, opened on its own.
|
||||
//
|
||||
// The launch ledger lives in this store, and admitting a launch is the first thing every replay-safe
|
||||
// `agent.launch` does — a terminal launch included. Reaching the store used to mean installing the
|
||||
// whole chat host (provider adapters, the host, the model catalog) first. The store is opened here
|
||||
// instead and the chat host, when something needs it, is built on this same instance: the store is
|
||||
// a single writer, so a second copy of it would diverge.
|
||||
|
||||
import { AgentSessionRecordStore } from './agent-session-record-store'
|
||||
import type { JournalHostDatabase } from '../native-chat/agent-session-journal/journal-host-database'
|
||||
import { openStructuredAgentSessionJournalDatabase } from './structured-agent-session-journal-open'
|
||||
import type { StructuredAgentSessionLogger } from '../native-chat/agent-session-wire/structured-agent-session-logger'
|
||||
|
||||
export type AgentSessionRecordStoreLocation = {
|
||||
stateDirectory: string
|
||||
hostId: string
|
||||
logger: StructuredAgentSessionLogger
|
||||
}
|
||||
|
||||
export type OpenedAgentSessionRecordStore = {
|
||||
journalDatabase: JournalHostDatabase
|
||||
store: AgentSessionRecordStore
|
||||
}
|
||||
|
||||
let opening: { stateDirectory: string; opened: Promise<OpenedAgentSessionRecordStore> } | null =
|
||||
null
|
||||
|
||||
export function openAgentSessionRecordStoreOnce(
|
||||
location: AgentSessionRecordStoreLocation
|
||||
): Promise<OpenedAgentSessionRecordStore> {
|
||||
if (opening) {
|
||||
// Loud rather than a silent second store: one process serves one profile.
|
||||
if (opening.stateDirectory !== location.stateDirectory) {
|
||||
return Promise.reject(new Error('agent_session_record_store_location_changed'))
|
||||
}
|
||||
return opening.opened
|
||||
}
|
||||
const slot = {
|
||||
stateDirectory: location.stateDirectory,
|
||||
opened: openRecordStore(location).catch((error: unknown) => {
|
||||
// A failed open must not poison the slot: the next caller retries.
|
||||
if (opening === slot) {
|
||||
opening = null
|
||||
}
|
||||
throw error
|
||||
})
|
||||
}
|
||||
opening = slot
|
||||
return slot.opened
|
||||
}
|
||||
|
||||
async function openRecordStore(
|
||||
location: AgentSessionRecordStoreLocation
|
||||
): Promise<OpenedAgentSessionRecordStore> {
|
||||
const journalDatabase = await openStructuredAgentSessionJournalDatabase(location)
|
||||
try {
|
||||
return {
|
||||
journalDatabase,
|
||||
store: AgentSessionRecordStore.open({ journalDatabase, hostId: location.hostId })
|
||||
}
|
||||
} catch (error) {
|
||||
journalDatabase.close()
|
||||
throw error
|
||||
}
|
||||
}
|
||||
|
||||
/** Empties the slot and hands back what it held, for whoever closes its database: the chat host's
|
||||
* teardown when one was built on it, the caller otherwise. */
|
||||
export async function releaseAgentSessionRecordStore(): Promise<OpenedAgentSessionRecordStore | null> {
|
||||
const released = opening
|
||||
opening = null
|
||||
return released ? released.opened.catch(() => null) : null
|
||||
}
|
||||
@@ -14,7 +14,12 @@ import { compareWorktreePs } from './runtime-worktree-status-projection'
|
||||
import type { Repo } from '../../shared/repo-types'
|
||||
import { enrichMissingRepoGitRemoteIdentities } from '../repo-git-remote-identity-enrichment'
|
||||
import { ensureStructuredAgentSessionHost as installStructuredAgentSessionHost } from './structured-agent-session-runtime'
|
||||
import { createStructuredAgentSessionLogger } from '../native-chat/agent-session-wire/structured-agent-session-logger'
|
||||
import {
|
||||
createStructuredAgentSessionLogger,
|
||||
neverThrowingStructuredAgentSessionLogger
|
||||
} from '../native-chat/agent-session-wire/structured-agent-session-logger'
|
||||
import { openAgentSessionRecordStoreOnce } from './agent-session-record-store-slot'
|
||||
import type { AgentSessionRecordStore } from './agent-session-record-store'
|
||||
import { maybeAutoRenameWorkspaceOnFirstStructuredTurn } from '../agent-hooks/first-work-structured-session-rename'
|
||||
import { firstWorkRenameDeps } from '../agent-hooks/first-work-rename-runtime'
|
||||
import { getProfileUserDataPath } from '../orca-profiles/profile-storage-paths'
|
||||
@@ -132,6 +137,17 @@ export class OrcaRuntimeWithGetWorktreePs extends OrcaRuntimeWithStartTuiIdleVis
|
||||
})
|
||||
}
|
||||
|
||||
/** The durable record store alone, without the chat host: launch admission needs only its
|
||||
* ledger. A host installed later is built on this same store. */
|
||||
async openAgentSessionRecordStore(): Promise<AgentSessionRecordStore> {
|
||||
const { store } = await openAgentSessionRecordStoreOnce({
|
||||
stateDirectory: getProfileUserDataPath(),
|
||||
hostId: LOCAL_EXECUTION_HOST_ID,
|
||||
logger: neverThrowingStructuredAgentSessionLogger(createStructuredAgentSessionLogger())
|
||||
})
|
||||
return store
|
||||
}
|
||||
|
||||
/**
|
||||
* Installs the structured agent-session host on first use. Lazy for the same
|
||||
* reason the orchestration DB is: the profile's user-data path is not final
|
||||
|
||||
@@ -11,6 +11,7 @@ import type {
|
||||
import type { RuntimeCapability } from '../../../shared/protocol-version'
|
||||
import type { OrchestrationCompatibilityEvidence } from '../../../shared/orchestration-compatibility-evidence'
|
||||
import type { OrchestrationSessionCaller } from '../orchestration/orchestration-caller-identity'
|
||||
import type { RpcCallerIdentity } from './rpc-caller-identity'
|
||||
|
||||
export type PairingRpcContext = {
|
||||
getEndpoints(params: PairingGetEndpointsParams): Promise<PairingGetEndpointsResult>
|
||||
@@ -75,6 +76,8 @@ export type RpcContext = {
|
||||
clientId?: string
|
||||
// Why: navigation is keyed by revocable device identity, never by the bearer credential or transient socket id.
|
||||
pairedDeviceId?: string
|
||||
// Why: host-assigned from the transport (rpc-caller-identity.ts); absent when the transport could not name its caller.
|
||||
caller?: RpcCallerIdentity
|
||||
// Why: lets handlers gate mobile payload truncation to phones only; undefined for in-process callers → treat as full-class (no clip).
|
||||
clientKind?: 'mobile' | 'runtime'
|
||||
// Why: negotiation is bound to the authenticated socket, never asserted by a destructive request.
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import type { RuntimeCapability } from '../../../shared/protocol-version'
|
||||
import type { TerminalStreamFrame } from '../../../shared/terminal-stream-protocol'
|
||||
import type { PairingRpcContext } from './core'
|
||||
import type { RpcCallerIdentity } from './rpc-caller-identity'
|
||||
|
||||
export type RpcDispatchStreamingOptions = {
|
||||
authenticatedCallerFingerprint?: string
|
||||
@@ -8,6 +9,8 @@ export type RpcDispatchStreamingOptions = {
|
||||
signal?: AbortSignal
|
||||
clientId?: string
|
||||
pairedDeviceId?: string
|
||||
/** Set by a transport that knows its caller but carries no paired device (the desktop's IPC). */
|
||||
caller?: RpcCallerIdentity
|
||||
clientKind?: 'mobile' | 'runtime'
|
||||
clientCapabilities?: readonly RuntimeCapability[]
|
||||
updateClientCapabilities?: (capabilities: readonly RuntimeCapability[]) => void
|
||||
|
||||
@@ -24,6 +24,7 @@ import { mapDispatcherError } from './dispatcher-error-response'
|
||||
import { parseRpcRequestParams } from './dispatcher-request-parsing'
|
||||
import { RpcStreamingDispatcher } from './rpc-streaming-dispatcher'
|
||||
import { invokeDispatcherUnaryMethod } from './dispatcher-unary-method-invocation'
|
||||
import { resolveRpcCallerIdentity } from './rpc-caller-identity'
|
||||
import {
|
||||
needsOrchestrationCallerResolution,
|
||||
resolveOrchestrationSessionCaller,
|
||||
@@ -116,6 +117,7 @@ export class RpcDispatcher {
|
||||
: undefined,
|
||||
requestId: request.id,
|
||||
clientId: options?.clientId,
|
||||
caller: resolveRpcCallerIdentity(options),
|
||||
clientKind: options?.clientKind,
|
||||
clientCapabilities: options?.clientCapabilities,
|
||||
updateClientCapabilities: options?.updateClientCapabilities,
|
||||
|
||||
@@ -9,8 +9,6 @@ import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import { openTestAgentSessionRecordStore } from '../../agent-session-record-store-test-harness'
|
||||
import { setStructuredAgentSessionHost } from '../../../native-chat/agent-session-wire/structured-agent-session-registry'
|
||||
import type { StructuredAgentSessionHost } from '../../../native-chat/agent-session-wire/structured-agent-session-host'
|
||||
import type { RpcContext } from '../core'
|
||||
import {
|
||||
CAPABLE_CLIENT,
|
||||
@@ -18,6 +16,7 @@ import {
|
||||
methodNamed,
|
||||
rpcContext,
|
||||
runtimeStub,
|
||||
setAgentLaunchRecordStore,
|
||||
type AgentLaunchRuntimeStub
|
||||
} from './agent-launch.test-fixture'
|
||||
|
||||
@@ -209,12 +208,11 @@ describe('a replayed launch', () => {
|
||||
beforeEach(async () => {
|
||||
directory = await mkdtemp(join(tmpdir(), 'orca-agent-launch-caller-'))
|
||||
const store = await openTestAgentSessionRecordStore(directory)
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: `deps.store` is the only member `agent.launch` reads, and a member it omits throws on call.
|
||||
setStructuredAgentSessionHost({ deps: { store } } as unknown as StructuredAgentSessionHost)
|
||||
setAgentLaunchRecordStore(store)
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
setStructuredAgentSessionHost(null)
|
||||
setAgentLaunchRecordStore(null)
|
||||
await rm(directory, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
|
||||
@@ -17,12 +17,11 @@ import type { RpcContext } from '../core'
|
||||
/** The paired client whose view this launch should move, or null when it moves the host's. */
|
||||
export function agentLaunchCallerNavigationId(
|
||||
target: AgentLaunchTarget,
|
||||
context: Pick<RpcContext, 'clientKind' | 'pairedDeviceId'>
|
||||
context: Pick<RpcContext, 'caller'>
|
||||
): string | null {
|
||||
if (target.kind !== 'existing' || context.clientKind === undefined) {
|
||||
return null
|
||||
}
|
||||
return context.pairedDeviceId?.trim() || null
|
||||
return target.kind === 'existing' && context.caller?.kind === 'paired-device'
|
||||
? context.caller.deviceId
|
||||
: null
|
||||
}
|
||||
|
||||
/** Bookkeeping, never a gate: the agent already runs, so a failure here only leaves the view as it was. */
|
||||
|
||||
@@ -13,8 +13,6 @@ import { join } from 'node:path'
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { AgentSessionRecordStore } from '../../agent-session-record-store'
|
||||
import { openTestAgentSessionRecordStore } from '../../agent-session-record-store-test-harness'
|
||||
import { setStructuredAgentSessionHost } from '../../../native-chat/agent-session-wire/structured-agent-session-registry'
|
||||
import type { StructuredAgentSessionHost } from '../../../native-chat/agent-session-wire/structured-agent-session-host'
|
||||
import type { OrcaRuntimeService } from '../../orca-runtime'
|
||||
import type { RpcContext } from '../core'
|
||||
import { RpcDispatcher } from '../dispatcher'
|
||||
@@ -24,6 +22,7 @@ import {
|
||||
methodNamed,
|
||||
rpcContext,
|
||||
runtimeStub,
|
||||
setAgentLaunchRecordStore,
|
||||
type AgentLaunchRuntimeStub as RuntimeStub
|
||||
} from './agent-launch.test-fixture'
|
||||
|
||||
@@ -222,12 +221,11 @@ describe('a live-pane refusal under a named operation', () => {
|
||||
beforeEach(async () => {
|
||||
directory = await mkdtemp(join(tmpdir(), 'orca-agent-launch-pane-'))
|
||||
store = await openTestAgentSessionRecordStore(directory)
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: `deps.store` is the only member `agent.launch` reads, and a member it omits throws on call.
|
||||
setStructuredAgentSessionHost({ deps: { store } } as unknown as StructuredAgentSessionHost)
|
||||
setAgentLaunchRecordStore(store)
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
setStructuredAgentSessionHost(null)
|
||||
setAgentLaunchRecordStore(null)
|
||||
await rm(directory, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
|
||||
@@ -10,11 +10,14 @@ import { join } from 'node:path'
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { AgentSessionRecordStore } from '../../agent-session-record-store'
|
||||
import { openTestAgentSessionRecordStore } from '../../agent-session-record-store-test-harness'
|
||||
import { setStructuredAgentSessionHost } from '../../../native-chat/agent-session-wire/structured-agent-session-registry'
|
||||
import type { StructuredAgentSessionHost } from '../../../native-chat/agent-session-wire/structured-agent-session-host'
|
||||
import type { OrcaRuntimeService } from '../../orca-runtime'
|
||||
import { RpcDispatcher } from '../dispatcher'
|
||||
import { methodNamed, runtimeStub, type AgentLaunchRuntimeStub } from './agent-launch.test-fixture'
|
||||
import {
|
||||
methodNamed,
|
||||
runtimeStub,
|
||||
setAgentLaunchRecordStore,
|
||||
type AgentLaunchRuntimeStub
|
||||
} from './agent-launch.test-fixture'
|
||||
|
||||
vi.mock('./structured-agent-session-create', () => ({
|
||||
createStructuredAgentSessionForWorktree: async () => ({
|
||||
@@ -57,12 +60,11 @@ describe('a launch whose terminal fails', () => {
|
||||
beforeEach(async () => {
|
||||
directory = await mkdtemp(join(tmpdir(), 'orca-agent-launch-prestart-'))
|
||||
store = await openTestAgentSessionRecordStore(directory)
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: `deps.store` is the only member `agent.launch` reads, and a member it omits throws on call.
|
||||
setStructuredAgentSessionHost({ deps: { store } } as unknown as StructuredAgentSessionHost)
|
||||
setAgentLaunchRecordStore(store)
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
setStructuredAgentSessionHost(null)
|
||||
setAgentLaunchRecordStore(null)
|
||||
await rm(directory, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
|
||||
@@ -30,8 +30,6 @@ import {
|
||||
openTestAgentSessionRecordStore,
|
||||
readPersistedTestAgentSessionStore
|
||||
} from '../../agent-session-record-store-test-harness'
|
||||
import { setStructuredAgentSessionHost } from '../../../native-chat/agent-session-wire/structured-agent-session-registry'
|
||||
import type { StructuredAgentSessionHost } from '../../../native-chat/agent-session-wire/structured-agent-session-host'
|
||||
import type { RpcContext } from '../core'
|
||||
import type { OrcaRuntimeService } from '../../orca-runtime'
|
||||
import { RpcDispatcher } from '../dispatcher'
|
||||
@@ -39,6 +37,7 @@ import {
|
||||
methodNamed,
|
||||
rpcContext,
|
||||
runtimeStub,
|
||||
setAgentLaunchRecordStore,
|
||||
type AgentLaunchRuntimeStub
|
||||
} from './agent-launch.test-fixture'
|
||||
|
||||
@@ -125,14 +124,12 @@ beforeEach(async () => {
|
||||
createStructuredSession.mockClear()
|
||||
directory = await mkdtemp(join(tmpdir(), 'orca-agent-launch-replay-'))
|
||||
store = await openTestAgentSessionRecordStore(directory)
|
||||
// The launch reaches the ledger through the installed host; nothing else on the host is used,
|
||||
// because the structured create below it is mocked out.
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: `deps.store` is the only member `agent.launch` reads, and a member it omits throws on call.
|
||||
setStructuredAgentSessionHost({ deps: { store } } as unknown as StructuredAgentSessionHost)
|
||||
// The ledger alone, as admission opens it: no chat host is installed.
|
||||
setAgentLaunchRecordStore(store)
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
setStructuredAgentSessionHost(null)
|
||||
setAgentLaunchRecordStore(null)
|
||||
await rm(directory, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
@@ -289,10 +286,7 @@ describe('a replay answers from the record', () => {
|
||||
const first = await launch(params, runtime)
|
||||
|
||||
const reopened = await openTestAgentSessionRecordStore(directory)
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: see the setup above.
|
||||
setStructuredAgentSessionHost({
|
||||
deps: { store: reopened }
|
||||
} as unknown as StructuredAgentSessionHost)
|
||||
setAgentLaunchRecordStore(reopened)
|
||||
const afterRestart = runtimeStub()
|
||||
|
||||
expect(await launch(params, afterRestart)).toEqual(first)
|
||||
@@ -503,10 +497,7 @@ describe('an unreadable launch payload costs one replay, never the store', () =>
|
||||
await rewriteRecordedLaunch({ outcome: { kind: 'structured' }, worktreeId: 'wt-1' })
|
||||
|
||||
const reopened = await openTestAgentSessionRecordStore(directory)
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: see the setup above.
|
||||
setStructuredAgentSessionHost({
|
||||
deps: { store: reopened }
|
||||
} as unknown as StructuredAgentSessionHost)
|
||||
setAgentLaunchRecordStore(reopened)
|
||||
|
||||
expect(reopened.listOperationRows()).toHaveLength(1)
|
||||
const retry = runtimeStub()
|
||||
|
||||
@@ -22,37 +22,19 @@ import type {
|
||||
AgentSessionOperationRefusalCode
|
||||
} from '../../../../shared/agent-session-operation-ledger'
|
||||
import { resolveAgentSessionReplayOutcome } from '../../../native-chat/agent-session-wire/structured-agent-session-replay-outcome'
|
||||
import { getStructuredAgentSessionHost } from '../../../native-chat/agent-session-wire/structured-agent-session-registry'
|
||||
import type { AgentSessionRecordStore } from '../../agent-session-record-store'
|
||||
import type { RpcContext } from '../core'
|
||||
import { rpcCallerOperationKey } from '../rpc-caller-identity'
|
||||
import type { AgentLaunchParams } from './agent-launch-schemas'
|
||||
|
||||
/** Remote replay needs the paired-device subject because its bearer credential can rotate. */
|
||||
export function agentLaunchOperationCallerKey(
|
||||
context: Pick<RpcContext, 'pairedDeviceId' | 'clientKind'>
|
||||
): string {
|
||||
if (context.clientKind === undefined) {
|
||||
return 'trusted-local:runtime'
|
||||
}
|
||||
const pairedDeviceId = context.pairedDeviceId?.trim()
|
||||
if (!pairedDeviceId) {
|
||||
/**
|
||||
* The ledger namespace of whoever the transport says is calling. A transport that could not name its
|
||||
* caller gets no replay safety at all, rather than a namespace shared with strangers.
|
||||
*/
|
||||
export function agentLaunchOperationCallerKey(context: Pick<RpcContext, 'caller'>): string {
|
||||
if (!context.caller) {
|
||||
throw new Error('agent_session_identity_required')
|
||||
}
|
||||
return pairedDeviceId
|
||||
}
|
||||
|
||||
/**
|
||||
* The store is owned by the structured session host, so reaching it installs that host — already
|
||||
* true of any structured launch, which installs it to attach. The change is that a terminal-bound
|
||||
* launch now opens the record store too, when and only when its caller asked for replay safety.
|
||||
*/
|
||||
async function requireLaunchOperationStore(context: RpcContext): Promise<AgentSessionRecordStore> {
|
||||
await context.runtime.ensureStructuredAgentSessionHost()
|
||||
const host = getStructuredAgentSessionHost()
|
||||
if (!host) {
|
||||
throw new Error('structured_agent_session_unsupported')
|
||||
}
|
||||
return host.deps.store
|
||||
return rpcCallerOperationKey(context.caller)
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -67,6 +49,9 @@ export type AgentLaunchAdmission =
|
||||
/** This caller owns the operation. It alone runs the effect, and it must settle the row. */
|
||||
| {
|
||||
decision: 'execute'
|
||||
/** The surface exists: records the launch as it stands, so a restart before `settle` replays
|
||||
* the running agent instead of refusing an unknown outcome. */
|
||||
record: (provisional: AgentLaunchResult) => Promise<void>
|
||||
settle: (result: AgentLaunchResult) => Promise<void>
|
||||
fail: (code: string) => Promise<void>
|
||||
/** Distinct from the launch id: the inner attach reserves in this same ledger. */
|
||||
@@ -134,8 +119,9 @@ export async function admitAgentLaunchOperation(
|
||||
if (!attachOperationId) {
|
||||
return refusal(operationId, 'agent_session_operation_invalid', 'is not a durable operation id')
|
||||
}
|
||||
const store = await requireLaunchOperationStore(context)
|
||||
const callerKey = agentLaunchOperationCallerKey(context)
|
||||
// The ledger alone: admitting a terminal launch has no use for the chat host.
|
||||
const store = await context.runtime.openAgentSessionRecordStore()
|
||||
const { decision: admitted, claim } = await store.admitAndClaimOperation(
|
||||
{ callerKey, operationId, fingerprint, now },
|
||||
// A fresh row, or a replayed one no one has answered yet, leaves the right to run open.
|
||||
@@ -169,21 +155,24 @@ export async function admitAgentLaunchOperation(
|
||||
refusal(operationId, 'agent_session_operation_unknown', 'is claimed but unsettled')
|
||||
)
|
||||
}
|
||||
const succeeded = (result: AgentLaunchResult) =>
|
||||
store.recordOperationOutcome({
|
||||
callerKey,
|
||||
operationId,
|
||||
outcome: {
|
||||
status: 'succeeded',
|
||||
// A terminal surface has a handle, not a session id; `launch` carries whichever it is.
|
||||
sessionId: result.outcome.kind === 'structured' ? result.outcome.sessionId : '',
|
||||
launch: result
|
||||
}
|
||||
})
|
||||
return {
|
||||
decision: 'execute',
|
||||
attachOperationId,
|
||||
callerKey,
|
||||
settle: (result) =>
|
||||
store.recordOperationOutcome({
|
||||
callerKey,
|
||||
operationId,
|
||||
outcome: {
|
||||
status: 'succeeded',
|
||||
// A terminal surface has a handle, not a session id; `launch` carries whichever it is.
|
||||
sessionId: result.outcome.kind === 'structured' ? result.outcome.sessionId : '',
|
||||
launch: result
|
||||
}
|
||||
}),
|
||||
// The same row shape twice: a build that predates the first write reads either one.
|
||||
record: succeeded,
|
||||
settle: succeeded,
|
||||
fail: (code) =>
|
||||
store.recordOperationOutcome({
|
||||
callerKey,
|
||||
|
||||
@@ -0,0 +1,375 @@
|
||||
/**
|
||||
* A launch's record across a host restart, against the real durable ledger.
|
||||
*
|
||||
* The record is written twice: once when the surface exists and once when the prompt's fate is
|
||||
* known. A host that dies at any point leaves the replay a truthful answer — the running agent once
|
||||
* its surface is recorded, an honest "unknown" before — and never a second agent. A "restart" here
|
||||
* is what a new process sees: the store reopened from disk and a runtime with no in-flight launches.
|
||||
*/
|
||||
|
||||
import { mkdtemp, rm } from 'node:fs/promises'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import { AGENT_LAUNCH_RUNTIME_CAPABILITY } from '../../../../shared/agent-launch-runtime-capability'
|
||||
import type { AgentLaunchResult } from '../../../../shared/agent-launch-intent'
|
||||
import type { AgentSessionOperationRow } from '../../../../shared/agent-session-operation-ledger'
|
||||
import type { AgentSessionRecordStore } from '../../agent-session-record-store'
|
||||
import { openTestAgentSessionRecordStore } from '../../agent-session-record-store-test-harness'
|
||||
import type { OrcaRuntimeService } from '../../orca-runtime'
|
||||
import type { RpcContext } from '../core'
|
||||
import { RpcDispatcher } from '../dispatcher'
|
||||
import { DESKTOP_RPC_CALLER } from '../rpc-caller-identity'
|
||||
import {
|
||||
methodNamed,
|
||||
rpcContext,
|
||||
runtimeStub,
|
||||
setAgentLaunchRecordStore,
|
||||
type AgentLaunchRuntimeStub
|
||||
} from './agent-launch.test-fixture'
|
||||
|
||||
const deliverTerminalPrompt = vi.hoisted(() =>
|
||||
vi.fn(async (_args: { handle: string }): Promise<boolean> => true)
|
||||
)
|
||||
vi.mock('./agent-launch-terminal-prompt', () => ({
|
||||
deliverTerminalAgentLaunchPrompt: deliverTerminalPrompt
|
||||
}))
|
||||
|
||||
const { AGENT_LAUNCH_METHODS } = await import('./agent-launch')
|
||||
const AGENT_LAUNCH_REPLAY = methodNamed(AGENT_LAUNCH_METHODS, 'agent.launchReplay')
|
||||
|
||||
// The ledger admits against `Date.now()`, so the id must be dated now.
|
||||
const OPERATION_ID = `${Date.now()}-000000000000000000000000000000ee`
|
||||
const PANE_KEY = '9b1deb4d-3b7d-4bad-9bdd-2b0d7b3dcb6d:3f2504e0-4f89-41d3-9a0c-0305e82c3301'
|
||||
// No structured preference: every launch here is a terminal agent, the surface a restart outlives.
|
||||
const TERMINAL_ONLY = {}
|
||||
const PROMPTED_LAUNCH = {
|
||||
agent: 'claude',
|
||||
target: { kind: 'existing', worktree: 'id:wt-7' },
|
||||
prompt: { text: 'fix the failing test', delivery: 'submit' },
|
||||
operationId: OPERATION_ID
|
||||
}
|
||||
const PHONE: Partial<RpcContext> = {
|
||||
clientKind: 'mobile',
|
||||
pairedDeviceId: 'device-1',
|
||||
clientCapabilities: [AGENT_LAUNCH_RUNTIME_CAPABILITY]
|
||||
}
|
||||
/** Exactly what the desktop's `runtime:call` handler hands the dispatcher. */
|
||||
const DESKTOP_IPC = {
|
||||
clientId: 'desktop-renderer',
|
||||
caller: DESKTOP_RPC_CALLER,
|
||||
clientKind: 'runtime' as const,
|
||||
clientCapabilities: [AGENT_LAUNCH_RUNTIME_CAPABILITY]
|
||||
}
|
||||
|
||||
let directory: string
|
||||
let store: AgentSessionRecordStore
|
||||
|
||||
function hostRuntime() {
|
||||
// A phone's launch into an existing workspace also moves the phone's own view to the new tab.
|
||||
return Object.assign(runtimeStub({ settings: TERMINAL_ONLY, terminalPaneKey: PANE_KEY }), {
|
||||
selectCreatedMobileSessionTabForClient: vi.fn(() => true)
|
||||
})
|
||||
}
|
||||
|
||||
function launch(
|
||||
runtime: AgentLaunchRuntimeStub,
|
||||
params: unknown = PROMPTED_LAUNCH,
|
||||
context: Partial<RpcContext> = PHONE
|
||||
): Promise<AgentLaunchResult> {
|
||||
return AGENT_LAUNCH_REPLAY.handler(
|
||||
AGENT_LAUNCH_REPLAY.params.parse(params),
|
||||
rpcContext(runtime, context)
|
||||
)
|
||||
}
|
||||
|
||||
/** A new process: the store reread from disk, and nothing in flight. */
|
||||
async function restartHost(): Promise<void> {
|
||||
store = await openTestAgentSessionRecordStore(directory)
|
||||
setAgentLaunchRecordStore(store)
|
||||
}
|
||||
|
||||
function row(callerKey = 'device-1'): AgentSessionOperationRow | null {
|
||||
return store.getOperationRow(callerKey, OPERATION_ID)
|
||||
}
|
||||
|
||||
/** The first write is fired, not awaited, so it lands a moment after the tab is published. */
|
||||
async function untilRecorded(status: AgentSessionOperationRow['outcome']['status']): Promise<void> {
|
||||
const deadline = Date.now() + 2_000
|
||||
while (row()?.outcome.status !== status) {
|
||||
if (Date.now() > deadline) {
|
||||
throw new Error(`the launch row never reached ${status}`)
|
||||
}
|
||||
await new Promise((resolve) => setTimeout(resolve, 5))
|
||||
}
|
||||
}
|
||||
|
||||
/** The store commits in order, so a write queued now lands after every write already fired. */
|
||||
async function ledgerWritesQueuedBefore(ledger: AgentSessionRecordStore): Promise<void> {
|
||||
await ledger.recordOperationOutcome({
|
||||
callerKey: 'no-such-caller',
|
||||
operationId: OPERATION_ID,
|
||||
outcome: { status: 'unknown' }
|
||||
})
|
||||
}
|
||||
|
||||
function dispatcherFor(runtime: AgentLaunchRuntimeStub): RpcDispatcher {
|
||||
return new RpcDispatcher({
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the fixture implements every runtime method agent.launch reaches, plus the id the dispatcher stamps on replies.
|
||||
runtime: { ...runtime, getRuntimeId: () => 'runtime-1' } as unknown as OrcaRuntimeService,
|
||||
methods: AGENT_LAUNCH_METHODS
|
||||
})
|
||||
}
|
||||
|
||||
beforeEach(async () => {
|
||||
deliverTerminalPrompt.mockReset()
|
||||
deliverTerminalPrompt.mockResolvedValue(true)
|
||||
directory = await mkdtemp(join(tmpdir(), 'orca-agent-launch-restart-'))
|
||||
store = await openTestAgentSessionRecordStore(directory)
|
||||
setAgentLaunchRecordStore(store)
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
setAgentLaunchRecordStore(null)
|
||||
await rm(directory, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
describe('a host restart mid-launch', () => {
|
||||
it('finds the running agent when the host died while its prompt waited for readiness', async () => {
|
||||
// The paste waits for the agent forever: the host dies first.
|
||||
let waiting: () => void = () => {}
|
||||
const readinessWait = new Promise<void>((resolve) => {
|
||||
waiting = resolve
|
||||
})
|
||||
deliverTerminalPrompt.mockImplementationOnce(() => {
|
||||
waiting()
|
||||
return new Promise<boolean>(() => {})
|
||||
})
|
||||
const dying = hostRuntime()
|
||||
void launch(dying)
|
||||
await readinessWait
|
||||
await ledgerWritesQueuedBefore(store)
|
||||
|
||||
await restartHost()
|
||||
const restarted = hostRuntime()
|
||||
const replayed = await launch(restarted)
|
||||
|
||||
expect(replayed).toEqual({
|
||||
outcome: { kind: 'terminal', handle: 'term_1', paneKey: PANE_KEY },
|
||||
worktreeId: 'wt-7',
|
||||
receipt: expect.objectContaining({ mode: 'terminal' }),
|
||||
// The dead host never pasted it, and saying so lets the caller keep the text.
|
||||
prompt: { delivery: 'submit', outcome: 'not-delivered' }
|
||||
})
|
||||
expect(restarted.createTerminal).not.toHaveBeenCalled()
|
||||
expect(dying.createTerminal).toHaveBeenCalledOnce()
|
||||
expect(deliverTerminalPrompt).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('stays unknown when the host died before the surface was recorded', async () => {
|
||||
const dying = hostRuntime()
|
||||
let spawnRequested: () => void = () => {}
|
||||
const spawning = new Promise<void>((resolve) => {
|
||||
spawnRequested = resolve
|
||||
})
|
||||
dying.createTerminal.mockImplementationOnce(() => {
|
||||
spawnRequested()
|
||||
return new Promise(() => {})
|
||||
})
|
||||
void launch(dying)
|
||||
await spawning
|
||||
|
||||
await restartHost()
|
||||
const restarted = hostRuntime()
|
||||
|
||||
// A pane may exist that nothing recorded; "unknown" is the truthful answer, never a relaunch.
|
||||
await expect(launch(restarted)).rejects.toThrow('agent_session_operation_unknown')
|
||||
expect(row()?.outcome.status).toBe('unknown')
|
||||
expect(restarted.createTerminal).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('replays the final answer when the host died after the prompt landed', async () => {
|
||||
const first = await launch(hostRuntime())
|
||||
expect(first.prompt).toEqual({ delivery: 'submit', outcome: 'handed-to-terminal' })
|
||||
|
||||
await restartHost()
|
||||
const restarted = hostRuntime()
|
||||
|
||||
await expect(launch(restarted)).resolves.toEqual(first)
|
||||
expect(restarted.createTerminal).not.toHaveBeenCalled()
|
||||
expect(deliverTerminalPrompt).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('lets the final write replace the first, never the other way round', async () => {
|
||||
let releasePaste: (pasted: boolean) => void = () => {}
|
||||
deliverTerminalPrompt.mockImplementationOnce(
|
||||
() =>
|
||||
new Promise<boolean>((resolve) => {
|
||||
releasePaste = resolve
|
||||
})
|
||||
)
|
||||
const running = launch(hostRuntime())
|
||||
await untilRecorded('succeeded')
|
||||
releasePaste(true)
|
||||
await running
|
||||
|
||||
await restartHost()
|
||||
const outcome = row()?.outcome
|
||||
expect(outcome?.status === 'succeeded' && outcome.launch).toMatchObject({
|
||||
prompt: { outcome: 'handed-to-terminal' }
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
describe('a reply lost three times', () => {
|
||||
it('starts one agent, and every retry gets its answer', async () => {
|
||||
const host = hostRuntime()
|
||||
const first = await launch(host)
|
||||
// Lost once and twice on the same host, then a third time across a restart.
|
||||
const second = await launch(host)
|
||||
const third = await launch(host)
|
||||
await restartHost()
|
||||
const afterRestart = hostRuntime()
|
||||
const fourth = await launch(afterRestart)
|
||||
|
||||
expect([second, third, fourth]).toEqual([first, first, first])
|
||||
expect(host.createTerminal).toHaveBeenCalledOnce()
|
||||
expect(afterRestart.createTerminal).not.toHaveBeenCalled()
|
||||
expect(deliverTerminalPrompt).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('starts one agent when three retries arrive while the first is still running', async () => {
|
||||
let releasePaste: (pasted: boolean) => void = () => {}
|
||||
deliverTerminalPrompt.mockImplementationOnce(
|
||||
() =>
|
||||
new Promise<boolean>((resolve) => {
|
||||
releasePaste = resolve
|
||||
})
|
||||
)
|
||||
const host = hostRuntime()
|
||||
const attempts = [launch(host), launch(host), launch(host)]
|
||||
await untilRecorded('succeeded')
|
||||
releasePaste(true)
|
||||
const [first, ...retries] = await Promise.all(attempts)
|
||||
|
||||
expect(retries).toEqual([first, first])
|
||||
expect(host.createTerminal).toHaveBeenCalledOnce()
|
||||
})
|
||||
})
|
||||
|
||||
describe('the desktop launches replay-safely', () => {
|
||||
it('admits the desktop window under its own host-assigned identity and replays its id', async () => {
|
||||
const host = hostRuntime()
|
||||
const request = {
|
||||
id: 'request-1',
|
||||
authToken: 'desktop-ipc',
|
||||
method: 'agent.launchReplay',
|
||||
params: PROMPTED_LAUNCH
|
||||
}
|
||||
|
||||
const first = await dispatcherFor(host).dispatch(request, DESKTOP_IPC)
|
||||
const replayed = await dispatcherFor(host).dispatch(request, DESKTOP_IPC)
|
||||
|
||||
expect(first).toMatchObject({ ok: true, result: { outcome: { kind: 'terminal' } } })
|
||||
expect(replayed).toMatchObject({ ok: true, result: first.ok ? first.result : null })
|
||||
expect(host.createTerminal).toHaveBeenCalledOnce()
|
||||
expect(row('trusted-local:desktop')?.outcome.status).toBe('succeeded')
|
||||
})
|
||||
|
||||
it('opens the ledger alone, never the chat host, to admit a terminal launch', async () => {
|
||||
const host = hostRuntime()
|
||||
|
||||
await launch(host)
|
||||
|
||||
expect(host.openAgentSessionRecordStore).toHaveBeenCalled()
|
||||
expect(host.ensureStructuredAgentSessionHost).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
|
||||
describe('a caller cannot claim an identity', () => {
|
||||
it('keys a paired device by its authenticated subject whatever its params say', async () => {
|
||||
const host = hostRuntime()
|
||||
const response = await new Promise<string>((resolve) => {
|
||||
void dispatcherFor(host).dispatchStreaming(
|
||||
{
|
||||
id: 'request-1',
|
||||
authToken: 'device-token',
|
||||
method: 'agent.launchReplay',
|
||||
params: {
|
||||
...PROMPTED_LAUNCH,
|
||||
caller: DESKTOP_RPC_CALLER,
|
||||
callerKey: 'trusted-local:desktop',
|
||||
pairedDeviceId: 'device-2'
|
||||
}
|
||||
},
|
||||
resolve,
|
||||
{ ...PHONE, clientId: 'device-token' }
|
||||
)
|
||||
})
|
||||
|
||||
expect(JSON.parse(response)).toMatchObject({ ok: true })
|
||||
expect(store.listOperationRows().map((entry) => entry.callerKey)).toEqual(['device-1'])
|
||||
})
|
||||
|
||||
it('refuses replay safety to a transport that cannot name its caller', async () => {
|
||||
const host = hostRuntime()
|
||||
|
||||
const response = await dispatcherFor(host).dispatch(
|
||||
{
|
||||
id: 'request-1',
|
||||
authToken: 'token',
|
||||
method: 'agent.launchReplay',
|
||||
params: { ...PROMPTED_LAUNCH, caller: DESKTOP_RPC_CALLER }
|
||||
},
|
||||
{ clientKind: 'runtime', clientCapabilities: [AGENT_LAUNCH_RUNTIME_CAPABILITY] }
|
||||
)
|
||||
|
||||
expect(response).toMatchObject({
|
||||
ok: false,
|
||||
error: { code: 'agent_session_identity_required' }
|
||||
})
|
||||
expect(host.createTerminal).not.toHaveBeenCalled()
|
||||
expect(store.listOperationRows()).toHaveLength(0)
|
||||
})
|
||||
})
|
||||
|
||||
describe('the ledger stays bounded', () => {
|
||||
it('keeps one row per launch, retained from admission, across both writes', async () => {
|
||||
let releasePaste: (pasted: boolean) => void = () => {}
|
||||
deliverTerminalPrompt.mockImplementationOnce(
|
||||
() =>
|
||||
new Promise<boolean>((resolve) => {
|
||||
releasePaste = resolve
|
||||
})
|
||||
)
|
||||
const running = launch(hostRuntime())
|
||||
await untilRecorded('succeeded')
|
||||
const afterFirstWrite = row()
|
||||
releasePaste(true)
|
||||
await running
|
||||
|
||||
expect(store.listOperationRows()).toHaveLength(1)
|
||||
expect(row()?.expiresAt).toBe(afterFirstWrite?.expiresAt)
|
||||
expect(row()?.recordedAt).toBe(afterFirstWrite?.recordedAt)
|
||||
})
|
||||
|
||||
it('holds the desktop to the same per-caller cap as every other caller', async () => {
|
||||
const now = Date.now()
|
||||
for (let index = 0; index < 512; index += 1) {
|
||||
await store.admitOperation({
|
||||
callerKey: 'trusted-local:desktop',
|
||||
operationId: `${now}-${index.toString(16).padStart(32, '0')}`,
|
||||
fingerprint: 'fp-seeded',
|
||||
now
|
||||
})
|
||||
}
|
||||
const host = hostRuntime()
|
||||
|
||||
await expect(launch(host, PROMPTED_LAUNCH, { ...DESKTOP_IPC })).rejects.toThrow(
|
||||
'agent_session_operation_capacity'
|
||||
)
|
||||
expect(host.createTerminal).not.toHaveBeenCalled()
|
||||
// Another caller's namespace is not starved by it.
|
||||
await expect(launch(host)).resolves.toMatchObject({ outcome: { kind: 'terminal' } })
|
||||
})
|
||||
})
|
||||
@@ -12,8 +12,6 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import { computeAgentLaunchFingerprint } from '../../../../shared/agent-launch-operation'
|
||||
import type { AgentSessionRecordStore } from '../../agent-session-record-store'
|
||||
import { openTestAgentSessionRecordStore } from '../../agent-session-record-store-test-harness'
|
||||
import { setStructuredAgentSessionHost } from '../../../native-chat/agent-session-wire/structured-agent-session-registry'
|
||||
import type { StructuredAgentSessionHost } from '../../../native-chat/agent-session-wire/structured-agent-session-host'
|
||||
import type { OrcaRuntimeService } from '../../orca-runtime'
|
||||
import type { RpcContext } from '../core'
|
||||
import { RpcDispatcher } from '../dispatcher'
|
||||
@@ -23,6 +21,7 @@ import {
|
||||
methodNamed,
|
||||
rpcContext,
|
||||
runtimeStub,
|
||||
setAgentLaunchRecordStore,
|
||||
type AgentLaunchRuntimeStub as RuntimeStub
|
||||
} from './agent-launch.test-fixture'
|
||||
|
||||
@@ -210,12 +209,11 @@ describe('a taken session id under a named operation', () => {
|
||||
beforeEach(async () => {
|
||||
directory = await mkdtemp(join(tmpdir(), 'orca-agent-launch-session-'))
|
||||
store = await openTestAgentSessionRecordStore(directory)
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: `deps.store` is the only member `agent.launch` reads, and a member it omits throws on call.
|
||||
setStructuredAgentSessionHost({ deps: { store } } as unknown as StructuredAgentSessionHost)
|
||||
setAgentLaunchRecordStore(store)
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
setStructuredAgentSessionHost(null)
|
||||
setAgentLaunchRecordStore(null)
|
||||
await rm(directory, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
|
||||
@@ -9,6 +9,8 @@ import { vi } from 'vitest'
|
||||
import { AGENT_LAUNCH_RUNTIME_CAPABILITY } from '../../../../shared/agent-launch-runtime-capability'
|
||||
import { AgentLaunchPaneAlreadyLiveError } from '../../../../shared/agent-launch-pane-already-live'
|
||||
import type { RpcContext } from '../core'
|
||||
import { resolveRpcCallerIdentity } from '../rpc-caller-identity'
|
||||
import type { AgentSessionRecordStore } from '../../agent-session-record-store'
|
||||
|
||||
export const STRUCTURED_PREFERENCE = {
|
||||
experimentalNativeChat: true,
|
||||
@@ -49,6 +51,13 @@ function reportPromptCarry(
|
||||
}
|
||||
}
|
||||
|
||||
let launchRecordStore: AgentSessionRecordStore | null = null
|
||||
|
||||
/** The ledger every stub's `openAgentSessionRecordStore` opens, as a process opens its one store. */
|
||||
export function setAgentLaunchRecordStore(store: AgentSessionRecordStore | null): void {
|
||||
launchRecordStore = store
|
||||
}
|
||||
|
||||
export function runtimeStub(options: AgentLaunchRuntimeStubOptions = {}) {
|
||||
const worktreeCreateResults = new Map<string, Promise<unknown>>()
|
||||
const waitForSetupTerminalCompletion = vi.fn(
|
||||
@@ -119,6 +128,12 @@ export function runtimeStub(options: AgentLaunchRuntimeStubOptions = {}) {
|
||||
folderWorkspace: null
|
||||
})),
|
||||
ensureStructuredAgentSessionHost: vi.fn(async () => {}),
|
||||
openAgentSessionRecordStore: vi.fn(async (): Promise<AgentSessionRecordStore> => {
|
||||
if (!launchRecordStore) {
|
||||
throw new Error('agent_session_record_store_unavailable')
|
||||
}
|
||||
return launchRecordStore
|
||||
}),
|
||||
waitForSetupTerminalCompletion
|
||||
}
|
||||
}
|
||||
@@ -139,12 +154,14 @@ export function methodNamed<TMethod extends { name: string }, TName extends stri
|
||||
}
|
||||
|
||||
// The one call the stub cannot satisfy structurally; every method it does implement is asserted.
|
||||
// The caller is stamped the way the dispatcher stamps it from the same transport fields.
|
||||
export function rpcContext(
|
||||
runtime: AgentLaunchRuntimeStub,
|
||||
context: Partial<RpcContext>
|
||||
): RpcContext {
|
||||
const caller = context.caller ?? resolveRpcCallerIdentity(context)
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the stub implements only the runtime surface these methods reach, so a method it omits throws on call rather than reading a wrong value.
|
||||
return { runtime, ...context } as unknown as RpcContext
|
||||
return { runtime, ...context, ...(caller ? { caller } : {}) } as unknown as RpcContext
|
||||
}
|
||||
|
||||
export const CAPABLE_CLIENT: Partial<RpcContext> = {
|
||||
|
||||
@@ -142,12 +142,20 @@ async function resolveUnlaunchedIntent(
|
||||
return intent
|
||||
}
|
||||
|
||||
/** What a launch admitted under an operation id carries into its execution. */
|
||||
type ReplaySafeLaunch = {
|
||||
attachOperationId: string
|
||||
callerKey: string
|
||||
terminalSpawn: TerminalSpawnDispatch
|
||||
/** Records the surface the moment it exists. Fired, never awaited: the ledger's transactions run
|
||||
* in order, so the final settle still lands after it, and the prompt never waits on bookkeeping. */
|
||||
recordSurface: (provisional: AgentLaunchResult) => void
|
||||
}
|
||||
|
||||
async function runAgentLaunch(
|
||||
intent: AgentLaunchIntent,
|
||||
context: RpcContext,
|
||||
attachOperationId?: string,
|
||||
operationCallerKey?: string,
|
||||
terminalSpawn?: TerminalSpawnDispatch
|
||||
replaySafe?: ReplaySafeLaunch
|
||||
): Promise<AgentLaunchResult> {
|
||||
const callerNavigationId = agentLaunchCallerNavigationId(intent.target, context)
|
||||
return executeAgentLaunch({
|
||||
@@ -155,19 +163,19 @@ async function runAgentLaunch(
|
||||
intent,
|
||||
surfaces: agentLaunchSurfaceFactory(
|
||||
context,
|
||||
attachOperationId,
|
||||
operationCallerKey,
|
||||
replaySafe?.attachOperationId,
|
||||
replaySafe?.callerKey,
|
||||
callerNavigationId !== null,
|
||||
terminalSpawn
|
||||
replaySafe?.terminalSpawn
|
||||
),
|
||||
workspaces: agentLaunchWorkspaceFactory(context, intent.agent),
|
||||
// The tab is shown as it is published, not after a prompt that can take a minute to land.
|
||||
...(callerNavigationId !== null
|
||||
? {
|
||||
onSurfacePublished: (surface) =>
|
||||
selectAgentLaunchTabForCaller(context.runtime, surface, callerNavigationId)
|
||||
}
|
||||
: {})
|
||||
onSurfacePublished: (surface) => {
|
||||
replaySafe?.recordSurface(surface)
|
||||
if (callerNavigationId !== null) {
|
||||
selectAgentLaunchTabForCaller(context.runtime, surface, callerNavigationId)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
@@ -255,13 +263,12 @@ async function executeReplaySafeAgentLaunch(
|
||||
const terminalSpawn = trackTerminalSpawnDispatch()
|
||||
let result: AgentLaunchResult
|
||||
try {
|
||||
result = await runAgentLaunch(
|
||||
intent,
|
||||
context,
|
||||
admission.attachOperationId,
|
||||
admission.callerKey,
|
||||
terminalSpawn
|
||||
)
|
||||
result = await runAgentLaunch(intent, context, {
|
||||
attachOperationId: admission.attachOperationId,
|
||||
callerKey: admission.callerKey,
|
||||
terminalSpawn,
|
||||
recordSurface: (provisional) => void settleQuietly(admission.record(provisional))
|
||||
})
|
||||
} catch (error) {
|
||||
const failedWithoutEffects = launchFailureWithoutEffectsCode(
|
||||
error,
|
||||
|
||||
@@ -0,0 +1,57 @@
|
||||
/**
|
||||
* Who is calling, as the host's own transport established it.
|
||||
*
|
||||
* One identity per caller, stamped by the dispatcher from what the connection proved — the runtime
|
||||
* socket, the desktop's IPC channel, or a paired device's authenticated socket — and never read
|
||||
* from request params, so a caller cannot claim to be another. Durable per-caller state (the launch
|
||||
* ledger's namespace) and per-caller behaviour (whose view a launch moves) both read this.
|
||||
*/
|
||||
|
||||
import type { RpcDispatchStreamingOptions } from './dispatcher-stream-options'
|
||||
|
||||
export type RpcCallerIdentity =
|
||||
/** The `orca` CLI over the runtime socket, and the in-process bridges that relay it. */
|
||||
| { kind: 'local-cli' }
|
||||
/** This host's own desktop app over IPC. One identity for every window, so a reload or an app
|
||||
* restart still names the same caller. */
|
||||
| { kind: 'desktop' }
|
||||
/** A paired phone or remote desktop, by its device subject, which survives credential rotation. */
|
||||
| { kind: 'paired-device'; deviceId: string }
|
||||
|
||||
export const DESKTOP_RPC_CALLER: RpcCallerIdentity = { kind: 'desktop' }
|
||||
|
||||
/**
|
||||
* A transport that names its caller is believed; one that declares no client at all is the
|
||||
* in-process runtime socket, trusted as it always was; one that declares a client it cannot name has
|
||||
* no identity, and anything that needs one refuses it.
|
||||
*/
|
||||
export function resolveRpcCallerIdentity(
|
||||
transport:
|
||||
| Pick<RpcDispatchStreamingOptions, 'caller' | 'clientKind' | 'pairedDeviceId'>
|
||||
| undefined
|
||||
): RpcCallerIdentity | undefined {
|
||||
if (transport?.caller) {
|
||||
return transport.caller
|
||||
}
|
||||
const deviceId = transport?.pairedDeviceId?.trim()
|
||||
if (deviceId) {
|
||||
return { kind: 'paired-device', deviceId }
|
||||
}
|
||||
return transport?.clientKind === undefined ? { kind: 'local-cli' } : undefined
|
||||
}
|
||||
|
||||
/**
|
||||
* The caller's namespace in the durable operation ledger. Rows outlive the process, so these strings
|
||||
* are permanent: the CLI and paired-device keys are the ones rows were written under before this
|
||||
* identity existed.
|
||||
*/
|
||||
export function rpcCallerOperationKey(caller: RpcCallerIdentity): string {
|
||||
switch (caller.kind) {
|
||||
case 'local-cli':
|
||||
return 'trusted-local:runtime'
|
||||
case 'desktop':
|
||||
return 'trusted-local:desktop'
|
||||
case 'paired-device':
|
||||
return caller.deviceId
|
||||
}
|
||||
}
|
||||
@@ -21,6 +21,7 @@ import { parseRpcRequestParams } from './dispatcher-request-parsing'
|
||||
import { routeDispatcherClientHostedBrowserRpc } from './dispatcher-client-browser-routing'
|
||||
import { needsLocalCallerFingerprint } from './dispatcher-caller-fingerprint'
|
||||
import { createDispatcherStreamingFeatureEmitter } from './dispatcher-streaming-feature-emitter'
|
||||
import { resolveRpcCallerIdentity } from './rpc-caller-identity'
|
||||
import {
|
||||
needsOrchestrationCallerResolution,
|
||||
resolveOrchestrationSessionCaller,
|
||||
@@ -138,6 +139,7 @@ export class RpcStreamingDispatcher {
|
||||
subscriptionRegistrationVersion,
|
||||
clientId: options?.clientId,
|
||||
pairedDeviceId: options?.pairedDeviceId,
|
||||
caller: resolveRpcCallerIdentity(options),
|
||||
clientKind: options?.clientKind,
|
||||
clientCapabilities: options?.clientCapabilities,
|
||||
updateClientCapabilities: options?.updateClientCapabilities,
|
||||
@@ -192,6 +194,7 @@ export class RpcStreamingDispatcher {
|
||||
connectionId: options?.connectionId,
|
||||
clientId: options?.clientId,
|
||||
pairedDeviceId: options?.pairedDeviceId,
|
||||
caller: resolveRpcCallerIdentity(options),
|
||||
clientKind: options?.clientKind,
|
||||
clientCapabilities: options?.clientCapabilities,
|
||||
updateClientCapabilities: options?.updateClientCapabilities,
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
/**
|
||||
* "This machine holds a structured chat" against a real host and record store. Session history,
|
||||
* resume preparation and replay-safe phone launches all build the host for a user who never had a
|
||||
* chat; only a chat record may turn the signal on, and the first one must turn it on at once.
|
||||
* resume preparation and replay-safe phone launches all open the record store for a user who never
|
||||
* had a chat; only a chat record may turn the signal on, and the first one must turn it on at once.
|
||||
*/
|
||||
|
||||
import { mkdtemp, rm } from 'node:fs/promises'
|
||||
@@ -21,6 +21,7 @@ import {
|
||||
stopStructuredAgentSessionRuntime
|
||||
} from './structured-agent-session-runtime'
|
||||
import { createStructuredAgentSessionLogger } from '../native-chat/agent-session-wire/structured-agent-session-logger'
|
||||
import { openAgentSessionRecordStoreOnce } from './agent-session-record-store-slot'
|
||||
|
||||
vi.mock('../ai-vault/session-scanner-worker-spawn', () => ({
|
||||
scanAiVaultSessionsInWorker: vi.fn(),
|
||||
@@ -53,6 +54,14 @@ const runtime = {
|
||||
ensureStructuredAgentSessionHost: async () => {
|
||||
await installHost()
|
||||
},
|
||||
openAgentSessionRecordStore: async () =>
|
||||
(
|
||||
await openAgentSessionRecordStoreOnce({
|
||||
stateDirectory,
|
||||
hostId: 'local',
|
||||
logger: createStructuredAgentSessionLogger()
|
||||
})
|
||||
).store,
|
||||
listAiVaultSessions: async () => ({
|
||||
sessions: [],
|
||||
issues: [],
|
||||
@@ -110,7 +119,12 @@ describe('whether this machine holds a structured chat', () => {
|
||||
|
||||
it('stays false when a replay-safe phone launch records its operation', async () => {
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: admission reads only these fields.
|
||||
const context = { runtime, clientKind: 'mobile', pairedDeviceId: 'phone-1' } as RpcContext
|
||||
const context = {
|
||||
runtime,
|
||||
clientKind: 'mobile',
|
||||
pairedDeviceId: 'phone-1',
|
||||
caller: { kind: 'paired-device', deviceId: 'phone-1' }
|
||||
} as RpcContext
|
||||
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: admission reads only the operation id.
|
||||
const params = { operationId: operationId('a1') } as Parameters<
|
||||
typeof admitAgentLaunchOperation
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
// An install that fails after it opened the chat journal closes that connection, so the next
|
||||
// install does not leave a second one open in the same process.
|
||||
// An install that fails after it opened the chat journal leaves that one connection to the record
|
||||
// store slot, which launch admission may already be using, and the next install builds on it: the
|
||||
// process never holds a second connection.
|
||||
|
||||
import { mkdtemp, rm } from 'node:fs/promises'
|
||||
import { tmpdir } from 'node:os'
|
||||
@@ -54,18 +55,20 @@ afterEach(async () => {
|
||||
})
|
||||
|
||||
describe('an install that fails after opening the chat journal', () => {
|
||||
it('closes the journal connection it opened, and the next install opens its own', async () => {
|
||||
it('keeps the one connection it opened, and the next install builds on it', async () => {
|
||||
const open = vi.spyOn(JournalHostDatabase, 'open')
|
||||
mocks.failWiring.mockReturnValueOnce(true)
|
||||
|
||||
await expect(install()).rejects.toThrow('model catalog wiring failed')
|
||||
const failed = await open.mock.results[0]?.value
|
||||
expect(failed).toBeInstanceOf(JournalHostDatabase)
|
||||
expect(failed.isClosed).toBe(true)
|
||||
const opened = await open.mock.results[0]?.value
|
||||
expect(opened).toBeInstanceOf(JournalHostDatabase)
|
||||
expect(opened.isClosed).toBe(false)
|
||||
|
||||
await expect(install()).resolves.toBeDefined()
|
||||
expect(open).toHaveBeenCalledTimes(2)
|
||||
const reopened = await open.mock.results[1]?.value
|
||||
expect(reopened?.isClosed).toBe(false)
|
||||
expect(open).toHaveBeenCalledOnce()
|
||||
expect(opened.isClosed).toBe(false)
|
||||
|
||||
await stopStructuredAgentSessionRuntime()
|
||||
expect(opened.isClosed).toBe(true)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -37,9 +37,11 @@ import {
|
||||
readClaudeManagedAccountGateSettings,
|
||||
type ClaudeManagedAccountGateSettings
|
||||
} from '../native-chat/claude-structured-managed-account-support'
|
||||
import { AgentSessionRecordStore } from './agent-session-record-store'
|
||||
import type { JournalHostDatabase } from '../native-chat/agent-session-journal/journal-host-database'
|
||||
import { openStructuredAgentSessionJournalDatabase } from './structured-agent-session-journal-open'
|
||||
import {
|
||||
openAgentSessionRecordStoreOnce,
|
||||
releaseAgentSessionRecordStore,
|
||||
type OpenedAgentSessionRecordStore
|
||||
} from './agent-session-record-store-slot'
|
||||
import { legacyAgentSessionStorePath } from './agent-session-record-store-file'
|
||||
import { journalDatabasePath } from '../native-chat/agent-session-journal/journal-host-database'
|
||||
import { journalDatabaseHoldsAgentSessions } from '../native-chat/agent-session-journal/journal-database'
|
||||
@@ -192,6 +194,14 @@ export async function stopStructuredAgentSessionRuntime(options?: {
|
||||
if (installed) {
|
||||
outstanding.push(installed)
|
||||
}
|
||||
// A store admission opened with no host built on it has no teardown to close its database.
|
||||
const recordStore = await releaseAgentSessionRecordStore()
|
||||
if (
|
||||
recordStore &&
|
||||
!outstanding.some((runtime) => runtime.journalDatabase === recordStore.journalDatabase)
|
||||
) {
|
||||
recordStore.journalDatabase.close()
|
||||
}
|
||||
const failures: unknown[] = []
|
||||
for (const runtime of outstanding) {
|
||||
try {
|
||||
@@ -222,27 +232,21 @@ async function install(deps: StructuredAgentSessionRuntimeDeps): Promise<Install
|
||||
throw new Error(STRUCTURED_AGENT_SESSION_LOGGER_REQUIRED)
|
||||
}
|
||||
const logger = neverThrowingStructuredAgentSessionLogger(deps.logger)
|
||||
const journalDatabase = await openStructuredAgentSessionJournalDatabase({
|
||||
// The store launch admission may already have opened; a failed install leaves it to that slot.
|
||||
const recordStore = await openAgentSessionRecordStoreOnce({
|
||||
stateDirectory: deps.stateDirectory,
|
||||
hostId: deps.hostId,
|
||||
logger
|
||||
})
|
||||
try {
|
||||
return await installOnJournal({ ...deps, logger }, journalDatabase)
|
||||
} catch (error) {
|
||||
// Nothing else holds the connection yet, and the next install opens its own.
|
||||
journalDatabase.close()
|
||||
throw error
|
||||
}
|
||||
return await installOnJournal({ ...deps, logger }, recordStore)
|
||||
}
|
||||
|
||||
async function installOnJournal(
|
||||
deps: StructuredAgentSessionRuntimeDeps,
|
||||
journalDatabase: JournalHostDatabase
|
||||
{ journalDatabase, store }: OpenedAgentSessionRecordStore
|
||||
): Promise<InstalledRuntime> {
|
||||
const envResolvers = createStructuredAgentEnvironmentResolvers(deps)
|
||||
const { resolveCodexEnvironment, resolveClaudeInheritedEnv } = envResolvers
|
||||
const store = AgentSessionRecordStore.open({ journalDatabase, hostId: deps.hostId })
|
||||
let host: StructuredAgentSessionHost | null = null
|
||||
const lifecycle = createStructuredAgentSessionLifecycleDelivery({
|
||||
handle: (event) => host?.handleAdapterEvent(event),
|
||||
|
||||
Reference in New Issue
Block a user