mirror of
https://github.com/stablyai/orca.git
synced 2026-10-03 00:02:19 +00:00
Fix Codex status after Ctrl+C copy and side-chat navigation (#24339)
* fix: preserve Codex status on ambiguous Ctrl+C input * fix: confirm Codex turn cancellations from host rollout records * fix: keep ephemeral side hooks separate from the Codex main turn * perf: watch active Codex rollouts and skip unrelated records * fix: retain confirmed Codex cancellation across late relay events
This commit is contained in:
@@ -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()
|
||||
})
|
||||
})
|
||||
@@ -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()
|
||||
}
|
||||
})
|
||||
})
|
||||
@@ -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<string, unknown>,
|
||||
@@ -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' })
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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<TranscriptPoll>(
|
||||
private readonly transcriptPollScheduler = new AgentTranscriptPollScheduler<TranscriptPoll>(
|
||||
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
|
||||
}
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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'
|
||||
|
||||
@@ -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<string, ReturnType<typeof setTimeout>>()
|
||||
private transcriptPollScheduler: CodexSubagentPollScheduler<TranscriptPoll>
|
||||
private transcriptPollScheduler: AgentTranscriptPollScheduler<TranscriptPoll>
|
||||
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
|
||||
}
|
||||
|
||||
@@ -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 })
|
||||
}
|
||||
})
|
||||
@@ -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,
|
||||
|
||||
@@ -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(
|
||||
|
||||
+1
-1
@@ -861,7 +861,7 @@ describe('connectPanePty', () => {
|
||||
prompt: 'stop quickly',
|
||||
updatedAt: 1_000,
|
||||
stateStartedAt: 900,
|
||||
agentType: 'codex',
|
||||
agentType: 'custom-agent',
|
||||
terminalTitle: 'Codex',
|
||||
stateHistory: []
|
||||
}
|
||||
|
||||
+84
-75
@@ -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<typeof vi.fn> }
|
||||
).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<typeof vi.fn> }
|
||||
).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'
|
||||
}
|
||||
|
||||
@@ -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'
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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<T extends AgentHookEventPayload>(
|
||||
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
|
||||
}
|
||||
@@ -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<T extends AgentHookEventPayload>(
|
||||
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<T extends AgentHookEventPayload>(
|
||||
}
|
||||
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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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<typeof vi.fn>
|
||||
dispose: ReturnType<typeof vi.fn>
|
||||
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<number>(cadence, (key, value) => {
|
||||
before++
|
||||
baseline.schedule(key, value)
|
||||
})
|
||||
const optimized = new AgentTranscriptPollScheduler<number>(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<string>(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<number>(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<number>(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<number>(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)
|
||||
})
|
||||
})
|
||||
@@ -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<T> = {
|
||||
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<T> {
|
||||
private readonly entries = new Map<string, TranscriptWakeup<T>>()
|
||||
private readonly deadlines: CodexSubagentPollScheduler<TranscriptWakeup<T>>
|
||||
|
||||
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()
|
||||
}
|
||||
}
|
||||
@@ -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 })
|
||||
}
|
||||
})
|
||||
@@ -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)
|
||||
}
|
||||
@@ -21,9 +21,9 @@ export class CodexSubagentPollScheduler<T> {
|
||||
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()
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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<string, TrackedTranscriptSubagent>
|
||||
/** Incremental reviewer cursors for child rollouts, which must not replace the parent cursor. */
|
||||
reviewerCursorsByPath: Map<string, JsonlCursor>
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -0,0 +1,32 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { reconcileCodexTranscriptTurn, type CodexTranscriptTurn } from './codex-turn-transcript'
|
||||
|
||||
const event = (payload: unknown): Record<string, unknown> => ({ 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)
|
||||
})
|
||||
})
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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 })
|
||||
}
|
||||
})
|
||||
@@ -57,6 +57,8 @@ export async function emitCodexHookStatus(
|
||||
state: 'working' | 'done'
|
||||
prompt?: string
|
||||
lastAssistantMessage?: string
|
||||
transcriptPath?: string | null
|
||||
sessionId?: string
|
||||
}
|
||||
): Promise<void> {
|
||||
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) {
|
||||
|
||||
Reference in New Issue
Block a user