Files
orca/src/cli/runtime/websocket-transport.test.ts
T
Brennan BensonandMerge Sim 5868fdc9e3 feat(native-chat): report Codex background tasks in the chat strip (#19346)
* feat(native-chat): report Codex background tasks in the chat strip

The background-tasks strip works for Claude only; a structured Codex
session shows nothing in it. Feed it from the Codex app-server stream.

The strip stands for work that OUTLIVED a turn, which is what the
monitoring header, Claude's foreground suppression, and the conversation
command gate all already assume. Codex has no `is_backgrounded` flag, so
that fact is derived from the turn boundary: a `subAgentActivity` child or
a primary-thread `commandExecution` becomes visible once the turn it
belongs to completes and it is still unsettled.

`turn/completed` only reveals a task here, never settles one — measured on
`codex app-server` 0.153.4, a spawn_agent child reported `completed` 95.8s
after its parent turn ended. Only a child's own activity kind settles it.

Codex exposes no honest stop: `turn/interrupt` on a child ends its turn
without emitting a terminal activity item and leaves its shell running. So
the state carries a new optional `supportsStopAll: false`, the strip hides
a control that could not act, and the blocked-command message asks the user
to wait rather than to press a button that does not exist.

* refactor(codex): move session teardown out of the structured adapter

Merging main crossed the 300-line cap on
`codex-structured-session-adapter.ts`: the rewind backend (#19235) and this
branch's close-time strip clear both landed in it. The four close paths move
verbatim into `codex-structured-session-teardown.ts`, where they funnel
through one `settled` helper instead of repeating the notification-retry and
background-task cleanup at each call site. No ratchet bump.

Also normalize a background task's description once at receipt rather than on
every projection; the roster is re-projected on each observed frame.

* fix(codex): drop the shell row the journal already settles

A `commandExecution` still `inProgress` when its turn ends was reported as a
`command` task. But `settleCodexJournalTurn` writes exactly those items to the
journal as `state: 'failed'` on `turn/completed` and forgets them, so the strip
row would have claimed a shell was still running at the same instant Orca
recorded that it was not — two surfaces contradicting each other about the same
process.

A subagent is the opposite case and stays: the roster pointedly does not sweep
at a turn boundary, because children measurably outlive it. That leaves the
producer making exactly one claim — these spawn_agent children are still live
after their turn — which the durable roster row corroborates.

* fix(native-chat): track Codex background execution lifetimes

* fix(native-chat): keep running tool groups from claiming completion

* Fix runtime catalog and capability expectation

* fix(codex): keep a child's name on the command row that outlives it

A child agent's commands stay hidden behind its agent row while the child
works. Once the child's turn settles with a command still running, that
command surfaces as its own row labelled from the raw command string, so
'long_probe' became "/bin/zsh -lc 'ping -c 300 127.0.0.1 > /dev/null'"
at the moment that row was the only remaining signal for the work.

Qualify a child's command row with the child's label. Resolved on read,
so a label registered after the command still lands, and bounded by the
existing description cap so admission accounting stays valid. Primary-
thread commands are left unqualified: they have no child to name.

---------

Co-authored-by: Merge Sim <sim@local>
2026-09-09 00:13:20 -07:00

396 lines
13 KiB
TypeScript

import { createServer, type Server } from 'node:http'
import { mkdtempSync } from 'node:fs'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterEach, describe, expect, it, vi } from 'vitest'
import { WebSocketServer } from 'ws'
import { encodePairingOffer, type PairingOffer } from '../../shared/pairing'
import type { RuntimeStatus } from '../../shared/runtime-types'
import {
decrypt,
deriveSharedKey,
encrypt,
generateKeyPair,
publicKeyToBase64
} from '../../shared/e2ee-crypto'
import { RuntimeClient } from './client'
import { launchOrcaApp } from './launch'
import { addEnvironmentFromPairingCode } from './environments'
import { RuntimeClientError } from './types'
import {
AGENT_SESSION_BACKGROUND_TASK_STOP_CAPABILITY,
AGENT_SESSION_BOUNDARY_RUNTIME_CAPABILITY,
AUTOMATION_OWNER_FENCING_RUNTIME_CAPABILITY,
MIN_COMPATIBLE_RUNTIME_CLIENT_VERSION,
RUNTIME_PROTOCOL_VERSION,
SESSION_TABS_AUTHORITATIVE_INVENTORY_RUNTIME_CAPABILITY,
SESSION_TAB_CLOSE_INTENT_RUNTIME_CAPABILITY,
SKILL_INSTALL_RESULT_V2_CAPABILITY,
WORKTREE_GITHUB_PR_SUPPRESSION_RUNTIME_CAPABILITY,
WORKTREE_VISIBILITY_DEFAULTS_RUNTIME_CAPABILITY,
WORKTREE_VISIBILITY_SOURCE_DEFAULTS_RUNTIME_CAPABILITY
} from '../../shared/protocol-version'
vi.mock('./launch', () => ({
launchOrcaApp: vi.fn()
}))
type TestRuntime = {
endpoint: string
publicKeyB64: string
deviceToken: string
authFrames: Record<string, unknown>[]
requestMethods: string[]
connectionCount: () => number
close: () => Promise<void>
}
describe('CLI remote WebSocket transport', () => {
const servers: TestRuntime[] = []
afterEach(async () => {
vi.mocked(launchOrcaApp).mockClear()
await Promise.all(servers.splice(0).map((server) => server.close()))
})
it('calls a remote runtime through a mobile pairing offer', async () => {
const runtime = await startTestRuntime('runtime-ws-1')
servers.push(runtime)
const pairingUrl = encodePairingOffer({
v: 2,
endpoint: runtime.endpoint,
deviceToken: runtime.deviceToken,
publicKeyB64: runtime.publicKeyB64
})
const client = new RuntimeClient('/tmp/unused', 5_000, pairingUrl)
const response = await client.call<{ runtimeId: string }>('status.get')
expect(response.ok).toBe(true)
expect(response.result.runtimeId).toBe('runtime-ws-1')
expect(runtime.authFrames).toContainEqual(
expect.objectContaining({
clientCapabilities: [
AGENT_SESSION_BACKGROUND_TASK_STOP_CAPABILITY,
SESSION_TAB_CLOSE_INTENT_RUNTIME_CAPABILITY,
SESSION_TABS_AUTHORITATIVE_INVENTORY_RUNTIME_CAPABILITY,
AGENT_SESSION_BOUNDARY_RUNTIME_CAPABILITY,
SKILL_INSTALL_RESULT_V2_CAPABILITY,
WORKTREE_GITHUB_PR_SUPPRESSION_RUNTIME_CAPABILITY,
WORKTREE_VISIBILITY_DEFAULTS_RUNTIME_CAPABILITY,
WORKTREE_VISIBILITY_SOURCE_DEFAULTS_RUNTIME_CAPABILITY,
AUTOMATION_OWNER_FENCING_RUNTIME_CAPABILITY
]
})
)
})
it('rejects malformed remote pairing codes before local runtime lookup', () => {
expect(() => new RuntimeClient('/tmp/unused', 5_000, 'not-a-pairing-code')).toThrow(
RuntimeClientError
)
})
it('accepts a bare pairing payload as well as the orca URL wrapper', async () => {
const runtime = await startTestRuntime('runtime-ws-2', {
appVersion: '1.5.0',
remoteUpdateSupport: {
installMode: 'unsupported-headless-serve',
automatic: false,
reason: 'manual-service-update-required'
},
capabilities: ['updater.remote-control.v1'],
degradations: [
{
code: 'browser_unavailable',
capability: 'browser.headless.v1',
message: 'Browser automation is unavailable.'
}
]
})
servers.push(runtime)
const offer: PairingOffer = {
v: 2,
endpoint: runtime.endpoint,
deviceToken: runtime.deviceToken,
publicKeyB64: runtime.publicKeyB64
}
const pairingUrl = encodePairingOffer(offer)
const barePayload = new URLSearchParams(pairingUrl.slice(pairingUrl.indexOf('?') + 1)).get(
'code'
)!
const client = new RuntimeClient('/tmp/unused', 5_000, barePayload)
const status = await client.getCliStatus()
expect(status.result.app).toEqual({ running: false, pid: null })
expect(status.result.runtime.reachable).toBe(true)
expect(status.result.runtime.runtimeId).toBe('runtime-ws-2')
expect(status.result.runtime).toMatchObject({
appVersion: '1.5.0',
remoteUpdateSupport: { automatic: false, reason: 'manual-service-update-required' },
capabilities: ['updater.remote-control.v1'],
degradations: [expect.objectContaining({ code: 'browser_unavailable' })]
})
})
it('does not launch a local desktop app for remote-paired open', async () => {
const runtime = await startTestRuntime('runtime-remote-headless', {
desktopWindowStatus: 'initializing'
})
servers.push(runtime)
const client = new RuntimeClient(
'/tmp/unused',
5_000,
encodePairingOffer({
v: 2,
endpoint: runtime.endpoint,
deviceToken: runtime.deviceToken,
publicKeyB64: runtime.publicKeyB64
})
)
const status = await client.openOrca()
expect(status.result.app.desktopWindowStatus).toBe('initializing')
expect(launchOrcaApp).not.toHaveBeenCalled()
})
it('connects through a saved environment selector', async () => {
const runtime = await startTestRuntime('runtime-env-1')
servers.push(runtime)
const userDataPath = mkdtempSync(join(tmpdir(), 'orca-cli-env-'))
addEnvironmentFromPairingCode(userDataPath, {
name: 'remote-dev',
pairingCode: encodePairingOffer({
v: 2,
endpoint: runtime.endpoint,
deviceToken: runtime.deviceToken,
publicKeyB64: runtime.publicKeyB64
})
})
const client = new RuntimeClient(userDataPath, 5_000, null, 'remote-dev')
const status = await client.getCliStatus()
expect(status.result.app).toEqual({ running: false, pid: null })
expect(status.result.runtime.reachable).toBe(true)
expect(status.result.runtime.runtimeId).toBe('runtime-env-1')
})
it('blocks remote RPCs when the server protocol is too old', async () => {
const runtime = await startTestRuntime('runtime-old', { runtimeProtocolVersion: 1 })
servers.push(runtime)
const client = new RuntimeClient(
'/tmp/unused',
5_000,
encodePairingOffer({
v: 2,
endpoint: runtime.endpoint,
deviceToken: runtime.deviceToken,
publicKeyB64: runtime.publicKeyB64
})
)
await expect(client.call('repo.list')).rejects.toMatchObject({
code: 'incompatible_runtime',
message: expect.stringContaining('server is too old')
})
expect(runtime.connectionCount()).toBe(1)
expect(runtime.authFrames).toHaveLength(1)
expect(runtime.requestMethods).toEqual(['status.get'])
})
it('preflights and dispatches through one authenticated connection', async () => {
const runtime = await startTestRuntime('runtime-single-auth')
servers.push(runtime)
const client = new RuntimeClient(
'/tmp/unused',
5_000,
encodePairingOffer({
v: 2,
endpoint: runtime.endpoint,
deviceToken: runtime.deviceToken,
publicKeyB64: runtime.publicKeyB64
})
)
await expect(client.call('repo.list')).rejects.toMatchObject({ code: 'method_not_found' })
expect(runtime.connectionCount()).toBe(1)
expect(runtime.authFrames).toHaveLength(1)
expect(runtime.requestMethods).toEqual(['status.get', 'repo.list'])
})
it('blocks orchestration mutations when a remote runtime lacks the contract capability', async () => {
const runtime = await startTestRuntime('runtime-old-orchestration', { capabilities: [] })
servers.push(runtime)
const client = new RuntimeClient(
'/tmp/unused',
5_000,
encodePairingOffer({
v: 2,
endpoint: runtime.endpoint,
deviceToken: runtime.deviceToken,
publicKeyB64: runtime.publicKeyB64
})
)
await expect(client.call('orchestration.send', { subject: 'hello' })).rejects.toMatchObject({
code: 'orchestration_migration_required',
data: {
reason: 'runtime_capability_missing',
effectsApplied: false
}
})
})
})
async function startTestRuntime(
runtimeId: string,
statusOverrides: {
runtimeProtocolVersion?: number
minCompatibleRuntimeClientVersion?: number
desktopWindowStatus?: 'available' | 'openable' | 'initializing' | 'blocked'
appVersion?: string
remoteUpdateSupport?: {
installMode: 'unsupported-headless-serve'
automatic: false
reason: 'manual-service-update-required'
}
capabilities?: string[]
degradations?: RuntimeStatus['degradations']
} = {}
): Promise<TestRuntime> {
const serverKeyPair = generateKeyPair()
const deviceToken = `token-${runtimeId}`
const httpServer = createServer()
const wss = new WebSocketServer({ server: httpServer })
const authFrames: Record<string, unknown>[] = []
const requestMethods: string[] = []
let connectionCount = 0
wss.on('connection', (ws) => {
connectionCount += 1
let sharedKey: Uint8Array | null = null
let authenticated = false
ws.on('message', (data) => {
const frame = data.toString()
if (!sharedKey) {
const hello = JSON.parse(frame) as Record<string, unknown> & {
type?: string
publicKeyB64?: string
}
const clientPublicKey = Buffer.from(hello.publicKeyB64 ?? '', 'base64')
sharedKey = deriveSharedKey(serverKeyPair.secretKey, clientPublicKey)
ws.send(JSON.stringify({ type: 'e2ee_ready' }))
return
}
const plaintext = decrypt(frame, sharedKey)
if (!plaintext) {
ws.close(4003, 'decrypt failed')
return
}
if (!authenticated) {
const auth = JSON.parse(plaintext) as Record<string, unknown> & {
type?: string
deviceToken?: string
}
authFrames.push(auth)
if (auth.type !== 'e2ee_auth' || auth.deviceToken !== deviceToken) {
ws.send(encrypt(JSON.stringify({ type: 'e2ee_error' }), sharedKey))
ws.close(4001, 'auth failed')
return
}
authenticated = true
ws.send(encrypt(JSON.stringify({ type: 'e2ee_authenticated' }), sharedKey))
return
}
const request = JSON.parse(plaintext) as { id: string; method: string }
requestMethods.push(request.method)
const response =
request.method === 'status.get'
? {
id: request.id,
ok: true,
result: {
runtimeId,
rendererGraphEpoch: 1,
graphStatus: 'ready',
authoritativeWindowId: null,
desktopWindowStatus: statusOverrides.desktopWindowStatus,
liveTabCount: 0,
liveLeafCount: 0,
runtimeProtocolVersion:
statusOverrides.runtimeProtocolVersion ?? RUNTIME_PROTOCOL_VERSION,
minCompatibleRuntimeClientVersion:
statusOverrides.minCompatibleRuntimeClientVersion ??
MIN_COMPATIBLE_RUNTIME_CLIENT_VERSION,
appVersion: statusOverrides.appVersion,
remoteUpdateSupport: statusOverrides.remoteUpdateSupport,
capabilities: statusOverrides.capabilities,
degradations: statusOverrides.degradations
},
_meta: { runtimeId }
}
: {
id: request.id,
ok: false,
error: { code: 'method_not_found', message: 'Unknown method' },
_meta: { runtimeId }
}
ws.send(encrypt(JSON.stringify(response), sharedKey))
})
})
await listen(httpServer)
const address = httpServer.address()
if (!address || typeof address === 'string') {
throw new Error('Expected TCP test server')
}
return {
endpoint: `ws://127.0.0.1:${address.port}`,
publicKeyB64: publicKeyToBase64(serverKeyPair.publicKey),
deviceToken,
authFrames,
requestMethods,
connectionCount: () => connectionCount,
close: async () => {
await new Promise<void>((resolve) => {
wss.close(() => resolve())
for (const client of wss.clients) {
client.close()
}
})
await closeHttpServer(httpServer)
}
}
}
async function listen(server: Server): Promise<void> {
await new Promise<void>((resolve, reject) => {
server.once('error', reject)
server.listen(0, '127.0.0.1', () => {
server.off('error', reject)
resolve()
})
})
}
async function closeHttpServer(server: Server): Promise<void> {
await new Promise<void>((resolve, reject) => {
server.close((error) => {
if (error) {
reject(error)
return
}
resolve()
})
})
}