feat(daemon): hold composer startup until an owner-fenced release

This commit is contained in:
Neil
2026-09-05 03:31:37 -07:00
parent eae9f17236
commit f75781902d
14 changed files with 758 additions and 65 deletions
+76
View File
@@ -10,6 +10,82 @@
}
},
"gates": [
{
"id": "agent-session.deferred-composer-startup",
"title": "Composer command release belongs to one terminal incarnation",
"maturity": "experimental",
"protection": "partial",
"owner": "terminal-runtime",
"layer": "daemon-session-and-subprocess-contract",
"surfaces": ["retained composer shell", "deferred startup command", "shell readiness queue"],
"platforms": ["macos", "linux", "windows"],
"providers": ["daemon", "local", "wsl", "ssh", "remote-runtime"],
"coveredPlatforms": ["macos"],
"coveredProviders": ["daemon"],
"coverageNotes": "Actual Session and TerminalHost contracts run on macOS. Windows and WSL argv branches are mocked. No renderer or wire caller enables deferred startup yet; direct local, remote, and other native platform release paths remain unimplemented.",
"motivatingLinks": ["docs/reference/worktree-create-retained-drafts.md"],
"invariant": "Shell readiness cannot authorize deferred agent execution. Only the same operation and terminal incarnation may release once. Manual input retires an unreleased command; exit or termination prevents a queued command from reaching the subprocess. Ambiguous writes are never replayed.",
"oracle": "Record actual Session subprocess writes while emitting readiness, advancing its timeout, reattaching clients, injecting query replies and manual input, terminating, and throwing after writing bytes. Require no command before release, at most one after it, and no write from wrong identity or after teardown.",
"commands": [
"pnpm test src/main/daemon/session-deferred-startup.test.ts src/main/daemon/terminal-host-startup.test.ts src/main/daemon/session.test.ts src/main/daemon/terminal-host-agent-session.test.ts src/main/daemon/pty-subprocess-windows-shell-launch.test.ts src/main/daemon/pty-subprocess-wsl-launch.test.ts src/main/daemon/pty-subprocess-managed-agent-env.test.ts"
],
"testFiles": [
"src/main/daemon/session-deferred-startup.test.ts",
"src/main/daemon/terminal-host-startup.test.ts",
"src/main/daemon/session.test.ts",
"src/main/daemon/terminal-host-agent-session.test.ts",
"src/main/daemon/pty-subprocess-windows-shell-launch.test.ts",
"src/main/daemon/pty-subprocess-wsl-launch.test.ts",
"src/main/daemon/pty-subprocess-managed-agent-env.test.ts"
],
"assertionRefs": [
{
"file": "src/main/daemon/session-deferred-startup.test.ts",
"assertions": [
"does not launch on the readiness timeout without Create",
"accepts Create before readiness and queues exactly one command through the existing gate"
]
}
],
"evidenceRuns": [
{
"date": "2026-09-05",
"runner": "local",
"platform": "macos",
"command": "pnpm test src/main/daemon/session-deferred-startup.test.ts src/main/daemon/terminal-host-startup.test.ts src/main/daemon/session.test.ts src/main/daemon/terminal-host-agent-session.test.ts src/main/daemon/pty-subprocess-windows-shell-launch.test.ts src/main/daemon/pty-subprocess-wsl-launch.test.ts src/main/daemon/pty-subprocess-managed-agent-env.test.ts",
"result": "passed",
"durationSeconds": 0.546,
"summary": "145 tests passed across seven suites; Session integration tests exposed automatic CPR/DA replies retiring pending startup, fixed with the existing terminal reply parser."
}
],
"runtimeBudget": {
"p95Seconds": 15,
"scope": "Focused owner and shell argv contracts"
},
"flakeHistory": {
"status": "unknown",
"evidence": "Local focused runs only; no CI soak history."
},
"redGreenEvidence": {
"status": "partial",
"evidence": "CPR/DA integration tests failed before reply exclusion and passed after it. No complete intentional-break matrix or real deferred composer journey yet."
},
"performanceBudget": {
"required": true,
"evidence": "One bounded startup record per Session, no polling or subprocess added. Existing readiness queue carries one accepted callback. Input classification runs only while that Session has an unreleased command."
},
"promotionCriteria": [
"Implement capability-negotiated owner release and the rendered composer integration.",
"Verify live local, daemon, Windows/WSL and SSH release and crash recovery.",
"Meet manifest CI and soak policy."
],
"knownGaps": [
"No wire capability or renderer caller; this is an internal prerequisite.",
"No live shell-profile or selected-agent readiness timing for deferred launch.",
"Direct local provider and SSH/remote release parity remain unimplemented."
],
"demotionRule": "Keep experimental until topology, fault-injection and rendered proof exist; demote on premature or repeated launch, stale identity acceptance, or unexplained flakes."
},
{
"id": "terminal-output.prestarted-shell-snapshot-adoption",
"title": "Prestarted shell adoption paints covered output once",
@@ -0,0 +1,277 @@
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
import { Session } from './session'
import type { SubprocessHandle } from './session-subprocess-handle'
const operationId = 'composer-create-operation'
const submission = "codex 'explain this project'\r"
const readyMarker = '\x1b]777;orca-shell-ready\x07$ '
function recordingSubprocess() {
let onData: ((data: string) => void) | undefined
let onExit: ((code: number) => void) | undefined
const written: string[] = []
const write = vi.fn((data: string) => {
written.push(data)
})
const handle: SubprocessHandle = {
pid: 12345,
getForegroundProcess: () => 'bash',
write,
resize: vi.fn(),
kill: vi.fn(),
forceKill: vi.fn(),
signal: vi.fn(),
terminateOwnedTree: () => 'unavailable',
onData: (callback) => {
onData = callback
},
onExit: (callback) => {
onExit = callback
},
dispose: vi.fn()
}
return {
handle,
written,
write,
emit: (data: string) => onData?.(data),
exit: () => onExit?.(0)
}
}
describe('Session deferred startup command', () => {
let session: Session
let subprocess: ReturnType<typeof recordingSubprocess>
beforeEach(() => {
vi.useFakeTimers()
subprocess = recordingSubprocess()
})
afterEach(() => {
session?.dispose()
vi.useRealTimers()
})
function createSession(shellReadySupported = true, deferred = true): void {
session = new Session({
sessionId: 'composer-shell',
cols: 80,
rows: 24,
subprocess: subprocess.handle,
shellReadySupported,
shellReadyTimeoutMs: 1_000,
...(deferred ? { deferredStartup: { operationId, submission } } : {})
})
}
function release() {
return session.releaseStartupCommand(session.incarnationId, operationId)
}
async function becomeReady(): Promise<void> {
subprocess.emit(readyMarker)
await vi.advanceTimersByTimeAsync(250)
}
it('keeps the command held after shell readiness until Create releases it', async () => {
createSession()
await becomeReady()
expect(session.shellState).toBe('ready')
expect(subprocess.written).toEqual([])
expect(release()).toBe('accepted')
expect(subprocess.written).toEqual([submission])
expect(release()).toBe('accepted')
expect(subprocess.written).toEqual([submission])
})
it('accepts Create before readiness and queues exactly one command through the existing gate', async () => {
createSession()
expect(release()).toBe('accepted')
expect(release()).toBe('accepted')
expect(subprocess.written).toEqual([])
await becomeReady()
expect(subprocess.written).toEqual([submission])
expect(release()).toBe('accepted')
subprocess.emit(readyMarker)
await vi.advanceTimersByTimeAsync(2_000)
expect(subprocess.written).toEqual([submission])
})
it('does not launch on the readiness timeout without Create', async () => {
createSession()
await vi.advanceTimersByTimeAsync(2_000)
expect(session.shellState).toBe('timed_out')
expect(subprocess.written).toEqual([])
expect(release()).toBe('accepted')
expect(subprocess.written).toEqual([submission])
})
it('flushes an already accepted Create once when the readiness marker is missing', async () => {
createSession()
expect(release()).toBe('accepted')
await vi.advanceTimersByTimeAsync(2_000)
expect(subprocess.written).toEqual([submission])
expect(release()).toBe('accepted')
expect(subprocess.written).toEqual([submission])
})
it('releases immediately on shells without a readiness marker', () => {
createSession(false)
expect(subprocess.written).toEqual([])
expect(release()).toBe('accepted')
expect(release()).toBe('accepted')
expect(subprocess.written).toEqual([submission])
})
it('rejects stale incarnation and operation identities without consuming the valid release', async () => {
createSession()
await becomeReady()
expect(session.releaseStartupCommand('previous-incarnation', operationId)).toBe(
'identity-mismatch'
)
expect(session.releaseStartupCommand(session.incarnationId, 'previous-operation')).toBe(
'identity-mismatch'
)
expect(subprocess.written).toEqual([])
expect(release()).toBe('accepted')
expect(subprocess.written).toEqual([submission])
})
it('returns unavailable for an ordinary shell without a deferred command', () => {
createSession(false, false)
expect(release()).toBe('unavailable')
expect(subprocess.written).toEqual([])
})
it.each([false, true])(
'manual input retires the unreleased command (ready=%s)',
async (ready) => {
createSession()
if (ready) {
await becomeReady()
}
session.write('echo manual\r')
expect(release()).toBe('retired')
if (!ready) {
await becomeReady()
}
expect(subprocess.written).toEqual(['echo manual\r'])
expect(release()).toBe('retired')
}
)
it('empty input and terminal query replies do not retire the command', async () => {
createSession()
session.write('')
session.write('\x1b]10;rgb:ffff/ffff/ffff\x07')
await becomeReady()
expect(release()).toBe('accepted')
expect(subprocess.written.filter((data) => data === submission)).toEqual([submission])
})
it.each(['\x1b[1;1R', '\x1b[?1;2c'])(
'a non-user terminal reply %j does not retire the command',
async (reply) => {
createSession()
session.write(reply)
await becomeReady()
expect(release()).toBe('accepted')
expect(subprocess.written).toEqual([reply, submission])
}
)
it('retains the command across resize, detach and reattach without running it', async () => {
createSession()
const client = { onData: vi.fn(), onExit: vi.fn() }
const first = session.attachClient(client)
session.resize(120, 40)
session.detachClient(first)
await becomeReady()
session.attachClient(client)
session.detachAllClients()
session.attachClient(client)
expect(subprocess.written).toEqual([])
expect(release()).toBe('accepted')
expect(subprocess.written).toEqual([submission])
})
it.each(['exit', 'dispose', 'termination'] as const)(
'prevents release after %s',
async (stop) => {
createSession()
if (stop === 'exit') {
subprocess.exit()
}
if (stop === 'dispose') {
session.dispose()
}
if (stop === 'termination') {
session.beginTermination()
}
expect(release()).not.toBe('accepted')
await becomeReady()
await vi.advanceTimersByTimeAsync(2_000)
expect(subprocess.written).toEqual([])
}
)
it.each(['exit', 'dispose', 'termination'] as const)(
'does not send an accepted but queued command after %s',
async (stop) => {
createSession()
expect(release()).toBe('accepted')
if (stop === 'exit') {
subprocess.exit()
}
if (stop === 'dispose') {
session.dispose()
}
if (stop === 'termination') {
session.beginTermination()
}
await becomeReady()
await vi.advanceTimersByTimeAsync(2_000)
expect(subprocess.written).toEqual([])
}
)
it('claims release before a subprocess write can synchronously reenter', () => {
createSession(false)
const replies: string[] = []
subprocess.write.mockImplementation((data) => {
subprocess.written.push(data)
replies.push(release())
})
expect(release()).toBe('accepted')
expect(replies).toEqual(['accepted'])
expect(subprocess.written).toEqual([submission])
})
it('does not retry a command when the subprocess writes and then throws', async () => {
createSession()
await becomeReady()
subprocess.write.mockImplementation((data) => {
subprocess.written.push(data)
throw new Error('transport reply lost after delivery')
})
expect(release()).toBe('unverifiable')
expect(release()).toBe('unverifiable')
expect(subprocess.written).toEqual([submission])
})
it('records an ambiguous queued write without throwing from the readiness callback or retrying', async () => {
createSession()
expect(release()).toBe('accepted')
subprocess.write.mockImplementation((data) => {
subprocess.written.push(data)
throw new Error('transport reply lost after queued delivery')
})
await becomeReady()
expect(release()).toBe('unverifiable')
expect(release()).toBe('unverifiable')
expect(subprocess.written).toEqual([submission])
})
})
@@ -0,0 +1,61 @@
export type DeferredSessionStartup = {
operationId: string
submission: string
}
export type StartupCommandReleaseResult =
| 'accepted'
| 'unverifiable'
| 'retired'
| 'identity-mismatch'
| 'unavailable'
/** Keeps Create authorization separate from the shell readiness timeout. */
export class SessionDeferredStartup {
private state: 'pending' | 'accepted' | 'unverifiable' | 'retired' = 'pending'
private submission: string | null
private readonly operationId: string
constructor(startup: DeferredSessionStartup) {
this.operationId = startup.operationId
this.submission = startup.submission
}
get isPending(): boolean {
return this.state === 'pending'
}
retire(): void {
if (this.state === 'pending') {
this.state = 'retired'
this.submission = null
}
}
markUnverifiable(): void {
if (this.state === 'accepted') {
this.state = 'unverifiable'
}
}
release(operationId: string, write: (submission: string) => void): StartupCommandReleaseResult {
if (!operationId || operationId !== this.operationId) {
return 'identity-mismatch'
}
if (this.state !== 'pending') {
return this.state
}
const submission = this.submission
this.submission = null
this.state = 'accepted'
try {
if (submission) {
write(submission)
}
} catch {
// A throwing write can already have delivered bytes; never replay it.
this.state = 'unverifiable'
}
return this.state
}
}
+2
View File
@@ -2,6 +2,7 @@ import type { SubprocessHandle } from './session-subprocess-handle'
import type { TuiAgent } from '../../shared/tui-agent'
import type { PtyStartupIngressIntent } from '../../shared/pty-startup-ingress'
import type { PtyOwnerBackend } from '../../shared/pty-owner-backend'
import type { DeferredSessionStartup } from './session-deferred-startup'
export type SessionOptions = {
sessionId: string
@@ -12,6 +13,7 @@ export type SessionOptions = {
subprocess: SubprocessHandle
shellReadySupported: boolean
shellReadyTimeoutMs?: number
deferredStartup?: DeferredSessionStartup
/** Reports a readiness outcome worth diagnosing to the daemon's file log.
* Why not console: the detached daemon runs with stdio 'ignore'. */
reportReadinessEvent?: (event: string, details: Record<string, unknown>) => void
@@ -46,7 +46,7 @@ export class SessionShellReadyBarrier {
private promptReadinessProbe: ShellPromptReadinessProbe | null = null
private readyTimer: ReturnType<typeof setTimeout> | null = null
private releaseDeviceAttributesResponder: (() => void) | null = null
private preReadyStdinQueue: string[] = []
private preReadyStdinQueue: (string | (() => void))[] = []
private readonly postReadyFlushGate: PostReadyFlushGate
constructor(private readonly deps: SessionShellReadyBarrierDeps) {
@@ -100,7 +100,7 @@ export class SessionShellReadyBarrier {
}
/** Queues `data` when the gate is closed; false means the caller must write it through. */
tryEnqueue(data: string): boolean {
tryEnqueue(data: string | (() => void)): boolean {
if (!this.isGatingWrites) {
return false
}
@@ -244,7 +244,11 @@ export class SessionShellReadyBarrier {
const queued = this.preReadyStdinQueue
this.preReadyStdinQueue = []
for (const data of queued) {
this.deps.subprocess.write(data)
if (typeof data === 'string') {
this.deps.subprocess.write(data)
} else {
data()
}
}
}
}
+119
View File
@@ -0,0 +1,119 @@
import type { SessionOptions } from './session-options'
import type { SessionOutputPlane } from './session-output-plane'
import type { TerminalShellRecoveryBarrier } from './terminal-shell-recovery-barrier'
import { PtyStartupIngress } from '../../shared/pty-startup-ingress'
import { extractOnlyTerminalQueryReplies } from '../../shared/terminal-query-reply'
import {
SessionDeferredStartup,
type DeferredSessionStartup,
type StartupCommandReleaseResult
} from './session-deferred-startup'
import { SessionShellReadyBarrier } from './session-shell-ready-barrier'
import type { SubprocessHandle } from './session-subprocess-handle'
type SessionStartupInputOptions = {
incarnationId: string
isAlive(): boolean
isTerminating(): boolean
subprocess: Pick<SubprocessHandle, 'write'>
ingress: Pick<PtyStartupIngress, 'answerLiveQueryReply'>
shellReady: Pick<SessionShellReadyBarrier, 'tryEnqueue'>
deferredStartup?: DeferredSessionStartup
}
export class SessionStartupInput {
private readonly deferred: SessionDeferredStartup | undefined
private readonly options: Omit<SessionStartupInputOptions, 'deferredStartup'>
constructor({ deferredStartup, ...options }: SessionStartupInputOptions) {
this.options = options
this.deferred = deferredStartup ? new SessionDeferredStartup(deferredStartup) : undefined
}
write(data: string): void {
if (!this.options.isAlive() || this.options.ingress.answerLiveQueryReply(data)) {
return
}
if (this.deferred?.isPending && data.length > 0 && !extractOnlyTerminalQueryReplies(data)) {
this.deferred.retire()
}
// Preserve the post-marker queue until its flush gate opens.
if (!this.options.shellReady.tryEnqueue(data)) {
this.options.subprocess.write(data)
}
}
retire(): void {
this.deferred?.retire()
}
release(expectedIncarnationId: string, operationId: string): StartupCommandReleaseResult {
if (expectedIncarnationId !== this.options.incarnationId) {
return 'identity-mismatch'
}
if (!this.options.isAlive() || this.options.isTerminating()) {
return 'unavailable'
}
return this.deferred?.release(operationId, (data) => this.deliver(data)) ?? 'unavailable'
}
private deliver(data: string): void {
const write = (): void => {
if (!this.options.isAlive() || this.options.isTerminating()) {
return
}
try {
this.options.subprocess.write(data)
} catch {
this.deferred?.markUnverifiable()
}
}
if (!this.options.shellReady.tryEnqueue(write)) {
write()
}
}
}
export function createSessionStartupInput(args: {
opts: SessionOptions
output: SessionOutputPlane
recoveryBarrier: TerminalShellRecoveryBarrier
isAlive(): boolean
isTerminating(): boolean
incarnationId: string
}): {
shellReady: SessionShellReadyBarrier
startupIngress: PtyStartupIngress
input: SessionStartupInput
} {
const { opts, output, recoveryBarrier } = args
const subprocess = opts.subprocess
const shellReady = new SessionShellReadyBarrier({
sessionId: opts.sessionId,
subprocess,
responderParser: output.responderParser,
shellReadySupported: opts.shellReadySupported,
...(opts.reportReadinessEvent ? { reportReadinessEvent: opts.reportReadinessEvent } : {}),
shellReadyTimeoutMs: opts.shellReadyTimeoutMs,
installDeviceAttributesFilter: () => output.installDeviceAttributesFilter(),
releaseDeviceAttributesFilter: () => output.releaseDeviceAttributesFilter(),
acceptStartupIngress: (data) => startupIngress.accept(data)
})
const startupIngress = new PtyStartupIngress({
...(opts.startupIngress ? { intent: opts.startupIngress } : {}),
...(opts.ownerBackend ? { ownerBackend: opts.ownerBackend } : {}),
write: (data) => subprocess.write(data),
onEmission: (emission) => recoveryBarrier.accept(emission)
})
const input = new SessionStartupInput({
incarnationId: args.incarnationId,
isAlive: args.isAlive,
isTerminating: args.isTerminating,
subprocess,
ingress: startupIngress,
shellReady,
deferredStartup: opts.deferredStartup
})
return { shellReady, startupIngress, input }
}
+26 -36
View File
@@ -2,7 +2,7 @@ import { isValidPtySize } from './daemon-pty-size'
import type { SessionOutputPlane, AttachedClient } from './session-output-plane'
import { createSessionOutputPipeline } from './session-output-pipeline'
import { SessionProducerPause } from './session-producer-pause'
import { SessionShellReadyBarrier } from './session-shell-ready-barrier'
import type { SessionShellReadyBarrier } from './session-shell-ready-barrier'
import type { TerminalShellRecoveryBarrier } from './terminal-shell-recovery-barrier'
import {
SessionTerminationController,
@@ -13,7 +13,9 @@ import type { JobTerminationOutcome } from '../windows/windows-pty-job'
import type { SessionOptions } from './session-options'
import type { TuiAgent } from '../../shared/tui-agent'
import { randomUUID } from 'node:crypto'
import { PtyStartupIngress } from '../../shared/pty-startup-ingress'
import type { PtyStartupIngress } from '../../shared/pty-startup-ingress'
import type { StartupCommandReleaseResult } from './session-deferred-startup'
import { createSessionStartupInput, type SessionStartupInput } from './session-startup-input'
import type {
SessionState,
@@ -40,6 +42,7 @@ export class Session {
private readonly termination: SessionTerminationController
private readonly startupIngress: PtyStartupIngress
private readonly recoveryBarrier: TerminalShellRecoveryBarrier
private readonly input: SessionStartupInput
constructor(opts: SessionOptions) {
this.sessionId = opts.sessionId
@@ -68,24 +71,17 @@ export class Session {
releaseProducerPause: (pauseOpts) => this.producerPause.release(pauseOpts)
})
this.shellReady = new SessionShellReadyBarrier({
sessionId: this.sessionId,
subprocess: this.subprocess,
responderParser: this.output.responderParser,
shellReadySupported: opts.shellReadySupported,
...(opts.reportReadinessEvent ? { reportReadinessEvent: opts.reportReadinessEvent } : {}),
shellReadyTimeoutMs: opts.shellReadyTimeoutMs,
installDeviceAttributesFilter: () => this.output.installDeviceAttributesFilter(),
releaseDeviceAttributesFilter: () => this.output.releaseDeviceAttributesFilter(),
acceptStartupIngress: (data) => this.startupIngress.accept(data)
})
this.startupIngress = new PtyStartupIngress({
...(opts.startupIngress ? { intent: opts.startupIngress } : {}),
...(opts.ownerBackend ? { ownerBackend: opts.ownerBackend } : {}),
write: (data) => this.subprocess.write(data),
onEmission: (emission) => this.recoveryBarrier.accept(emission)
const startup = createSessionStartupInput({
opts,
output: this.output,
recoveryBarrier: this.recoveryBarrier,
isAlive: () => !this._disposed && this.isAlive,
isTerminating: () => this.isTerminating,
incarnationId: this.incarnationId
})
this.shellReady = startup.shellReady
this.startupIngress = startup.startupIngress
this.input = startup.input
this.shellReady.startPromptReadinessProbe()
this.subprocess.onData((data) => {
if (!this._disposed) {
@@ -140,23 +136,14 @@ export class Session {
}
write(data: string): void {
if (this._state === 'exited' || this._disposed) {
return
}
this.input.write(data)
}
// Daemon POSIX PTYs need the local provider's cooked-echo containment (#13137).
// DA1/CPR stay immediate unless an echo-risk reply is already held (#13892, #15559).
if (this.startupIngress.answerLiveQueryReply(data)) {
return
}
// Why: keep queuing during the post-ready flush-gate window ('ready' but not yet flushed); a
// direct write would race fresh input ahead of the buffered startup command.
if (this.shellReady.tryEnqueue(data)) {
return
}
this.subprocess.write(data)
releaseStartupCommand(
expectedIncarnationId: string,
operationId: string
): StartupCommandReleaseResult {
return this.input.release(expectedIncarnationId, operationId)
}
resize(cols: number, rows: number): void {
@@ -201,6 +188,7 @@ export class Session {
}
signal(sig: string): void {
this.input.retire()
this.termination.signal(sig)
}
@@ -346,6 +334,7 @@ export class Session {
return
}
this._disposed = true
this.input.retire()
this.output.markDisposed()
// Why: never leave a paused fd behind on teardown; the handle's dead-guard makes this a no-op once the child is reaped.
this.producerPause.release({ resume: true })
@@ -355,6 +344,7 @@ export class Session {
}
private handleSubprocessExit(code: number, cause?: TerminalExitCause): void {
this.input.retire()
this.termination.markPhysicalExit()
if (this._disposed) {
return
@@ -377,7 +367,7 @@ export class Session {
this.termination.cancelForceKillFallback()
this.shellReady.clearReadyTimer()
this.shellReady.clearFlushGate()
this.shellReady.clearPendingWrites()
// Why: release the ptmx fd here or node-pty's _socket leaks the master fd until GC (docs/fix-pty-fd-leak.md).
// Not via #teardownSubprocess: it flips `_disposed`, short-circuiting the later Session.dispose() reaper.
@@ -3,6 +3,7 @@ import type { AgentSessionOwnerBinding } from '../../shared/agent-session-host-a
import type { CreateOrAttachOptions, CreateOrAttachResult } from './terminal-host-create-contract'
export type InternalCreateOrAttachOptions = CreateOrAttachOptions & {
deferredStartupOperationId?: string
agentSessionGeneration?: string
isCanceled?: () => boolean
cancelSignal?: AbortSignal
@@ -38,6 +39,7 @@ export async function createOrAttachClaimedAgentSession(args: {
...args.options,
sessionId: ensured.owner.ptyId,
command: undefined,
deferredStartupOperationId: undefined,
agentSessionEnsure: undefined,
attachOnly: true
})
@@ -60,11 +60,12 @@ describe('TerminalHost agent-session claims', () => {
await host.dispose()
})
it('adopts one claimed provider session across different requested daemon ids', async () => {
it.each([false, true])('adopts claimed sessions with deferred=%s', async (deferred) => {
const first = await host.createOrAttach({
sessionId: 'session-claimed-first',
cols: 80,
rows: 24,
...(deferred ? { command: 'codex', deferredStartupOperationId: 'operation' } : {}),
streamClient: { onData: vi.fn(), onExit: vi.fn() },
agentSessionEnsure: { claim, surface }
})
@@ -72,6 +73,7 @@ describe('TerminalHost agent-session claims', () => {
sessionId: 'session-claimed-retry',
cols: 80,
rows: 24,
...(deferred ? { command: 'codex', deferredStartupOperationId: 'operation' } : {}),
streamClient: { onData: vi.fn(), onExit: vi.fn() },
agentSessionEnsure: {
claim,
@@ -88,6 +90,13 @@ describe('TerminalHost agent-session claims', () => {
owner: { ptyId: 'session-claimed-first', surface }
})
expect(spawnSubprocess).toHaveBeenCalledOnce()
expect(subprocess?.write).not.toHaveBeenCalled()
if (deferred) {
expect(
host.releaseStartupCommand('session-claimed-first', second.incarnationId, 'operation')
).toBe('accepted')
expect(subprocess?.write).toHaveBeenCalledOnce()
}
})
it('cannot adopt a live session that predates provider-session claims', async () => {
+1
View File
@@ -12,6 +12,7 @@ export type TerminalHostOptions = {
env?: Record<string, string>
envToDelete?: string[]
command?: string
deferStartupCommand?: boolean
startupCommandDelivery?: StartupCommandDelivery
launchAgent?: TuiAgent
shellOverride?: string
+21 -10
View File
@@ -119,6 +119,7 @@ async function spawnAndPublishSession(
env: opts.env,
envToDelete: opts.envToDelete,
command: opts.command,
...(opts.deferredStartupOperationId ? { deferStartupCommand: true } : {}),
startupCommandDelivery: opts.startupCommandDelivery,
...(opts.launchAgent ? { launchAgent: opts.launchAgent } : {}),
shellOverride: opts.shellOverride,
@@ -133,6 +134,12 @@ async function spawnAndPublishSession(
const shellReadySupported =
(opts.shellReadySupported ?? false) &&
(subprocess.shellPath === undefined || shellPathSupportsPtyStartupBarrier(subprocess.shellPath))
const startupSubmission = opts.command
? buildStartupCommandSubmission(opts.command, {
submit: process.platform === 'win32' ? '\r' : '\n',
bracketedPasteSafe: shellReadySupported
})
: undefined
const session = new Session({
sessionId: opts.sessionId,
cols: size.cols,
@@ -146,6 +153,14 @@ async function spawnAndPublishSession(
wslDistro
}),
shellReadySupported,
...(opts.deferredStartupOperationId && startupSubmission
? {
deferredStartup: {
operationId: opts.deferredStartupOperationId,
submission: startupSubmission
}
}
: {}),
scrollback: resolveDaemonSessionScrollbackRows(),
historySeedChunks: opts.historySeedChunks,
...(opts.startupIngress ? { startupIngress: opts.startupIngress } : {}),
@@ -174,13 +189,16 @@ async function spawnAndPublishSession(
const token = session.attachClient(opts.streamClient)
const startupCommandWritten =
Boolean(opts.command) && !subprocess.startupCommandDeliveredInShellArgs
Boolean(opts.command) &&
!opts.deferredStartupOperationId &&
!subprocess.startupCommandDeliveredInShellArgs
// Why: without this, a missing command and a lost one log identically.
// Length, never the text -- launches can carry credentials.
try {
deps.reportReadinessEvent?.('startup-command-delivery', {
sessionId: opts.sessionId,
written: startupCommandWritten,
...(opts.deferredStartupOperationId ? { deferred: true } : {}),
hasCommand: Boolean(opts.command),
commandLength: opts.command?.length ?? 0,
viaShellArgs: subprocess.startupCommandDeliveredInShellArgs === true,
@@ -189,15 +207,8 @@ async function spawnAndPublishSession(
} catch {
// Diagnostics must never turn a live PTY into a failed create.
}
if (startupCommandWritten && opts.command) {
const submit = process.platform === 'win32' ? '\r' : '\n'
// Why: only Orca-wrapped shells advertise the paste-safe startup barrier.
session.write(
buildStartupCommandSubmission(opts.command, {
submit,
bracketedPasteSafe: shellReadySupported
})
)
if (startupCommandWritten && startupSubmission) {
session.write(startupSubmission)
}
return {
@@ -0,0 +1,37 @@
import type { Session } from './session'
import type { StartupCommandReleaseResult } from './session-deferred-startup'
import { TerminalAttachCanceledError } from './daemon-errors'
import type { InternalCreateOrAttachOptions } from './terminal-host-agent-session-claim'
export function assertTerminalHostCreateAllowed(
opts: InternalCreateOrAttachOptions,
creationFenced: boolean
): void {
if (
opts.deferredStartupOperationId !== undefined &&
(!opts.deferredStartupOperationId || !opts.command)
) {
throw new Error('Deferred startup requires an operation identity and command')
}
if (creationFenced) {
throw new Error('Terminal host is shutting down')
}
if (opts.isCanceled?.()) {
throw new TerminalAttachCanceledError(opts.sessionId)
}
}
export function createTerminalHostStartupReleaser(
sessions: ReadonlyMap<string, Session>,
isCreationFenced: () => boolean
) {
return (
sessionId: string,
expectedIncarnationId: string,
operationId: string
): StartupCommandReleaseResult =>
isCreationFenced()
? 'unavailable'
: (sessions.get(sessionId)?.releaseStartupCommand(expectedIncarnationId, operationId) ??
'unavailable')
}
@@ -130,3 +130,107 @@ describe('TerminalHost startup command delivery logging', () => {
expect(sub.write).toHaveBeenCalledWith(`codex${process.platform === 'win32' ? '\r' : '\n'}`)
})
})
describe('TerminalHost deferred command ownership', () => {
const command = 'codex DEFERRED_STARTUP_MARKER'
const operationId = 'composer-reservation'
let sub: SubprocessHandle
let host: TerminalHost
let spawn: ReturnType<typeof vi.fn<() => SubprocessHandle>>
beforeEach(() => {
sub = mockSubprocess()
let onExit: ((code: number) => void) | undefined
sub.onExit = (callback) => {
onExit = callback
}
sub.forceKill = vi.fn(() => {
onExit?.(0)
})
spawn = vi.fn(() => sub)
host = new TerminalHost({ spawnSubprocess: spawn })
})
afterEach(async () => {
await host.dispose()
})
async function create() {
return host.createOrAttach({
sessionId: 'retained-shell',
cols: 80,
rows: 24,
command,
deferredStartupOperationId: operationId,
streamClient: { onData: vi.fn(), onExit: vi.fn() }
})
}
it('retains the original command for planning and holds all execution until release', async () => {
const created = await create()
expect(spawn).toHaveBeenCalledWith(
expect.objectContaining({ command, deferStartupCommand: true })
)
expect(sub.write).not.toHaveBeenCalled()
expect(host.releaseStartupCommand('retained-shell', created.incarnationId, operationId)).toBe(
'accepted'
)
expect(sub.write).toHaveBeenCalledOnce()
expect(sub.write).toHaveBeenCalledWith(
`${command}${process.platform === 'win32' ? '\r' : '\n'}`
)
expect(host.releaseStartupCommand('retained-shell', created.incarnationId, operationId)).toBe(
'accepted'
)
expect(sub.write).toHaveBeenCalledOnce()
})
it('never spawns on unknown release and rejects another operation or incarnation', async () => {
expect(host.releaseStartupCommand('missing', 'old', operationId)).toBe('unavailable')
expect(spawn).not.toHaveBeenCalled()
const created = await create()
expect(host.releaseStartupCommand('retained-shell', 'old', operationId)).toBe(
'identity-mismatch'
)
expect(host.releaseStartupCommand('retained-shell', created.incarnationId, 'other')).toBe(
'identity-mismatch'
)
expect(sub.write).not.toHaveBeenCalled()
})
it('reattaches without releasing or replacing the original pending command', async () => {
const original = await create()
const attached = await create()
expect(attached.isNew).toBe(false)
expect(attached.incarnationId).toBe(original.incarnationId)
expect(spawn).toHaveBeenCalledOnce()
expect(sub.write).not.toHaveBeenCalled()
expect(host.releaseStartupCommand('retained-shell', attached.incarnationId, operationId)).toBe(
'accepted'
)
expect(sub.write).toHaveBeenCalledOnce()
})
it('retires an unused launch after manual input into the retained shell', async () => {
const created = await create()
host.write('retained-shell', 'vim\r')
expect(host.releaseStartupCommand('retained-shell', created.incarnationId, operationId)).toBe(
'retired'
)
expect(sub.write).toHaveBeenCalledExactlyOnceWith('vim\r')
})
it('rejects an empty operation identity before spawning', async () => {
await expect(
host.createOrAttach({
sessionId: 'invalid',
cols: 80,
rows: 24,
command,
deferredStartupOperationId: '',
streamClient: { onData: vi.fn(), onExit: vi.fn() }
})
).rejects.toThrow('Deferred startup requires')
expect(spawn).not.toHaveBeenCalled()
})
})
+15 -15
View File
@@ -19,7 +19,10 @@ import { resolveTerminalHostSessionCwd } from './terminal-host-session-cwd'
import { TerminalHostTombstones } from './terminal-host-tombstones'
import { listLiveTerminalHostSessions } from './terminal-host-session-listing'
import { createOrAttachTerminalSession } from './terminal-host-session-create'
import { TerminalAttachCanceledError } from './daemon-errors'
import {
assertTerminalHostCreateAllowed,
createTerminalHostStartupReleaser
} from './terminal-host-startup-operations'
import { rejectOnAbort } from './terminal-attach-cancellation'
import { randomUUID } from 'node:crypto'
import { pruneRetiredPtyIncarnations } from '../../shared/retired-pty-incarnations'
@@ -76,7 +79,7 @@ export class TerminalHost {
}
async createOrAttach(opts: InternalCreateOrAttachOptions): Promise<CreateOrAttachResult> {
this.assertCreateOrAttachAllowed(opts)
assertTerminalHostCreateAllowed(opts, this.creationFenced)
for (
let inFlight = this.pendingCreations.get(opts.sessionId);
inFlight !== undefined;
@@ -86,9 +89,9 @@ export class TerminalHost {
// minutes. Waiting unconditionally is what let one dead path strand every
// later create and attach for the session, so a canceled caller leaves.
await Promise.race([inFlight, rejectOnAbort(opts.cancelSignal, opts.sessionId)])
this.assertCreateOrAttachAllowed(opts)
assertTerminalHostCreateAllowed(opts, this.creationFenced)
}
this.assertCreateOrAttachAllowed(opts)
assertTerminalHostCreateAllowed(opts, this.creationFenced)
let settleCreation: () => void = () => {}
this.pendingCreations.set(
@@ -107,13 +110,14 @@ export class TerminalHost {
Boolean(this.sessions.get(owner.ptyId)?.isAlive)
),
createOrAttach: async (options) => {
this.assertCreateOrAttachAllowed(options)
assertTerminalHostCreateAllowed(options, this.creationFenced)
if (options.agentSessionGeneration && this.sessions.get(options.sessionId)?.isAlive) {
throw new Error('agent_session_claim_unavailable')
}
return await createOrAttachTerminalSession(options, {
sessions: this.sessions,
assertCreateAllowed: () => this.assertCreateOrAttachAllowed(options),
assertCreateAllowed: () =>
assertTerminalHostCreateAllowed(options, this.creationFenced),
sessionTeardown: this.sessionTeardown,
killedTombstones: this.killedTombstones,
spawnSubprocess: this.spawnSubprocess,
@@ -146,19 +150,15 @@ export class TerminalHost {
}
}
private assertCreateOrAttachAllowed(opts: InternalCreateOrAttachOptions): void {
if (this.creationFenced) {
throw new Error('Terminal host is shutting down')
}
if (opts.isCanceled?.()) {
throw new TerminalAttachCanceledError(opts.sessionId)
}
}
write(sessionId: string, data: string): void {
this.getAliveSession(sessionId).write(data)
}
readonly releaseStartupCommand = createTerminalHostStartupReleaser(
this.sessions,
() => this.creationFenced
)
closeStartupQueryAuthority(sessionId: string): number {
return this.getAliveSession(sessionId).closeStartupQueryAuthority()
}