feat(native-chat): add Claude structured adapter

This commit is contained in:
Merge Sim
2026-08-31 17:44:30 -07:00
parent 26031ca317
commit f289fee894
35 changed files with 4665 additions and 7 deletions
@@ -0,0 +1,29 @@
import { describe, expect, it } from 'vitest'
import { buildClaudeChildProcessEnv } from './claude-child-process-environment'
describe('Claude child process environment', () => {
it('strips case-insensitive inherited auth and session stamps on Windows', () => {
const env = buildClaudeChildProcessEnv(
{
ANTHROPIC_AUTH_TOKEN: 'configured-token',
Claude_Code_Session_Id: 'configured-session'
},
{
platform: 'win32',
inheritedEnv: {
anthropic_api_key: 'inherited-key',
Anthropic_Custom_Headers: 'Authorization: inherited',
claude_code_child_session: '1',
CLAUDE_CODE_SESSION_ID: 'inherited-session',
SAFE_VALUE: 'preserved'
}
}
)
expect(env).toEqual({
ANTHROPIC_AUTH_TOKEN: 'configured-token',
Claude_Code_Session_Id: 'configured-session',
SAFE_VALUE: 'preserved'
})
})
})
@@ -0,0 +1,47 @@
import { CLAUDE_AUTH_ENV_VARS, applyClaudeEnvPatch } from '../claude-accounts/environment'
const CLAUDE_CHILD_SESSION_STAMP_ENV_KEYS = [
'CLAUDE_CODE_CHILD_SESSION',
'CLAUDE_CODE_SESSION_ID',
'CLAUDE_CODE_BRIDGE_SESSION_ID'
] as const
function cloneProcessEnv(source: NodeJS.ProcessEnv): Record<string, string> {
const env: Record<string, string> = {}
for (const [key, value] of Object.entries(source)) {
if (value !== undefined) {
env[key] = value
}
}
return env
}
export function buildClaudeChildProcessEnv(
configuredEnv: Record<string, string> = {},
options: { inheritedEnv?: NodeJS.ProcessEnv; platform?: NodeJS.Platform } = {}
): Record<string, string> {
const inheritedEnv = options.inheritedEnv ?? process.env
const platform = options.platform ?? process.platform
const env = applyClaudeEnvPatch(cloneProcessEnv(inheritedEnv), {}, { stripAuthEnv: true })
if (platform === 'win32') {
const authKeys = new Set(CLAUDE_AUTH_ENV_VARS.map((key) => key.toUpperCase()))
for (const [key, value] of Object.entries(env)) {
const normalized = key.toUpperCase()
if (
authKeys.has(normalized) ||
(normalized === 'ANTHROPIC_CUSTOM_HEADERS' &&
/authorization|x-api-key|api-key|bearer/i.test(value))
) {
delete env[key]
}
}
}
for (const key of CLAUDE_CHILD_SESSION_STAMP_ENV_KEYS) {
for (const inheritedKey of Object.keys(env)) {
if (inheritedKey === key || (platform === 'win32' && inheritedKey.toUpperCase() === key)) {
delete env[inheritedKey]
}
}
}
return Object.assign(env, configuredEnv)
}
@@ -0,0 +1,197 @@
import { EventEmitter } from 'node:events'
import { PassThrough } from 'node:stream'
import { afterEach, describe, expect, it, vi } from 'vitest'
import type { spawnProcess } from '../../shared/child-process/run-process'
import {
openClaudeStreamJsonConnection,
type ClaudeControlRequest
} from './claude-stream-json-connection'
type FakeChild = EventEmitter & {
pid: number
stdin: PassThrough
stdout: PassThrough
stderr: PassThrough
kill: ReturnType<typeof vi.fn>
}
function fakeSpawn() {
const child = new EventEmitter() as FakeChild
child.pid = 4321
child.stdin = new PassThrough()
child.stdout = new PassThrough()
child.stderr = new PassThrough()
child.kill = vi.fn(() => true)
const spawnMock = vi.fn((_spec: Parameters<typeof spawnProcess>[0]) => child)
const spawnImpl = spawnMock as unknown as typeof spawnProcess
return { child, spawnImpl, spawnMock }
}
function writtenFrames(child: FakeChild): Promise<Record<string, unknown>[]> {
return new Promise((resolve) => {
setImmediate(() => {
const text = child.stdin.read()?.toString('utf8') ?? ''
resolve(
text
.trim()
.split('\n')
.filter(Boolean)
.map((line) => JSON.parse(line) as Record<string, unknown>)
)
})
})
}
describe('Claude stream-json connection', () => {
afterEach(() => vi.unstubAllEnvs())
it('spawns in the pinned workspace and routes acknowledged control requests', async () => {
const process = fakeSpawn()
const connection = await openClaudeStreamJsonConnection(
{
command: 'claude',
args: ['-p'],
cwd: '/work/repo',
env: { CLAUDE_CONFIG_DIR: '/accounts/one' }
},
{},
process.spawnImpl
)
const listing = connection.request('list_models')
const outbound = await writtenFrames(process.child)
const requestId = (outbound[0] as { request_id: string }).request_id
process.child.stdout.write(
`${JSON.stringify({
type: 'control_response',
response: {
subtype: 'success',
request_id: requestId,
response: { models: [{ value: 'sonnet' }] }
}
})}\n`
)
await expect(listing).resolves.toEqual({ models: [{ value: 'sonnet' }] })
expect(process.spawnImpl).toHaveBeenCalledWith(
expect.objectContaining({
program: 'claude',
args: ['-p'],
cwd: '/work/repo',
env: expect.objectContaining({ CLAUDE_CONFIG_DIR: '/accounts/one' })
})
)
})
it('routes provider permission controls and writes their response envelope', async () => {
const process = fakeSpawn()
let inbound: ClaudeControlRequest | null = null
const connection = await openClaudeStreamJsonConnection(
{ command: 'claude', args: [], cwd: '/work' },
{ onControlRequest: (request) => (inbound = request) },
process.spawnImpl
)
process.child.stdout.write(
`${JSON.stringify({
type: 'control_request',
request_id: 'permission-1',
request: { subtype: 'can_use_tool', tool_name: 'Bash' }
})}\n`
)
await new Promise((resolve) => setImmediate(resolve))
expect(inbound).toMatchObject({ request_id: 'permission-1' })
await connection.respond('permission-1', { behavior: 'deny', message: 'No' })
expect(await writtenFrames(process.child)).toEqual([
{
type: 'control_response',
response: {
subtype: 'success',
request_id: 'permission-1',
response: { behavior: 'deny', message: 'No' }
}
}
])
})
it('forwards unmatched control responses to the provider-frame path', async () => {
const process = fakeSpawn()
const onMessage = vi.fn()
await openClaudeStreamJsonConnection(
{ command: 'claude', args: [], cwd: '/work' },
{ onMessage },
process.spawnImpl
)
const frame = {
type: 'control_response',
response: { subtype: 'success', request_id: 'provider-owned', response: { opaque: true } }
}
process.child.stdout.write(`${JSON.stringify(frame)}\n`)
await new Promise((resolve) => setImmediate(resolve))
expect(onMessage).toHaveBeenCalledWith(frame)
})
it('fails closed on malformed provider output', async () => {
const process = fakeSpawn()
const onExit = vi.fn()
const onMessage = vi.fn()
await openClaudeStreamJsonConnection(
{ command: 'claude', args: [], cwd: '/work' },
{ onExit, onMessage },
process.spawnImpl
)
process.child.stdout.write('{not-json}\n{"type":"assistant"}\n')
await new Promise((resolve) => setImmediate(resolve))
expect(onMessage).not.toHaveBeenCalled()
expect(onExit).toHaveBeenCalledWith(expect.objectContaining({ message: expect.any(String) }))
})
it('uses configured auth and strips inherited auth and child-session stamps', async () => {
vi.stubEnv('ANTHROPIC_API_KEY', 'inherited-key')
vi.stubEnv('CLAUDE_CODE_CHILD_SESSION', '1')
vi.stubEnv('CLAUDE_CODE_SESSION_ID', 'inherited-session')
const process = fakeSpawn()
await openClaudeStreamJsonConnection(
{
command: 'claude',
args: [],
cwd: '/work',
env: {
ANTHROPIC_AUTH_TOKEN: 'configured-token',
ANTHROPIC_BASE_URL: 'https://gateway.example.test',
CLAUDE_CODE_SESSION_ID: 'configured-session'
}
},
{},
process.spawnImpl
)
const env = process.spawnMock.mock.calls[0]?.[0]?.env
expect(env).toMatchObject({
ANTHROPIC_AUTH_TOKEN: 'configured-token',
ANTHROPIC_BASE_URL: 'https://gateway.example.test',
CLAUDE_CODE_SESSION_ID: 'configured-session'
})
expect(env?.ANTHROPIC_API_KEY).toBeUndefined()
expect(env?.CLAUDE_CODE_CHILD_SESSION).toBeUndefined()
})
it('settles concurrent and repeated closes after a spawn failure', async () => {
const process = fakeSpawn()
const connection = await openClaudeStreamJsonConnection(
{ command: 'missing-claude', args: [], cwd: '/work' },
{},
process.spawnImpl
)
process.child.emit('error', new Error('spawn missing-claude ENOENT'))
await expect(Promise.all([connection.close(), connection.close()])).resolves.toEqual([
true,
true
])
await expect(connection.close()).resolves.toBe(true)
})
})
@@ -0,0 +1,286 @@
import { forceTerminateProcessTree } from '../../shared/child-process/process-tree-termination'
import { spawnProcess } from '../../shared/child-process/run-process'
import { waitForProcessExitUntil } from '../codex/codex-process-exit-deadline'
import { buildClaudeChildProcessEnv } from './claude-child-process-environment'
import { attachClaudeStreamJsonStdout } from './claude-stream-json-stdout'
export type ClaudeStreamJsonLaunch = {
command: string
args: string[]
cwd: string
env?: Record<string, string>
}
export type ClaudeControlRequest = {
type: 'control_request'
request_id: string
request: Record<string, unknown> & { subtype: string }
}
export type ClaudeControlCancelRequest = {
type: 'control_cancel_request'
request_id: string
}
export type ClaudeControlResponder = Pick<
ClaudeStreamJsonConnection,
'respond' | 'respondWithError'
>
export type ClaudeStreamJsonConnectionHandlers = {
onMessage?: (message: Record<string, unknown>) => void
onControlRequest?: (request: ClaudeControlRequest, responder?: ClaudeControlResponder) => void
onControlCancelRequest?: (request: ClaudeControlCancelRequest) => void
onExit?: (error: Error) => void
}
export type ClaudeStreamJsonConnection = {
readonly pid: number | undefined
readonly closed: boolean
send: (message: Record<string, unknown>) => Promise<void>
request: (
subtype: string,
params?: Record<string, unknown>,
options?: { timeoutMs?: number }
) => Promise<unknown>
respond: (requestId: string, response: unknown) => Promise<void>
respondWithError: (requestId: string, error: string) => Promise<void>
close: () => Promise<boolean>
}
export class ClaudeControlRequestError extends Error {
constructor(
readonly subtype: string,
message: string
) {
super(message)
this.name = 'ClaudeControlRequestError'
}
}
const DEFAULT_REQUEST_TIMEOUT_MS = 30_000
const GRACEFUL_EXIT_MS = 1_500
const FORCED_EXIT_MS = 1_000
const STDERR_TAIL_MAX_BYTES = 8192
const STDOUT_LINE_MAX_BYTES = 8 * 1024 * 1024
type PendingRequest = {
subtype: string
resolve: (result: unknown) => void
reject: (error: Error) => void
timer: ReturnType<typeof setTimeout>
}
function isRecord(value: unknown): value is Record<string, unknown> {
return typeof value === 'object' && value !== null && !Array.isArray(value)
}
function exitError(stderrTail: string, cause?: Error): Error {
const detail = stderrTail.trim()
const message = detail ? `claude stream-json exited: ${detail}` : 'claude stream-json exited'
return cause ? new Error(message, { cause }) : new Error(message)
}
export async function openClaudeStreamJsonConnection(
launch: ClaudeStreamJsonLaunch,
handlers: ClaudeStreamJsonConnectionHandlers = {},
spawnImpl: typeof spawnProcess = spawnProcess
): Promise<ClaudeStreamJsonConnection> {
const child = spawnImpl({
program: launch.command,
args: launch.args,
cwd: launch.cwd,
env: buildClaudeChildProcessEnv(launch.env),
stdio: ['pipe', 'pipe', 'pipe'],
detached: process.platform !== 'win32'
})
const pending = new Map<string, PendingRequest>()
let nextRequestId = 1
let stderrTail = ''
let exited = false
let closing = false
let terminalError: Error | null = null
let writeChain: Promise<void> = Promise.resolve()
let closePromise: Promise<boolean> | null = null
let settleExit = (): void => {}
const exitPromise = new Promise<void>((resolve) => {
settleExit = resolve
})
const markExited = (): void => {
exited = true
settleExit()
}
child.on('exit', markExited)
const failPending = (error: Error): void => {
for (const waiter of pending.values()) {
clearTimeout(waiter.timer)
waiter.reject(error)
}
pending.clear()
}
const handleUnexpectedEnd = (cause?: Error): void => {
if (terminalError) {
return
}
terminalError = exitError(stderrTail, cause)
failPending(terminalError)
if (!closing) {
handlers.onExit?.(terminalError)
}
}
const writeLine = (payload: Record<string, unknown>): Promise<void> => {
const write = async (): Promise<void> => {
if (closing || exited || terminalError || child.stdin.destroyed || !child.stdin.writable) {
throw terminalError ?? new Error('claude stream-json connection is closed')
}
const line = `${JSON.stringify(payload)}\n`
await new Promise<void>((resolve, reject) => {
child.stdin.write(line, (error) => (error ? reject(error) : resolve()))
})
}
const queued = writeChain.then(write)
writeChain = queued.catch(() => {})
return queued
}
const respond = (requestId: string, response: unknown): Promise<void> =>
writeLine({
type: 'control_response',
response: { subtype: 'success', request_id: requestId, response }
})
const respondWithError = (requestId: string, error: string): Promise<void> =>
writeLine({
type: 'control_response',
response: { subtype: 'error', request_id: requestId, error }
})
const dispatchMessage = (message: Record<string, unknown>): void => {
if (message.type === 'control_response') {
const response = isRecord(message.response) ? message.response : null
const requestId = typeof response?.request_id === 'string' ? response.request_id : null
const waiter = requestId ? pending.get(requestId) : undefined
if (!waiter || !requestId || !response) {
handlers.onMessage?.(message)
return
}
pending.delete(requestId)
clearTimeout(waiter.timer)
if (response.subtype === 'success') {
waiter.resolve(response.response)
} else {
waiter.reject(
new ClaudeControlRequestError(
waiter.subtype,
typeof response.error === 'string'
? response.error
: `claude ${waiter.subtype} request failed`
)
)
}
return
}
if (
message.type === 'control_request' &&
typeof message.request_id === 'string' &&
isRecord(message.request) &&
typeof message.request.subtype === 'string'
) {
handlers.onControlRequest?.(message as ClaudeControlRequest, { respond, respondWithError })
return
}
if (message.type === 'control_cancel_request' && typeof message.request_id === 'string') {
handlers.onControlCancelRequest?.(message as ClaudeControlCancelRequest)
return
}
handlers.onMessage?.(message)
}
const stdout = attachClaudeStreamJsonStdout({
stdout: child.stdout,
maxLineBytes: STDOUT_LINE_MAX_BYTES,
onMessage: dispatchMessage,
onFailure: (error) => {
void forceTerminateProcessTree(child)
handleUnexpectedEnd(error)
}
})
child.stderr.setEncoding('utf8').on('data', (chunk: string) => {
stderrTail = (stderrTail + chunk).slice(-STDERR_TAIL_MAX_BYTES)
})
child.on('error', (error) => {
markExited()
handleUnexpectedEnd(error)
})
child.on('close', () => {
markExited()
stdout.flush()
handleUnexpectedEnd()
})
child.stdin.on('error', (error) => {
if (!closing) {
void forceTerminateProcessTree(child)
handleUnexpectedEnd(error)
}
})
const request = (
subtype: string,
params: Record<string, unknown> = {},
options: { timeoutMs?: number } = {}
): Promise<unknown> => {
const requestId = `orca-${nextRequestId++}`
const promise = new Promise<unknown>((resolve, reject) => {
const timer = setTimeout(() => {
pending.delete(requestId)
reject(new Error(`claude ${subtype} request timed out`))
}, options.timeoutMs ?? DEFAULT_REQUEST_TIMEOUT_MS)
timer.unref?.()
pending.set(requestId, { subtype, resolve, reject, timer })
})
void writeLine({
type: 'control_request',
request_id: requestId,
request: { subtype, ...params }
}).catch((error) => {
const waiter = pending.get(requestId)
if (waiter) {
pending.delete(requestId)
clearTimeout(waiter.timer)
waiter.reject(error as Error)
}
})
return promise
}
const close = (): Promise<boolean> => {
closePromise ??= (async () => {
closing = true
try {
child.stdin.end()
} catch {
// The reap below still owns the process.
}
await waitForProcessExitUntil(exitPromise, GRACEFUL_EXIT_MS)
const descendantsExited = await forceTerminateProcessTree(child)
await waitForProcessExitUntil(exitPromise, FORCED_EXIT_MS)
failPending(new Error('claude stream-json connection closed'))
return descendantsExited && exited
})()
return closePromise
}
return {
get pid() {
return child.pid
},
get closed() {
return closing || exited || terminalError !== null
},
send: writeLine,
request,
respond,
respondWithError,
close
}
}
@@ -0,0 +1,80 @@
import type { Readable } from 'node:stream'
export type ClaudeStreamJsonStdout = {
flush: () => void
}
function isRecord(value: unknown): value is Record<string, unknown> {
return typeof value === 'object' && value !== null && !Array.isArray(value)
}
export function attachClaudeStreamJsonStdout(input: {
stdout: Readable
maxLineBytes: number
onMessage: (message: Record<string, unknown>) => void
onFailure: (error: Error) => void
}): ClaudeStreamJsonStdout {
let buffer = ''
let stopped = false
const stop = (): void => {
if (stopped) {
return
}
stopped = true
buffer = ''
input.stdout.removeListener('data', onData)
input.stdout.pause()
}
const fail = (error: Error): void => {
stop()
input.onFailure(error)
}
const parseLine = (line: string): void => {
if (!line.trim()) {
return
}
try {
const parsed: unknown = JSON.parse(line)
if (isRecord(parsed)) {
input.onMessage(parsed)
}
} catch (error) {
fail(new Error('claude emitted invalid stream-json', { cause: error }))
}
}
const onData = (chunk: string): void => {
if (stopped) {
return
}
buffer += chunk
let newline = buffer.indexOf('\n')
while (newline !== -1) {
const line = buffer.slice(0, newline)
buffer = buffer.slice(newline + 1)
if (Buffer.byteLength(line, 'utf8') > input.maxLineBytes) {
fail(new Error('claude emitted an oversized stream-json line'))
return
}
parseLine(line)
if (stopped) {
return
}
newline = buffer.indexOf('\n')
}
if (Buffer.byteLength(buffer, 'utf8') > input.maxLineBytes) {
fail(new Error('claude emitted an oversized stream-json line'))
}
}
input.stdout.setEncoding('utf8').on('data', onData)
return {
flush: () => {
if (!stopped && buffer.trim()) {
const line = buffer
buffer = ''
parseLine(line)
}
}
}
}
@@ -0,0 +1,79 @@
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import type { ClaudeControlRequest, ClaudeControlResponder } from './claude-stream-json-connection'
import { resolveClaudeReplayWaiter } from './claude-structured-dispatch'
import {
handleClaudeInboundControl,
handleClaudeInboundControlCancel
} from './claude-structured-inbound-control'
import type { ClaudeInitDeadline } from './claude-structured-init-deadline'
import { readClaudeInit } from './claude-structured-init-proof'
import type { ClaudeStructuredProviderEvents } from './claude-structured-provider-events'
import {
observeClaudeTopLevelLeaf,
type ClaudeAcquisitionAttempt,
type ClaudeSession,
type ClaudeStructuredSessionEvent
} from './claude-structured-session-state'
export class ClaudeStructuredAcquisitionEvents {
private session: ClaudeSession | null = null
private leafUuid: string | null = null
constructor(
private readonly sessionId: string,
private readonly attempt: ClaudeAcquisitionAttempt,
private readonly sink: StructuredAgentSessionEventSink | undefined,
private readonly initDeadline: ClaudeInitDeadline,
private readonly providerEvents: ClaudeStructuredProviderEvents
) {}
observeResumeLeaf(leafUuid: string | null): void {
this.leafUuid = leafUuid
}
publish(session: ClaudeSession): void {
this.session = session
}
observedLeaf(): string | null {
return this.leafUuid
}
emit(event: ClaudeStructuredSessionEvent): void {
this.providerEvents.deliver(this.attempt, this.sessionId, () =>
this.providerEvents.emit(this.session, this.sink, event)
)
}
onMessage = (message: Record<string, unknown>): void => {
const init = readClaudeInit(message)
if (init) {
this.initDeadline.resolve(init)
}
this.leafUuid = observeClaudeTopLevelLeaf(this.leafUuid, message)
if (this.session) {
this.session.leafUuid = this.leafUuid
resolveClaudeReplayWaiter(this.session, message)
}
this.emit({ type: 'message', sessionId: this.sessionId, message })
}
onControlRequest = (request: ClaudeControlRequest, responder?: ClaudeControlResponder): void => {
handleClaudeInboundControl({
sessionId: this.sessionId,
attempt: this.attempt,
request,
responder,
emit: (event) => this.emit(event)
})
}
onControlCancelRequest = ({ request_id: requestId }: { request_id: string }): void => {
handleClaudeInboundControlCancel({
sessionId: this.sessionId,
attempt: this.attempt,
requestId,
emit: (event) => this.emit(event)
})
}
}
@@ -0,0 +1,34 @@
import { applyClaudePromptAnswer } from './claude-structured-prompt-replies'
import { ClaudeControlRequestError } from './claude-stream-json-connection'
import type { ClaudeSession } from './claude-structured-session-state'
export async function cancelClaudeTurn(
session: ClaudeSession,
timeoutMs: number | undefined
): Promise<{ cancelled: boolean }> {
try {
await session.connection.request('interrupt', {}, { timeoutMs })
return { cancelled: true }
} catch (error) {
if (error instanceof ClaudeControlRequestError) {
return { cancelled: false }
}
throw error
}
}
export async function answerClaudePrompt(
session: ClaudeSession,
input: { itemId: string; kind: 'approval' | 'question'; optionId: string }
): Promise<void> {
const found = session.prompts.find(input.itemId)
if (!found || found.prompt.kind !== input.kind) {
throw new Error(`claude is no longer waiting on ${input.itemId}`)
}
const response = applyClaudePromptAnswer(found, input.optionId)
if (response === null) {
return
}
session.prompts.forget(found.prompt)
await session.connection.respond(found.prompt.requestId, response)
}
@@ -0,0 +1,145 @@
import { mkdtemp, rm, writeFile } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { describe, expect, it, vi } from 'vitest'
import type { AgentJournalMessageItem } from '../../shared/agent-session-journal-types'
import { dispatchClaudeTurn, resolveClaudeReplayWaiter } from './claude-structured-dispatch'
import type { ClaudeSession } from './claude-structured-session-state'
function sessionFor(send = vi.fn().mockResolvedValue(undefined)): ClaudeSession {
return {
connection: { send } as unknown as ClaudeSession['connection'],
providerSessionId: 'provider-session',
leafUuid: null,
fence: 1,
prompts: {} as ClaudeSession['prompts'],
dispatchWaiters: [],
options: new Map(),
reportedOptions: {},
events: undefined,
translator: null
}
}
function userMessage(blocks: AgentJournalMessageItem['blocks']): AgentJournalMessageItem {
return { kind: 'message', role: 'user', blocks }
}
describe('Claude structured text dispatch', () => {
it('accepts a slash command from its result receipt when Claude omits the user replay', async () => {
const session = sessionFor()
const dispatched = dispatchClaudeTurn(
session,
{ clientMessageId: 'client-1', body: userMessage([{ type: 'text', text: '/permissions' }]) },
100
)
await vi.waitFor(() => expect(session.dispatchWaiters).toHaveLength(1))
resolveClaudeReplayWaiter(session, {
type: 'result',
subtype: 'success',
session_id: 'provider-session',
uuid: 'command-result-uuid'
})
await expect(dispatched).resolves.toEqual({
state: 'accepted',
providerIdentity: {
provider: 'claude',
sessionId: 'provider-session',
uuid: 'command-result-uuid'
}
})
})
it('does not mistake a normal turn result for its missing user replay', async () => {
const session = sessionFor()
const dispatched = dispatchClaudeTurn(
session,
{ clientMessageId: 'client-1', body: userMessage([{ type: 'text', text: 'hello' }]) },
100
)
await vi.waitFor(() => expect(session.dispatchWaiters).toHaveLength(1))
resolveClaudeReplayWaiter(session, {
type: 'result',
session_id: 'provider-session',
uuid: 'unrelated-result-uuid'
})
expect(session.dispatchWaiters).toHaveLength(1)
resolveClaudeReplayWaiter(session, {
type: 'user',
parent_tool_use_id: null,
session_id: 'provider-session',
uuid: 'user-replay-uuid'
})
await expect(dispatched).resolves.toMatchObject({
state: 'accepted',
providerIdentity: { uuid: 'user-replay-uuid' }
})
})
it('leaves image dispatch explicitly unavailable for slice 2', async () => {
const session = sessionFor()
const body = userMessage(
Array.from({ length: 21 }, (_, index) => ({
type: 'image-ref' as const,
url: `https://example.test/${index}.png`
}))
)
await expect(
dispatchClaudeTurn(session, { clientMessageId: 'client-1', body }, 1)
).resolves.toEqual({
state: 'rejected',
reason: 'Claude structured image dispatch is not available yet'
})
expect(session.connection.send).not.toHaveBeenCalled()
})
it('rejects local images whose aggregate size exceeds twenty MiB', async () => {
const directory = await mkdtemp(join(tmpdir(), 'orca-claude-images-'))
try {
const paths = await Promise.all(
Array.from({ length: 5 }, async (_, index) => {
const path = join(directory, `${index}.png`)
await writeFile(path, Buffer.alloc(5 * 1024 * 1024))
return path
})
)
const session = sessionFor()
const body = userMessage(paths.map((path) => ({ type: 'image-ref' as const, path })))
await expect(
dispatchClaudeTurn(session, { clientMessageId: 'client-1', body }, 1)
).resolves.toEqual({
state: 'rejected',
reason: 'Claude structured image dispatch is not available yet'
})
expect(session.connection.send).not.toHaveBeenCalled()
} finally {
await rm(directory, { recursive: true, force: true })
}
})
it('rejects a local image by actual bytes read beyond the per-image cap', async () => {
const directory = await mkdtemp(join(tmpdir(), 'orca-claude-image-'))
try {
const path = join(directory, 'oversized.png')
await writeFile(path, Buffer.alloc(5 * 1024 * 1024 + 1))
const session = sessionFor()
const body = userMessage([{ type: 'image-ref', path }])
await expect(
dispatchClaudeTurn(session, { clientMessageId: 'client-1', body }, 1)
).resolves.toEqual({
state: 'rejected',
reason: 'Claude structured image dispatch is not available yet'
})
expect(session.connection.send).not.toHaveBeenCalled()
} finally {
await rm(directory, { recursive: true, force: true })
}
})
})
@@ -0,0 +1,108 @@
import type { AgentJournalMessageItem } from '../../shared/agent-session-journal-types'
import type { NativeChatBlock } from '../../shared/native-chat-types'
import type { AgentSessionDispatchOutcome } from '../native-chat/agent-session-wire/structured-agent-session-adapter'
import type { ClaudeSession } from './claude-structured-session-state'
import { readClaudeFrameString } from './claude-structured-init-proof'
export function resolveClaudeReplayWaiter(
session: ClaudeSession,
message: Record<string, unknown>
): void {
const isUserReplay = message.type === 'user' && message.parent_tool_use_id === null
const isCompletedCommand = message.type === 'result'
if (
(!isUserReplay && !isCompletedCommand) ||
readClaudeFrameString(message, 'session_id') !== session.providerSessionId
) {
return
}
const uuid = readClaudeFrameString(message, 'uuid')
const current = session.dispatchWaiters[0]
if (isCompletedCommand && !current?.acceptsResult) {
return
}
const waiter = uuid ? session.dispatchWaiters.shift() : undefined
if (waiter && uuid) {
clearTimeout(waiter.timer)
waiter.resolve(uuid)
}
}
async function messageContent(body: AgentJournalMessageItem): Promise<unknown[]> {
if (body.role !== 'user') {
throw new Error('Claude dispatch accepts only user messages')
}
const content: unknown[] = []
for (const block of body.blocks as NativeChatBlock[]) {
if (block.type === 'text' && block.text.length > 0) {
content.push({ type: 'text', text: block.text })
} else if (block.type === 'image-ref') {
throw new Error('Claude structured image dispatch is not available yet')
}
}
if (content.length === 0) {
throw new Error('Claude dispatch requires text')
}
return content
}
function waitForReplay(
session: ClaudeSession,
timeoutMs: number,
acceptsResult: boolean
): Promise<string | null> {
return new Promise((resolve) => {
const waiter = {
acceptsResult,
resolve,
timer: setTimeout(() => {
const index = session.dispatchWaiters.indexOf(waiter)
if (index !== -1) {
session.dispatchWaiters.splice(index, 1)
}
resolve(null)
}, timeoutMs)
}
waiter.timer.unref?.()
session.dispatchWaiters.push(waiter)
})
}
export async function dispatchClaudeTurn(
session: ClaudeSession,
input: { clientMessageId: string; body: AgentJournalMessageItem },
timeoutMs: number
): Promise<AgentSessionDispatchOutcome> {
let content: unknown[]
try {
content = await messageContent(input.body)
} catch (error) {
return { state: 'rejected', reason: (error as Error).message }
}
const acceptsResult = input.body.blocks.some(
(block) => block.type === 'text' && block.text.trimStart().startsWith('/')
)
const replayed = waitForReplay(session, timeoutMs, acceptsResult)
try {
await session.connection.send({
type: 'user',
message: { role: 'user', content },
parent_tool_use_id: null,
session_id: session.providerSessionId
})
} catch (error) {
const waiter = session.dispatchWaiters.shift()
if (waiter) {
clearTimeout(waiter.timer)
waiter.resolve(null)
}
return { state: 'unknown', reason: (error as Error).message }
}
const uuid = await replayed
return uuid
? {
state: 'accepted',
providerIdentity: { provider: 'claude', sessionId: session.providerSessionId, uuid }
}
: { state: 'unknown', reason: 'claude accepted a message but did not replay its uuid in time' }
}
@@ -0,0 +1,105 @@
import { describe, expect, it, vi } from 'vitest'
import { createClaudeAcquisitionAttempt } from './claude-structured-session-state'
import { ClaudePromptRegistry } from './claude-structured-prompt-replies'
import {
CLAUDE_BLOCKING_CONTROL_REQUEST_SUBTYPES,
CLAUDE_CAN_USE_TOOL_SUBTYPE,
CLAUDE_REQUEST_USER_DIALOG_SUBTYPE,
handleClaudeInboundControl
} from './claude-structured-inbound-control'
function rejectingAttempt() {
const attempt = createClaudeAcquisitionAttempt(new ClaudePromptRegistry())
const respond = vi.fn().mockRejectedValue(new Error('connection closed'))
const respondWithError = vi.fn().mockRejectedValue(new Error('connection closed'))
attempt.connection = { respond, respondWithError } as never
return { attempt, respond, respondWithError }
}
describe('Claude inbound control', () => {
it('routes a valid permission request to the durable prompt registry', () => {
const control = rejectingAttempt()
const emit = vi.fn()
expect(
handleClaudeInboundControl({
sessionId: 'session-1',
attempt: control.attempt,
request: {
type: 'control_request',
request_id: 'permission-1',
request: {
subtype: CLAUDE_CAN_USE_TOOL_SUBTYPE,
tool_use_id: 'tool-1',
tool_name: 'Bash',
input: { command: 'git status' }
}
},
emit
})
).toEqual({ kind: 'prompt' })
expect(control.respond).not.toHaveBeenCalled()
expect(control.respondWithError).not.toHaveBeenCalled()
expect(emit).toHaveBeenCalledWith(expect.objectContaining({ type: 'prompt' }))
})
it('denies a malformed permission request with a protocol response', () => {
const control = rejectingAttempt()
expect(
handleClaudeInboundControl({
sessionId: 'session-1',
attempt: control.attempt,
request: {
type: 'control_request',
request_id: 'permission-2',
request: { subtype: CLAUDE_CAN_USE_TOOL_SUBTYPE, tool_use_id: 'tool-2' }
},
emit: vi.fn()
})
).toEqual({ kind: 'responded', subtype: CLAUDE_CAN_USE_TOOL_SUBTYPE })
expect(control.respond).toHaveBeenCalledWith('permission-2', {
behavior: 'deny',
message: 'Orca could not decode this permission request.',
toolUseID: 'tool-2'
})
})
it('consumes rejected writes while declining unsupported controls', async () => {
const dialog = rejectingAttempt()
handleClaudeInboundControl({
sessionId: 'session-1',
attempt: dialog.attempt,
request: {
type: 'control_request',
request_id: 'dialog-1',
request: { subtype: 'request_user_dialog' }
},
emit: vi.fn()
})
const unsupported = rejectingAttempt()
handleClaudeInboundControl({
sessionId: 'session-1',
attempt: unsupported.attempt,
request: {
type: 'control_request',
request_id: 'control-1',
request: { subtype: 'future_control' }
},
emit: vi.fn()
})
await new Promise((resolve) => setImmediate(resolve))
expect(dialog.respond).toHaveBeenCalledWith('dialog-1', { behavior: 'cancelled' })
expect(unsupported.respondWithError).toHaveBeenCalledWith(
'control-1',
'Orca does not handle future_control'
)
})
it('enumerates every blocking control request in the stable stream-json surface', () => {
expect(new Set(CLAUDE_BLOCKING_CONTROL_REQUEST_SUBTYPES)).toEqual(
new Set([CLAUDE_CAN_USE_TOOL_SUBTYPE, CLAUDE_REQUEST_USER_DIALOG_SUBTYPE])
)
})
})
@@ -0,0 +1,90 @@
import type { ClaudeControlRequest, ClaudeControlResponder } from './claude-stream-json-connection'
import type {
ClaudeAcquisitionAttempt,
ClaudeStructuredSessionEvent
} from './claude-structured-session-state'
export const CLAUDE_CAN_USE_TOOL_SUBTYPE = 'can_use_tool'
export const CLAUDE_REQUEST_USER_DIALOG_SUBTYPE = 'request_user_dialog'
export const CLAUDE_BLOCKING_CONTROL_REQUEST_SUBTYPES = [
CLAUDE_CAN_USE_TOOL_SUBTYPE,
CLAUDE_REQUEST_USER_DIALOG_SUBTYPE
] as const
export type ClaudeInboundControlDisposition =
| { kind: 'prompt' }
| { kind: 'responded'; subtype: string }
export function handleClaudeInboundControlCancel(input: {
sessionId: string
attempt: ClaudeAcquisitionAttempt
requestId: string
emit: (event: ClaudeStructuredSessionEvent) => void
}): void {
const prompt = input.attempt.prompts.cancel(input.requestId)
input.emit(
prompt
? {
type: 'prompt-cancelled',
sessionId: input.sessionId,
promptKey: prompt.promptKey
}
: {
type: 'provider-frame',
sessionId: input.sessionId,
kind: 'control_cancel_request',
payload: { request_id: input.requestId }
}
)
}
/** Every Claude control request becomes a durable prompt or receives a safe reply. */
export function handleClaudeInboundControl(input: {
sessionId: string
attempt: ClaudeAcquisitionAttempt
request: ClaudeControlRequest
responder?: ClaudeControlResponder
emit: (event: ClaudeStructuredSessionEvent) => void
}): ClaudeInboundControlDisposition {
const responder = input.responder ?? input.attempt.connection
const subtype = input.request.request.subtype
if (subtype === CLAUDE_REQUEST_USER_DIALOG_SUBTYPE) {
void responder?.respond(input.request.request_id, { behavior: 'cancelled' }).catch(() => {})
input.emit({
type: 'provider-frame',
sessionId: input.sessionId,
kind: `control_request:${subtype}`,
payload: input.request.request
})
return { kind: 'responded', subtype }
}
const prompt = input.attempt.prompts.register(input.request)
if (prompt) {
input.emit({ type: 'prompt', sessionId: input.sessionId, prompt })
return { kind: 'prompt' }
}
if (subtype === CLAUDE_CAN_USE_TOOL_SUBTYPE) {
const toolUseId =
typeof input.request.request.tool_use_id === 'string'
? input.request.request.tool_use_id
: undefined
void responder
?.respond(input.request.request_id, {
behavior: 'deny',
message: 'Orca could not decode this permission request.',
...(toolUseId ? { toolUseID: toolUseId } : {})
})
.catch(() => {})
} else {
void responder
?.respondWithError(input.request.request_id, `Orca does not handle ${subtype}`)
.catch(() => {})
}
input.emit({
type: 'provider-frame',
sessionId: input.sessionId,
kind: `control_request:${subtype}`,
payload: input.request.request
})
return { kind: 'responded', subtype }
}
@@ -0,0 +1,72 @@
import type { ClaudeInitObservation } from './claude-structured-init-proof'
import { claudeInitializationAuthError } from './claude-structured-init-proof'
import type { ClaudeStreamJsonConnection } from './claude-stream-json-connection'
import { AgentSessionAcquisitionRefusal } from '../native-chat/agent-session-wire/structured-agent-session-adapter'
export type ClaudeInitDeadline = {
promise: Promise<ClaudeInitObservation>
resolve: (init: ClaudeInitObservation) => void
reject: (error: Error) => void
start: () => void
clear: () => void
}
export function claudeInitTimeoutError(
sessionId: string,
timeoutMs: number
): AgentSessionAcquisitionRefusal {
return new AgentSessionAcquisitionRefusal(
`Claude did not finish starting session ${sessionId} within ${Math.ceil(timeoutMs / 1000)} seconds. Verify the selected Claude account is signed in and CLAUDE_CONFIG_DIR contains valid credentials, then retry; no SessionStart or system/init proof arrived.`
)
}
export async function requestClaudeInitialization(
connection: ClaudeStreamJsonConnection,
sessionId: string,
timeoutMs: number
): Promise<unknown> {
try {
const result = await connection.request(
'initialize',
{ supportedDialogKinds: [] },
{ timeoutMs }
)
const authError = claudeInitializationAuthError(result)
if (authError) {
throw authError
}
return result
} catch (error) {
if (error instanceof Error && error.message === 'claude initialize request timed out') {
throw claudeInitTimeoutError(sessionId, timeoutMs)
}
throw error
}
}
export function createClaudeInitDeadline(sessionId: string, timeoutMs: number): ClaudeInitDeadline {
let resolve = (_init: ClaudeInitObservation): void => {}
let reject = (_error: Error): void => {}
const promise = new Promise<ClaudeInitObservation>((resolvePromise, rejectPromise) => {
resolve = resolvePromise
reject = rejectPromise
})
void promise.catch(() => {})
let timer: ReturnType<typeof setTimeout> | null = null
return {
promise,
resolve,
reject,
start: () => {
timer = setTimeout(() => reject(claudeInitTimeoutError(sessionId, timeoutMs)), timeoutMs)
timer.unref?.()
},
clear: () => {
if (timer) {
clearTimeout(timer)
timer = null
}
}
}
}
@@ -0,0 +1,74 @@
import { CLAUDE_DEFAULT_SETTING_SOURCES } from './claude-structured-launch-resolution'
import type { ClaudeAuthDiagnostic } from './claude-structured-session-state'
import { AgentSessionAcquisitionRefusal } from '../native-chat/agent-session-wire/structured-agent-session-adapter'
export type ClaudeInitObservation = {
providerSessionId: string
uuid: string | null
message: Record<string, unknown>
}
function isRecord(value: unknown): value is Record<string, unknown> {
return typeof value === 'object' && value !== null && !Array.isArray(value)
}
export function readClaudeFrameString(source: Record<string, unknown>, key: string): string | null {
const value = source[key]
return typeof value === 'string' && value.length > 0 ? value : null
}
export function readClaudeInit(message: Record<string, unknown>): ClaudeInitObservation | null {
const hookName = readClaudeFrameString(message, 'hook_name')
const isInit = message.type === 'system' && message.subtype === 'init'
const isSessionStart =
message.type === 'system' &&
(message.subtype === 'hook_started' || message.subtype === 'hook_response') &&
hookName?.startsWith('SessionStart:') === true
if (!isInit && !isSessionStart) {
return null
}
const providerSessionId = readClaudeFrameString(message, 'session_id')
return providerSessionId
? {
providerSessionId,
uuid: isInit ? readClaudeFrameString(message, 'uuid') : null,
message
}
: null
}
export function readClaudeModels(initialization: unknown): unknown[] {
return isRecord(initialization) && Array.isArray(initialization.models)
? initialization.models
: []
}
export function claudeInitializationAuthError(
initialization: unknown
): AgentSessionAcquisitionRefusal | null {
const account =
isRecord(initialization) && isRecord(initialization.account) ? initialization.account : null
return readClaudeFrameString(account ?? {}, 'tokenSource') === 'none'
? new AgentSessionAcquisitionRefusal(
'Claude is not signed in for the selected account. Sign in with the Claude CLI for this CLAUDE_CONFIG_DIR, then retry.'
)
: null
}
export function claudeAuthDiagnostic(
init: ClaudeInitObservation,
settings: unknown
): ClaudeAuthDiagnostic {
const env = isRecord(settings) && isRecord(settings.env) ? settings.env : {}
const apiKeySource = readClaudeFrameString(init.message, 'apiKeySource')
const configured = (key: string): boolean =>
(typeof env[key] === 'string' && (env[key] as string).trim().length > 0) ||
Boolean(process.env[key]?.trim())
return {
apiKeySourceConfigured: apiKeySource !== null && apiKeySource !== 'none',
baseUrlConfigured: configured('ANTHROPIC_BASE_URL'),
authTokenConfigured: configured('ANTHROPIC_AUTH_TOKEN'),
apiKeyConfigured: configured('ANTHROPIC_API_KEY'),
settingSources: CLAUDE_DEFAULT_SETTING_SOURCES
}
}
@@ -0,0 +1,158 @@
import type {
AgentJournalItemBody,
AgentJournalItemIdentity,
AgentJournalMessageItem
} from '../../shared/agent-session-journal-types'
import type { NativeChatBlock } from '../../shared/native-chat-types'
import {
boundInlineText,
DEFAULT_JOURNAL_PAYLOAD_LIMITS
} from '../native-chat/agent-session-journal/journal-payload-bounds'
export type ClaudeMessageEnvelope = {
sessionId: string
uuid: string
role: 'assistant' | 'user'
content: unknown[]
}
export type ClaudeToolUse = { id: string; name: string; input: unknown }
export type ClaudeToolResult = { toolUseId: string; output: string; failed: boolean }
export function claudeRecord(value: unknown): Record<string, unknown> | null {
return typeof value === 'object' && value !== null && !Array.isArray(value)
? (value as Record<string, unknown>)
: null
}
export function claudeText(value: unknown): string | null {
return typeof value === 'string' && value.length > 0 ? value : null
}
export function readClaudeMessageEnvelope(
frame: Record<string, unknown>
): ClaudeMessageEnvelope | null {
if (frame.type !== 'assistant' && frame.type !== 'user') {
return null
}
const message = claudeRecord(frame.message)
const sessionId = claudeText(frame.session_id)
const uuid = claudeText(frame.uuid)
const role = message?.role
return sessionId && uuid && (role === 'assistant' || role === 'user')
? {
sessionId,
uuid,
role,
content: Array.isArray(message?.content) ? message.content : []
}
: null
}
export function claudeMessageIdentity(
envelope: Pick<ClaudeMessageEnvelope, 'sessionId' | 'uuid'>
): AgentJournalItemIdentity {
return { provider: 'claude', sessionId: envelope.sessionId, uuid: envelope.uuid }
}
function messageBlocks(envelope: ClaudeMessageEnvelope): NativeChatBlock[] {
const blocks: NativeChatBlock[] = []
for (const value of envelope.content) {
const part = claudeRecord(value)
const text = claudeText(part?.text)
if (part?.type === 'text' && text) {
blocks.push({ type: 'text', text })
continue
}
const source = claudeRecord(part?.source)
const url = claudeText(source?.url)
if (part?.type === 'image' && source?.type === 'url' && url) {
blocks.push({ type: 'image-ref', url })
}
}
return blocks
}
export function claudeMessageBody(envelope: ClaudeMessageEnvelope): AgentJournalMessageItem | null {
const blocks = messageBlocks(envelope)
return blocks.length > 0 ? { kind: 'message', role: envelope.role, blocks } : null
}
export function claudeToolUses(envelope: ClaudeMessageEnvelope): ClaudeToolUse[] {
return envelope.content.flatMap((value) => {
const part = claudeRecord(value)
const id = claudeText(part?.id)
const name = claudeText(part?.name)
return part?.type === 'tool_use' && id && name ? [{ id, name, input: part.input ?? null }] : []
})
}
function resultText(value: unknown): string {
if (typeof value === 'string') {
return value
}
if (!Array.isArray(value)) {
return value === undefined ? '' : JSON.stringify(value)
}
return value
.flatMap((entry) => {
if (typeof entry === 'string') {
return [entry]
}
const part = claudeRecord(entry)
return part?.type === 'text' && typeof part.text === 'string' ? [part.text] : []
})
.join('\n')
}
export function claudeToolResults(envelope: ClaudeMessageEnvelope): ClaudeToolResult[] {
return envelope.content.flatMap((value) => {
const part = claudeRecord(value)
const toolUseId = claudeText(part?.tool_use_id)
return part?.type === 'tool_result' && toolUseId
? [
{
toolUseId,
output: resultText(part.content),
failed: part.is_error === true
}
]
: []
})
}
export function claudeThinkingText(envelope: ClaudeMessageEnvelope): string | null {
const parts = envelope.content.flatMap((value) => {
const part = claudeRecord(value)
const thinking = claudeText(part?.thinking)
return part?.type === 'thinking' && thinking ? [thinking] : []
})
return parts.length > 0 ? parts.join('\n') : null
}
export function claudeToolBody(input: {
tool: ClaudeToolUse
result?: ClaudeToolResult
}): AgentJournalItemBody {
return {
kind: 'tool-call',
name: input.tool.name,
input: input.tool.input,
state: input.result ? (input.result.failed ? 'failed' : 'completed') : 'running',
...(input.result
? { output: boundInlineText(input.result.output, DEFAULT_JOURNAL_PAYLOAD_LIMITS).bounded }
: {})
}
}
export function claudeStreamingMessageBody(text: string): AgentJournalMessageItem {
return { kind: 'message', role: 'assistant', blocks: [{ type: 'text', text }] }
}
export function claudeToolIdentity(sessionId: string, toolUseId: string): AgentJournalItemIdentity {
return { provider: 'orca', clientMessageId: `claude-tool:${sessionId}:${toolUseId}` }
}
export function claudeThinkingIdentity(sessionId: string, uuid: string): AgentJournalItemIdentity {
return { provider: 'orca', clientMessageId: `claude-thinking:${sessionId}:${uuid}` }
}
@@ -0,0 +1,349 @@
import { describe, expect, it, vi } from 'vitest'
import type {
AgentJournalItemBody,
AgentJournalItemIdentity
} from '../../shared/agent-session-journal-types'
import { agentJournalItemKey } from '../../shared/agent-session-journal-item-key'
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import {
boundInlineText,
DEFAULT_JOURNAL_PAYLOAD_LIMITS
} from '../native-chat/agent-session-journal/journal-payload-bounds'
import type { ClaudePendingPrompt } from './claude-structured-prompt-replies'
import { createClaudeJournalTranslator } from './claude-structured-journal-translation'
function sinkState() {
const items: { identity: AgentJournalItemIdentity; body: AgentJournalItemBody }[] = []
const tombstones: AgentJournalItemIdentity[] = []
const sink: StructuredAgentSessionEventSink = {
appendItem: (identity, body) => items.push({ identity, body }),
appendTombstone: (identity) => tombstones.push(identity),
publish: vi.fn()
}
return { sink, items, tombstones }
}
function message(
type: 'assistant' | 'user',
uuid: string,
content: unknown[],
parentToolUseId: string | null = null
) {
return {
type: 'message' as const,
sessionId: 'orca-session',
message: {
type,
uuid,
session_id: 'claude-session',
parent_tool_use_id: parentToolUseId,
message: { role: type, content }
}
}
}
describe('Claude structured journal translation', () => {
it('reuses the shared coalescer and finalizes the provider-keyed message row', () => {
const state = sinkState()
let scheduled: (() => void) | null = null
const translator = createClaudeJournalTranslator({
sink: state.sink,
schedule: (run, delay) => {
expect(delay).toBe(60)
scheduled = run
return () => {
scheduled = null
}
}
})
translator.handle({
type: 'message',
sessionId: 'orca-session',
message: {
type: 'stream_event',
uuid: 'assistant-1',
session_id: 'claude-session',
event: { type: 'content_block_delta', delta: { type: 'text_delta', text: 'Hel' } }
}
})
translator.handle({
type: 'message',
sessionId: 'orca-session',
message: {
type: 'stream_event',
uuid: 'assistant-1',
session_id: 'claude-session',
event: { type: 'content_block_delta', delta: { type: 'text_delta', text: 'lo' } }
}
})
expect(state.items).toEqual([])
const run = scheduled as (() => void) | null
run?.()
expect(state.items.at(-1)).toEqual({
identity: { provider: 'claude', sessionId: 'claude-session', uuid: 'assistant-1' },
body: { kind: 'message', role: 'assistant', blocks: [{ type: 'text', text: 'Hello' }] }
})
translator.handle(message('assistant', 'assistant-1', [{ type: 'text', text: 'Hello!' }]))
expect(state.items.at(-1)).toEqual({
identity: { provider: 'claude', sessionId: 'claude-session', uuid: 'assistant-1' },
body: { kind: 'message', role: 'assistant', blocks: [{ type: 'text', text: 'Hello!' }] }
})
})
it('journals turn lifecycle and updates one tool row through its result', () => {
const state = sinkState()
const translator = createClaudeJournalTranslator({ sink: state.sink })
translator.handle(message('user', 'user-1', [{ type: 'text', text: 'List files' }]))
translator.handle(
message('assistant', 'assistant-tool', [
{ type: 'tool_use', id: 'tool-1', name: 'Bash', input: { command: 'ls' } }
])
)
translator.handle(
message(
'user',
'tool-result-1',
[{ type: 'tool_result', tool_use_id: 'tool-1', content: 'a.ts\nb.ts' }],
'tool-1'
)
)
const keyed = new Map(
state.items.map((item) => [agentJournalItemKey(item.identity), item.body])
)
expect(keyed.get('claude:claude-session:user-1')).toMatchObject({
kind: 'message',
role: 'user'
})
expect(keyed.get('orca:claude-tool%3Aclaude-session%3Atool-1')).toMatchObject({
kind: 'tool-call',
name: 'Bash',
state: 'completed',
output: { head: 'a.ts\nb.ts', truncated: false }
})
expect(
state.items.some(
(item) => item.body.kind === 'status' && item.body.turnLifecycle?.turnId === 'user-1'
)
).toBe(true)
translator.handle({
type: 'message',
sessionId: 'orca-session',
message: { type: 'result', session_id: 'claude-session', uuid: 'result-1' }
})
expect(state.tombstones.at(-1)).toMatchObject({
provider: 'legacy',
agent: 'claude',
recordId: 'turn-lifecycle:user-1'
})
})
it('bounds persisted thinking text to the shared journal payload limit', () => {
const state = sinkState()
const translator = createClaudeJournalTranslator({ sink: state.sink })
const thinking = 'considering '.repeat(20_000)
translator.handle(message('assistant', 'assistant-thinking', [{ type: 'thinking', thinking }]))
expect(state.items.at(-1)?.body).toEqual({
kind: 'status',
text: boundInlineText(thinking, DEFAULT_JOURNAL_PAYLOAD_LIMITS).text
})
})
it('starts a cancellable lifecycle for image-only root user replays', () => {
const state = sinkState()
const translator = createClaudeJournalTranslator({ sink: state.sink })
translator.handle(
message('user', 'user-image', [
{ type: 'image', source: { type: 'base64', media_type: 'image/png', data: 'AA==' } }
])
)
expect(state.items.at(-1)?.body).toEqual({
kind: 'status',
text: 'Claude is working…',
turnLifecycle: { turnId: 'user-image', state: 'running' }
})
})
it('renders empty user frames through the provider fallback', () => {
const state = sinkState()
const translator = createClaudeJournalTranslator({ sink: state.sink })
translator.handle(message('user', 'control-only', []))
expect(state.items).toEqual(
expect.arrayContaining([
expect.objectContaining({
body: expect.objectContaining({
kind: 'status',
providerFrame: expect.objectContaining({ kind: 'message:user:empty' })
})
})
])
)
})
it('renders every unmodeled Claude frame family as a bounded provider row', () => {
const state = sinkState()
const translator = createClaudeJournalTranslator({ sink: state.sink })
translator.handle({
type: 'message',
sessionId: 'orca-session',
message: { type: 'system', subtype: 'compact_boundary', summary: 'x'.repeat(100_000) }
})
translator.handle({
type: 'message',
sessionId: 'orca-session',
message: { type: 'system', subtype: 'hook_response', hook_name: 'PostToolUse' }
})
translator.handle({
type: 'message',
sessionId: 'orca-session',
message: { type: 'system', subtype: 'command_started', command: '/compact' }
})
translator.handle({
type: 'message',
sessionId: 'orca-session',
message: { type: 'result', usage: { input_tokens: 12 }, total_cost_usd: 0.01 }
})
translator.handle({
type: 'message',
sessionId: 'orca-session',
message: { type: 'tool_progress', tool_use_id: 'tool-1', elapsed_time_seconds: 2 }
})
translator.handle({
type: 'message',
sessionId: 'orca-session',
message: { type: 'prompt_suggestion', suggestion: '/compact' }
})
translator.handle(
message('user', 'attachment-1', [
{ type: 'document', source: { type: 'base64', media_type: 'application/pdf' } }
])
)
translator.handle({
type: 'provider-frame',
sessionId: 'orca-session',
kind: 'control_request:future_control',
payload: { subtype: 'future_control' }
})
const frames = state.items.flatMap((item) =>
item.body.kind === 'status' && item.body.providerFrame ? [item.body.providerFrame] : []
)
expect(frames.map((frame) => frame.kind)).toEqual([
'message:system:command_started',
'message:result',
'message:user:content:document',
'control_request:future_control'
])
})
it('preserves grouped questions as addressable items and cancels them durably', () => {
const state = sinkState()
const bindings: unknown[][] = []
const translator = createClaudeJournalTranslator({
sink: state.sink,
bindPromptItemId: (...args) => bindings.push(args)
})
const approval = prompt({
requestId: 'permission-1',
promptKey: 'permission-1',
toolUseId: 'tool-1',
toolName: 'Bash',
kind: 'approval',
input: { command: 'git status' },
questionIds: []
})
translator.handle({ type: 'prompt', sessionId: 'orca-session', prompt: approval })
expect(state.items.at(-1)?.body).toMatchObject({
kind: 'approval',
title: 'Allow Bash?',
options: expect.arrayContaining([{ id: 'allow', label: 'Allow' }])
})
expect(bindings[0]).toEqual([
'orca:claude-prompt%3Aorca-session%3Apermission-1',
'permission-1'
])
const questions = prompt({
requestId: 'questions-1',
promptKey: 'questions-1',
toolUseId: 'tool-q',
toolName: 'AskUserQuestion',
kind: 'question',
input: {
questions: [
{ question: 'Library?', options: [{ label: 'Luxon' }] },
{ question: 'Ship?', options: [{ label: 'Yes' }] }
]
},
questionIds: ['Library?', 'Ship?']
})
translator.handle({ type: 'prompt', sessionId: 'orca-session', prompt: questions })
expect(state.items.filter((item) => item.body.kind === 'question')).toHaveLength(2)
expect(state.items.at(-1)?.body).toMatchObject({
kind: 'question',
question: 'Ship?'
})
expect(bindings.at(-1)).toEqual([
'orca:claude-prompt%3Aorca-session%3Aquestions-1%3Aq2',
'questions-1'
])
const multiSelect = prompt({
requestId: 'questions-multi',
promptKey: 'questions-multi',
toolUseId: 'tool-multi',
toolName: 'AskUserQuestion',
kind: 'question',
input: {
questions: [
{
question: 'Libraries?',
multiSelect: true,
options: [{ label: 'Luxon' }, { label: 'Temporal' }]
}
]
},
questionIds: ['Libraries?']
})
translator.handle({ type: 'prompt', sessionId: 'orca-session', prompt: multiSelect })
expect(state.items.at(-1)?.body).toMatchObject({
kind: 'question',
question: 'Libraries?',
options: [{ label: 'Luxon' }, { label: 'Temporal' }],
freeTextQuestionId: 'q1'
})
translator.handle({
type: 'prompt-cancelled',
sessionId: 'orca-session',
promptKey: 'questions-1'
})
expect(state.tombstones).toHaveLength(2)
})
})
function prompt(
input: Pick<
ClaudePendingPrompt,
'requestId' | 'promptKey' | 'toolUseId' | 'toolName' | 'kind' | 'input' | 'questionIds'
>
): ClaudePendingPrompt {
return {
...input,
suggestions: [],
answers: new Map(),
request: { subtype: 'can_use_tool' }
}
}
@@ -0,0 +1,301 @@
import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types'
import { agentJournalItemKey } from '../../shared/agent-session-journal-item-key'
import {
createAgentSessionDeltaCoalescer,
type AgentSessionDeltaCoalescerDeps
} from '../native-chat/agent-session-wire/agent-session-delta-coalescer'
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import {
boundInlineText,
DEFAULT_JOURNAL_PAYLOAD_LIMITS
} from '../native-chat/agent-session-journal/journal-payload-bounds'
import type { ClaudeStructuredSessionEvent } from './claude-structured-session-state'
import {
claudeMessageBody,
claudeMessageIdentity,
claudeRecord,
claudeStreamingMessageBody,
claudeText,
claudeThinkingIdentity,
claudeThinkingText,
claudeToolBody,
claudeToolIdentity,
claudeToolResults,
claudeToolUses,
readClaudeMessageEnvelope,
type ClaudeToolUse
} from './claude-structured-item-translation'
import {
claudeApprovalItem,
claudePromptIdentity,
claudeQuestionItems
} from './claude-structured-prompt-items'
import type { ClaudePromptRegistry } from './claude-structured-prompt-replies'
import {
claudeProviderFrameKind,
createClaudeProviderFrameFallback,
isModeledClaudeContent
} from './claude-structured-provider-fallback'
export type ClaudeJournalTranslatorDeps = {
sink: StructuredAgentSessionEventSink
bindPromptItemId?: (journalItemId: string, promptKey: string, questionId?: string) => void
coalesceMs?: number
schedule?: AgentSessionDeltaCoalescerDeps['schedule']
}
export type ClaudeJournalTranslator = {
handle: (event: ClaudeStructuredSessionEvent) => void
flush: () => void
dispose: () => void
}
export function createClaudeSessionJournalTranslator(
sink: StructuredAgentSessionEventSink | undefined,
prompts: ClaudePromptRegistry
): ClaudeJournalTranslator | null {
return sink
? createClaudeJournalTranslator({
sink,
bindPromptItemId: (itemId, promptKey, questionId) =>
prompts.bindJournalItemId(itemId, promptKey, questionId)
})
: null
}
function lifecycleIdentity(sessionId: string, turnId: string): AgentJournalItemIdentity {
return {
provider: 'legacy',
agent: 'claude',
sessionId,
recordId: `turn-lifecycle:${turnId}`
}
}
function streamDelta(message: Record<string, unknown>): string | null {
if (message.type !== 'stream_event') {
return null
}
const event = claudeRecord(message.event)
if (event?.type === 'content_block_start') {
return claudeText(claudeRecord(event.content_block)?.text)
}
if (event?.type !== 'content_block_delta') {
return null
}
const delta = claudeRecord(event.delta)
return delta?.type === 'text_delta' ? claudeText(delta.text) : null
}
export function createClaudeJournalTranslator(
deps: ClaudeJournalTranslatorDeps
): ClaudeJournalTranslator {
const tools = new Map<string, ClaudeToolUse>()
const promptItems = new Map<string, AgentJournalItemIdentity[]>()
const streamIdentities = new Map<string, AgentJournalItemIdentity>()
const latestStreamText = new Map<string, string>()
const checkpointLengths = new Map<string, number>()
let currentTurn: { sessionId: string; turnId: string } | null = null
const providerFallback = createClaudeProviderFrameFallback(deps.sink)
const publishLifecycle = (sessionId: string, turnId: string, running: boolean): void => {
const identity = lifecycleIdentity(sessionId, turnId)
if (running) {
deps.sink.appendItem(identity, {
kind: 'status',
text: 'Claude is working…',
turnLifecycle: { turnId, state: 'running' }
})
} else {
deps.sink.appendTombstone(identity)
}
deps.sink.publish()
}
const persistStream = (key: string, text: string, force: boolean): void => {
latestStreamText.set(key, text)
const checkpointLength = checkpointLengths.get(key) ?? 0
const nextLength = Math.max(checkpointLength + 32, Math.ceil(checkpointLength * 1.125))
if (!force && checkpointLength > 0 && text.length < nextLength) {
return
}
const identity = streamIdentities.get(key)
if (!identity) {
return
}
checkpointLengths.set(key, text.length)
deps.sink.appendItem(identity, claudeStreamingMessageBody(text))
deps.sink.publish()
}
const coalescer = createAgentSessionDeltaCoalescer({
windowMs: deps.coalesceMs,
schedule: deps.schedule,
emit: (key, text) => persistStream(key, text, false)
})
const flushStreams = (): void => {
coalescer.flushAll()
for (const [key, text] of latestStreamText) {
if (checkpointLengths.get(key) !== text.length) {
persistStream(key, text, true)
}
}
}
const handleStream = (message: Record<string, unknown>): boolean => {
const delta = streamDelta(message)
const sessionId = claudeText(message.session_id)
const uuid = claudeText(message.uuid)
if (!delta || !sessionId || !uuid) {
return false
}
const key = agentJournalItemKey({ provider: 'claude', sessionId, uuid })
streamIdentities.set(key, { provider: 'claude', sessionId, uuid })
coalescer.append(key, delta)
return true
}
const handleMessage = (message: Record<string, unknown>): boolean => {
const envelope = readClaudeMessageEnvelope(message)
if (!envelope) {
return false
}
const identity = claudeMessageIdentity(envelope)
const key = agentJournalItemKey(identity)
coalescer.forget(key)
latestStreamText.delete(key)
checkpointLengths.delete(key)
streamIdentities.delete(key)
let changed = false
const body = claudeMessageBody(envelope)
if (body) {
deps.sink.appendItem(identity, body)
changed = true
}
for (const tool of claudeToolUses(envelope)) {
tools.set(tool.id, tool)
deps.sink.appendItem(
claudeToolIdentity(envelope.sessionId, tool.id),
claudeToolBody({ tool })
)
changed = true
}
for (const result of claudeToolResults(envelope)) {
const tool = tools.get(result.toolUseId) ?? {
id: result.toolUseId,
name: 'tool',
input: null
}
deps.sink.appendItem(
claudeToolIdentity(envelope.sessionId, result.toolUseId),
claudeToolBody({ tool, result })
)
changed = true
}
const thinking = claudeThinkingText(envelope)
if (thinking) {
deps.sink.appendItem(claudeThinkingIdentity(envelope.sessionId, envelope.uuid), {
kind: 'status',
text: boundInlineText(thinking, DEFAULT_JOURNAL_PAYLOAD_LIMITS).text
})
changed = true
}
const unhandledContent = envelope.content.filter((part) => !isModeledClaudeContent(part))
for (const part of unhandledContent) {
const partType = claudeText(claudeRecord(part)?.type) ?? 'unknown'
providerFallback.append(`message:${envelope.role}:content:${partType}`, part)
changed = true
}
if (envelope.content.length === 0) {
providerFallback.append(`message:${envelope.role}:empty`, message)
changed = true
}
if (
envelope.role === 'user' &&
envelope.content.length > 0 &&
message.parent_tool_use_id === null
) {
if (currentTurn) {
publishLifecycle(currentTurn.sessionId, currentTurn.turnId, false)
}
currentTurn = { sessionId: envelope.sessionId, turnId: envelope.uuid }
publishLifecycle(envelope.sessionId, envelope.uuid, true)
}
if (changed) {
deps.sink.publish()
}
return true
}
const handlePrompt = (event: Extract<ClaudeStructuredSessionEvent, { type: 'prompt' }>): void => {
const identities: AgentJournalItemIdentity[] = []
if (event.prompt.kind === 'question') {
for (const question of claudeQuestionItems({
sessionId: event.sessionId,
prompt: event.prompt
})) {
identities.push(question.identity)
deps.sink.appendItem(question.identity, question.body)
deps.bindPromptItemId?.(agentJournalItemKey(question.identity), event.prompt.promptKey)
}
} else {
const identity = claudePromptIdentity({
sessionId: event.sessionId,
promptKey: event.prompt.promptKey
})
identities.push(identity)
deps.sink.appendItem(identity, claudeApprovalItem(event.prompt))
deps.bindPromptItemId?.(agentJournalItemKey(identity), event.prompt.promptKey)
}
promptItems.set(event.prompt.promptKey, identities)
deps.sink.publish()
}
return {
handle: (event) => {
if (event.type === 'ended') {
flushStreams()
if (currentTurn) {
publishLifecycle(currentTurn.sessionId, currentTurn.turnId, false)
currentTurn = null
}
return
}
if (event.type === 'message' && handleStream(event.message)) {
return
}
flushStreams()
if (event.type === 'prompt') {
handlePrompt(event)
} else if (event.type === 'prompt-cancelled') {
for (const identity of promptItems.get(event.promptKey) ?? []) {
deps.sink.appendTombstone(identity)
}
promptItems.delete(event.promptKey)
deps.sink.publish()
} else if (event.type === 'message' && event.message.type === 'result') {
if (currentTurn) {
publishLifecycle(currentTurn.sessionId, currentTurn.turnId, false)
currentTurn = null
}
providerFallback.append(claudeProviderFrameKind(event.message), event.message)
} else if (event.type === 'message') {
if (!handleMessage(event.message)) {
providerFallback.append(claudeProviderFrameKind(event.message), event.message)
}
} else if (event.type === 'provider-frame') {
providerFallback.append(event.kind, event.payload)
}
},
flush: flushStreams,
dispose: () => {
coalescer.dispose()
tools.clear()
promptItems.clear()
streamIdentities.clear()
latestStreamText.clear()
checkpointLengths.clear()
}
}
}
@@ -0,0 +1,130 @@
import { describe, expect, it } from 'vitest'
import type { AgentSessionRecord } from '../../shared/agent-session-record'
import { LOCAL_EXECUTION_HOST_ID } from '../../shared/execution-host'
import type { AgentSessionRecordStore } from '../runtime/agent-session-record-store'
import {
CLAUDE_DEFAULT_SETTING_SOURCES,
CLAUDE_STRUCTURED_BASE_ARGS,
claudeSessionIdForOrcaSession,
createClaudeStructuredLaunchResolver
} from './claude-structured-launch-resolution'
const SESSION_ID = 'orca-session-1'
const IDENTITY = { sessionId: SESSION_ID } as Parameters<
ReturnType<typeof createClaudeStructuredLaunchResolver>
>[0]['identity']
function record(overrides: Partial<AgentSessionRecord> = {}): AgentSessionRecord {
return {
sessionId: SESSION_ID,
provider: 'claude',
location: {
executionHostId: LOCAL_EXECUTION_HOST_ID,
wslDistro: null,
workspaceId: 'workspace-1',
workspaceKind: 'folder'
},
accountHome: { variable: 'CLAUDE_CONFIG_DIR', path: '/home/work/.claude' },
providerHandleChain: [],
...overrides
} as AgentSessionRecord
}
function resolverFor(value: AgentSessionRecord | null, resolveEnv?: () => Record<string, string>) {
return createClaudeStructuredLaunchResolver({
store: { getRecord: () => value } as unknown as AgentSessionRecordStore,
resolveWorkspacePath: async (id) => `/repos/${id}`,
resolveCommand: () => '/usr/local/bin/claude',
...(resolveEnv ? { resolveEnvironment: async () => resolveEnv() } : {})
})
}
describe('claude structured launch resolution', () => {
it('pre-mints a stable provider id and pins interactive setting sources', async () => {
const first = await resolverFor(record())({ identity: IDENTITY })
const second = await resolverFor(record())({ identity: IDENTITY })
expect(first.providerSessionId).toBe(claudeSessionIdForOrcaSession(SESSION_ID))
expect(second.providerSessionId).toBe(first.providerSessionId)
expect(first).toMatchObject({
command: '/usr/local/bin/claude',
cwd: '/repos/workspace-1',
claudeConfigDir: '/home/work/.claude',
resumeLeafUuid: null,
resumed: false
})
expect(first.args).toContain('--session-id')
expect(first.args).toContain(first.providerSessionId)
expect(first.args).toContain('--permission-prompt-tool')
expect(first.args).toContain('stdio')
expect(first.args).toContain('--setting-sources')
expect(first.args).toContain(CLAUDE_DEFAULT_SETTING_SOURCES.join(','))
expect(CLAUDE_STRUCTURED_BASE_ARGS).toContain('--verbose')
})
it('resumes the session and leaf at the durable chain head', async () => {
const launch = await resolverFor(
record({
providerHandleChain: [
{ handle: { provider: 'claude', sessionId: 'provider-old', leafUuid: 'leaf-old' } },
{
handle: {
provider: 'claude',
sessionId: 'provider-current',
leafUuid: 'leaf-current'
}
}
] as AgentSessionRecord['providerHandleChain']
})
)({ identity: IDENTITY })
expect(launch).toMatchObject({
providerSessionId: 'provider-current',
resumeLeafUuid: 'leaf-current',
resumed: true
})
expect(launch.args.slice(-2)).toEqual(['--resume', 'provider-current'])
})
it('uses the runtime environment instead of the scrubbed legacy launchEnv', async () => {
const pinned = record({
launchEnv: {
ANTHROPIC_AUTH_TOKEN: 'first-token',
ANTHROPIC_BASE_URL: 'https://gateway.example.test'
}
})
const resolver = resolverFor(pinned, () => ({
ANTHROPIC_AUTH_TOKEN: 'rotated-token',
ANTHROPIC_BASE_URL: 'https://gateway.example.test'
}))
expect((await resolver({ identity: IDENTITY })).env).toEqual({
ANTHROPIC_AUTH_TOKEN: 'rotated-token',
ANTHROPIC_BASE_URL: 'https://gateway.example.test'
})
expect((await resolver({ identity: IDENTITY })).env?.ANTHROPIC_AUTH_TOKEN).toBe('rotated-token')
})
it('refuses other hosts, WSL, providers, and account-home variables', async () => {
await expect(
resolverFor(record({ location: { ...record().location, executionHostId: 'ssh:build' } }))({
identity: IDENTITY
})
).rejects.toThrow(/local host/)
await expect(
resolverFor(record({ location: { ...record().location, wslDistro: 'Ubuntu' } }))({
identity: IDENTITY
})
).rejects.toThrow(/local host/)
await expect(
resolverFor(record({ provider: 'codex' } as Partial<AgentSessionRecord>))({
identity: IDENTITY
})
).rejects.toThrow(/codex session/)
await expect(
resolverFor(record({ accountHome: { variable: 'CODEX_HOME', path: '/tmp/codex' } }))({
identity: IDENTITY
})
).rejects.toThrow(/CLAUDE_CONFIG_DIR/)
})
})
@@ -0,0 +1,104 @@
import { createHash } from 'node:crypto'
import type { AgentSessionJournalIdentity } from '../../shared/agent-session-journal-types'
import { agentSessionProviderHandleChainHead } from '../../shared/agent-session-provider-handle'
import { LOCAL_EXECUTION_HOST_ID } from '../../shared/execution-host'
import { resolveClaudeCommand } from '../codex-cli/command'
import type { AgentSessionRecordStore } from '../runtime/agent-session-record-store'
export const CLAUDE_DEFAULT_SETTING_SOURCES = ['user', 'project', 'local'] as const
export const CLAUDE_STRUCTURED_BASE_ARGS = [
'-p',
'--input-format',
'stream-json',
'--output-format',
'stream-json',
'--include-partial-messages',
'--verbose',
'--replay-user-messages',
'--permission-prompt-tool',
'stdio',
'--setting-sources',
CLAUDE_DEFAULT_SETTING_SOURCES.join(',')
]
export type ClaudeStructuredLaunch = {
command: string
args: string[]
cwd: string
env?: Record<string, string>
claudeConfigDir: string
providerSessionId: string
resumeLeafUuid: string | null
resumed: boolean
}
export type ClaudeStructuredLaunchResolverDeps = {
store: AgentSessionRecordStore
resolveWorkspacePath: (workspaceId: string) => Promise<string>
resolveCommand?: (options?: { pathEnv?: string | null; homePath?: string }) => string
resolveEnvironment?: () => Promise<NodeJS.ProcessEnv>
platform?: NodeJS.Platform
}
export function claudeSessionIdForOrcaSession(sessionId: string): string {
const bytes = createHash('sha256').update(`orca-claude:${sessionId}`).digest().subarray(0, 16)
bytes[6] = ((bytes[6] ?? 0) & 0x0f) | 0x40
bytes[8] = ((bytes[8] ?? 0) & 0x3f) | 0x80
const hex = bytes.toString('hex')
return `${hex.slice(0, 8)}-${hex.slice(8, 12)}-${hex.slice(12, 16)}-${hex.slice(16, 20)}-${hex.slice(20)}`
}
export function createClaudeStructuredLaunchResolver(
deps: ClaudeStructuredLaunchResolverDeps
): (input: { identity: AgentSessionJournalIdentity }) => Promise<ClaudeStructuredLaunch> {
return async ({ identity }) => {
const record = deps.store.getRecord(identity.sessionId)
if (!record) {
throw new Error(`no durable agent-session record for ${identity.sessionId}`)
}
if (record.provider !== 'claude') {
throw new Error(`session ${identity.sessionId} is a ${record.provider} session`)
}
if (
record.location.executionHostId !== LOCAL_EXECUTION_HOST_ID ||
record.location.wslDistro !== null
) {
throw new Error(
`claude structured sessions run on the local host, not ${record.location.executionHostId}`
)
}
if ((deps.platform ?? process.platform) === 'win32') {
throw new Error('claude structured sessions are unavailable on win32')
}
if (record.accountHome.variable !== 'CLAUDE_CONFIG_DIR') {
throw new Error(`claude sessions pin CLAUDE_CONFIG_DIR, not ${record.accountHome.variable}`)
}
const head = agentSessionProviderHandleChainHead(record.providerHandleChain)
const providerSessionId =
head?.handle.provider === 'claude'
? head.handle.sessionId
: claudeSessionIdForOrcaSession(identity.sessionId)
const providerArgs =
head?.handle.provider === 'claude'
? ['--resume', providerSessionId]
: ['--session-id', providerSessionId]
const environment = await deps.resolveEnvironment?.()
const pathEnv = environment?.PATH ?? environment?.Path ?? null
const homePath = environment?.HOME ?? environment?.USERPROFILE
const command = (deps.resolveCommand ?? resolveClaudeCommand)({
pathEnv,
...(homePath ? { homePath } : {})
})
return {
command,
args: [...CLAUDE_STRUCTURED_BASE_ARGS, ...providerArgs],
cwd: await deps.resolveWorkspacePath(record.location.workspaceId),
...(environment ? { env: { ...environment } as Record<string, string> } : {}),
claudeConfigDir: record.accountHome.path,
providerSessionId,
resumeLeafUuid: head?.handle.provider === 'claude' ? head.handle.leafUuid : null,
resumed: head?.handle.provider === 'claude'
}
}
}
@@ -0,0 +1,57 @@
import { ClaudeControlRequestError } from './claude-stream-json-connection'
import { AgentSessionOptionRejectedError } from '../native-chat/agent-session-wire/structured-agent-session-option-error'
import type { ClaudeSession } from './claude-structured-session-state'
const OPTION_ORDER = ['model', 'effort', 'permissionMode'] as const
export function restoredClaudeStructuredSessionOptions(
options: Readonly<Record<string, string>> | undefined
): Map<string, string> {
return new Map(
OPTION_ORDER.flatMap((key) => {
const value = options?.[key]
return value ? [[key, value] as const] : []
})
)
}
export async function setClaudeStructuredOption(
session: ClaudeSession,
input: { key: string; value: string },
timeoutMs: number | undefined
): Promise<Readonly<Record<string, string>>> {
const request =
input.key === 'model'
? { subtype: 'set_model', params: { model: input.value } }
: input.key === 'permissionMode'
? { subtype: 'set_permission_mode', params: { mode: input.value } }
: input.key === 'effort'
? { subtype: 'apply_flag_settings', params: { settings: { effortLevel: input.value } } }
: null
if (!request) {
throw new AgentSessionOptionRejectedError(
`claude stream-json has no session option named ${input.key}`
)
}
try {
await session.connection.request(request.subtype, request.params, { timeoutMs })
} catch (error) {
if (error instanceof ClaudeControlRequestError) {
throw new AgentSessionOptionRejectedError(error)
}
throw error
}
session.options.set(input.key, input.value)
return Object.fromEntries(session.options)
}
export async function restoreClaudeStructuredSessionOptions(
session: ClaudeSession,
timeoutMs: number | undefined
): Promise<void> {
const options = [...session.options.entries()]
session.options.clear()
for (const [key, value] of options) {
await setClaudeStructuredOption(session, { key, value }, timeoutMs)
}
}
@@ -0,0 +1,46 @@
import { describe, expect, it, vi } from 'vitest'
import { claudeProcessIdentity } from './claude-structured-owner-identity'
const IDENTITY = {
sessionId: 'session-identity',
workspaceId: 'workspace-1',
hostId: 'local',
agent: 'claude' as const,
providerHandle: { kind: 'claude' as const, sessionId: 'provider-1', leafUuid: null }
}
describe('claude process identity', () => {
it('records start time and spawn token', async () => {
await expect(
claudeProcessIdentity(
{ identity: IDENTITY, spawnToken: 'spawn-a', pid: 4242 },
async () => 123
)
).resolves.toEqual({
hostId: 'local',
pid: 4242,
processStartTimeMs: 123,
spawnToken: 'spawn-a'
})
})
it('retries unreadable start times before refusing', async () => {
const readStartTime = vi
.fn<(pid: number) => Promise<number | null>>()
.mockResolvedValueOnce(null)
.mockResolvedValueOnce(null)
.mockResolvedValueOnce(456)
await expect(
claudeProcessIdentity({ identity: IDENTITY, spawnToken: 'spawn-a', pid: 4242 }, readStartTime)
).resolves.toMatchObject({ processStartTimeMs: 456 })
expect(readStartTime).toHaveBeenCalledTimes(3)
})
it('refuses when start time remains unreadable', async () => {
const readStartTime = vi.fn(async () => null)
await expect(
claudeProcessIdentity({ identity: IDENTITY, spawnToken: 'spawn-a', pid: 4242 }, readStartTime)
).rejects.toThrow('start time')
expect(readStartTime).toHaveBeenCalledTimes(3)
})
})
@@ -1,4 +1,40 @@
import type { AgentSessionJournalIdentity } from '../../shared/agent-session-journal-types'
import type { AgentSessionProviderHandleLink } from '../../shared/agent-session-provider-handle'
import type { AgentSessionProcessIdentity } from '../../shared/agent-session-record'
import { readProcessStartTimeMs } from '../runtime/agent-session-process-identity-probe'
export const CLAUDE_SPAWN_TOKEN_ENV = 'ORCA_AGENT_SESSION_SPAWN_TOKEN'
const START_TIME_READ_ATTEMPTS = 3
export async function claudeProcessIdentity(
input: {
identity: AgentSessionJournalIdentity
spawnToken: string
pid: number | undefined
},
readStartTime: (pid: number) => Promise<number | null> = readProcessStartTimeMs
): Promise<AgentSessionProcessIdentity> {
if (input.pid === undefined) {
throw new Error('claude stream-json started without a pid')
}
let processStartTimeMs: number | null = null
for (
let attempt = 0;
attempt < START_TIME_READ_ATTEMPTS && processStartTimeMs === null;
attempt += 1
) {
processStartTimeMs = await readStartTime(input.pid)
}
if (processStartTimeMs === null) {
throw new Error(`claude stream-json start time for pid ${input.pid} could not be read`)
}
return {
hostId: input.identity.hostId,
pid: input.pid,
processStartTimeMs,
spawnToken: input.spawnToken
}
}
export function claudeProviderHandleLink(input: {
sessionId: string
@@ -0,0 +1,114 @@
import type {
AgentJournalApprovalItem,
AgentJournalItemIdentity,
AgentJournalPromptOption,
AgentJournalQuestionItem
} from '../../shared/agent-session-journal-types'
import {
boundInlineText,
DEFAULT_JOURNAL_PAYLOAD_LIMITS
} from '../native-chat/agent-session-journal/journal-payload-bounds'
import { claudeRecord, claudeText } from './claude-structured-item-translation'
import {
CLAUDE_APPROVAL_DECISIONS,
encodeClaudeQuestionOptionId,
type ClaudeApprovalDecision,
type ClaudePendingPrompt
} from './claude-structured-prompt-replies'
const APPROVAL_LABELS: Record<ClaudeApprovalDecision, string> = {
allow: 'Allow',
allowForSession: 'Allow for this session',
deny: 'Deny',
cancel: 'Stop'
}
const PENDING = {
state: 'pending',
selectedOptionId: null,
resolvedBy: null,
resolvedAt: null
} as const
export function claudePromptIdentity(input: {
sessionId: string
promptKey: string
questionId?: string
}): AgentJournalItemIdentity {
const suffix = input.questionId ? `:${input.questionId}` : ''
return {
provider: 'orca',
clientMessageId: `claude-prompt:${input.sessionId}:${input.promptKey}${suffix}`
}
}
export function claudeApprovalItem(prompt: ClaudePendingPrompt): AgentJournalApprovalItem {
const serialized = JSON.stringify(prompt.input)
return {
kind: 'approval',
title: `Allow ${prompt.toolName}?`,
detail: serialized ? boundInlineText(serialized, DEFAULT_JOURNAL_PAYLOAD_LIMITS).text : null,
options: CLAUDE_APPROVAL_DECISIONS.map((decision) => ({
id: decision,
label: APPROVAL_LABELS[decision]
})),
resolution: { ...PENDING }
}
}
export type ClaudeQuestionItem = {
identity: AgentJournalItemIdentity
body: AgentJournalQuestionItem
}
function questionOptions(
question: Record<string, unknown>,
questionAddress: string
): AgentJournalPromptOption[] {
if (!Array.isArray(question.options)) {
return []
}
return question.options.flatMap((value, index) => {
const option = claudeRecord(value)
const label = claudeText(option?.label)
return label
? [
{
id: encodeClaudeQuestionOptionId(questionAddress, `choice-${index + 1}`),
label
}
]
: []
})
}
export function claudeQuestionItems(input: {
sessionId: string
prompt: ClaudePendingPrompt
}): ClaudeQuestionItem[] {
const values = Array.isArray(input.prompt.input.questions) ? input.prompt.input.questions : []
return values.flatMap((value, index): ClaudeQuestionItem[] => {
const question = claudeRecord(value)
const questionAddress = `q${index + 1}`
const text = claudeText(question?.question) ?? claudeText(question?.header)
const header = claudeText(question?.header)
return question && input.prompt.questionIds[index] && text
? [
{
identity: claudePromptIdentity({
sessionId: input.sessionId,
promptKey: input.prompt.promptKey,
questionId: questionAddress
}),
body: {
kind: 'question',
question: header ? `${header}: ${text}` : text,
options: questionOptions(question, questionAddress),
freeTextQuestionId: questionAddress,
resolution: { ...PENDING }
}
}
]
: []
})
}
@@ -0,0 +1,233 @@
import type { ClaudeControlRequest } from './claude-stream-json-connection'
export const CLAUDE_APPROVAL_DECISIONS = ['allow', 'allowForSession', 'deny', 'cancel'] as const
export type ClaudeApprovalDecision = (typeof CLAUDE_APPROVAL_DECISIONS)[number]
export type ClaudePendingPrompt = {
requestId: string
promptKey: string
toolUseId: string
toolName: string
kind: 'approval' | 'question'
input: Record<string, unknown>
suggestions: unknown[]
questionIds: readonly string[]
answers: Map<string, string | readonly string[]>
request: ClaudeControlRequest['request']
}
type PromptBinding = {
address: string
questionId?: string
}
function isRecord(value: unknown): value is Record<string, unknown> {
return typeof value === 'object' && value !== null && !Array.isArray(value)
}
function readString(value: unknown): string | null {
return typeof value === 'string' && value.trim().length > 0 ? value : null
}
function questionsFrom(input: Record<string, unknown>): Record<string, unknown>[] {
return Array.isArray(input.questions) ? input.questions.filter(isRecord) : []
}
function questionIdFromAddress(prompt: ClaudePendingPrompt, address: string): string | null {
const match = /^q([1-9]\d*)$/.exec(address)
const index = match ? Number(match[1]) - 1 : -1
return index >= 0 ? (prompt.questionIds[index] ?? null) : null
}
function questionAnswer(prompt: ClaudePendingPrompt, questionId: string, optionId: string): string {
const decoded = decodeClaudeQuestionOptionId(optionId)
if (!decoded) {
return optionId
}
const questionIndex = prompt.questionIds.indexOf(questionId)
if (questionIndex === -1) {
return optionId
}
const choice = /^choice-([1-9]\d*)$/.exec(decoded.answer)
const optionIndex = choice ? Number(choice[1]) - 1 : -1
const question = questionsFrom(prompt.input)[questionIndex]
const options = Array.isArray(question?.options) ? question.options : []
const option = options[optionIndex]
const label = isRecord(option) ? readString(option.label) : null
if (decoded.questionId === `q${questionIndex + 1}` && label) {
return label
}
if (decoded.questionId === `q${questionIndex + 1}`) {
return decoded.answer
}
const legacyChoice = options.some(
(candidate) => isRecord(candidate) && readString(candidate.label) === decoded.answer
)
return decoded.questionId === questionId && (legacyChoice || decoded.answer.trim().length > 0)
? decoded.answer
: optionId
}
function questionId(question: Record<string, unknown>, index: number): string {
return readString(question.question) ?? readString(question.header) ?? `question-${index + 1}`
}
export function encodeClaudeQuestionOptionId(questionId: string, answer: string): string {
return `${encodeURIComponent(questionId)}:${encodeURIComponent(answer)}`
}
export function decodeClaudeQuestionOptionId(
optionId: string
): { questionId: string; answer: string } | null {
const separator = optionId.indexOf(':')
if (separator <= 0) {
return null
}
try {
return {
questionId: decodeURIComponent(optionId.slice(0, separator)),
answer: decodeURIComponent(optionId.slice(separator + 1))
}
} catch {
return null
}
}
export class ClaudePromptRegistry {
private readonly prompts = new Map<string, ClaudePendingPrompt>()
private readonly journalBindings = new Map<string, PromptBinding>()
register(control: ClaudeControlRequest): ClaudePendingPrompt | null {
if (control.request.subtype !== 'can_use_tool') {
return null
}
const toolUseId = readString(control.request.tool_use_id)
const toolName = readString(control.request.tool_name)
const input = isRecord(control.request.input) ? control.request.input : null
if (!toolUseId || !toolName || !input) {
return null
}
const questions = toolName === 'AskUserQuestion' ? questionsFrom(input) : []
const prompt: ClaudePendingPrompt = {
requestId: control.request_id,
promptKey: control.request_id,
toolUseId,
toolName,
kind: questions.length > 0 ? 'question' : 'approval',
input,
suggestions: Array.isArray(control.request.permission_suggestions)
? control.request.permission_suggestions
: [],
questionIds: questions.map(questionId),
answers: new Map(),
request: control.request
}
this.prompts.set(prompt.promptKey, prompt)
return prompt
}
bindJournalItemId(journalItemId: string, promptKey: string, questionIdForItem?: string): void {
this.journalBindings.set(journalItemId, {
address: promptKey,
...(questionIdForItem ? { questionId: questionIdForItem } : {})
})
}
find(itemId: string): { prompt: ClaudePendingPrompt; questionId?: string } | null {
const binding = this.journalBindings.get(itemId)
const prompt = this.prompts.get(binding?.address ?? itemId)
return prompt
? { prompt, ...(binding?.questionId ? { questionId: binding.questionId } : {}) }
: null
}
cancel(requestId: string): ClaudePendingPrompt | null {
const prompt = this.prompts.get(requestId) ?? null
if (prompt) {
this.forget(prompt)
}
return prompt
}
forget(prompt: ClaudePendingPrompt): void {
this.prompts.delete(prompt.promptKey)
for (const [itemId, binding] of this.journalBindings) {
if (binding.address === prompt.promptKey) {
this.journalBindings.delete(itemId)
}
}
}
clear(): ClaudePendingPrompt[] {
const pending = [...this.prompts.values()]
this.prompts.clear()
this.journalBindings.clear()
return pending
}
}
function approvalResponse(prompt: ClaudePendingPrompt, optionId: string): Record<string, unknown> {
if (!(CLAUDE_APPROVAL_DECISIONS as readonly string[]).includes(optionId)) {
throw new Error(`${optionId} is not a Claude approval decision`)
}
const decision = optionId as ClaudeApprovalDecision
if (decision === 'allow' || decision === 'allowForSession') {
return {
behavior: 'allow',
updatedInput: prompt.input,
...(decision === 'allowForSession' && prompt.suggestions.length > 0
? { updatedPermissions: prompt.suggestions }
: {}),
toolUseID: prompt.toolUseId
}
}
return {
behavior: 'deny',
message: decision === 'cancel' ? 'User stopped this turn.' : 'User denied this action.',
...(decision === 'cancel' ? { interrupt: true } : {}),
toolUseID: prompt.toolUseId
}
}
function questionResponse(
prompt: ClaudePendingPrompt,
optionId: string,
boundQuestionId?: string
): Record<string, unknown> | null {
const decoded = decodeClaudeQuestionOptionId(optionId)
const decodedQuestionId = decoded
? (questionIdFromAddress(prompt, decoded.questionId) ??
(prompt.questionIds.includes(decoded.questionId) ? decoded.questionId : null))
: null
const selectedQuestionId =
boundQuestionId ??
decodedQuestionId ??
(prompt.questionIds.length === 1 ? prompt.questionIds[0] : null)
if (!selectedQuestionId || !prompt.questionIds.includes(selectedQuestionId)) {
throw new Error(`${optionId} does not name a question on Claude prompt ${prompt.promptKey}`)
}
const answer = questionAnswer(prompt, selectedQuestionId, optionId)
prompt.answers.set(selectedQuestionId, answer)
if (prompt.questionIds.some((id) => !prompt.answers.has(id))) {
return null
}
const answers: Record<string, string | readonly string[]> = {}
for (const id of prompt.questionIds) {
answers[id] = prompt.answers.get(id) as string
}
return {
behavior: 'allow',
updatedInput: { ...prompt.input, answers },
toolUseID: prompt.toolUseId
}
}
export function applyClaudePromptAnswer(
found: { prompt: ClaudePendingPrompt; questionId?: string },
optionId: string
): Record<string, unknown> | null {
if (found.prompt.kind === 'approval') {
return approvalResponse(found.prompt, optionId)
}
return questionResponse(found.prompt, optionId, found.questionId)
}
@@ -0,0 +1,43 @@
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import { settleClaudeExitedSession } from './claude-structured-session-close'
import type {
ClaudeAcquisitionAttempt,
ClaudeSession,
ClaudeStructuredSessionEvent
} from './claude-structured-session-state'
export class ClaudeStructuredProviderEvents {
constructor(
private readonly sessions: Map<string, ClaudeSession>,
private readonly onEvent?: (event: ClaudeStructuredSessionEvent) => void
) {}
deliver(attempt: ClaudeAcquisitionAttempt, sessionId: string, event: () => void): void {
if (!attempt.published) {
attempt.buffered.push(event)
return
}
if (this.sessions.get(sessionId)?.connection === attempt.connection) {
event()
}
}
handleExit(sessionId: string, attempt: ClaudeAcquisitionAttempt, error: Error): void {
const session = this.sessions.get(sessionId)
if (!session || session.connection !== attempt.connection) {
return
}
this.sessions.delete(sessionId)
this.emit(session, session.events, { type: 'ended', sessionId, reason: error.message })
settleClaudeExitedSession(session)
}
emit(
session: ClaudeSession | null,
_events: StructuredAgentSessionEventSink | undefined,
event: ClaudeStructuredSessionEvent
): void {
session?.translator?.handle(event)
this.onEvent?.(event)
}
}
@@ -0,0 +1,52 @@
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import { unhandledProviderFrameJournalItem } from '../native-chat/agent-session-wire/unhandled-provider-frame'
import { claudeRecord, claudeText } from './claude-structured-item-translation'
export function claudeProviderFrameKind(message: Record<string, unknown>): string {
const type = claudeText(message.type) ?? 'unknown'
const subtype = claudeText(message.subtype)
const eventType = claudeText(claudeRecord(message.event)?.type)
return ['message', type, subtype ?? eventType].filter(Boolean).join(':')
}
export function isModeledClaudeContent(value: unknown): boolean {
const part = claudeRecord(value)
if (!part) {
return false
}
if (part.type === 'text') {
return claudeText(part.text) !== null
}
if (part.type === 'image') {
const source = claudeRecord(part.source)
return source?.type === 'url' && claudeText(source.url) !== null
}
if (part.type === 'tool_use') {
return claudeText(part.id) !== null && claudeText(part.name) !== null
}
if (part.type === 'tool_result') {
return claudeText(part.tool_use_id) !== null
}
return part.type === 'thinking' && claudeText(part.thinking) !== null
}
export function createClaudeProviderFrameFallback(sink: StructuredAgentSessionEventSink): {
append: (kind: string, payload: unknown) => void
} {
let sequence = 0
return {
append: (kind, payload) => {
const translated = unhandledProviderFrameJournalItem('claude', kind, payload)
if (!translated) {
return
}
sequence += 1
sink.appendItem(
{ provider: 'orca', clientMessageId: `provider-frame:claude:${sequence}` },
translated.body,
translated.blobs
)
sink.publish()
}
}
}
@@ -0,0 +1,667 @@
import { describe, expect, it } from 'vitest'
import type {
AgentJournalMessageItem,
AgentSessionJournalIdentity
} from '../../shared/agent-session-journal-types'
import { AgentSessionAcquisitionRefusal } from '../native-chat/agent-session-wire/structured-agent-session-adapter'
import type {
ClaudeStreamJsonConnection,
ClaudeStreamJsonConnectionHandlers,
ClaudeStreamJsonLaunch,
openClaudeStreamJsonConnection
} from './claude-stream-json-connection'
import { ClaudeControlRequestError } from './claude-stream-json-connection'
import { CLAUDE_SPAWN_TOKEN_ENV } from './claude-structured-owner-identity'
import { encodeClaudeQuestionOptionId } from './claude-structured-prompt-replies'
import {
CLAUDE_STRUCTURED_INIT_TIMEOUT_MS,
ClaudeStructuredSessionAdapter,
type ClaudeStructuredLaunch,
type ClaudeStructuredSessionEvent
} from './claude-structured-session-adapter'
const PROVIDER_SESSION_ID = '819cf9f8-e43c-4ad7-b50f-54aa158a726a'
const USER_MESSAGE: AgentJournalMessageItem = {
kind: 'message',
role: 'user',
blocks: [{ type: 'text', text: 'ship it' }]
}
function identityFor(sessionId = 'session-1'): AgentSessionJournalIdentity {
return {
sessionId,
workspaceId: 'workspace-1',
hostId: 'host-1',
agent: 'claude',
providerHandle: { kind: 'claude', sessionId: PROVIDER_SESSION_ID, leafUuid: null }
}
}
type Route = (params: Record<string, unknown> | undefined) => unknown
type FakeConnection = Omit<ClaudeStreamJsonConnection, 'closed'> & {
closed: boolean
launch: ClaudeStreamJsonLaunch
handlers: ClaudeStreamJsonConnectionHandlers
calls: { subtype: string; params?: Record<string, unknown> }[]
sent: Record<string, unknown>[]
replies: { requestId: string; response?: unknown; error?: string }[]
closeCount: number
}
function fakeClaude(
options: {
initSessionId?: string
initUuid?: string
initModel?: string
initEffort?: string
initProof?: 'init' | 'session-start' | 'none'
initAccount?: unknown
exitBeforeInit?: string
settings?: unknown
replayUuid?: string | null
routes?: Record<string, Route>
} = {}
): {
connections: FakeConnection[]
openConnection: typeof openClaudeStreamJsonConnection
routes: Record<string, Route>
} {
const connections: FakeConnection[] = []
const routes = options.routes ?? {}
const openConnection = (async (launch, handlers = {}) => {
const connection: FakeConnection = {
launch,
handlers,
calls: [],
sent: [],
replies: [],
closeCount: 0,
pid: 4321,
closed: false,
request: async (subtype, params) => {
connection.calls.push({ subtype, params })
if (subtype === 'initialize') {
if (options.exitBeforeInit) {
handlers.onExit?.(new Error(options.exitBeforeInit))
return { models: [] }
}
if (options.initProof === 'session-start') {
handlers.onMessage?.({
type: 'system',
subtype: 'hook_started',
hook_name: 'SessionStart:startup',
session_id: options.initSessionId ?? PROVIDER_SESSION_ID,
uuid: options.initUuid ?? 'init-uuid'
})
} else if (options.initProof !== 'none') {
handlers.onMessage?.({
type: 'system',
subtype: 'init',
session_id: options.initSessionId ?? PROVIDER_SESSION_ID,
uuid: options.initUuid ?? 'init-uuid',
model: options.initModel ?? 'claude-sonnet-5',
effortLevel: options.initEffort ?? 'high',
apiKeySource: 'none'
})
}
return {
models: [{ value: 'claude-sonnet', displayName: 'Sonnet' }],
...(options.initAccount === undefined ? {} : { account: options.initAccount })
}
}
if (subtype === 'get_settings') {
return options.settings ?? { env: {} }
}
const route = routes[subtype]
return route ? route(params) : {}
},
send: async (message) => {
connection.sent.push(message)
if (message.type === 'user' && options.replayUuid !== null) {
handlers.onMessage?.({
...message,
uuid: options.replayUuid ?? 'user-uuid'
})
}
},
respond: async (requestId, response) => {
connection.replies.push({ requestId, response })
},
respondWithError: async (requestId, error) => {
connection.replies.push({ requestId, error })
},
close: async () => {
connection.closeCount += 1
connection.closed = true
return true
}
}
connections.push(connection)
return connection
}) as typeof openClaudeStreamJsonConnection
return { connections, openConnection, routes }
}
function adapterFor(
claude: ReturnType<typeof fakeClaude>,
launch: Partial<ClaudeStructuredLaunch> = {},
events: ClaudeStructuredSessionEvent[] = [],
persistedHandles: unknown[] = [],
initTimeoutMs?: number
): ClaudeStructuredSessionAdapter {
return new ClaudeStructuredSessionAdapter({
resolveLaunch: async () => ({
command: 'claude',
args: ['-p'],
cwd: '/work/repo',
claudeConfigDir: '/accounts/claude',
providerSessionId: PROVIDER_SESSION_ID,
resumeLeafUuid: null,
resumed: false,
...launch
}),
onEvent: (event) => events.push(event),
openConnection: claude.openConnection,
readProcessStartTime: async () => 1_700_000_000_000,
now: () => 1_700_000_000_500,
...(initTimeoutMs === undefined ? {} : { initTimeoutMs }),
dispatchAckTimeoutMs: 10,
persistHandle: async (handle) => {
persistedHandles.push(handle)
}
})
}
async function acquired(
claude: ReturnType<typeof fakeClaude>,
launch: Partial<ClaudeStructuredLaunch> = {},
events: ClaudeStructuredSessionEvent[] = []
): Promise<ClaudeStructuredSessionAdapter> {
const adapter = adapterFor(claude, launch, events)
await adapter.acquire({ identity: identityFor(), fence: 7, spawnToken: 'spawn-9' })
return adapter
}
describe('ClaudeStructuredSessionAdapter.acquire', () => {
it('finishes its startup deadline before the paired mobile request deadline', () => {
expect(CLAUDE_STRUCTURED_INIT_TIMEOUT_MS).toBeLessThan(30_000)
})
it('pins the account, proves init, and reports the process and chain leaf', async () => {
const claude = fakeClaude()
const events: ClaudeStructuredSessionEvent[] = []
const adapter = adapterFor(claude, {}, events)
const acquisition = await adapter.acquire({
identity: identityFor(),
fence: 7,
spawnToken: 'spawn-9'
})
expect(claude.connections[0].launch).toMatchObject({
cwd: '/work/repo',
env: {
[CLAUDE_SPAWN_TOKEN_ENV]: 'spawn-9',
CLAUDE_CONFIG_DIR: '/accounts/claude'
}
})
expect(claude.connections[0].calls.slice(0, 2)).toEqual([
{ subtype: 'initialize', params: { supportedDialogKinds: [] } },
{ subtype: 'get_settings', params: {} }
])
expect(acquisition.process).toEqual({
hostId: 'host-1',
pid: 4321,
processStartTimeMs: 1_700_000_000_000,
spawnToken: 'spawn-9'
})
expect(acquisition.link).toMatchObject({
handle: { provider: 'claude', sessionId: PROVIDER_SESSION_ID, leafUuid: null },
origin: 'created',
mintedAtFence: 7,
observedAt: 1_700_000_000_500
})
expect(events[0]).toMatchObject({ type: 'message', message: { subtype: 'init' } })
})
it('restores persisted model and effort before publishing a reacquired session', async () => {
const claude = fakeClaude()
const adapter = adapterFor(claude, { resumed: true })
await adapter.acquire({
identity: identityFor(),
fence: 7,
spawnToken: 'spawn-9',
options: { model: 'opus', effort: 'high' }
})
expect(claude.connections[0].calls.slice(-2)).toEqual([
{ subtype: 'set_model', params: { model: 'opus' } },
{ subtype: 'apply_flag_settings', params: { settings: { effortLevel: 'high' } } }
])
await expect(adapter.readOptions({ sessionId: 'session-1', fence: 7 })).resolves.toMatchObject({
current: { model: 'opus', effort: 'high' }
})
})
it('forwards configured launch environment while keeping ownership pins authoritative', async () => {
const claude = fakeClaude()
const adapter = adapterFor(claude, {
env: {
ANTHROPIC_AUTH_TOKEN: 'configured-token',
ANTHROPIC_BASE_URL: 'https://gateway.example.test',
CLAUDE_CONFIG_DIR: '/wrong/account',
[CLAUDE_SPAWN_TOKEN_ENV]: 'wrong-token'
}
})
await adapter.acquire({ identity: identityFor(), fence: 7, spawnToken: 'spawn-9' })
expect(claude.connections[0].launch.env).toEqual({
ANTHROPIC_AUTH_TOKEN: 'configured-token',
ANTHROPIC_BASE_URL: 'https://gateway.example.test',
CLAUDE_CONFIG_DIR: '/accounts/claude',
[CLAUDE_SPAWN_TOKEN_ENV]: 'spawn-9'
})
})
it('accepts SessionStart as the pre-turn session proof from the real CLI protocol', async () => {
const claude = fakeClaude({ initProof: 'session-start', initUuid: 'session-start-uuid' })
const events: ClaudeStructuredSessionEvent[] = []
const adapter = adapterFor(claude, {}, events)
const acquisition = await adapter.acquire({
identity: identityFor(),
fence: 7,
spawnToken: 'spawn-9'
})
expect(acquisition.link.handle).toEqual({
provider: 'claude',
sessionId: PROVIDER_SESSION_ID,
leafUuid: null
})
expect(events[0]).toMatchObject({
type: 'message',
message: { subtype: 'hook_started', hook_name: 'SessionStart:startup' }
})
})
it('records only non-secret effective auth-lane diagnostics', async () => {
const claude = fakeClaude({
settings: {
env: {
ANTHROPIC_BASE_URL: 'https://gateway.example.test',
ANTHROPIC_AUTH_TOKEN: 'secret'
}
}
})
const events: ClaudeStructuredSessionEvent[] = []
await acquired(claude, {}, events)
const diagnostic = events.find((event) => event.type === 'auth-diagnostic')
expect(diagnostic).toEqual({
type: 'auth-diagnostic',
sessionId: 'session-1',
diagnostic: {
apiKeySourceConfigured: false,
baseUrlConfigured: true,
authTokenConfigured: true,
apiKeyConfigured: false,
settingSources: ['user', 'project', 'local']
}
})
expect(JSON.stringify(diagnostic)).not.toContain('secret')
expect(JSON.stringify(diagnostic)).not.toContain('gateway.example.test')
})
it('resumes the same provider id and refuses an init proof for another session', async () => {
const resumedClaude = fakeClaude()
const resumed = adapterFor(resumedClaude, {
resumed: true,
resumeLeafUuid: 'leaf-before'
})
const acquisition = await resumed.acquire({
identity: identityFor(),
fence: 9,
spawnToken: 'spawn-9'
})
expect(acquisition.link.origin).toBe('resumed')
expect(acquisition.link.handle).toEqual({
provider: 'claude',
sessionId: PROVIDER_SESSION_ID,
leafUuid: 'leaf-before'
})
const wrongClaude = fakeClaude({ initSessionId: 'different-session' })
const wrong = adapterFor(wrongClaude)
await expect(
wrong.acquire({ identity: identityFor(), fence: 7, spawnToken: 'spawn-9' })
).rejects.toThrow(/expected/)
expect(wrongClaude.connections[0].closeCount).toBe(1)
})
it('surfaces a CLI startup failure instead of waiting for the init deadline', async () => {
const claude = fakeClaude({ exitBeforeInit: 'Claude login required' })
const adapter = adapterFor(claude)
await expect(
adapter.acquire({ identity: identityFor(), fence: 7, spawnToken: 'spawn-9' })
).rejects.toThrow('Claude login required')
expect(claude.connections[0].closeCount).toBe(1)
})
it('closes a silent unauthenticated startup with actionable account guidance', async () => {
const claude = fakeClaude({ initProof: 'none' })
const adapter = adapterFor(claude, {}, [], [], 20)
const error = await adapter
.acquire({ identity: identityFor(), fence: 7, spawnToken: 'spawn-9' })
.catch((cause: unknown) => cause)
expect(error).toBeInstanceOf(AgentSessionAcquisitionRefusal)
expect(error).toMatchObject({
message: expect.stringMatching(/selected Claude account is signed in.*CLAUDE_CONFIG_DIR/s)
})
expect(claude.connections[0].calls[0]).toEqual({
subtype: 'initialize',
params: { supportedDialogKinds: [] }
})
expect(claude.connections[0].closeCount).toBe(1)
})
it('refuses an unauthenticated initialize response even when SessionStart runs', async () => {
const claude = fakeClaude({
initProof: 'session-start',
initAccount: { apiProvider: 'firstParty', tokenSource: 'none' }
})
const adapter = adapterFor(claude)
await expect(
adapter.acquire({ identity: identityFor(), fence: 7, spawnToken: 'spawn-9' })
).rejects.toThrow(/not signed in.*Claude CLI.*CLAUDE_CONFIG_DIR/s)
expect(claude.connections[0].closeCount).toBe(1)
})
})
describe('ClaudeStructuredSessionAdapter turns and controls', () => {
it('accepts a dispatch only after Claude replays its provider uuid', async () => {
const claude = fakeClaude({ replayUuid: 'user-provider-uuid' })
const adapter = await acquired(claude)
const result = await adapter.dispatch({
sessionId: 'session-1',
clientMessageId: 'client-1',
body: USER_MESSAGE,
fence: 7
})
expect(result).toEqual({
state: 'accepted',
providerIdentity: {
provider: 'claude',
sessionId: PROVIDER_SESSION_ID,
uuid: 'user-provider-uuid'
}
})
expect(claude.connections[0].sent[0]).toMatchObject({
type: 'user',
message: { role: 'user', content: [{ type: 'text', text: 'ship it' }] },
session_id: PROVIDER_SESSION_ID
})
})
it('leaves delivery unconfirmed when no replay uuid arrives', async () => {
const adapter = await acquired(fakeClaude({ replayUuid: null }))
await expect(
adapter.dispatch({
sessionId: 'session-1',
clientMessageId: 'client-1',
body: USER_MESSAGE,
fence: 7
})
).resolves.toMatchObject({ state: 'unknown' })
})
it('requires an acknowledged interrupt and supports controlled options', async () => {
const claude = fakeClaude()
const adapter = await acquired(claude)
await expect(
adapter.cancelTurn({ sessionId: 'session-1', turnId: 'turn-1', fence: 7 })
).resolves.toEqual({ cancelled: true })
await expect(
adapter.setOption({ sessionId: 'session-1', key: 'model', value: 'sonnet', fence: 7 })
).resolves.toEqual({ model: 'sonnet' })
expect(claude.connections[0].calls.slice(-2)).toEqual([
{ subtype: 'interrupt', params: {} },
{ subtype: 'set_model', params: { model: 'sonnet' } }
])
claude.routes.interrupt = () => {
throw new ClaudeControlRequestError('interrupt', 'not running')
}
await expect(
adapter.cancelTurn({ sessionId: 'session-1', turnId: 'turn-2', fence: 7 })
).resolves.toEqual({ cancelled: false })
claude.routes.interrupt = () => {
throw new Error('claude interrupt request timed out')
}
await expect(
adapter.cancelTurn({ sessionId: 'session-1', turnId: 'turn-3', fence: 7 })
).rejects.toThrow('timed out')
})
it('classifies provider-declined options without treating timeouts as settled', async () => {
const claude = fakeClaude({
routes: {
set_model: () => {
throw new ClaudeControlRequestError('set_model', 'model unavailable')
}
}
})
const adapter = await acquired(claude)
await expect(
adapter.setOption({ sessionId: 'session-1', key: 'model', value: 'fable', fence: 7 })
).rejects.toMatchObject({ name: 'AgentSessionOptionRejectedError' })
claude.routes.set_model = () => {
throw new Error('claude set_model request timed out')
}
await expect(
adapter.setOption({ sessionId: 'session-1', key: 'model', value: 'opus', fence: 7 })
).rejects.toThrow('timed out')
})
it('hydrates live model choices and maps the resolved current model to its CLI id', async () => {
const claude = fakeClaude({
initModel: 'claude-sonnet-5',
routes: {
list_models: () => ({
models: [
{ value: 'default', resolvedModel: 'claude-opus-5', displayName: 'Default' },
{
value: 'opus',
resolvedModel: 'claude-opus-5',
displayName: 'Opus',
supportsEffort: true,
supportedEffortLevels: ['low', 'high']
},
{
value: 'sonnet',
resolvedModel: 'claude-sonnet-5',
displayName: 'Sonnet'
}
]
})
}
})
const adapter = await acquired(claude)
await expect(adapter.readOptions({ sessionId: 'session-1', fence: 7 })).resolves.toEqual({
models: [
{
id: 'opus',
label: 'Opus',
isDefault: true,
efforts: [
{ value: 'low', label: 'Low' },
{ value: 'high', label: 'High' }
]
},
{ id: 'sonnet', label: 'Sonnet', isDefault: false, efforts: [] }
],
current: { model: 'sonnet', effort: 'high' }
})
})
it('keeps the shared Claude seed when live model discovery is unavailable', async () => {
const claude = fakeClaude({
initModel: 'custom-model',
routes: {
list_models: () => {
throw new Error('unsupported')
}
}
})
const adapter = await acquired(claude)
const result = await adapter.readOptions({ sessionId: 'session-1', fence: 7 })
expect(result.models.map((model) => model.id)).toEqual([
'fable',
'opus',
'sonnet',
'haiku',
'custom-model'
])
expect(result.current).toEqual({ model: 'custom-model', effort: 'high' })
})
})
describe('ClaudeStructuredSessionAdapter prompts', () => {
it('turns can_use_tool into an addressable durable approval callback', async () => {
const claude = fakeClaude()
const events: ClaudeStructuredSessionEvent[] = []
const adapter = await acquired(claude, {}, events)
claude.connections[0].handlers.onControlRequest?.({
type: 'control_request',
request_id: 'permission-1',
request: {
subtype: 'can_use_tool',
tool_name: 'Bash',
tool_use_id: 'tool-1',
input: { command: 'git status' },
permission_suggestions: [{ type: 'addRules' }]
}
})
expect(events.at(-1)).toMatchObject({
type: 'prompt',
prompt: { kind: 'approval', toolName: 'Bash', promptKey: 'permission-1' }
})
adapter.bindPromptItemId('session-1', 'journal-approval', 'permission-1')
await adapter.answerPrompt({
sessionId: 'session-1',
itemId: 'journal-approval',
kind: 'approval',
optionId: 'allowForSession',
fence: 7
})
expect(claude.connections[0].replies).toEqual([
{
requestId: 'permission-1',
response: {
behavior: 'allow',
updatedInput: { command: 'git status' },
updatedPermissions: [{ type: 'addRules' }],
toolUseID: 'tool-1'
}
}
])
})
it('collects every AskUserQuestion card before answering the one callback', async () => {
const claude = fakeClaude()
const adapter = await acquired(claude)
claude.connections[0].handlers.onControlRequest?.({
type: 'control_request',
request_id: 'question-1',
request: {
subtype: 'can_use_tool',
tool_name: 'AskUserQuestion',
tool_use_id: 'tool-question',
input: {
questions: [
{ question: 'Library?', options: [{ label: 'Luxon' }] },
{ question: 'Ship now?', options: [{ label: 'Yes' }] }
]
}
}
})
adapter.bindPromptItemId('session-1', 'journal-q1', 'question-1', 'Library?')
adapter.bindPromptItemId('session-1', 'journal-q2', 'question-1', 'Ship now?')
await adapter.answerPrompt({
sessionId: 'session-1',
itemId: 'journal-q1',
kind: 'question',
optionId: encodeClaudeQuestionOptionId('Library?', 'Luxon'),
fence: 7
})
expect(claude.connections[0].replies).toEqual([])
await adapter.answerPrompt({
sessionId: 'session-1',
itemId: 'journal-q2',
kind: 'question',
optionId: encodeClaudeQuestionOptionId('Ship now?', 'Yes'),
fence: 7
})
expect(claude.connections[0].replies[0]).toMatchObject({
requestId: 'question-1',
response: {
behavior: 'allow',
updatedInput: { answers: { 'Library?': 'Luxon', 'Ship now?': 'Yes' } },
toolUseID: 'tool-question'
}
})
})
it('does not move the persisted cursor from a result-frame uuid', async () => {
const claude = fakeClaude()
const events: ClaudeStructuredSessionEvent[] = []
const persistedHandles: unknown[] = []
const adapter = adapterFor(
claude,
{ resumeLeafUuid: 'previous-leaf', resumed: true },
events,
persistedHandles
)
await adapter.acquire({ identity: identityFor(), fence: 7, spawnToken: 'spawn-9' })
claude.connections[0].handlers.onMessage?.({
type: 'result',
session_id: PROVIDER_SESSION_ID,
uuid: 'result-uuid'
})
await adapter.closeSession('session-1')
expect(persistedHandles).toEqual([
{
sessionId: 'session-1',
providerSessionId: PROVIDER_SESSION_ID,
leafUuid: 'previous-leaf',
fence: 7
}
])
expect(events.at(-2)).toEqual({
type: 'handle',
sessionId: 'session-1',
providerSessionId: PROVIDER_SESSION_ID,
leafUuid: 'previous-leaf',
fence: 7
})
expect(claude.connections[0].closeCount).toBe(1)
})
})
@@ -0,0 +1,272 @@
import type {
AgentSessionAcquisition,
StructuredAgentSessionAcquireInput,
StructuredAgentSessionAdapter
} from '../native-chat/agent-session-wire/structured-agent-session-adapter'
import {
AgentSessionPreSpawnError as PreSpawnError,
rethrowAfterAgentSessionAcquisitionCleanup as rethrowAfterCleanup
} from '../native-chat/agent-session-wire/structured-agent-session-adapter'
import type { AgentSessionExecutionLocation } from '../../shared/agent-session-record'
import { LOCAL_EXECUTION_HOST_ID } from '../../shared/execution-host'
import { openClaudeStreamJsonConnection } from './claude-stream-json-connection'
import { answerClaudePrompt, cancelClaudeTurn } from './claude-structured-control-actions'
import { dispatchClaudeTurn } from './claude-structured-dispatch'
import { claudeAuthDiagnostic, readClaudeModels } from './claude-structured-init-proof'
import {
createClaudeInitDeadline,
requestClaudeInitialization
} from './claude-structured-init-deadline'
import { CLAUDE_SPAWN_TOKEN_ENV, claudeProcessIdentity } from './claude-structured-owner-identity'
import { ClaudeStructuredAcquisitionEvents } from './claude-structured-acquisition-events'
import {
restoreClaudeStructuredSessionOptions,
restoredClaudeStructuredSessionOptions,
setClaudeStructuredOption
} from './claude-structured-options'
import { ClaudePromptRegistry } from './claude-structured-prompt-replies'
import { readClaudeStructuredSessionOptions } from './claude-structured-session-options'
import { createClaudeSessionPublication } from './claude-structured-session-publication'
import { createClaudeSessionJournalTranslator } from './claude-structured-journal-translation'
import { ClaudeStructuredProviderEvents } from './claude-structured-provider-events'
import {
cancelClaudeAcquisitionAttempt,
ClaudeAcquisitionRegistry,
type ClaudeSession,
type ClaudeStructuredSessionAdapterDeps
} from './claude-structured-session-state'
import { closeClaudePublishedSession } from './claude-structured-session-close'
export type { ClaudeStructuredLaunch } from './claude-structured-launch-resolution'
export type {
ClaudeAuthDiagnostic,
ClaudeStructuredSessionAdapterDeps,
ClaudeStructuredSessionEvent
} from './claude-structured-session-state'
export const CLAUDE_STRUCTURED_INIT_TIMEOUT_MS = 10_000
const DISPATCH_ACK_TIMEOUT_MS = 10_000
export class ClaudeStructuredSessionAdapter implements StructuredAgentSessionAdapter {
private readonly sessions = new Map<string, ClaudeSession>()
private readonly acquisitions = new ClaudeAcquisitionRegistry()
private readonly events: ClaudeStructuredProviderEvents
constructor(private readonly deps: ClaudeStructuredSessionAdapterDeps) {
this.events = new ClaudeStructuredProviderEvents(this.sessions, deps.onEvent)
}
supportsCreate = (location: AgentSessionExecutionLocation, agent: string): boolean =>
agent === 'claude' &&
location.executionHostId === LOCAL_EXECUTION_HOST_ID &&
location.wslDistro === null &&
(process.platform === 'darwin' || process.platform === 'linux')
supportsLocation = (location: AgentSessionExecutionLocation): boolean =>
this.supportsCreate(location, 'claude')
async acquire(input: StructuredAgentSessionAcquireInput): Promise<AgentSessionAcquisition> {
const sessionId = input.identity.sessionId
const prompts = new ClaudePromptRegistry()
const translator = createClaudeSessionJournalTranslator(input.events, prompts)
const { previous, attempt } = this.acquisitions.start(sessionId, prompts)
const initTimeoutMs = this.deps.initTimeoutMs ?? CLAUDE_STRUCTURED_INIT_TIMEOUT_MS
const initDeadline = createClaudeInitDeadline(sessionId, initTimeoutMs)
const acquisitionEvents = new ClaudeStructuredAcquisitionEvents(
sessionId,
attempt,
input.events,
initDeadline,
this.events
)
try {
if (!(await cancelClaudeAcquisitionAttempt(previous))) {
throw new Error(`superseded claude session ${sessionId} did not prove provider exit`)
}
this.acquisitions.assertCurrent(sessionId, attempt)
await this.closePublishedSession(sessionId)
this.acquisitions.assertCurrent(sessionId, attempt)
const launch = await this.deps
.resolveLaunch({ identity: input.identity })
.catch((error: unknown) => {
throw new PreSpawnError(error)
})
acquisitionEvents.observeResumeLeaf(launch.resumeLeafUuid)
this.acquisitions.assertCurrent(sessionId, attempt)
const open = this.deps.openConnection ?? openClaudeStreamJsonConnection
const connection = await open(
{
command: launch.command,
args: launch.args,
cwd: launch.cwd,
env: {
...launch.env,
[CLAUDE_SPAWN_TOKEN_ENV]: input.spawnToken,
CLAUDE_CONFIG_DIR: launch.claudeConfigDir
}
},
{
onMessage: acquisitionEvents.onMessage,
onControlRequest: acquisitionEvents.onControlRequest,
onControlCancelRequest: acquisitionEvents.onControlCancelRequest,
onExit: (error) => {
if (!attempt.published) {
initDeadline.reject(error)
}
this.events.handleExit(sessionId, attempt, error)
}
}
)
attempt.connection = connection
this.acquisitions.assertCurrent(sessionId, attempt)
initDeadline.start()
const [initialization, init] = await Promise.all([
requestClaudeInitialization(connection, sessionId, initTimeoutMs),
initDeadline.promise
])
const models = readClaudeModels(initialization)
acquisitionEvents.emit({ type: 'options', sessionId, models })
initDeadline.clear()
this.acquisitions.assertCurrent(sessionId, attempt)
if (init.providerSessionId !== launch.providerSessionId) {
throw new Error(
`claude proved session ${init.providerSessionId}, expected ${launch.providerSessionId}`
)
}
const settings = await connection
.request('get_settings', {}, { timeoutMs: this.deps.requestTimeoutMs })
.catch(() => null)
acquisitionEvents.emit({
type: 'auth-diagnostic',
sessionId,
diagnostic: claudeAuthDiagnostic(init, settings)
})
const process = await claudeProcessIdentity(
{ ...input, pid: connection.pid },
this.deps.readProcessStartTime
)
this.acquisitions.assertCurrent(sessionId, attempt)
if (connection.closed) {
throw new Error(`claude stream-json for session ${sessionId} exited while being acquired`)
}
const publication = createClaudeSessionPublication({
connection,
init,
leafUuid: acquisitionEvents.observedLeaf(),
fence: input.fence,
resumed: launch.resumed,
prompts,
translator,
events: input.events,
process,
options: restoredClaudeStructuredSessionOptions(input.options),
...(this.deps.mintLinkId ? { linkId: this.deps.mintLinkId() } : {}),
observedAt: this.deps.now?.() ?? Date.now()
})
const acquired: AgentSessionAcquisition = publication.acquisition
const liveSession = publication.session
acquisitionEvents.publish(liveSession)
await restoreClaudeStructuredSessionOptions(liveSession, this.deps.requestTimeoutMs)
this.acquisitions.assertCurrent(sessionId, attempt)
this.acquisitions.deleteIfCurrent(sessionId, attempt)
this.sessions.set(sessionId, liveSession)
attempt.published = true
for (const event of attempt.buffered.splice(0)) {
event()
}
return acquired
} catch (error) {
initDeadline.clear()
this.acquisitions.deleteIfCurrent(sessionId, attempt)
if (this.sessions.get(sessionId)?.connection !== attempt.connection) {
translator?.dispose()
prompts.clear()
await rethrowAfterCleanup(
{
releaseAcquisition: async () =>
attempt.connection ? await attempt.connection.close() : true
},
sessionId,
error
)
}
throw error
} finally {
attempt.finish()
}
}
bindPromptItemId(
sessionId: string,
journalItemId: string,
promptKey: string,
questionId?: string
): void {
this.sessions.get(sessionId)?.prompts.bindJournalItemId(journalItemId, promptKey, questionId)
}
dispatch: StructuredAgentSessionAdapter['dispatch'] = (input) =>
dispatchClaudeTurn(
this.session(input.sessionId),
input,
this.deps.dispatchAckTimeoutMs ?? DISPATCH_ACK_TIMEOUT_MS
)
cancelTurn: StructuredAgentSessionAdapter['cancelTurn'] = (input) =>
cancelClaudeTurn(this.session(input.sessionId), this.deps.requestTimeoutMs)
answerPrompt: StructuredAgentSessionAdapter['answerPrompt'] = (input) =>
answerClaudePrompt(this.session(input.sessionId), input)
setOption: StructuredAgentSessionAdapter['setOption'] = (input) =>
setClaudeStructuredOption(this.session(input.sessionId), input, this.deps.requestTimeoutMs)
readOptions = (input: { sessionId: string; fence: number }) =>
readClaudeStructuredSessionOptions(this.session(input.sessionId), this.deps.requestTimeoutMs)
releaseAcquisition(input: { sessionId: string }): Promise<boolean> {
return this.closeSession(input.sessionId)
}
async closeSession(sessionId: string): Promise<boolean> {
const attempt = this.acquisitions.get(sessionId)
if (attempt) {
attempt.cancelled = true
const exited = attempt.connection ? await attempt.connection.close() : true
await attempt.finished
if (!exited) {
return false
}
}
return this.closePublishedSession(sessionId)
}
disposeSession = (sessionId: string): Promise<boolean> => this.closeSession(sessionId)
private closePublishedSession(sessionId: string): Promise<boolean> {
return closeClaudePublishedSession({
sessions: this.sessions,
sessionId,
...(this.deps.persistHandle ? { persistHandle: this.deps.persistHandle } : {}),
...(this.deps.onEvent ? { onEvent: this.deps.onEvent } : {})
})
}
async closeAll(): Promise<void> {
this.acquisitions.close()
const ids = new Set([...this.sessions.keys(), ...this.acquisitions.sessionIds()])
const outcomes = await Promise.all([...ids].map((sessionId) => this.closeSession(sessionId)))
if (outcomes.some((exited) => !exited)) {
throw new Error('claude structured session teardown could not prove provider-child exit')
}
}
private session(sessionId: string): ClaudeSession {
const session = this.sessions.get(sessionId)
if (!session) {
throw new Error(`no live claude stream-json session for ${sessionId}`)
}
return session
}
}
@@ -0,0 +1,79 @@
import type { ClaudeSession, ClaudeStructuredSessionEvent } from './claude-structured-session-state'
export function settleClaudeDispatchWaiters(session: ClaudeSession): void {
for (const waiter of session.dispatchWaiters.splice(0)) {
clearTimeout(waiter.timer)
waiter.resolve(null)
}
}
export function settleClaudeExitedSession(session: ClaudeSession): void {
settleClaudeDispatchWaiters(session)
session.prompts.clear()
session.translator?.dispose()
}
export async function closeClaudePublishedSession(input: {
sessions: Map<string, ClaudeSession>
sessionId: string
persistHandle?: (handle: {
sessionId: string
providerSessionId: string
leafUuid: string | null
fence: number
}) => Promise<void>
onEvent?: (event: ClaudeStructuredSessionEvent) => void
}): Promise<boolean> {
const session = input.sessions.get(input.sessionId)
if (!session) {
return true
}
settleClaudeDispatchWaiters(session)
const pending = session.prompts.clear()
await Promise.allSettled(
pending.map((prompt) =>
session.connection.respond(prompt.requestId, {
behavior: 'deny',
message: 'Structured Claude session closed.',
interrupt: true,
toolUseID: prompt.toolUseId
})
)
)
session.translator?.flush()
if (!(await session.connection.close())) {
return false
}
input.sessions.delete(input.sessionId)
let persistenceError: unknown
try {
await input.persistHandle?.({
sessionId: input.sessionId,
providerSessionId: session.providerSessionId,
leafUuid: session.leafUuid,
fence: session.fence
})
input.onEvent?.({
type: 'handle',
sessionId: input.sessionId,
providerSessionId: session.providerSessionId,
leafUuid: session.leafUuid,
fence: session.fence
})
} catch (error) {
persistenceError = error
} finally {
const ended = {
type: 'ended',
sessionId: input.sessionId,
reason: 'claude session closed'
} as const
session.translator?.handle(ended)
input.onEvent?.(ended)
session.translator?.dispose()
}
if (persistenceError) {
throw persistenceError
}
return true
}
@@ -0,0 +1,112 @@
import type {
AgentSessionModelOption,
AgentSessionOptionChoice,
AgentSessionOptionsResult
} from '../../shared/agent-session-wire'
import { CLAUDE_SESSION_OPTION_CATALOG } from '../../shared/agent-session-option-catalog-claude-codex'
import type { CatalogModel } from '../../shared/agent-session-option-catalog-types'
import type { ClaudeSession } from './claude-structured-session-state'
type ListedModel = AgentSessionModelOption & { resolvedModel: string | null }
function record(value: unknown): Record<string, unknown> | null {
return typeof value === 'object' && value !== null && !Array.isArray(value)
? (value as Record<string, unknown>)
: null
}
function text(value: unknown): string | null {
return typeof value === 'string' && value.trim() ? value : null
}
function effortLabel(value: string): string {
return value === 'xhigh' ? 'Extra high' : `${value.charAt(0).toUpperCase()}${value.slice(1)}`
}
function listedEfforts(row: Record<string, unknown>): AgentSessionOptionChoice[] {
return row.supportsEffort === true && Array.isArray(row.supportedEffortLevels)
? row.supportedEffortLevels.flatMap((value) => {
const effort = text(value)
return effort ? [{ value: effort, label: effortLabel(effort) }] : []
})
: []
}
function listedModels(value: unknown): ListedModel[] {
const response = record(value)
const rows = Array.isArray(response?.models)
? response.models.map(record).filter((row): row is Record<string, unknown> => row !== null)
: []
const defaultRow = rows.find((row) => text(row.value) === 'default')
const defaultResolvedModel = text(defaultRow?.resolvedModel)
const seen = new Set<string>()
return rows.flatMap((row) => {
const id = text(row.value)
if (!id || id === 'default' || seen.has(id)) {
return []
}
seen.add(id)
const resolvedModel = text(row.resolvedModel)
const description = text(row.description)
return [
{
id,
label: text(row.displayName) ?? id,
...(description ? { description } : {}),
isDefault: resolvedModel !== null && resolvedModel === defaultResolvedModel,
efforts: listedEfforts(row),
resolvedModel
}
]
})
}
function seedEfforts(model: CatalogModel): AgentSessionOptionChoice[] {
const effort = model.options.find((option) => option.id === 'effort')
return effort?.kind.type === 'select' ? effort.kind.choices : []
}
function seedModels(): ListedModel[] {
return CLAUDE_SESSION_OPTION_CATALOG.models.map((model) => ({
id: model.id,
label: model.label,
...(model.description ? { description: model.description } : {}),
isDefault: model.isDefault === true,
efforts: seedEfforts(model),
resolvedModel: null
}))
}
function currentModelId(models: ListedModel[], reportedModel: string | undefined): string {
const matched = reportedModel
? models.find((model) => model.id === reportedModel || model.resolvedModel === reportedModel)
: undefined
return (
matched?.id ?? reportedModel ?? models.find((model) => model.isDefault)?.id ?? models[0]!.id
)
}
export async function readClaudeStructuredSessionOptions(
session: ClaudeSession,
timeoutMs: number | undefined
): Promise<AgentSessionOptionsResult> {
const response = await session.connection
.request('list_models', {}, { timeoutMs })
.catch(() => null)
const discovered = listedModels(response)
const models = discovered.length > 0 ? discovered : seedModels()
const reportedModel = session.options.get('model') ?? session.reportedOptions.model
const model = currentModelId(models, reportedModel)
if (!models.some((entry) => entry.id === model)) {
models.push({ id: model, label: model, isDefault: false, efforts: [], resolvedModel: null })
}
const effort = session.options.get('effort') ?? session.reportedOptions.effort
return {
models: models.map((entry) => {
const { resolvedModel, ...withoutResolvedModel } = entry
void resolvedModel
return withoutResolvedModel
}),
current: { model, ...(effort ? { effort } : {}) }
}
}
@@ -0,0 +1,52 @@
import type { AgentSessionAcquisition } from '../native-chat/agent-session-wire/structured-agent-session-adapter'
import { readClaudeFrameString, type ClaudeInitObservation } from './claude-structured-init-proof'
import { claudeProviderHandleLink } from './claude-structured-owner-identity'
import type { ClaudePromptRegistry } from './claude-structured-prompt-replies'
import type { ClaudeJournalTranslator } from './claude-structured-journal-translation'
import type { ClaudeSession } from './claude-structured-session-state'
export function createClaudeSessionPublication(input: {
connection: ClaudeSession['connection']
init: ClaudeInitObservation
leafUuid: string | null
fence: number
resumed: boolean
prompts: ClaudePromptRegistry
translator: ClaudeJournalTranslator | null
events: ClaudeSession['events']
process: AgentSessionAcquisition['process']
linkId?: string
observedAt: number
options?: ReadonlyMap<string, string>
}): { acquisition: AgentSessionAcquisition; session: ClaudeSession } {
const model = readClaudeFrameString(input.init.message, 'model')
const effort = readClaudeFrameString(input.init.message, 'effortLevel')
return {
acquisition: {
process: input.process,
link: claudeProviderHandleLink({
sessionId: input.init.providerSessionId,
leafUuid: input.leafUuid,
resumed: input.resumed,
fence: input.fence,
...(input.linkId ? { linkId: input.linkId } : {}),
observedAt: input.observedAt
})
},
session: {
connection: input.connection,
providerSessionId: input.init.providerSessionId,
leafUuid: input.leafUuid,
fence: input.fence,
prompts: input.prompts,
dispatchWaiters: [],
options: new Map(input.options),
reportedOptions: {
...(model ? { model } : {}),
...(effort ? { effort } : {})
},
translator: input.translator,
events: input.events
}
}
}
@@ -0,0 +1,177 @@
import type { AgentSessionJournalIdentity } from '../../shared/agent-session-journal-types'
import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import type {
ClaudeStreamJsonConnection,
openClaudeStreamJsonConnection
} from './claude-stream-json-connection'
import type { ClaudeStructuredLaunch } from './claude-structured-launch-resolution'
import type { ClaudeJournalTranslator } from './claude-structured-journal-translation'
import type { ClaudePendingPrompt, ClaudePromptRegistry } from './claude-structured-prompt-replies'
export type ClaudeAuthDiagnostic = {
apiKeySourceConfigured: boolean
baseUrlConfigured: boolean
authTokenConfigured: boolean
apiKeyConfigured: boolean
settingSources: readonly string[]
}
export type ClaudeStructuredSessionEvent =
| { type: 'message'; sessionId: string; message: Record<string, unknown> }
| { type: 'provider-frame'; sessionId: string; kind: string; payload: unknown }
| { type: 'prompt'; sessionId: string; prompt: ClaudePendingPrompt }
| { type: 'prompt-cancelled'; sessionId: string; promptKey: string }
| { type: 'options'; sessionId: string; models: unknown[] }
| {
type: 'handle'
sessionId: string
providerSessionId: string
leafUuid: string | null
fence: number
}
| { type: 'auth-diagnostic'; sessionId: string; diagnostic: ClaudeAuthDiagnostic }
| { type: 'ended'; sessionId: string; reason: string }
export type ClaudeStructuredSessionAdapterDeps = {
resolveLaunch: (input: {
identity: AgentSessionJournalIdentity
}) => Promise<ClaudeStructuredLaunch>
onEvent?: (event: ClaudeStructuredSessionEvent) => void
openConnection?: typeof openClaudeStreamJsonConnection
readProcessStartTime?: (pid: number) => Promise<number | null>
mintLinkId?: () => string
now?: () => number
requestTimeoutMs?: number
initTimeoutMs?: number
dispatchAckTimeoutMs?: number
persistHandle?: (input: {
sessionId: string
providerSessionId: string
leafUuid: string | null
fence: number
}) => Promise<void>
}
/** The only stream-frame path allowed to advance Claude's advisory transcript cursor. */
export function observeClaudeTopLevelLeaf(
currentLeaf: string | null,
message: Record<string, unknown>
): string | null {
if (
(message.type !== 'user' && message.type !== 'assistant') ||
message.parent_tool_use_id !== null ||
typeof message.uuid !== 'string' ||
message.uuid.length === 0
) {
return currentLeaf
}
return message.uuid
}
export type ClaudeDispatchWaiter = {
resolve: (uuid: string | null) => void
timer: ReturnType<typeof setTimeout>
acceptsResult: boolean
}
export type ClaudeSession = {
connection: ClaudeStreamJsonConnection
providerSessionId: string
leafUuid: string | null
fence: number
prompts: ClaudePromptRegistry
dispatchWaiters: ClaudeDispatchWaiter[]
options: Map<string, string>
reportedOptions: { model?: string; effort?: string }
translator: ClaudeJournalTranslator | null
events: StructuredAgentSessionEventSink | undefined
}
export type ClaudeAcquisitionAttempt = {
connection: ClaudeStreamJsonConnection | null
prompts: ClaudePromptRegistry
buffered: (() => void)[]
published: boolean
cancelled: boolean
finished: Promise<void>
finish: () => void
}
export function createClaudeAcquisitionAttempt(
prompts: ClaudePromptRegistry
): ClaudeAcquisitionAttempt {
let finish = (): void => {}
const finished = new Promise<void>((resolve) => {
finish = resolve
})
return {
connection: null,
prompts,
buffered: [],
published: false,
cancelled: false,
finished,
finish
}
}
export class ClaudeAcquisitionRegistry {
private readonly attempts = new Map<string, ClaudeAcquisitionAttempt>()
private closing = false
get size(): number {
return this.attempts.size
}
start(
sessionId: string,
prompts: ClaudePromptRegistry
): {
previous: ClaudeAcquisitionAttempt | undefined
attempt: ClaudeAcquisitionAttempt
} {
if (this.closing) {
throw new Error('claude structured session adapter is closing')
}
const previous = this.attempts.get(sessionId)
const attempt = createClaudeAcquisitionAttempt(prompts)
this.attempts.set(sessionId, attempt)
return { previous, attempt }
}
assertCurrent(sessionId: string, attempt: ClaudeAcquisitionAttempt): void {
if (this.closing || attempt.cancelled || this.attempts.get(sessionId) !== attempt) {
throw new Error(`claude session ${sessionId} was superseded while being acquired`)
}
}
get(sessionId: string): ClaudeAcquisitionAttempt | undefined {
return this.attempts.get(sessionId)
}
deleteIfCurrent(sessionId: string, attempt: ClaudeAcquisitionAttempt): void {
if (this.attempts.get(sessionId) === attempt) {
this.attempts.delete(sessionId)
}
}
sessionIds(): IterableIterator<string> {
return this.attempts.keys()
}
close(): void {
this.closing = true
}
}
export async function cancelClaudeAcquisitionAttempt(
attempt: ClaudeAcquisitionAttempt | undefined
): Promise<boolean> {
if (!attempt) {
return true
}
attempt.cancelled = true
const exited = attempt.connection ? await attempt.connection.close() : true
await attempt.finished
return exited
}
@@ -189,6 +189,12 @@ const CODEX_ITEM_CLASSIFICATIONS: Record<string, ProviderFrameClassification> =
contextCompaction: 'status-chrome'
}
const CLAUDE_CONTROL_FRAME_CLASSIFICATIONS: Record<string, ProviderFrameClassification> = {
'control_request:can_use_tool': 'timeline-substantive',
'control_request:request_user_dialog': 'timeline-substantive',
control_cancel_request: 'timeline-substantive'
}
function notificationKind(kind: string): string {
return kind.startsWith('notification:') ? kind.slice('notification:'.length) : kind
}
@@ -215,7 +221,10 @@ function catalogClassification(
]
}
if (provider === 'claude') {
return PROVIDER_FRAME_CLASSIFICATIONS.claude[kind as ClaudeStreamJsonFrameKind]
return (
PROVIDER_FRAME_CLASSIFICATIONS.claude[kind as ClaudeStreamJsonFrameKind] ??
CLAUDE_CONTROL_FRAME_CLASSIFICATIONS[kind]
)
}
return undefined
}
@@ -0,0 +1,152 @@
import { describe, expect, it, vi } from 'vitest'
import type { AgentSessionJournalIdentity } from '../../../shared/agent-session-journal-types'
import type { AgentSessionExecutionLocation } from '../../../shared/agent-session-record'
import type { StructuredAgentSessionAdapter } from './structured-agent-session-adapter'
import { StructuredAgentSessionAdapterRouter } from './structured-agent-session-adapter-router'
import { adapterSupportsCreate } from './structured-agent-session-provider-support'
const LOCATION: AgentSessionExecutionLocation = {
executionHostId: 'local',
wslDistro: null,
workspaceId: 'workspace-1',
workspaceKind: 'folder'
}
function identity(agent: 'claude' | 'codex', sessionId: string): AgentSessionJournalIdentity {
return {
sessionId,
workspaceId: 'workspace-1',
hostId: 'local',
agent,
providerHandle:
agent === 'claude'
? { kind: 'claude', sessionId: 'provider-1', leafUuid: null }
: { kind: 'codex', threadId: 'thread-1' }
}
}
function fakeAdapter(agent: 'claude' | 'codex') {
const dispatch = vi.fn(async () => ({ state: 'rejected' as const, reason: agent }))
const closeSession = vi.fn(async () => true)
const releaseAcquisition = vi.fn(async () => true)
const adapter = {
supportsCreate: vi.fn(
(_location: AgentSessionExecutionLocation, candidate: string) => candidate === agent
),
acquire: vi.fn(async () => ({
process: {
hostId: 'local',
pid: agent === 'claude' ? 101 : 102,
processStartTimeMs: 1,
spawnToken: 'spawn'
},
link:
agent === 'claude'
? {
linkId: 'claude-link',
handle: { provider: 'claude' as const, sessionId: 'provider-1', leafUuid: null },
origin: 'created' as const,
mintedAtFence: 1,
observedAt: 1
}
: {
linkId: 'codex-link',
handle: { provider: 'codex' as const, threadId: 'thread-1' },
origin: 'created' as const,
mintedAtFence: 1,
observedAt: 1
}
})),
releaseAcquisition,
dispatch,
cancelTurn: vi.fn(async () => ({ cancelled: false })),
answerPrompt: vi.fn(async () => {}),
setOption: vi.fn(async () => {}),
readOptions: vi.fn(async () => ({
models: [],
current: { model: `${agent}-model` }
})),
historyFilePath: vi.fn(async () => null),
closeSession,
disposeSession: closeSession
} satisfies StructuredAgentSessionAdapter
return { adapter, dispatch, closeSession, releaseAcquisition }
}
function dispatchInput(sessionId: string) {
return {
sessionId,
clientMessageId: 'client-1',
body: { kind: 'message' as const, role: 'user' as const, blocks: [] },
fence: 1
}
}
describe('StructuredAgentSessionAdapterRouter', () => {
it('routes every acquired session through its provider owner, including Codex', async () => {
const claude = fakeAdapter('claude')
const codex = fakeAdapter('codex')
const router = new StructuredAgentSessionAdapterRouter(
{ claude: claude.adapter, codex: codex.adapter },
async () => {}
)
await router.acquire({
identity: identity('codex', 'codex-session'),
fence: 1,
spawnToken: 's'
})
await router.acquire({
identity: identity('claude', 'claude-session'),
fence: 1,
spawnToken: 's'
})
await router.dispatch(dispatchInput('codex-session'))
await router.dispatch(dispatchInput('claude-session'))
expect(codex.dispatch).toHaveBeenCalledTimes(1)
expect(claude.dispatch).toHaveBeenCalledTimes(1)
})
it('makes Claude available through router capability, not the Codex fallback', () => {
const claude = fakeAdapter('claude')
const codex = fakeAdapter('codex')
const router = new StructuredAgentSessionAdapterRouter(
{ claude: claude.adapter, codex: codex.adapter },
async () => {}
)
expect(adapterSupportsCreate(router, LOCATION, 'claude')).toBe(true)
expect(claude.adapter.supportsCreate).toHaveBeenCalledWith(LOCATION, 'claude')
expect(codex.adapter.supportsCreate).toHaveBeenCalledTimes(0)
})
it('retains the owner until provider-child exit is proven', async () => {
const claude = fakeAdapter('claude')
const codex = fakeAdapter('codex')
codex.closeSession.mockResolvedValueOnce(false)
const router = new StructuredAgentSessionAdapterRouter(
{ claude: claude.adapter, codex: codex.adapter },
async () => {}
)
await router.acquire({ identity: identity('codex', 'session-1'), fence: 1, spawnToken: 's' })
expect(await router.closeSession('session-1')).toBe(false)
await expect(router.dispatch(dispatchInput('session-1'))).resolves.toMatchObject({
reason: 'codex'
})
})
it('proves unknown-owner cleanup only when every adapter proves absence', async () => {
const claude = fakeAdapter('claude')
const codex = fakeAdapter('codex')
claude.releaseAcquisition.mockResolvedValueOnce(false)
const router = new StructuredAgentSessionAdapterRouter(
{ claude: claude.adapter, codex: codex.adapter },
async () => {}
)
expect(await router.releaseAcquisition({ sessionId: 'unknown-1' })).toBe(false)
expect(await router.releaseAcquisition({ sessionId: 'unknown-2' })).toBe(true)
})
})
@@ -0,0 +1,127 @@
import type { AgentSessionJournalIdentity } from '../../../shared/agent-session-journal-types'
import type { AgentSessionExecutionLocation } from '../../../shared/agent-session-record'
import type { StructuredAgentSessionAdapter } from './structured-agent-session-adapter'
type RoutedAgent = 'claude' | 'codex'
export class StructuredAgentSessionAdapterRouter implements StructuredAgentSessionAdapter {
private readonly owners = new Map<string, StructuredAgentSessionAdapter>()
constructor(
private readonly adapters: Record<RoutedAgent, StructuredAgentSessionAdapter>,
private readonly closeAdapters: () => Promise<void>
) {}
supportsCreate = (location: AgentSessionExecutionLocation, agent: string): boolean => {
const adapter = this.adapterForAgent(agent)
return adapter
? (adapter.supportsCreate?.(location, agent) ?? adapter.supportsLocation?.(location) ?? false)
: false
}
supportsLocation = (location: AgentSessionExecutionLocation): boolean =>
Object.values(this.adapters).some((adapter) => adapter.supportsLocation?.(location) ?? false)
async acquire(input: Parameters<StructuredAgentSessionAdapter['acquire']>[0]) {
const adapter = this.requireAgent(input.identity)
const acquired = await adapter.acquire(input)
this.owners.set(input.identity.sessionId, adapter)
return acquired
}
async releaseAcquisition(input: { sessionId: string }): Promise<boolean> {
const adapter = this.owners.get(input.sessionId)
if (adapter) {
const exited = (await adapter.releaseAcquisition?.(input)) === true
if (exited) {
this.owners.delete(input.sessionId)
}
return exited
}
return this.cleanupUnknownOwner((candidate) => candidate.releaseAcquisition?.(input))
}
dispatch: StructuredAgentSessionAdapter['dispatch'] = (input) =>
this.owner(input.sessionId).dispatch(input)
cancelTurn: StructuredAgentSessionAdapter['cancelTurn'] = (input) =>
this.owner(input.sessionId).cancelTurn(input)
answerPrompt: StructuredAgentSessionAdapter['answerPrompt'] = (input) =>
this.owner(input.sessionId).answerPrompt(input)
setOption: StructuredAgentSessionAdapter['setOption'] = (input) =>
this.owner(input.sessionId).setOption(input)
readOptions = (input: { sessionId: string; fence: number }) => {
const reader = this.owner(input.sessionId).readOptions
if (!reader) {
throw new Error(`structured session ${input.sessionId} does not report options`)
}
return reader(input)
}
historyFilePath = (input: { identity: AgentSessionJournalIdentity }) =>
this.requireAgent(input.identity).historyFilePath?.(input) ?? Promise.resolve(null)
async closeSession(sessionId: string): Promise<boolean> {
const adapter = this.owners.get(sessionId)
if (!adapter) {
return this.cleanupUnknownOwner((candidate) => candidate.closeSession?.(sessionId))
}
const exited = (await adapter.closeSession?.(sessionId)) === true
if (exited) {
this.owners.delete(sessionId)
}
return exited
}
async disposeSession(sessionId: string): Promise<boolean> {
const adapter = this.owners.get(sessionId)
if (!adapter) {
return this.cleanupUnknownOwner((candidate) =>
(candidate.disposeSession ?? candidate.closeSession)?.(sessionId)
)
}
const dispose = adapter.disposeSession ?? adapter.closeSession
const exited = (await dispose?.call(adapter, sessionId)) === true
if (exited) {
this.owners.delete(sessionId)
}
return exited
}
async closeAll(): Promise<void> {
this.owners.clear()
await this.closeAdapters()
}
private owner(sessionId: string): StructuredAgentSessionAdapter {
const adapter = this.owners.get(sessionId)
if (!adapter) {
throw new Error(`no live structured adapter owns ${sessionId}`)
}
return adapter
}
private requireAgent(identity: AgentSessionJournalIdentity): StructuredAgentSessionAdapter {
const adapter = this.adapterForAgent(identity.agent)
if (!adapter) {
throw new Error(`structured sessions do not support ${identity.agent}`)
}
return adapter
}
private adapterForAgent(agent: string): StructuredAgentSessionAdapter | null {
return agent === 'claude' || agent === 'codex' ? this.adapters[agent] : null
}
private async cleanupUnknownOwner(
cleanup: (adapter: StructuredAgentSessionAdapter) => Promise<boolean> | undefined
): Promise<boolean> {
const outcomes = await Promise.allSettled(
Object.values(this.adapters).map(async (adapter) => (await cleanup(adapter)) === true)
)
return outcomes.every((outcome) => outcome.status === 'fulfilled' && outcome.value === true)
}
}
@@ -11,12 +11,19 @@ import { existsSync } from 'node:fs'
import { join } from 'node:path'
import type { AgentSessionOwnerProbe } from '../../shared/agent-session-lease-adjudication'
import type { AgentSessionRecord } from '../../shared/agent-session-record'
import { createClaudeStructuredLaunchResolver } from '../claude/claude-structured-launch-resolution'
import {
ClaudeStructuredSessionAdapter,
type ClaudeStructuredSessionAdapterDeps
} from '../claude/claude-structured-session-adapter'
import { claudeProviderHandleLink } from '../claude/claude-structured-owner-identity'
import { createCodexStructuredLaunchResolver } from '../codex/codex-structured-launch-resolution'
import {
CodexStructuredSessionAdapter,
type CodexStructuredSessionAdapterDeps
} from '../codex/codex-structured-session-adapter'
import { StructuredAgentSessionHost } from '../native-chat/agent-session-wire/structured-agent-session-host'
import { StructuredAgentSessionAdapterRouter } from '../native-chat/agent-session-wire/structured-agent-session-adapter-router'
import type { StructuredAgentSessionHandoffTransport } from '../native-chat/agent-session-wire/structured-agent-session-handoff-types'
import { setStructuredAgentSessionHost } from '../native-chat/agent-session-wire/structured-agent-session-registry'
import { AgentSessionRecordStore } from './agent-session-record-store'
@@ -55,15 +62,19 @@ export type StructuredAgentSessionRuntimeDeps = {
claimKeyId: string
resolveWorkspacePath: (workspaceId: string) => Promise<string>
resolveCodexCommand?: (options?: { pathEnv?: string | null; homePath?: string }) => string
resolveClaudeCommand?: (options?: { pathEnv?: string | null; homePath?: string }) => string
/** Provider transports are overridden only to drive the runtime against scripted children. */
openCodexConnection?: CodexStructuredSessionAdapterDeps['openConnection']
openClaudeConnection?: ClaudeStructuredSessionAdapterDeps['openConnection']
/** Scripted app-servers carry fake pids the real start-time read cannot answer for. */
readProcessStartTime?: CodexStructuredSessionAdapterDeps['readProcessStartTime']
readClaudeProcessStartTime?: ClaudeStructuredSessionAdapterDeps['readProcessStartTime']
resolveLaunchArgs?: (provider: AgentSessionRecord['provider']) => Promise<string[]> | string[]
resolveLaunchEnv?: () => Promise<NodeJS.ProcessEnv>
resolveLaunchEnvOverlay?: () => Promise<Record<string, string>> | Record<string, string>
resolveEnvironment?: () => Promise<NodeJS.ProcessEnv>
resolveCodexOverrides?: () => NodeJS.ProcessEnv
resolveClaudeOverrides?: () => NodeJS.ProcessEnv
onError?: (input: { scope: string; error: unknown }) => void
handoffTransport?: StructuredAgentSessionHandoffTransport
reapOrphanChildren?: typeof stopOrphanAgentSessionChildren
@@ -71,7 +82,7 @@ export type StructuredAgentSessionRuntimeDeps = {
type InstalledRuntime = {
host: StructuredAgentSessionHost
adapter: CodexStructuredSessionAdapter
adapter: StructuredAgentSessionAdapterRouter
/** Resolves after every adapter-exit recovery callback has settled. */
waitForRecovery: () => Promise<void>
}
@@ -89,7 +100,7 @@ export function ensureStructuredAgentSessionHost(
return installing.then((installed) => installed.host)
}
/** Drops the host and reaps every Codex child under it. Runtime teardown and
/** Drops the host and reaps every structured provider child under it. Runtime teardown and
* test isolation take the same path, so neither can leave a live app-server. */
export async function stopStructuredAgentSessionRuntime(): Promise<void> {
const pending = installing
@@ -118,11 +129,13 @@ export async function stopStructuredAgentSessionRuntime(): Promise<void> {
async function install(deps: StructuredAgentSessionRuntimeDeps): Promise<InstalledRuntime> {
const bootEnvironment = (deps.resolveEnvironment ?? resolveLoginShellEnvironment)()
const resolveEnvironment = async (): Promise<NodeJS.ProcessEnv> => ({
const resolveEnvironment = async (
provider: AgentSessionRecord['provider']
): Promise<NodeJS.ProcessEnv> => ({
...(await bootEnvironment),
...(await deps.resolveLaunchEnv?.()),
...(await deps.resolveLaunchEnvOverlay?.()),
...deps.resolveCodexOverrides?.()
...(provider === 'claude' ? deps.resolveClaudeOverrides?.() : deps.resolveCodexOverrides?.())
})
const store = await AgentSessionRecordStore.open({
directory: join(deps.stateDirectory, RECORD_STORE_DIR_NAME),
@@ -151,7 +164,7 @@ async function install(deps: StructuredAgentSessionRuntimeDeps): Promise<Install
resolveLaunch: createCodexStructuredLaunchResolver({
store,
resolveWorkspacePath: deps.resolveWorkspacePath,
resolveEnvironment,
resolveEnvironment: () => resolveEnvironment('codex'),
...(deps.resolveCodexCommand ? { resolveCommand: deps.resolveCodexCommand } : {})
}),
...(deps.openCodexConnection ? { openConnection: deps.openCodexConnection } : {}),
@@ -172,7 +185,36 @@ async function install(deps: StructuredAgentSessionRuntimeDeps): Promise<Install
})
}
})
const adapter = codex
const claude = new ClaudeStructuredSessionAdapter({
resolveLaunch: createClaudeStructuredLaunchResolver({
store,
resolveWorkspacePath: deps.resolveWorkspacePath,
resolveEnvironment: () => resolveEnvironment('claude'),
...(deps.resolveClaudeCommand ? { resolveCommand: deps.resolveClaudeCommand } : {})
}),
...(deps.openClaudeConnection ? { openConnection: deps.openClaudeConnection } : {}),
readProcessStartTime: deps.readClaudeProcessStartTime ?? deps.readProcessStartTime,
persistHandle: async (observed) => {
const now = Date.now()
await store.transitionHandoff(observed.sessionId, (record) =>
recordAgentSessionProviderHandle({
record,
fence: record.lease.runtimeFence,
link: claudeProviderHandleLink({
sessionId: observed.providerSessionId,
leafUuid: observed.leafUuid,
resumed: record.providerHandleChain.length > 0,
fence: record.lease.runtimeFence,
observedAt: now
}),
now
})
)
}
})
const adapter = new StructuredAgentSessionAdapterRouter({ claude, codex }, async () => {
await Promise.all([claude.closeAll(), codex.closeAll()])
})
host = new StructuredAgentSessionHost({
store,
adapter,