diff --git a/src/main/agent-hooks/server-codex-ctrl-c-inference.test.ts b/src/main/agent-hooks/server-codex-ctrl-c-inference.test.ts new file mode 100644 index 00000000000..9a7100fa031 --- /dev/null +++ b/src/main/agent-hooks/server-codex-ctrl-c-inference.test.ts @@ -0,0 +1,51 @@ +import { afterEach, describe, expect, it, vi } from 'vitest' +import { AgentHookServer, _internals } from './server' +import { PANE } from './server.test-fixtures' + +afterEach(() => { + _internals.resetCachesForTests() + vi.useRealTimers() +}) + +describe('Codex Ctrl+C status inference', () => { + it.each([null, 'ssh-connection'])('preserves working status on host %s', (connectionId) => { + vi.useFakeTimers() + vi.setSystemTime(1_000) + const server = new AgentHookServer() + server.ingestRemote( + { + paneKey: PANE, + tabId: 'tab-1', + worktreeId: 'folder-1', + hookEventName: 'UserPromptSubmit', + payload: { state: 'working', prompt: 'main task', agentType: 'codex' } + }, + connectionId + ) + const baseline = server.getStatusSnapshot()[0] + const request = { + paneKey: PANE, + baselineUpdatedAt: baseline.receivedAt, + baselineStateStartedAt: baseline.stateStartedAt, + baselinePrompt: baseline.prompt, + baselineAgentType: baseline.agentType, + intent: 'ctrl-c' as const + } + vi.setSystemTime(1_500) + expect(server.inferInterrupt(request)).toBe(false) + expect(server.getStatusSnapshot()).toEqual([baseline]) + + server.ingestRemote( + { + paneKey: PANE, + tabId: 'tab-1', + worktreeId: 'folder-1', + hookEventName: 'Stop', + payload: { state: 'done', prompt: 'main task', agentType: 'codex' } + }, + connectionId + ) + expect(server.getStatusSnapshot()[0]).toMatchObject({ state: 'done' }) + expect(server.getStatusSnapshot()[0].interrupted).toBeUndefined() + }) +}) diff --git a/src/main/agent-hooks/server-codex-turn-interruption.test.ts b/src/main/agent-hooks/server-codex-turn-interruption.test.ts new file mode 100644 index 00000000000..b5ab96d54f2 --- /dev/null +++ b/src/main/agent-hooks/server-codex-turn-interruption.test.ts @@ -0,0 +1,132 @@ +import { afterEach, describe, expect, it, vi } from 'vitest' +import { appendFileSync, mkdtempSync, rmSync, writeFileSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { AgentHookServer } from './server' +import { PANE } from './server.test-fixtures' + +const line = (payload: unknown): string => `${JSON.stringify({ type: 'event_msg', payload })}\n` + +describe('Codex recorded turn interruption', () => { + const dirs: string[] = [] + afterEach(() => { + for (const dir of dirs) { + rmSync(dir, { recursive: true, force: true }) + } + dirs.length = 0 + }) + + it.each([false, true])('settles the lead and preserves child work: %s', async (withChild) => { + const dir = mkdtempSync(join(tmpdir(), 'codex-turn-interruption-')) + dirs.push(dir) + const transcriptPath = join(dir, 'rollout-root.jsonl') + const childPath = join(dir, 'rollout-child-child-1.jsonl') + writeFileSync(transcriptPath, line({ type: 'task_started', turn_id: 'turn-1' })) + if (withChild) { + appendFileSync( + transcriptPath, + line({ + type: 'sub_agent_activity', + agent_thread_id: 'child-1', + kind: 'started', + agent_path: '/root/child', + occurred_at_ms: Date.now() + }) + ) + writeFileSync(childPath, line({ type: 'task_started' })) + } + const server = new AgentHookServer() + await server.start({ env: 'production' }) + try { + const env = server.buildPtyEnv() + const response = await fetch(`http://127.0.0.1:${env.ORCA_AGENT_HOOK_PORT}/hook/codex`, { + method: 'POST', + headers: { + 'Content-Type': 'application/json', + 'X-Orca-Agent-Hook-Token': env.ORCA_AGENT_HOOK_TOKEN + }, + body: JSON.stringify({ + paneKey: PANE, + tabId: 'tab-1', + worktreeId: 'folder-1', + payload: { + hook_event_name: 'UserPromptSubmit', + session_id: 'main-session', + prompt: 'main task', + transcript_path: transcriptPath + } + }) + }) + expect(response.status).toBe(204) + const beforeSide = server.getStatusSnapshot()[0] + for (const hookEventName of ['SessionStart', 'UserPromptSubmit', 'Stop']) { + const sideResponse = await fetch( + `http://127.0.0.1:${env.ORCA_AGENT_HOOK_PORT}/hook/codex`, + { + method: 'POST', + headers: { + 'Content-Type': 'application/json', + 'X-Orca-Agent-Hook-Token': env.ORCA_AGENT_HOOK_TOKEN + }, + body: JSON.stringify({ + paneKey: PANE, + tabId: 'tab-1', + worktreeId: 'folder-1', + payload: { + hook_event_name: hookEventName, + session_id: 'side-session', + transcript_path: null, + prompt: 'side chat' + } + }) + } + ) + expect(sideResponse.status).toBe(204) + expect(server.getStatusSnapshot()[0]).toEqual(beforeSide) + } + const baseline = server.getStatusSnapshot()[0] + expect( + server.inferInterrupt({ + paneKey: PANE, + baselineUpdatedAt: baseline.receivedAt, + baselineStateStartedAt: baseline.stateStartedAt, + baselinePrompt: baseline.prompt, + baselineAgentType: 'codex', + intent: 'ctrl-c' + }) + ).toBe(false) + await new Promise((resolve) => setTimeout(resolve, 600)) + expect(server.getStatusSnapshot()[0].state).toBe('working') + appendFileSync( + transcriptPath, + line({ type: 'turn_aborted', turn_id: 'turn-1', reason: 'interrupted' }) + ) + await vi.waitFor( + () => { + expect(server.getStatusSnapshot()[0]).toMatchObject({ + state: withChild ? 'working' : 'done', + mainAgent: { state: 'done', outcome: 'cancellation' } + }) + }, + { timeout: 2_000 } + ) + if (withChild) { + expect(server.getStatusSnapshot()[0].interrupted).toBeUndefined() + appendFileSync(childPath, line({ type: 'task_complete' })) + await vi.waitFor( + () => + expect(server.getStatusSnapshot()[0]).toMatchObject({ + state: 'done', + interrupted: true, + mainAgent: { state: 'done', outcome: 'cancellation' } + }), + { timeout: 2_000 } + ) + } else { + expect(server.getStatusSnapshot()[0].interrupted).toBe(true) + } + } finally { + server.stop() + } + }) +}) diff --git a/src/main/agent-hooks/server-main-agent-turn-verdicts.test.ts b/src/main/agent-hooks/server-main-agent-turn-verdicts.test.ts index 4f7d01903f1..d5f69600526 100644 --- a/src/main/agent-hooks/server-main-agent-turn-verdicts.test.ts +++ b/src/main/agent-hooks/server-main-agent-turn-verdicts.test.ts @@ -1,4 +1,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { appendFileSync, mkdtempSync, rmSync, writeFileSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' import { AgentHookServer, _internals } from './server' import { buildBody, postHookEvent, PANE } from './server.test-fixtures' @@ -44,20 +47,6 @@ describe('main agent turn verdicts and clocks', () => { expect(response.status).toBe(204) } - function inferCtrlC(agentType: 'codex'): void { - const baseline = server.getStatusSnapshot()[0] - expect( - server.inferInterrupt({ - paneKey: PANE, - baselineUpdatedAt: baseline.receivedAt, - baselineStateStartedAt: baseline.stateStartedAt, - baselinePrompt: baseline.prompt, - baselineAgentType: agentType, - intent: 'ctrl-c' - }) - ).toBe(true) - } - it('starts a new Claude session with its own main agent clock', async () => { await post('/hook/claude', { hook_event_name: 'UserPromptSubmit', prompt: 'first' }) await post('/hook/claude', { hook_event_name: 'Stop' }) @@ -97,10 +86,31 @@ describe('main agent turn verdicts and clocks', () => { } ) - it('keeps an inferred Codex cancellation across a late root Stop', async () => { - await post('/hook/codex', { hook_event_name: 'UserPromptSubmit', prompt: 'long task' }) - vi.setSystemTime(1_001_000) - inferCtrlC('codex') + it('keeps a recorded Codex cancellation across a late root Stop', async () => { + const directory = mkdtempSync(join(tmpdir(), 'codex-turn-verdict-')) + const transcriptPath = join(directory, 'rollout.jsonl') + try { + writeFileSync( + transcriptPath, + '{"type":"event_msg","payload":{"type":"task_started","turn_id":"turn-1"}}\n' + ) + await post('/hook/codex', { + hook_event_name: 'UserPromptSubmit', + prompt: 'long task', + transcript_path: transcriptPath + }) + vi.setSystemTime(1_001_000) + appendFileSync( + transcriptPath, + '{"type":"event_msg","payload":{"type":"turn_aborted","turn_id":"turn-1","reason":"interrupted"}}\n' + ) + await vi.waitFor( + () => expect(server.getStatusSnapshot()[0]?.mainAgent?.outcome).toBe('cancellation'), + { timeout: 2_000 } + ) + } finally { + rmSync(directory, { recursive: true, force: true }) + } const cancelled = server.getStatusSnapshot()[0]?.mainAgent expect(cancelled).toMatchObject({ state: 'done', outcome: 'cancellation' }) @@ -115,7 +125,7 @@ describe('main agent turn verdicts and clocks', () => { expect(server.getStatusSnapshot()[0]).toMatchObject({ state: 'working', mainAgent: cancelled }) }) - it('keeps an inferred Codex cancellation across a late relayed root Stop', () => { + it('keeps a host-confirmed Codex cancellation across a late relayed root Stop', () => { const relayed = ( hookEventName: string, payload: Record, @@ -134,7 +144,21 @@ describe('main agent turn verdicts and clocks', () => { ) relayed('UserPromptSubmit', { state: 'working' }, { hasExplicitPrompt: true }) vi.setSystemTime(1_001_000) - inferCtrlC('codex') + server.ingestRemote( + { + paneKey: PANE, + tabId: 'tab-1', + worktreeId: 'wt-1', + payload: { + prompt: 'long task', + agentType: 'codex', + state: 'done', + interrupted: true, + mainAgent: { state: 'done', outcome: 'cancellation', stateStartedAt: 1_001_000 } + } + }, + 'conn-1' + ) const cancelled = server.getStatusSnapshot()[0]?.mainAgent expect(cancelled).toMatchObject({ state: 'done', outcome: 'cancellation' }) diff --git a/src/main/agent-hooks/server/server-status-inference.ts b/src/main/agent-hooks/server/server-status-inference.ts index 7dbe2620890..06f94f619ec 100644 --- a/src/main/agent-hooks/server/server-status-inference.ts +++ b/src/main/agent-hooks/server/server-status-inference.ts @@ -7,6 +7,7 @@ import { isAgentInterruptInputIntent, isNavigationEscapeIntent, requiresDoubleEscapeInterrupt, + shouldIgnoreInterruptIntent, type AgentInterruptInferenceRequest } from '../../../shared/agent-interrupt-intent' import { @@ -42,8 +43,7 @@ export abstract class AgentHookServerStatusInference extends AgentHookServerRowO } const payload = existing.payload const agentType: AgentType | undefined = payload.agentType - // Why: Droid's Ctrl+C exits the CLI (handled by PTY lifecycle) rather than interrupting the current turn. - if (agentType === 'droid' && request.intent === 'ctrl-c') { + if (shouldIgnoreInterruptIntent(agentType, request.intent)) { return false } // Why: these agents use the first Escape as a TUI cancel that can leave the turn running; only a double Escape infers an interrupt. diff --git a/src/main/agent-hooks/server/server-status-retries.ts b/src/main/agent-hooks/server/server-status-retries.ts index 4d92b372a35..3007d862395 100644 --- a/src/main/agent-hooks/server/server-status-retries.ts +++ b/src/main/agent-hooks/server/server-status-retries.ts @@ -1,3 +1,4 @@ +import { pollCodexTranscriptStatus } from '../../../shared/agent-hook-listener/providers/codex-transcript-poll' import { normalizeHookPayload } from '../../../shared/agent-hook-listener' import { hasPendingAgentResultText, @@ -6,9 +7,10 @@ import { import type { AgentHookSource } from '../../../shared/agent-hook-relay' import { shouldPollHookTranscript, + hookTranscriptWatchPath, transcriptPollUpdate } from '../../../shared/agent-hook-listener/transcript-poll-policy' -import { CodexSubagentPollScheduler } from '../../../shared/codex-subagent-poll-scheduler' +import { AgentTranscriptPollScheduler } from '../../../shared/agent-transcript-poll-scheduler' import type { EnrichedAgentHookEventPayload } from './server-types' import { ASSISTANT_MESSAGE_RETRY_ATTEMPTS, @@ -24,7 +26,7 @@ type TranscriptPoll = { } export abstract class AgentHookServerStatusRetries extends AgentHookServerStatusUpdate { - private readonly transcriptPollScheduler = new CodexSubagentPollScheduler( + private readonly transcriptPollScheduler = new AgentTranscriptPollScheduler( CODEX_SUBAGENT_POLL_MS, (paneKey, poll) => this.runTranscriptPoll(paneKey, poll) ) @@ -55,11 +57,15 @@ export abstract class AgentHookServerStatusRetries extends AgentHookServerStatus if (source !== 'codex' && source !== 'muse') { return } - this.transcriptPollScheduler.clear(original.paneKey) if (!shouldPollHookTranscript(this.state, source, original)) { + this.transcriptPollScheduler.clear(original.paneKey) return } - this.transcriptPollScheduler.schedule(original.paneKey, { source, body, original }) + this.transcriptPollScheduler.schedule( + original.paneKey, + { source, body, original }, + hookTranscriptWatchPath(this.state, source, original.paneKey) + ) } private runTranscriptPoll(paneKey: string, poll: TranscriptPoll): void { @@ -73,7 +79,10 @@ export abstract class AgentHookServerStatusRetries extends AgentHookServerStatus ) { return } - const normalized = normalizeHookPayload(this.state, source, body, this.env) + const normalized = + source === 'codex' + ? pollCodexTranscriptStatus(this.state, original) + : normalizeHookPayload(this.state, source, body, this.env) if (!normalized) { return } diff --git a/src/main/agent-hooks/server/server-status-update.ts b/src/main/agent-hooks/server/server-status-update.ts index bb40c187134..ebb87c3c9c2 100644 --- a/src/main/agent-hooks/server/server-status-update.ts +++ b/src/main/agent-hooks/server/server-status-update.ts @@ -101,7 +101,7 @@ export abstract class AgentHookServerStatusUpdate extends AgentHookServerStatusA const stateReconciledPayload = terminalOwnedPayload.connectionId && terminalOwnedPayload.payload.agentType === 'codex' && - terminalOwnedPayload.hookEventName + (terminalOwnedPayload.hookEventName || terminalOwnedPayload.payload.mainAgent) ? { ...terminalOwnedPayload, payload: reconcileRemoteCodexState( diff --git a/src/main/native-chat/transcript-native-watcher.ts b/src/main/native-chat/transcript-native-watcher.ts index 8cccd022c57..5aa615fdb9a 100644 --- a/src/main/native-chat/transcript-native-watcher.ts +++ b/src/main/native-chat/transcript-native-watcher.ts @@ -1,87 +1,4 @@ -import { watch, type FSWatcher } from 'node:fs' -import { basename, dirname } from 'node:path' - -export type TranscriptNativeWatcher = { - /** Best-effort bind; false keeps the caller in reconciliation-only mode. */ - bind: () => boolean - /** Detach from an identity that may no longer represent the watched path. */ - invalidate: () => void - needsRebind: () => boolean - dispose: () => void -} - -/** - * Optional fs.watch acceleration for transcript reconciliation. Native watches - * can fail on otherwise-readable remote filesystems, so binding is retryable - * and never owns transcript liveness; the caller's polling loop does. - */ -export function createTranscriptNativeWatcher( - filePath: string, - onEvent: () => void, - onError: () => void -): TranscriptNativeWatcher { - const watchedName = basename(filePath) - let disposed = false - let watcher: FSWatcher | null = null - let rebindNeeded = true - - function invalidateCandidate(candidate: FSWatcher): void { - if (watcher !== candidate) { - return - } - watcher = null - rebindNeeded = true - candidate.close() - } - - return { - bind(): boolean { - if (disposed || watcher) { - return watcher !== null - } - let nextWatcher: FSWatcher - try { - // Why: watching the parent survives target-file replacement on macOS. - nextWatcher = watch(dirname(filePath), (event, changedName) => { - if (changedName !== null && changedName.toString() !== watchedName) { - return - } - // Why: a parent replacement may emit rename without a watcher error. - if (event === 'rename') { - invalidateCandidate(nextWatcher) - } - onEvent() - }) - } catch { - rebindNeeded = true - return false - } - // Why: an active tail should not keep a headless runtime alive during shutdown. - nextWatcher.unref?.() - nextWatcher.on('error', () => { - if (disposed || watcher !== nextWatcher) { - return - } - invalidateCandidate(nextWatcher) - onError() - }) - watcher = nextWatcher - rebindNeeded = false - return true - }, - invalidate(): void { - if (watcher) { - invalidateCandidate(watcher) - } else { - rebindNeeded = true - } - }, - needsRebind: () => rebindNeeded, - dispose(): void { - disposed = true - watcher?.close() - watcher = null - rebindNeeded = false - } - } -} +export { + createTranscriptNativeWatcher, + type TranscriptNativeWatcher +} from '../../shared/transcript-native-watcher' diff --git a/src/relay/agent-hook-result-retry-scheduler.ts b/src/relay/agent-hook-result-retry-scheduler.ts index 3a6c9cd07fb..79f0b876a96 100644 --- a/src/relay/agent-hook-result-retry-scheduler.ts +++ b/src/relay/agent-hook-result-retry-scheduler.ts @@ -1,3 +1,4 @@ +import { pollCodexTranscriptStatus } from '../shared/agent-hook-listener/providers/codex-transcript-poll' // Deferred re-normalization timers for late-arriving agent results: the transcript hadn't caught up // when the hook fired, so re-read the same body on a timer and re-apply only if it changed. Both // timer families live in one owner so pane teardown and server stop tear both down in one ordered @@ -12,9 +13,10 @@ import type { HookListenerState } from '../shared/agent-hook-listener/listener-s import type { AgentHookSource } from '../shared/agent-hook-relay' import { shouldPollHookTranscript, + hookTranscriptWatchPath, transcriptPollUpdate } from '../shared/agent-hook-listener/transcript-poll-policy' -import { CodexSubagentPollScheduler } from '../shared/codex-subagent-poll-scheduler' +import { AgentTranscriptPollScheduler } from '../shared/agent-transcript-poll-scheduler' const ASSISTANT_MESSAGE_RETRY_ATTEMPTS = 5 const ASSISTANT_MESSAGE_RETRY_MS = 50 @@ -43,12 +45,12 @@ export type AgentHookResultRetryHost = { export class AgentHookResultRetryScheduler { private assistantMessageRetryTimers = new Map>() - private transcriptPollScheduler: CodexSubagentPollScheduler + private transcriptPollScheduler: AgentTranscriptPollScheduler private host: AgentHookResultRetryHost constructor(host: AgentHookResultRetryHost) { this.host = host - this.transcriptPollScheduler = new CodexSubagentPollScheduler( + this.transcriptPollScheduler = new AgentTranscriptPollScheduler( CODEX_SUBAGENT_POLL_MS, (paneKey, poll) => this.runTranscriptPoll(paneKey, poll) ) @@ -86,17 +88,21 @@ export class AgentHookResultRetryScheduler { if (source !== 'codex' && source !== 'muse') { return } - this.transcriptPollScheduler.clear(original.paneKey) if (!shouldPollHookTranscript(this.host.state, source, original)) { + this.transcriptPollScheduler.clear(original.paneKey) return } - this.transcriptPollScheduler.schedule(original.paneKey, { - source, - body, - original, - env, - version - }) + this.transcriptPollScheduler.schedule( + original.paneKey, + { + source, + body, + original, + env, + version + }, + hookTranscriptWatchPath(this.host.state, source, original.paneKey) + ) } private runTranscriptPoll(paneKey: string, poll: TranscriptPoll): void { @@ -110,7 +116,10 @@ export class AgentHookResultRetryScheduler { ) { return } - const event = normalizeHookPayload(this.host.state, source, body, this.host.env) + const event = + source === 'codex' + ? pollCodexTranscriptStatus(this.host.state, original) + : normalizeHookPayload(this.host.state, source, body, this.host.env) if (!event) { return } diff --git a/src/relay/agent-hook-server-codex-turn-interruption.test.ts b/src/relay/agent-hook-server-codex-turn-interruption.test.ts new file mode 100644 index 00000000000..451b21e560b --- /dev/null +++ b/src/relay/agent-hook-server-codex-turn-interruption.test.ts @@ -0,0 +1,79 @@ +import { expect, it, vi } from 'vitest' +import { appendFileSync, mkdtempSync, rmSync, writeFileSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import type { AgentHookRelayEnvelope } from '../shared/agent-hook-relay' +import { RelayAgentHookServer } from './agent-hook-server' +import { AgentHookServer } from '../main/agent-hooks/server' + +const PANE_KEY = 'tab-1:11111111-1111-4111-8111-111111111111' +const line = (payload: unknown): string => `${JSON.stringify({ type: 'event_msg', payload })}\n` + +it('forwards host-confirmed Codex interruption without requiring a local rollout', async () => { + const dir = mkdtempSync(join(tmpdir(), 'relay-codex-interruption-')) + const transcriptPath = join(dir, 'rollout-root.jsonl') + writeFileSync(transcriptPath, line({ type: 'task_started', turn_id: 'turn-1' })) + const desktop = new AgentHookServer() + const forward = vi.fn((envelope: AgentHookRelayEnvelope) => + desktop.ingestRemote(envelope, 'ssh-connection') + ) + const server = new RelayAgentHookServer({ endpointDir: dir, forward }) + await server.start() + try { + const { port, token } = server.getCoordinates() + const response = await fetch(`http://127.0.0.1:${port}/hook/codex`, { + method: 'POST', + headers: { 'Content-Type': 'application/json', 'X-Orca-Agent-Hook-Token': token }, + body: JSON.stringify({ + paneKey: PANE_KEY, + tabId: 'tab-1', + worktreeId: 'folder-1', + payload: { + hook_event_name: 'UserPromptSubmit', + session_id: 'main-session', + prompt: 'remote main task', + transcript_path: transcriptPath + } + }) + }) + expect(response.status).toBe(204) + expect(desktop.getStatusSnapshot()[0].state).toBe('working') + const sideResponse = await fetch(`http://127.0.0.1:${port}/hook/codex`, { + method: 'POST', + headers: { 'Content-Type': 'application/json', 'X-Orca-Agent-Hook-Token': token }, + body: JSON.stringify({ + paneKey: PANE_KEY, + tabId: 'tab-1', + worktreeId: 'folder-1', + payload: { + hook_event_name: 'UserPromptSubmit', + session_id: 'side-session', + transcript_path: null, + prompt: 'side chat' + } + }) + }) + expect(sideResponse.status).toBe(204) + expect(forward).toHaveBeenCalledTimes(1) + + appendFileSync( + transcriptPath, + line({ type: 'turn_aborted', turn_id: 'turn-1', reason: 'interrupted' }) + ) + await vi.waitFor( + () => + expect(desktop.getStatusSnapshot()[0]).toMatchObject({ + connectionId: 'ssh-connection', + state: 'done', + interrupted: true, + mainAgent: { state: 'done', outcome: 'cancellation' } + }), + { timeout: 2_000 } + ) + expect(forward.mock.calls.at(-1)?.[0].hookEventName).toBeUndefined() + expect(forward).toHaveBeenCalledTimes(2) + } finally { + server.stop() + rmSync(dir, { recursive: true, force: true }) + } +}) diff --git a/src/renderer/src/components/terminal-pane/agent-interrupt-inference.test.ts b/src/renderer/src/components/terminal-pane/agent-interrupt-inference.test.ts index ebb50a41f26..8e2cd48de6e 100644 --- a/src/renderer/src/components/terminal-pane/agent-interrupt-inference.test.ts +++ b/src/renderer/src/components/terminal-pane/agent-interrupt-inference.test.ts @@ -170,9 +170,9 @@ describe('agent interrupt inference', () => { entry = undefined }) - it('does not infer Ctrl+C for Droid', () => { + it.each(['codex', 'droid'] as const)('does not infer Ctrl+C for %s', (agentType) => { vi.useFakeTimers() - let entry: AgentStatusEntry | undefined = makeEntry({ agentType: 'droid' }) + let entry: AgentStatusEntry | undefined = makeEntry({ agentType }) const inferInterrupt = vi.fn() const tracker = createAgentInterruptInference({ paneKey: PANE_KEY, diff --git a/src/renderer/src/components/terminal-pane/agent-interrupt-inference.ts b/src/renderer/src/components/terminal-pane/agent-interrupt-inference.ts index 0b39565b35c..bb9661c362a 100644 --- a/src/renderer/src/components/terminal-pane/agent-interrupt-inference.ts +++ b/src/renderer/src/components/terminal-pane/agent-interrupt-inference.ts @@ -6,6 +6,7 @@ import { AGENT_INTERRUPT_SETTLE_MS, isNavigationEscapeIntent, requiresDoubleEscapeInterrupt, + shouldIgnoreInterruptIntent, type AgentInterruptInferenceRequest, type AgentInterruptInputIntent } from '../../../../shared/agent-interrupt-intent' @@ -50,13 +51,6 @@ function shouldFlushInterruptImmediately( ) } -function shouldIgnoreInterruptIntent( - agentType: AgentStatusEntry['agentType'], - intent: AgentInterruptInputIntent -): boolean { - return agentType === 'droid' && intent === 'ctrl-c' -} - /** Why: skip a round-trip main will refuse anyway. Scoped to 'working' so Claude's * AskUserQuestion dismissal — a 'waiting' row — still reaches inferQuestionAnswered. */ function isIgnorableNavigationEscape( diff --git a/src/renderer/src/components/terminal-pane/pty-connection-command-finished-cleanup.test.ts b/src/renderer/src/components/terminal-pane/pty-connection-command-finished-cleanup.test.ts index 75b96f02684..9c9db1ba02c 100644 --- a/src/renderer/src/components/terminal-pane/pty-connection-command-finished-cleanup.test.ts +++ b/src/renderer/src/components/terminal-pane/pty-connection-command-finished-cleanup.test.ts @@ -861,7 +861,7 @@ describe('connectPanePty', () => { prompt: 'stop quickly', updatedAt: 1_000, stateStartedAt: 900, - agentType: 'codex', + agentType: 'custom-agent', terminalTitle: 'Codex', stateHistory: [] } diff --git a/src/renderer/src/components/terminal-pane/pty-connection-interrupt-inference.test.ts b/src/renderer/src/components/terminal-pane/pty-connection-interrupt-inference.test.ts index 7e9034c1d61..c0527b21c65 100644 --- a/src/renderer/src/components/terminal-pane/pty-connection-interrupt-inference.test.ts +++ b/src/renderer/src/components/terminal-pane/pty-connection-interrupt-inference.test.ts @@ -144,73 +144,82 @@ describe('connectPanePty', () => { await restoreTerminalTestGlobals() }) - it('infers interrupts only from the focused terminal key target', async () => { - const { connectPanePty } = await import('./pty-connection') - const transport = createMockTransport() - transportFactoryQueue.push(transport) - vi.useFakeTimers() - vi.setSystemTime(1_100) - const paneKey = makePaneKey('tab-1', LEAF_1) - mockStoreState.agentStatusByPaneKey[paneKey] = { - state: 'working', - prompt: 'stop this task', - updatedAt: 1_000, - stateStartedAt: 900, - agentType: 'codex', - paneKey, - terminalTitle: 'Codex', - stateHistory: [] + it.each(['custom-agent', 'codex'] as const)( + 'handles focused Ctrl+C for %s without guessing Codex cancellation', + async (agentType) => { + const { connectPanePty } = await import('./pty-connection') + const transport = createMockTransport() + transportFactoryQueue.push(transport) + vi.useFakeTimers() + vi.setSystemTime(1_100) + const paneKey = makePaneKey('tab-1', LEAF_1) + mockStoreState.agentStatusByPaneKey[paneKey] = { + state: 'working', + prompt: 'stop this task', + updatedAt: 1_000, + stateStartedAt: 900, + agentType, + paneKey, + terminalTitle: 'Codex', + stateHistory: [] + } + const terminalTarget = createKeyboardEventTarget() + const unrelatedTarget = createKeyboardEventTarget() + ;( + globalThis.window as unknown as { addEventListener?: ReturnType } + ).addEventListener = vi.fn() + const pane = createPane(1) + ;(pane.terminal as { element?: unknown }).element = terminalTarget.target + let onDataHandler: ((data: string) => void) | null = null + pane.terminal.onData = vi.fn(((handler: (data: string) => void) => { + onDataHandler = handler + return { dispose: vi.fn() } + }) as typeof pane.terminal.onData) + + connectPanePty(pane as never, createManager(1) as never, createDeps() as never) + if (!onDataHandler) { + throw new Error('expected onData handler to be registered') + } + unrelatedTarget.dispatch({ + key: 'c', + ctrlKey: true, + metaKey: false, + altKey: false, + shiftKey: false, + repeat: false + } as KeyboardEvent) + vi.advanceTimersByTime(500) + expect(window.api.agentStatus.inferInterrupt).not.toHaveBeenCalled() + expect(globalThis.window.addEventListener).not.toHaveBeenCalled() + + terminalTarget.dispatch({ + key: 'c', + ctrlKey: true, + metaKey: false, + altKey: false, + shiftKey: false, + repeat: false + } as KeyboardEvent) + ;(onDataHandler as unknown as (data: string) => void)('\x03') + await flushAsyncTicks() + vi.advanceTimersByTime(500) + + if (agentType === 'codex') { + expect(window.api.agentStatus.inferInterrupt).not.toHaveBeenCalled() + expect(mockStoreState.agentStatusByPaneKey[paneKey]).toMatchObject({ state: 'working' }) + expect(transport.sendInputAccepted).toHaveBeenCalledWith('\x03', 'query-reply') + return + } + expect(window.api.agentStatus.inferInterrupt).toHaveBeenCalledWith({ + paneKey, + baselineUpdatedAt: 1_000, + baselineStateStartedAt: 900, + baselinePrompt: 'stop this task', + baselineAgentType: agentType, + intent: 'ctrl-c' + }) } - const terminalTarget = createKeyboardEventTarget() - const unrelatedTarget = createKeyboardEventTarget() - ;( - globalThis.window as unknown as { addEventListener?: ReturnType } - ).addEventListener = vi.fn() - const pane = createPane(1) - ;(pane.terminal as { element?: unknown }).element = terminalTarget.target - let onDataHandler: ((data: string) => void) | null = null - pane.terminal.onData = vi.fn(((handler: (data: string) => void) => { - onDataHandler = handler - return { dispose: vi.fn() } - }) as typeof pane.terminal.onData) - - connectPanePty(pane as never, createManager(1) as never, createDeps() as never) - if (!onDataHandler) { - throw new Error('expected onData handler to be registered') - } - unrelatedTarget.dispatch({ - key: 'c', - ctrlKey: true, - metaKey: false, - altKey: false, - shiftKey: false, - repeat: false - } as KeyboardEvent) - vi.advanceTimersByTime(500) - expect(window.api.agentStatus.inferInterrupt).not.toHaveBeenCalled() - expect(globalThis.window.addEventListener).not.toHaveBeenCalled() - - terminalTarget.dispatch({ - key: 'c', - ctrlKey: true, - metaKey: false, - altKey: false, - shiftKey: false, - repeat: false - } as KeyboardEvent) - ;(onDataHandler as unknown as (data: string) => void)('\x03') - await flushAsyncTicks() - vi.advanceTimersByTime(500) - - expect(window.api.agentStatus.inferInterrupt).toHaveBeenCalledWith({ - paneKey, - baselineUpdatedAt: 1_000, - baselineStateStartedAt: 900, - baselinePrompt: 'stop this task', - baselineAgentType: 'codex', - intent: 'ctrl-c' - }) - }) + ) it('clears stale working pane title after inferred interrupt applies', async () => { const { connectPanePty } = await import('./pty-connection') @@ -231,7 +240,7 @@ describe('connectPanePty', () => { prompt: 'stop visible spinner', updatedAt: 1_000, stateStartedAt: 900, - agentType: 'codex', + agentType: 'custom-agent', paneKey, terminalTitle: 'Codex working', stateHistory: [] @@ -382,7 +391,7 @@ describe('connectPanePty', () => { prompt: 'stop from real terminal byte', updatedAt: 1_000, stateStartedAt: 900, - agentType: 'codex', + agentType: 'custom-agent', paneKey, terminalTitle: 'Codex working', stateHistory: [] @@ -408,7 +417,7 @@ describe('connectPanePty', () => { baselineUpdatedAt: 1_000, baselineStateStartedAt: 900, baselinePrompt: 'stop from real terminal byte', - baselineAgentType: 'codex', + baselineAgentType: 'custom-agent', intent: 'ctrl-c' }) }) @@ -454,7 +463,7 @@ describe('connectPanePty', () => { prompt: 'stop enhanced keyboard input', updatedAt: 1_000, stateStartedAt: 900, - agentType: 'codex', + agentType: 'custom-agent', paneKey, terminalTitle: 'Codex working', stateHistory: [] @@ -484,7 +493,7 @@ describe('connectPanePty', () => { baselineUpdatedAt: 1_000, baselineStateStartedAt: 900, baselinePrompt: 'stop enhanced keyboard input', - baselineAgentType: 'codex', + baselineAgentType: 'custom-agent', intent: 'ctrl-c' }) }) @@ -506,7 +515,7 @@ describe('connectPanePty', () => { prompt: 'stop after process exit', updatedAt: 1_000, stateStartedAt: 900, - agentType: 'codex', + agentType: 'custom-agent', paneKey, terminalTitle: 'Codex working', stateHistory: [] @@ -540,7 +549,7 @@ describe('connectPanePty', () => { baselineUpdatedAt: 1_000, baselineStateStartedAt: 900, baselinePrompt: 'stop after process exit', - baselineAgentType: 'codex', + baselineAgentType: 'custom-agent', intent: 'ctrl-c' }) expect(mockStoreState.dropAgentStatus).not.toHaveBeenCalled() @@ -563,7 +572,7 @@ describe('connectPanePty', () => { prompt: 'stop and leave shell', updatedAt: 1_000, stateStartedAt: 900, - agentType: 'codex', + agentType: 'custom-agent', paneKey, terminalTitle: 'Codex working', stateHistory: [] @@ -575,7 +584,7 @@ describe('connectPanePty', () => { interrupted: true, updatedAt: 1_100, stateStartedAt: 1_100, - agentType: 'codex', + agentType: 'custom-agent', paneKey, terminalTitle: 'Terminal 1' } diff --git a/src/shared/agent-hook-listener-relay-dependency.test.ts b/src/shared/agent-hook-listener-relay-dependency.test.ts index 62dc8ff3877..70caf4f6fdf 100644 --- a/src/shared/agent-hook-listener-relay-dependency.test.ts +++ b/src/shared/agent-hook-listener-relay-dependency.test.ts @@ -126,6 +126,7 @@ describe('agent hook listener relay dependency boundary', () => { 'agent-hook-listener/hook-envelope.ts', 'agent-hook-listener/listener-limits.ts', 'agent-hook-listener/listener-state.ts', + 'agent-hook-listener/providers/codex-transcript-poll.ts', 'agent-hook-listener/request-body.ts', 'agent-hook-listener/source-routing.ts', 'agent-hook-listener/transcript-poll-policy.ts' diff --git a/src/shared/agent-hook-listener/providers/codex-events.ts b/src/shared/agent-hook-listener/providers/codex-events.ts index e5663766c2d..6781ce0d425 100644 --- a/src/shared/agent-hook-listener/providers/codex-events.ts +++ b/src/shared/agent-hook-listener/providers/codex-events.ts @@ -150,6 +150,20 @@ export function normalizeCodexEvent( return normalizeCodexSubagentLifecycleEvent(state, eventName, paneKey, hookPayload) } + const sessionId = readString(hookPayload, 'session_id') + const currentSessionId = state.lastStatusByPaneKey.get(paneKey)?.providerSession?.id + // Ephemeral side chats share the pane but must not replace its recorded main turn. + if ( + hookPayload.transcript_path === null && + !readString(hookPayload, 'agent_id') && + state.codexSubagentTranscriptByPaneKey.get(paneKey)?.parent.filePath && + sessionId && + currentSessionId && + sessionId !== currentSessionId + ) { + return null + } + // Why: Codex's request_user_input (0.145+) is auto-allowed, so it fires PreToolUse while blocked on a human answer; map to waiting like grok's ask_user_question. const isUserInputPreTool = eventName === 'PreToolUse' && @@ -195,6 +209,12 @@ export function normalizeCodexEvent( transcriptPath ) } + if (!agentId && (eventName === 'UserPromptSubmit' || eventName === 'SessionStart')) { + const transcript = state.codexSubagentTranscriptByPaneKey.get(paneKey) + if (transcript) { + transcript.rootTurn.interrupted = false + } + } if (agentId) { // Why: reconcile the child rollout reviewer before classifying its approval, including after relay restart. const childState = resolveCodexApprovalOwnedState( diff --git a/src/shared/agent-hook-listener/providers/codex-state.ts b/src/shared/agent-hook-listener/providers/codex-state.ts index ebd6b362def..c13fdbf95ed 100644 --- a/src/shared/agent-hook-listener/providers/codex-state.ts +++ b/src/shared/agent-hook-listener/providers/codex-state.ts @@ -220,6 +220,13 @@ export function reconcileRemoteCodexState( } } + if (!eventName && !agentId && payload.mainAgent && payload.mainAgent.state !== 'blocked') { + setCodexMainAgentTurnState(state, paneKey, { + ...payload.mainAgent, + state: payload.mainAgent.state, + model: payload.model ?? state.codexLeadStateByPaneKey.get(paneKey)?.model + }) + } const lead = state.codexLeadStateByPaneKey.get(paneKey) if (!lead) { return payload diff --git a/src/shared/agent-hook-listener/providers/codex-transcript-poll.ts b/src/shared/agent-hook-listener/providers/codex-transcript-poll.ts new file mode 100644 index 00000000000..73a1b740d90 --- /dev/null +++ b/src/shared/agent-hook-listener/providers/codex-transcript-poll.ts @@ -0,0 +1,32 @@ +import { reconcileCodexSubagentTranscript } from '../../codex-subagent-transcript' +import type { AgentHookEventPayload } from '../listener-event' +import type { HookListenerState } from '../listener-state' +import { buildCodexChildDrivenStatusPayload } from './codex-events' +import { getOrCreateCodexSubagentRoster, markCodexLeadTurnInterrupted } from './codex-state' + +/** Polling reads new host-owned records without replaying the hook that started the turn. */ +export function pollCodexTranscriptStatus( + state: HookListenerState, + original: T +): T | undefined { + const transcript = state.codexSubagentTranscriptByPaneKey.get(original.paneKey) + if (!transcript?.parent.filePath) { + return undefined + } + const changed = reconcileCodexSubagentTranscript( + transcript, + getOrCreateCodexSubagentRoster(state, original.paneKey), + transcript.parent.filePath + ) + const interrupted = + transcript.rootTurn.interrupted && + state.codexLeadStateByPaneKey.get(original.paneKey)?.state !== 'done' + if (!changed && !interrupted) { + return original + } + if (interrupted) { + markCodexLeadTurnInterrupted(state, original.paneKey) + } + const payload = buildCodexChildDrivenStatusPayload(state, undefined, original.paneKey, {}) + return payload ? { ...original, payload } : undefined +} diff --git a/src/shared/agent-hook-listener/transcript-poll-policy.ts b/src/shared/agent-hook-listener/transcript-poll-policy.ts index d27814cf928..31cf93dac26 100644 --- a/src/shared/agent-hook-listener/transcript-poll-policy.ts +++ b/src/shared/agent-hook-listener/transcript-poll-policy.ts @@ -11,7 +11,11 @@ export function shouldPollHookTranscript( event: AgentHookEventPayload ): boolean { if (source === 'codex') { - return hasCodexTranscriptSubagents(state, event.paneKey) + return ( + hasCodexTranscriptSubagents(state, event.paneKey) || + (state.codexLeadStateByPaneKey.get(event.paneKey)?.state !== 'done' && + Boolean(state.codexSubagentTranscriptByPaneKey.get(event.paneKey)?.parent.filePath)) + ) } if (source === 'muse') { // Why: Muse's question tool fires no hook, so only its session log shows the wait and its answer. @@ -20,12 +24,26 @@ export function shouldPollHookTranscript( return false } +/** Child rollout discovery keeps its existing cadence; only a root-only pane can use one file watch. */ +export function hookTranscriptWatchPath( + state: HookListenerState, + source: AgentHookSource, + paneKey: string +): string | undefined { + return source === 'codex' && !hasCodexTranscriptSubagents(state, paneKey) + ? state.codexSubagentTranscriptByPaneKey.get(paneKey)?.parent.filePath + : undefined +} + /** Returns the poll result to publish, or undefined when it carries nothing new. */ export function transcriptPollUpdate( source: AgentHookSource, original: T, polled: T ): T | undefined { + if (original === polled) { + return undefined + } if (source === 'muse') { const changed = polled.payload.state !== original.payload.state || @@ -37,5 +55,10 @@ export function transcriptPollUpdate( } const subagentsChanged = JSON.stringify(polled.payload.subagents) !== JSON.stringify(original.payload.subagents) - return subagentsChanged ? polled : undefined + const leadChanged = + polled.payload.mainAgent?.state !== original.payload.mainAgent?.state || + polled.payload.mainAgent?.outcome !== original.payload.mainAgent?.outcome + return subagentsChanged || leadChanged + ? { ...polled, hasExplicitPrompt: undefined, hookEventName: undefined } + : undefined } diff --git a/src/shared/agent-interrupt-intent.ts b/src/shared/agent-interrupt-intent.ts index 24f6f179e73..70ce0d81b26 100644 --- a/src/shared/agent-interrupt-intent.ts +++ b/src/shared/agent-interrupt-intent.ts @@ -18,6 +18,14 @@ export function isAgentInterruptInputIntent(intent: unknown): intent is AgentInt return intent === 'plain-escape' || intent === 'ctrl-c' } +// Ctrl+C can copy text or leave a Codex side chat without cancelling the main turn. +export function shouldIgnoreInterruptIntent( + agentType: AgentType | undefined, + intent: AgentInterruptInputIntent +): boolean { + return intent === 'ctrl-c' && (agentType === 'codex' || agentType === 'droid') +} + // Why: these TUIs also close an overlay on a bare Escape (Claude's /btw composer, OMP/Pi's // focused-child and settings views). The keypress is ambiguous at the source and nothing outside // the TUI can disambiguate it, so it is never evidence a turn ended — only the provider's own diff --git a/src/shared/agent-transcript-poll-scheduler.test.ts b/src/shared/agent-transcript-poll-scheduler.test.ts new file mode 100644 index 00000000000..82ab4dda4c7 --- /dev/null +++ b/src/shared/agent-transcript-poll-scheduler.test.ts @@ -0,0 +1,142 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' + +const watchers = vi.hoisted(() => ({ + bound: true, + entries: Array<{ + wake: () => void + error: () => void + bind: ReturnType + dispose: ReturnType + rebind: boolean + }>() +})) +vi.mock('./transcript-native-watcher', () => ({ + createTranscriptNativeWatcher: (_path: string, wake: () => void, error: () => void) => { + const entry = { + wake, + error, + rebind: true, + bind: vi.fn(() => { + entry.rebind = !watchers.bound + return watchers.bound + }), + dispose: vi.fn() + } + watchers.entries.push(entry) + return { bind: entry.bind, dispose: entry.dispose, needsRebind: () => entry.rebind } + } +})) +import { AgentTranscriptPollScheduler } from './agent-transcript-poll-scheduler' +import { CodexSubagentPollScheduler } from './codex-subagent-poll-scheduler' + +describe('host transcript wakeup budget', () => { + beforeEach(() => { + vi.useFakeTimers() + watchers.bound = true + watchers.entries.length = 0 + }) + afterEach(() => vi.useRealTimers()) + + it.each([500, 1_000])( + 'bounds quiet root probes across 100 panes with one timer (cadence %i ms)', + (cadence) => { + let before = 0 + let after = 0 + const baseline = new CodexSubagentPollScheduler(cadence, (key, value) => { + before++ + baseline.schedule(key, value) + }) + const optimized = new AgentTranscriptPollScheduler(cadence, (key, value) => { + after++ + optimized.schedule(key, value, `/rollout-${key}.jsonl`) + }) + for (let i = 0; i < 100; i++) { + baseline.schedule(String(i), i) + optimized.schedule(String(i), i, `/rollout-${i}.jsonl`) + } + expect(vi.getTimerCount()).toBe(2) + vi.advanceTimersByTime(60_000) + expect(before).toBe(100 * (60_000 / cadence)) + expect(after).toBe(1_200) + optimized.clearAll() + baseline.clearAll() + expect(vi.getTimerCount()).toBe(0) + expect(watchers.entries.every((entry) => entry.dispose.mock.calls.length === 1)).toBe(true) + } + ) + + it('coalesces a stream without postponing reads and reuses the latest hook', () => { + const seen: string[] = [] + const scheduler = new AgentTranscriptPollScheduler(500, (key, value) => { + seen.push(value) + scheduler.schedule(key, value, '/rollout.jsonl') + }) + scheduler.schedule('pane', 'first', '/rollout.jsonl') + vi.advanceTimersByTime(500) + const watch = watchers.entries[0]! + for (let i = 0; i < 1_000; i++) { + watch.wake() + scheduler.schedule('pane', `hook-${i}`, '/rollout.jsonl') + } + expect(watchers.entries).toHaveLength(1) + expect(vi.getTimerCount()).toBe(1) + vi.advanceTimersByTime(500) + expect(seen).toEqual(['first', 'hook-999']) + scheduler.clear('pane') + watch.wake() + expect(vi.getTimerCount()).toBe(0) + expect(watch.dispose).toHaveBeenCalledOnce() + }) + + it('keeps unreadable watches on the old cadence but bounds bind retries', () => { + watchers.bound = false + let reads = 0 + const scheduler = new AgentTranscriptPollScheduler(500, (key, value) => { + reads++ + scheduler.schedule(key, value, '/missing.jsonl') + }) + scheduler.schedule('pane', 1, '/missing.jsonl') + vi.advanceTimersByTime(10_000) + expect(reads).toBe(20) + expect(watchers.entries[0]!.bind).toHaveBeenCalledTimes(3) + scheduler.clearAll() + }) + + it('rebinds after rename or error and drops replaced watches and stale callbacks', () => { + const seen: number[] = [] + const scheduler = new AgentTranscriptPollScheduler(500, (key, value) => { + seen.push(value) + scheduler.schedule(key, value, '/new.jsonl') + }) + scheduler.schedule('pane', 1, '/old.jsonl') + const old = watchers.entries[0]! + scheduler.schedule('pane', 2, '/new.jsonl') + old.wake() + expect(old.dispose).toHaveBeenCalledOnce() + vi.advanceTimersByTime(500) + const current = watchers.entries[1]! + current.rebind = true + current.wake() + vi.advanceTimersByTime(500) + expect(current.bind).toHaveBeenCalledTimes(2) + current.rebind = true + current.error() + vi.advanceTimersByTime(500) + expect(current.bind).toHaveBeenCalledTimes(3) + expect(seen).toEqual([2, 2, 2]) + scheduler.clearAll() + }) + + it('retains polling for children and WSL UNC and releases abandoned callbacks', () => { + const seen = vi.fn() + const scheduler = new AgentTranscriptPollScheduler(500, seen) + scheduler.schedule('child', 1) + scheduler.schedule('wsl', 2, '\\\\wsl.localhost\\Ubuntu\\rollout.jsonl') + scheduler.schedule('root', 3, '/rollout.jsonl') + vi.advanceTimersByTime(500) + expect(seen).toHaveBeenCalledTimes(3) + expect(watchers.entries).toHaveLength(1) + expect(watchers.entries[0]!.dispose).toHaveBeenCalledOnce() + expect(vi.getTimerCount()).toBe(0) + }) +}) diff --git a/src/shared/agent-transcript-poll-scheduler.ts b/src/shared/agent-transcript-poll-scheduler.ts new file mode 100644 index 00000000000..ccefe786783 --- /dev/null +++ b/src/shared/agent-transcript-poll-scheduler.ts @@ -0,0 +1,125 @@ +import { CodexSubagentPollScheduler } from './codex-subagent-poll-scheduler' +import { + createTranscriptNativeWatcher, + type TranscriptNativeWatcher +} from './transcript-native-watcher' +import { isWslUncPath } from './wsl-paths' + +const ROOT_RECONCILIATION_MS = 5_000 + +type TranscriptWakeup = { + value: T + filePath?: string + watcher?: TranscriptNativeWatcher + watcherBound: boolean + nextBindAt: number + scheduled: boolean + eventPending: boolean + initialCheck: boolean +} + +/** One deadline timer; root-only sessions use targeted file events with a sparse safety probe. */ +export class AgentTranscriptPollScheduler { + private readonly entries = new Map>() + private readonly deadlines: CodexSubagentPollScheduler> + + constructor( + private readonly delayMs: number, + private readonly onDue: (key: string, value: T) => void + ) { + this.deadlines = new CodexSubagentPollScheduler(delayMs, (key, entry) => { + if (this.entries.get(key) !== entry) { + return + } + entry.scheduled = false + entry.eventPending = false + entry.initialCheck = false + try { + this.onDue(key, entry.value) + } finally { + if (this.entries.get(key) === entry && !entry.scheduled) { + this.clear(key) + } + } + }) + } + + schedule(key: string, value: T, filePath?: string): void { + let entry = this.entries.get(key) + if (entry && entry.filePath !== filePath) { + this.clear(key) + entry = undefined + } + if (!entry) { + entry = { + value, + filePath, + watcherBound: false, + nextBindAt: 0, + scheduled: false, + eventPending: false, + initialCheck: true + } + this.entries.set(key, entry) + if (filePath && !isWslUncPath(filePath)) { + const watchedEntry = entry + const wake = (): void => { + if (this.entries.get(key) !== watchedEntry) { + return + } + if (watchedEntry.watcher?.needsRebind()) { + watchedEntry.watcherBound = false + watchedEntry.nextBindAt = 0 + } + if (watchedEntry.eventPending) { + return + } + watchedEntry.eventPending = true + watchedEntry.scheduled = true + this.deadlines.schedule(key, watchedEntry, this.delayMs) + } + entry.watcher = createTranscriptNativeWatcher( + filePath, + wake, + () => { + watchedEntry.watcherBound = false + watchedEntry.nextBindAt = 0 + wake() + }, + { watchParent: false } + ) + } + } + entry.value = value + // The pending callback reads this entry's latest hook, so repeated hooks need neither a new watch nor a new timer. + if (entry.scheduled) { + return + } + const now = performance.now() + if (entry.watcher?.needsRebind() && now >= entry.nextBindAt) { + entry.watcherBound = entry.watcher.bind() + entry.nextBindAt = now + ROOT_RECONCILIATION_MS + } + const delay = entry.watcherBound && !entry.initialCheck ? ROOT_RECONCILIATION_MS : this.delayMs + entry.scheduled = true + this.deadlines.schedule(key, entry, delay) + } + + clear(key: string): void { + const entry = this.entries.get(key) + if (!entry) { + return + } + this.entries.delete(key) + this.deadlines.clear(key) + entry.watcher?.dispose() + } + + clearAll(): void { + this.deadlines.clearAll() + for (const entry of this.entries.values()) { + entry.watcher?.dispose() + } + this.entries.clear() + } +} diff --git a/src/shared/codex-status-transcript-line.test.ts b/src/shared/codex-status-transcript-line.test.ts new file mode 100644 index 00000000000..ecfbcb4a0e2 --- /dev/null +++ b/src/shared/codex-status-transcript-line.test.ts @@ -0,0 +1,42 @@ +import { appendFileSync, mkdtempSync, rmSync, writeFileSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { expect, it, vi } from 'vitest' +import { + createCodexSubagentTranscriptState, + reconcileCodexSubagentTranscript +} from './codex-subagent-transcript' +import type { CodexSubagentRoster } from './codex-subagent-roster' + +it('skips decoding 1,000 message/token records and unchanged reads, but retains an abort', () => { + const directory = mkdtempSync(join(tmpdir(), 'codex-status-budget-')) + const path = join(directory, 'rollout.jsonl') + const state = createCodexSubagentTranscriptState() + const roster: CodexSubagentRoster = new Map() + writeFileSync(path, '') + reconcileCodexSubagentTranscript(state, roster, path) + const irrelevantRecords = Array.from({ length: 1_000 }, (_, i) => + JSON.stringify({ + type: 'event_msg', + payload: { type: i % 2 ? 'agent_message' : 'token_count', message: 'message' } + }) + ).join('\n') + appendFileSync(path, `${irrelevantRecords}\n`) + const parse = vi.spyOn(JSON, 'parse') + try { + expect(reconcileCodexSubagentTranscript(state, roster, path)).toBe(false) + expect(parse).not.toHaveBeenCalled() + appendFileSync( + path, + '{"type":"event_msg","payload":{"type":"turn_aborted","reason":"interrupted"}}\n' + ) + expect(reconcileCodexSubagentTranscript(state, roster, path)).toBe(true) + expect(parse).toHaveBeenCalledOnce() + expect(state.rootTurn.interrupted).toBe(true) + expect(reconcileCodexSubagentTranscript(state, roster, path)).toBe(false) + expect(parse).toHaveBeenCalledOnce() + } finally { + parse.mockRestore() + rmSync(directory, { recursive: true, force: true }) + } +}) diff --git a/src/shared/codex-status-transcript-line.ts b/src/shared/codex-status-transcript-line.ts new file mode 100644 index 00000000000..ad0b6b22004 --- /dev/null +++ b/src/shared/codex-status-transcript-line.ts @@ -0,0 +1,7 @@ +// Codex writes these lifecycle types as JSON strings; message and token records need no decoding. +const STATUS_RECORD_TYPE = + /"(?:turn_context|thread_settings_applied|sub_agent_activity|task_started|turn_started|task_complete|turn_complete|turn_aborted)"/ + +export function isCodexStatusTranscriptLine(line: string): boolean { + return STATUS_RECORD_TYPE.test(line) +} diff --git a/src/shared/codex-subagent-poll-scheduler.ts b/src/shared/codex-subagent-poll-scheduler.ts index b843ff034bb..9a753d80af3 100644 --- a/src/shared/codex-subagent-poll-scheduler.ts +++ b/src/shared/codex-subagent-poll-scheduler.ts @@ -21,9 +21,9 @@ export class CodexSubagentPollScheduler { private readonly now: () => number = monotonicNow ) {} - schedule(key: string, value: T): void { + schedule(key: string, value: T, delayMs = this.delayMs): void { this.entries.delete(key) - this.entries.set(key, { value, dueAt: this.now() + this.delayMs }) + this.entries.set(key, { value, dueAt: this.now() + delayMs }) this.arm() } diff --git a/src/shared/codex-subagent-reviewer.ts b/src/shared/codex-subagent-reviewer.ts index 9fff4db175b..946f59352c8 100644 --- a/src/shared/codex-subagent-reviewer.ts +++ b/src/shared/codex-subagent-reviewer.ts @@ -1,6 +1,7 @@ import { extname, isAbsolute } from 'node:path' import { readJsonlCursor, type JsonRecord } from './codex-rollout-jsonl-cursor' +import { isCodexStatusTranscriptLine } from './codex-status-transcript-line' import type { CodexSubagentTranscriptState } from './codex-subagent-transcript' const REVIEWER_CURSOR_MAX_PATHS = 64 @@ -78,7 +79,7 @@ export function reconcileCodexSubagentReviewer( cursor = { filePath: normalizedPath, offset: 0, carry: '' } state.reviewerCursorsByPath.set(normalizedPath, cursor) } - const records = readJsonlCursor(cursor) + const records = readJsonlCursor(cursor, isCodexStatusTranscriptLine) if (records === undefined) { state.reviewerCursorsByPath.delete(normalizedPath) state.reviewersByPath.delete(normalizedPath) diff --git a/src/shared/codex-subagent-transcript.ts b/src/shared/codex-subagent-transcript.ts index 5ce841ca252..9f3f2ec81dc 100644 --- a/src/shared/codex-subagent-transcript.ts +++ b/src/shared/codex-subagent-transcript.ts @@ -8,6 +8,10 @@ import { type JsonRecord } from './codex-rollout-jsonl-cursor' +import { reconcileCodexTranscriptTurn, type CodexTranscriptTurn } from './codex-turn-transcript' + +import { isCodexStatusTranscriptLine } from './codex-status-transcript-line' + import { readApprovalsReviewer } from './codex-subagent-reviewer' import type { CodexApprovalsReviewer } from './codex-subagent-reviewer' @@ -34,6 +38,7 @@ type TrackedTranscriptSubagent = JsonlCursor & { export type CodexSubagentTranscriptState = { parent: JsonlCursor + rootTurn: CodexTranscriptTurn subagents: Map /** Incremental reviewer cursors for child rollouts, which must not replace the parent cursor. */ reviewerCursorsByPath: Map @@ -175,6 +180,7 @@ function childIsComplete(records: JsonRecord[]): boolean { export function createCodexSubagentTranscriptState(): CodexSubagentTranscriptState { return { parent: { offset: 0, carry: '' }, + rootTurn: { interrupted: false }, subagents: new Map(), reviewerCursorsByPath: new Map(), reviewersByPath: new Map() @@ -191,23 +197,30 @@ export function reconcileCodexSubagentTranscript( state: CodexSubagentTranscriptState, roster: CodexSubagentRoster, transcriptPath: string | undefined -): void { +): boolean { const normalizedPath = normalizedTranscriptPath(transcriptPath) if (!normalizedPath) { - return + return false } + let changed = false if (state.parent.filePath !== normalizedPath) { for (const id of state.subagents.keys()) { finishCodexSubagent(roster, id) } + changed = true state.parent = { filePath: normalizedPath, offset: 0, carry: '' } + state.rootTurn = { interrupted: false } state.subagents.clear() state.reviewerCursorsByPath.clear() state.reviewersByPath.clear() // Why: a different rollout is a different session, so its predecessor's reviewer is void. state.approvalsReviewer = undefined } - const parentRecords = readJsonlCursor(state.parent) + const parentRecords = readJsonlCursor(state.parent, isCodexStatusTranscriptLine) + changed ||= + Boolean(parentRecords?.length) || + (parentRecords === undefined && state.approvalsReviewer !== undefined) + reconcileCodexTranscriptTurn(state.rootTurn, parentRecords ?? []) // A stale reviewer must never turn an unreadable rollout into a hidden prompt. state.approvalsReviewer = parentRecords === undefined @@ -248,7 +261,8 @@ export function reconcileCodexSubagentTranscript( entriesByDirectory ) } - const records = readJsonlCursor(tracked) + const records = readJsonlCursor(tracked, isCodexStatusTranscriptLine) + changed ||= Boolean(records?.length) if (!records) { // Why: a rollout that never appears (or is deleted) has no completion event, so time-box it instead of leaking a working row. tracked.filePath = undefined @@ -267,7 +281,9 @@ export function reconcileCodexSubagentTranscript( continue } } + changed = true finishCodexSubagent(roster, id) state.subagents.delete(id) } + return changed } diff --git a/src/shared/codex-turn-transcript.test.ts b/src/shared/codex-turn-transcript.test.ts new file mode 100644 index 00000000000..cddd8edadb3 --- /dev/null +++ b/src/shared/codex-turn-transcript.test.ts @@ -0,0 +1,32 @@ +import { describe, expect, it } from 'vitest' +import { reconcileCodexTranscriptTurn, type CodexTranscriptTurn } from './codex-turn-transcript' + +const event = (payload: unknown): Record => ({ type: 'event_msg', payload }) + +describe('Codex root turn evidence', () => { + it('ignores a stale abort for an older turn and a quoted interruption message', () => { + const turn: CodexTranscriptTurn = { interrupted: false } + reconcileCodexTranscriptTurn(turn, [ + event({ type: 'task_started', turn_id: 'new' }), + event({ type: 'turn_aborted', turn_id: 'old', reason: 'interrupted' }), + event({ type: 'agent_message', message: 'Conversation interrupted - use /feedback' }) + ]) + expect(turn).toEqual({ turnId: 'new', interrupted: false }) + reconcileCodexTranscriptTurn(turn, [ + event({ type: 'turn_aborted', turn_id: 'new', reason: 'interrupted' }) + ]) + expect(turn.interrupted).toBe(true) + reconcileCodexTranscriptTurn(turn, [event({ type: 'task_started', turn_id: 'next' })]) + expect(turn).toEqual({ turnId: 'next', interrupted: false }) + }) + + it('accepts older abort events without a turn id and ignores other abort reasons', () => { + const turn: CodexTranscriptTurn = { interrupted: false } + reconcileCodexTranscriptTurn(turn, [event({ type: 'turn_aborted', reason: 'budget_limited' })]) + expect(turn.interrupted).toBe(false) + reconcileCodexTranscriptTurn(turn, [event({ type: 'turn_aborted', reason: 'interrupted' })]) + expect(turn.interrupted).toBe(true) + reconcileCodexTranscriptTurn(turn, [event({ type: 'turn_complete' })]) + expect(turn.interrupted).toBe(false) + }) +}) diff --git a/src/shared/codex-turn-transcript.ts b/src/shared/codex-turn-transcript.ts new file mode 100644 index 00000000000..3e43454292c --- /dev/null +++ b/src/shared/codex-turn-transcript.ts @@ -0,0 +1,33 @@ +import { record, type JsonRecord } from './codex-rollout-jsonl-cursor' + +export type CodexTranscriptTurn = { + turnId?: string + interrupted: boolean +} + +/** The root rollout, rather than a terminal key or painted notice, confirms cancellation. */ +export function reconcileCodexTranscriptTurn( + turn: CodexTranscriptTurn, + records: readonly JsonRecord[] +): void { + for (const item of records) { + if (item.type !== 'event_msg') { + continue + } + const payload = record(item.payload) + if (!payload) { + continue + } + const turnId = typeof payload.turn_id === 'string' ? payload.turn_id : undefined + if (payload.type === 'task_started' || payload.type === 'turn_started') { + turn.turnId = turnId + turn.interrupted = false + } else if (!turnId || !turn.turnId || turnId === turn.turnId) { + if (payload.type === 'turn_aborted' && payload.reason === 'interrupted') { + turn.interrupted = true + } else if (payload.type === 'task_complete' || payload.type === 'turn_complete') { + turn.interrupted = false + } + } + } +} diff --git a/src/shared/transcript-native-watcher.ts b/src/shared/transcript-native-watcher.ts new file mode 100644 index 00000000000..393b39e9588 --- /dev/null +++ b/src/shared/transcript-native-watcher.ts @@ -0,0 +1,91 @@ +import { watch, type FSWatcher } from 'node:fs' +import { basename, dirname } from 'node:path' + +export type TranscriptNativeWatcher = { + /** Best-effort bind; false keeps the caller in reconciliation-only mode. */ + bind: () => boolean + /** Detach from an identity that may no longer represent the watched path. */ + invalidate: () => void + needsRebind: () => boolean + dispose: () => void +} + +/** + * Optional fs.watch acceleration for transcript reconciliation. Native watches + * can fail on otherwise-readable remote filesystems, so binding is retryable + * and never owns transcript liveness; the caller's polling loop does. + */ +export function createTranscriptNativeWatcher( + filePath: string, + onEvent: () => void, + onError: () => void, + options: { watchParent?: boolean } = {} +): TranscriptNativeWatcher { + const watchedName = basename(filePath) + let disposed = false + let watcher: FSWatcher | null = null + let rebindNeeded = true + + function invalidateCandidate(candidate: FSWatcher): void { + if (watcher !== candidate) { + return + } + watcher = null + rebindNeeded = true + candidate.close() + } + + return { + bind(): boolean { + if (disposed || watcher) { + return watcher !== null + } + let nextWatcher: FSWatcher + try { + // Parent watches survive replacement; append-only readers can watch one file to avoid sibling events. + nextWatcher = watch( + options.watchParent === false ? filePath : dirname(filePath), + (event, changedName) => { + if (changedName !== null && changedName.toString() !== watchedName) { + return + } + // Why: a parent replacement may emit rename without a watcher error. + if (event === 'rename') { + invalidateCandidate(nextWatcher) + } + onEvent() + } + ) + } catch { + rebindNeeded = true + return false + } + // Why: an active tail should not keep a headless runtime alive during shutdown. + nextWatcher.unref?.() + nextWatcher.on('error', () => { + if (disposed || watcher !== nextWatcher) { + return + } + invalidateCandidate(nextWatcher) + onError() + }) + watcher = nextWatcher + rebindNeeded = false + return true + }, + invalidate(): void { + if (watcher) { + invalidateCandidate(watcher) + } else { + rebindNeeded = true + } + }, + needsRebind: () => rebindNeeded, + dispose(): void { + disposed = true + watcher?.close() + watcher = null + rebindNeeded = false + } + } +} diff --git a/tests/e2e/codex-ctrl-c-status.spec.ts b/tests/e2e/codex-ctrl-c-status.spec.ts new file mode 100644 index 00000000000..4a258103d08 --- /dev/null +++ b/tests/e2e/codex-ctrl-c-status.spec.ts @@ -0,0 +1,89 @@ +import { appendFileSync, mkdtempSync, rmSync, writeFileSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { test, expect } from './helpers/orca-app' +import { emitCodexHookStatus, readHookEndpoint } from './helpers/agent-hook-endpoint' +import { ensureTerminalVisible, waitForActiveWorktree, waitForSessionReady } from './helpers/store' +import { + focusActiveTerminalInput, + waitForActivePaneHookDescriptor, + waitForActiveTerminalManager +} from './helpers/terminal' + +test('Codex Ctrl+C preserves working status until a confirmed interruption', async ({ + orcaPage, + electronApp +}, testInfo) => { + await waitForSessionReady(orcaPage) + await waitForActiveWorktree(orcaPage) + await ensureTerminalVisible(orcaPage) + await waitForActiveTerminalManager(orcaPage, 30_000) + const endpoint = await readHookEndpoint(electronApp) + const descriptor = await waitForActivePaneHookDescriptor(orcaPage) + await orcaPage.evaluate(() => { + const state = window.__store?.getState() + state?.setAgentActivityDisplayMode('full') + if (state && !state.worktreeCardProperties.includes('inline-agents')) { + state.setWorktreeCardProperties([...state.worktreeCardProperties, 'inline-agents']) + } + }) + const dir = mkdtempSync(join(tmpdir(), 'orca-codex-interruption-')) + const transcriptPath = join(dir, 'rollout-root.jsonl') + writeFileSync( + transcriptPath, + `${JSON.stringify({ type: 'event_msg', payload: { type: 'task_started', turn_id: 'turn-1' } })}\n` + ) + try { + await emitCodexHookStatus(endpoint, { + ...descriptor, + transcriptPath, + sessionId: 'main-session', + state: 'working', + prompt: 'Main task continues' + }) + const working = orcaPage.locator('[aria-label="Working"]') + const interrupted = orcaPage.locator('[aria-label="Interrupted"]') + await expect(working.first()).toBeVisible() + await emitCodexHookStatus(endpoint, { + ...descriptor, + state: 'working', + sessionId: 'side-session', + transcriptPath: null, + prompt: 'Side chat' + }) + await expect(orcaPage.getByText('Main task continues', { exact: true })).toBeVisible() + await focusActiveTerminalInput(orcaPage) + await orcaPage.keyboard.press('Control+c') + // Allow the old 500 ms inference timer to fire before recording the rendered result. + await orcaPage.waitForTimeout(1_000) + await orcaPage.screenshot({ + path: testInfo.outputPath('status-after-ctrl-c.png'), + clip: { x: 0, y: 180, width: 280, height: 240 } + }) + await expect(interrupted).toHaveCount(0) + await expect(working.first()).toBeVisible() + + appendFileSync( + transcriptPath, + `${JSON.stringify({ type: 'event_msg', payload: { type: 'turn_aborted', turn_id: 'turn-1', reason: 'interrupted' } })}\n` + ) + await expect(interrupted.first()).toBeVisible() + await expect(working).toHaveCount(0) + await orcaPage.screenshot({ + path: testInfo.outputPath('status-after-confirmed-interruption.png'), + clip: { x: 0, y: 180, width: 280, height: 240 } + }) + await emitCodexHookStatus(endpoint, { + ...descriptor, + state: 'working', + prompt: 'Next main task' + }) + await expect(working.first()).toBeVisible() + await expect(interrupted).toHaveCount(0) + await emitCodexHookStatus(endpoint, { ...descriptor, state: 'done' }) + await expect(working).toHaveCount(0) + await expect(interrupted).toHaveCount(0) + } finally { + rmSync(dir, { recursive: true, force: true }) + } +}) diff --git a/tests/e2e/helpers/agent-hook-endpoint.ts b/tests/e2e/helpers/agent-hook-endpoint.ts index ec530f8bb40..5364b93c8ef 100644 --- a/tests/e2e/helpers/agent-hook-endpoint.ts +++ b/tests/e2e/helpers/agent-hook-endpoint.ts @@ -57,6 +57,8 @@ export async function emitCodexHookStatus( state: 'working' | 'done' prompt?: string lastAssistantMessage?: string + transcriptPath?: string | null + sessionId?: string } ): Promise { const [tabId] = status.paneKey.split(':') @@ -82,7 +84,11 @@ export async function emitCodexHookStatus( worktreeId: status.worktreeId, env: endpoint.env, version: endpoint.version, - payload + payload: { + ...payload, + ...(status.transcriptPath !== undefined ? { transcript_path: status.transcriptPath } : {}), + ...(status.sessionId ? { session_id: status.sessionId } : {}) + } }) }) if (response.status !== 204) {