From f75781902db3ee0be844ebc3c9e9efa58dfb4395 Mon Sep 17 00:00:00 2001 From: Neil <4138956+nwparker@users.noreply.github.com> Date: Sat, 5 Sep 2026 03:31:37 -0700 Subject: [PATCH] feat(daemon): hold composer startup until an owner-fenced release --- config/reliability-gates.jsonc | 76 +++++ .../daemon/session-deferred-startup.test.ts | 277 ++++++++++++++++++ src/main/daemon/session-deferred-startup.ts | 61 ++++ src/main/daemon/session-options.ts | 2 + .../daemon/session-shell-ready-barrier.ts | 10 +- src/main/daemon/session-startup-input.ts | 119 ++++++++ src/main/daemon/session.ts | 62 ++-- .../terminal-host-agent-session-claim.ts | 2 + .../terminal-host-agent-session.test.ts | 11 +- src/main/daemon/terminal-host-options.ts | 1 + .../daemon/terminal-host-session-create.ts | 31 +- .../terminal-host-startup-operations.ts | 37 +++ src/main/daemon/terminal-host-startup.test.ts | 104 +++++++ src/main/daemon/terminal-host.ts | 30 +- 14 files changed, 758 insertions(+), 65 deletions(-) create mode 100644 src/main/daemon/session-deferred-startup.test.ts create mode 100644 src/main/daemon/session-deferred-startup.ts create mode 100644 src/main/daemon/session-startup-input.ts create mode 100644 src/main/daemon/terminal-host-startup-operations.ts diff --git a/config/reliability-gates.jsonc b/config/reliability-gates.jsonc index 1fed7b8315d..ce4d98541b3 100644 --- a/config/reliability-gates.jsonc +++ b/config/reliability-gates.jsonc @@ -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", diff --git a/src/main/daemon/session-deferred-startup.test.ts b/src/main/daemon/session-deferred-startup.test.ts new file mode 100644 index 00000000000..5af981e5821 --- /dev/null +++ b/src/main/daemon/session-deferred-startup.test.ts @@ -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 + + 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 { + 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]) + }) +}) diff --git a/src/main/daemon/session-deferred-startup.ts b/src/main/daemon/session-deferred-startup.ts new file mode 100644 index 00000000000..084a666aa10 --- /dev/null +++ b/src/main/daemon/session-deferred-startup.ts @@ -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 + } +} diff --git a/src/main/daemon/session-options.ts b/src/main/daemon/session-options.ts index 1e2d38814d4..17e39e75441 100644 --- a/src/main/daemon/session-options.ts +++ b/src/main/daemon/session-options.ts @@ -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) => void diff --git a/src/main/daemon/session-shell-ready-barrier.ts b/src/main/daemon/session-shell-ready-barrier.ts index 6fce89af89e..55b2fa97d58 100644 --- a/src/main/daemon/session-shell-ready-barrier.ts +++ b/src/main/daemon/session-shell-ready-barrier.ts @@ -46,7 +46,7 @@ export class SessionShellReadyBarrier { private promptReadinessProbe: ShellPromptReadinessProbe | null = null private readyTimer: ReturnType | 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() + } } } } diff --git a/src/main/daemon/session-startup-input.ts b/src/main/daemon/session-startup-input.ts new file mode 100644 index 00000000000..1427bb5ecb4 --- /dev/null +++ b/src/main/daemon/session-startup-input.ts @@ -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 + ingress: Pick + shellReady: Pick + deferredStartup?: DeferredSessionStartup +} + +export class SessionStartupInput { + private readonly deferred: SessionDeferredStartup | undefined + private readonly options: Omit + + 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 } +} diff --git a/src/main/daemon/session.ts b/src/main/daemon/session.ts index b0f1dfa538d..1d4e8726cbf 100644 --- a/src/main/daemon/session.ts +++ b/src/main/daemon/session.ts @@ -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. diff --git a/src/main/daemon/terminal-host-agent-session-claim.ts b/src/main/daemon/terminal-host-agent-session-claim.ts index 54251e40df2..9c5fc59764c 100644 --- a/src/main/daemon/terminal-host-agent-session-claim.ts +++ b/src/main/daemon/terminal-host-agent-session-claim.ts @@ -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 }) diff --git a/src/main/daemon/terminal-host-agent-session.test.ts b/src/main/daemon/terminal-host-agent-session.test.ts index 1cf929fb93f..02d33d9191a 100644 --- a/src/main/daemon/terminal-host-agent-session.test.ts +++ b/src/main/daemon/terminal-host-agent-session.test.ts @@ -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 () => { diff --git a/src/main/daemon/terminal-host-options.ts b/src/main/daemon/terminal-host-options.ts index 9bb5ff48c0f..d003c1158af 100644 --- a/src/main/daemon/terminal-host-options.ts +++ b/src/main/daemon/terminal-host-options.ts @@ -12,6 +12,7 @@ export type TerminalHostOptions = { env?: Record envToDelete?: string[] command?: string + deferStartupCommand?: boolean startupCommandDelivery?: StartupCommandDelivery launchAgent?: TuiAgent shellOverride?: string diff --git a/src/main/daemon/terminal-host-session-create.ts b/src/main/daemon/terminal-host-session-create.ts index 8f6833c3d9f..97866e9a3e5 100644 --- a/src/main/daemon/terminal-host-session-create.ts +++ b/src/main/daemon/terminal-host-session-create.ts @@ -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 { diff --git a/src/main/daemon/terminal-host-startup-operations.ts b/src/main/daemon/terminal-host-startup-operations.ts new file mode 100644 index 00000000000..85eb76312be --- /dev/null +++ b/src/main/daemon/terminal-host-startup-operations.ts @@ -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, + isCreationFenced: () => boolean +) { + return ( + sessionId: string, + expectedIncarnationId: string, + operationId: string + ): StartupCommandReleaseResult => + isCreationFenced() + ? 'unavailable' + : (sessions.get(sessionId)?.releaseStartupCommand(expectedIncarnationId, operationId) ?? + 'unavailable') +} diff --git a/src/main/daemon/terminal-host-startup.test.ts b/src/main/daemon/terminal-host-startup.test.ts index e7b7c7b4c09..c124b912acf 100644 --- a/src/main/daemon/terminal-host-startup.test.ts +++ b/src/main/daemon/terminal-host-startup.test.ts @@ -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 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() + }) +}) diff --git a/src/main/daemon/terminal-host.ts b/src/main/daemon/terminal-host.ts index 18f7b82a90f..4e90ecdfa42 100644 --- a/src/main/daemon/terminal-host.ts +++ b/src/main/daemon/terminal-host.ts @@ -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 { - 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() }