diff --git a/src/main/runtime/orca-runtime-on-pty-data.ts b/src/main/runtime/orca-runtime-on-pty-data.ts index 5d27da94c14..8d6024dc759 100644 --- a/src/main/runtime/orca-runtime-on-pty-data.ts +++ b/src/main/runtime/orca-runtime-on-pty-data.ts @@ -256,6 +256,8 @@ export class OrcaRuntimeWithOnPtyData extends OrcaRuntimeWithPreparePtyExecution ...(cwdChanged && cwd !== null ? { cwd } : {}), ...(sourceRanges && sourceRanges.length > 0 ? { sourceRanges } : {}) })) + // A resumed agent can publish its restored screen before it owns the PTY foreground. + this.orchestrationMailboxPointerDelivery.redriveAfterPtyOutput(ptyId) return outputSequence } } diff --git a/src/main/runtime/orchestration/mailbox-pointer-delivery.ts b/src/main/runtime/orchestration/mailbox-pointer-delivery.ts index 201cc61aa78..ca5e0fae48a 100644 --- a/src/main/runtime/orchestration/mailbox-pointer-delivery.ts +++ b/src/main/runtime/orchestration/mailbox-pointer-delivery.ts @@ -4,6 +4,7 @@ import type { OrchestrationMessageWaiter } from './mailbox-pointer-eligibility' import type { OrchestrationMailboxLeaf, OrchestrationMailboxOwner } from './mailbox-owner' import { getOrchestrationMailboxPointerCandidates } from './mailbox-pointer-candidates' import { OrchestrationMailboxStatuslessCodexProofCoordinator } from './mailbox-statusless-codex-proof-coordinator' +import { OrchestrationMailboxStatuslessCodexRedrive } from './mailbox-statusless-codex-redrive' import { isStatuslessIdleProofCurrent } from './mailbox-statusless-idle-proof' import { stageOrchestrationMailboxPointer } from './mailbox-pointer-stage' import type { SubmitStatuslessCodexPointer } from './mailbox-statusless-codex-submit' @@ -39,9 +40,13 @@ type PointerDeliveryDependencies = { export class OrchestrationMailboxPointerDelivery { private readonly state = new OrchestrationMailboxPointerState() private readonly statuslessCodexProofs: OrchestrationMailboxStatuslessCodexProofCoordinator + private readonly statuslessCodexRedrives: OrchestrationMailboxStatuslessCodexRedrive constructor(private readonly deps: PointerDeliveryDependencies) { this.statuslessCodexProofs = new OrchestrationMailboxStatuslessCodexProofCoordinator(deps) + this.statuslessCodexRedrives = new OrchestrationMailboxStatuslessCodexRedrive((mailboxHandle) => + this.redrive(mailboxHandle, true) + ) } deliverForHandle(handle: string, reservedTypes?: ReadonlySet): void { @@ -122,7 +127,7 @@ export class OrchestrationMailboxPointerDelivery this.statuslessCodexRedrives.schedule(ptyId, mailboxHandle, sequence), + clearDeferredOutputRedrive: ( + ptyId: string, + mailboxHandle: string, + sequence: number + ) => this.statuslessCodexRedrives.clear(ptyId, mailboxHandle, sequence) + } : {}), writePty: this.deps.writePty, settle: (settledPtyId, settledFlight) => this.settle(settledPtyId, settledFlight), @@ -224,6 +241,7 @@ export class OrchestrationMailboxPointerDelivery = { isLeafPtyProvenAbsent: (ptyId: string) => Promise requestSleepingRecipientWake?: (mailboxHandle: string) => void submitStatuslessCodexPointer?: SubmitStatuslessCodexPointer + deferRedriveUntilPtyOutput?: (ptyId: string, mailboxHandle: string, sequence: number) => boolean + clearDeferredOutputRedrive?: (ptyId: string, mailboxHandle: string, sequence: number) => void writePty: (ptyId: string, data: string) => boolean | Promise settle: (ptyId: string, flight: OrchestrationMailboxDeliveryFlight) => void redrive: (mailboxHandle: string, force?: boolean) => void @@ -63,7 +65,12 @@ export function stageOrchestrationMailboxPointer + const redriveMailbox = vi.fn((mailboxHandle: string) => { + delivery.deliverForHandle(mailboxHandle) + }) + delivery = new OrchestrationMailboxPointerDelivery({ mailboxOwner: { resolve: () => MAILBOX } as unknown as OrchestrationMailboxOwner, deliveryTarget: { resolveTerminalHandle: () => TERMINAL_HANDLE, @@ -51,7 +55,7 @@ function makeHarness( getTerminalProcessIncarnation: () => processIncarnation, isLeafPtyProvenAbsent: () => Promise.resolve(false), proveStatuslessCodexIdle, - redriveMailbox: () => undefined, + redriveMailbox, ...(options.submitStatuslessCodexPointer ? { submitStatuslessCodexPointer: options.submitStatuslessCodexPointer } : {}), @@ -67,6 +71,7 @@ function makeHarness( setProcessIncarnation: (next: string) => { processIncarnation = next }, + redriveMailbox, writePty } } @@ -230,6 +235,48 @@ describe('statusless Codex mailbox pointer delivery', () => { expect(harness.writePty).not.toHaveBeenCalled() }) + it('retries once when ready output arrived before foreground ownership settled', async () => { + const firstSubmission = deferred() + const submitStatuslessCodexPointer = vi + .fn() + .mockImplementationOnce(async () => firstSubmission.promise) + .mockResolvedValueOnce() + const harness = makeHarness(() => Promise.resolve('pty-1:incarnation-1'), { + submitStatuslessCodexPointer + }) + + harness.delivery.deliverForHandle(MAILBOX) + await vi.waitFor(() => expect(submitStatuslessCodexPointer).toHaveBeenCalledTimes(1)) + harness.delivery.redriveAfterPtyOutput('pty-1') + firstSubmission.reject(new Error('foreground_not_ready')) + await Promise.resolve() + await vi.advanceTimersByTimeAsync(1_000) + await vi.waitFor(() => expect(submitStatuslessCodexPointer).toHaveBeenCalledTimes(2)) + + expect(harness.redriveMailbox).toHaveBeenCalledTimes(1) + expect(harness.db.getUnreadMessages(MAILBOX)[0]?.delivered_at).toEqual(expect.any(String)) + }) + + it('bounds a permanent structured submission failure to one retry', async () => { + const submitStatuslessCodexPointer = vi + .fn() + .mockRejectedValue(new Error('foreground_not_ready')) + const harness = makeHarness(() => Promise.resolve('pty-1:incarnation-1'), { + submitStatuslessCodexPointer + }) + + harness.delivery.deliverForHandle(MAILBOX) + await vi.waitFor(() => expect(submitStatuslessCodexPointer).toHaveBeenCalledTimes(1)) + await vi.advanceTimersByTimeAsync(1_000) + await vi.waitFor(() => expect(submitStatuslessCodexPointer).toHaveBeenCalledTimes(2)) + harness.delivery.redriveAfterPtyOutput('pty-1') + await vi.advanceTimersByTimeAsync(10_000) + + expect(submitStatuslessCodexPointer).toHaveBeenCalledTimes(2) + expect(harness.redriveMailbox).toHaveBeenCalledTimes(1) + expect(harness.db.getUnreadMessages(MAILBOX)[0]?.delivered_at).toBeNull() + }) + it('requeues staged mail when the process is replaced before submit', async () => { const harness = makeHarness(() => Promise.resolve('pty-1:incarnation-1')) diff --git a/src/main/runtime/orchestration/mailbox-statusless-codex-redrive.ts b/src/main/runtime/orchestration/mailbox-statusless-codex-redrive.ts new file mode 100644 index 00000000000..ec170ad3645 --- /dev/null +++ b/src/main/runtime/orchestration/mailbox-statusless-codex-redrive.ts @@ -0,0 +1,95 @@ +const REDRIVE_FALLBACK_MS = 1_000 + +type DeferredRedrive = { + sequence: number + armed: boolean + timer: ReturnType | null +} + +export class OrchestrationMailboxStatuslessCodexRedrive { + private readonly redrivesByPtyId = new Map>() + private readonly armedPtyIds = new Set() + + constructor(private readonly redrive: (mailboxHandle: string) => void) {} + + schedule(ptyId: string, mailboxHandle: string, sequence: number): boolean { + const deferred = this.redrivesByPtyId.get(ptyId) ?? new Map() + const existing = deferred.get(mailboxHandle) + if (existing && existing.sequence >= sequence) { + return false + } + if (existing?.timer) { + clearTimeout(existing.timer) + } + const retry = { sequence, armed: true, timer: null as ReturnType | null } + retry.timer = setTimeout(() => { + const current = this.redrivesByPtyId.get(ptyId)?.get(mailboxHandle) + if (current === retry && retry.armed) { + this.consume(ptyId, mailboxHandle, retry) + } + }, REDRIVE_FALLBACK_MS) + retry.timer.unref?.() + deferred.set(mailboxHandle, retry) + this.redrivesByPtyId.set(ptyId, deferred) + this.armedPtyIds.add(ptyId) + return true + } + + clear(ptyId: string, mailboxHandle: string, sequence: number): void { + const deferred = this.redrivesByPtyId.get(ptyId) + const retry = deferred?.get(mailboxHandle) + if (!deferred || !retry || retry.sequence > sequence) { + return + } + if (retry.timer) { + clearTimeout(retry.timer) + } + deferred.delete(mailboxHandle) + if (deferred.size === 0) { + this.redrivesByPtyId.delete(ptyId) + this.armedPtyIds.delete(ptyId) + } + } + + retirePty(ptyId: string): void { + for (const retry of this.redrivesByPtyId.get(ptyId)?.values() ?? []) { + if (retry.timer) { + clearTimeout(retry.timer) + } + } + this.redrivesByPtyId.delete(ptyId) + this.armedPtyIds.delete(ptyId) + } + + handlePtyOutput(ptyId: string): void { + if (!this.armedPtyIds.delete(ptyId)) { + return + } + const deferred = this.redrivesByPtyId.get(ptyId) + for (const [mailboxHandle, retry] of deferred ?? []) { + if (retry.armed) { + this.consume(ptyId, mailboxHandle, retry) + } + } + } + + private consume(ptyId: string, mailboxHandle: string, retry: DeferredRedrive): void { + retry.armed = false + if (retry.timer) { + clearTimeout(retry.timer) + retry.timer = null + } + this.refreshArmedPty(ptyId) + this.redrive(mailboxHandle) + } + + private refreshArmedPty(ptyId: string): void { + for (const retry of this.redrivesByPtyId.get(ptyId)?.values() ?? []) { + if (retry.armed) { + this.armedPtyIds.add(ptyId) + return + } + } + this.armedPtyIds.delete(ptyId) + } +} diff --git a/src/main/runtime/orchestration/mailbox-statusless-codex-submit.ts b/src/main/runtime/orchestration/mailbox-statusless-codex-submit.ts index 50fefe53774..951b1d4f304 100644 --- a/src/main/runtime/orchestration/mailbox-statusless-codex-submit.ts +++ b/src/main/runtime/orchestration/mailbox-statusless-codex-submit.ts @@ -28,6 +28,8 @@ type StatuslessCodexSubmitDependencies ReadonlySet | undefined getTerminalProcessIncarnation: (terminalHandle: string) => string | null submitStatuslessCodexPointer: SubmitStatuslessCodexPointer + deferRedriveUntilPtyOutput: (ptyId: string, mailboxHandle: string, sequence: number) => boolean + clearDeferredOutputRedrive: (ptyId: string, mailboxHandle: string, sequence: number) => void settle: (ptyId: string, flight: OrchestrationMailboxDeliveryFlight) => void redrive: (mailboxHandle: string, force?: boolean) => void } @@ -61,6 +63,7 @@ export function submitStatuslessCodexMailboxPointer @@ -100,15 +103,21 @@ export function submitStatuslessCodexMailboxPointer { let released = false if (finalizeReservation) { released = - submitted || clearAndRedrive || releaseWithoutRedrive + submitted || clearAndRedrive || releaseWithoutRedrive || deferredUntilOutput ? deps.state.clearWatermark(input.mailboxHandle, input.newestSequence, ptyId) : deps.state.deactivateWatermark(input.mailboxHandle, input.newestSequence, ptyId) } + if (submitted) { + deps.clearDeferredOutputRedrive(ptyId, input.mailboxHandle, input.newestSequence) + } deps.settle(ptyId, flight) if (released && clearAndRedrive) { deps.redrive(input.mailboxHandle, true) diff --git a/src/main/runtime/slept-pane-mail-wake.test.ts b/src/main/runtime/slept-pane-mail-wake.test.ts index 60dca2715b3..4eed5a3328b 100644 --- a/src/main/runtime/slept-pane-mail-wake.test.ts +++ b/src/main/runtime/slept-pane-mail-wake.test.ts @@ -6,6 +6,7 @@ import type { WorkspaceSessionState } from '../../shared/workspace-session-state import { InMemoryOrchestrationMessages, TEST_WORKTREE_ID, + deferred, makeRuntimeStoreWithWorkspaceSession, setInMemoryOrchestrationMessages } from './orca-runtime-test-fixtures.spec' @@ -24,6 +25,12 @@ const TAB_ID = 'tab-slept' const LEAF_ID = '11111111-1111-4111-8111-111111111111' const PANE_KEY = `${TAB_ID}:${LEAF_ID}` +async function flushMicrotasks(): Promise { + for (let index = 0; index < 10; index += 1) { + await Promise.resolve() + } +} + function sleepingRecord( overrides: Partial = {} ): SleepingAgentSessionRecord { @@ -61,6 +68,7 @@ async function sleptPaneRuntime(record: SleepingAgentSessionRecord): Promise<{ connected: boolean tabMountSends: unknown[][] write: ReturnType + confirmForegroundProcess: ReturnType remountWithPty: (ptyId: string) => void remountStatuslessCodex: (ptyId: string) => void setForegroundProcess: (process: string | null) => void @@ -76,12 +84,14 @@ async function sleptPaneRuntime(record: SleepingAgentSessionRecord): Promise<{ let foregroundProcess: string | null = null let confirmedForegroundProcess: string | null | undefined let foregroundConfirmationSupported = true + const confirmForegroundProcess = vi.fn(async () => + confirmedForegroundProcess === undefined ? foregroundProcess : confirmedForegroundProcess + ) runtime.setPtyController({ write, kill: vi.fn(), getForegroundProcess: async () => foregroundProcess, - confirmForegroundProcess: async () => - confirmedForegroundProcess === undefined ? foregroundProcess : confirmedForegroundProcess, + confirmForegroundProcess, supportsForegroundProcessConfirmation: () => foregroundConfirmationSupported } as never) runtime.attachWindow(1) @@ -143,6 +153,7 @@ async function sleptPaneRuntime(record: SleepingAgentSessionRecord): Promise<{ connected: row!.connected, tabMountSends, write, + confirmForegroundProcess, setForegroundProcess: (process) => { foregroundProcess = process }, @@ -512,4 +523,65 @@ describe('mail addressed to a listed slept pane', () => { vi.useRealTimers() } }) + + it('retries after a resumed Codex takes foreground just after its restored screen appears', async () => { + vi.useFakeTimers() + try { + const { + runtime, + db, + handle, + write, + confirmForegroundProcess, + remountStatuslessCodex, + setForegroundProcess + } = await sleptPaneRuntime(sleepingRecord({ agent: 'codex' })) + db.setRun({ id: 'run_test', coordinator_handle: handle, coordinator_pane_key: PANE_KEY }) + const message = db.insertMessage({ + from: 'term_worker', + to: 'run:run_test', + subject: 'worker done', + type: 'worker_done' + }) + runtime.notifyMessageArrived('run:run_test', 'worker_done') + await Promise.resolve() + await Promise.resolve() + + const firstForegroundConfirmation = deferred() + confirmForegroundProcess + .mockImplementationOnce(() => firstForegroundConfirmation.promise) + .mockResolvedValue('codex') + setForegroundProcess('codex') + write.mockImplementation((ptyId: string, data: string) => { + if (data === '\r') { + runtime.onPtyData(ptyId, '\x1b]0;Codex working\x07', 102) + } + return true + }) + remountStatuslessCodex('pty-codex-woken') + await vi.advanceTimersByTimeAsync(2_000) + await vi.waitFor(() => expect(confirmForegroundProcess).toHaveBeenCalledTimes(1)) + runtime.onPtyData('pty-codex-woken', 'restored output\n', 101) + firstForegroundConfirmation.resolve('zsh') + await flushMicrotasks() + expect(write).not.toHaveBeenCalled() + + await vi.advanceTimersByTimeAsync(12_000) + expect(confirmForegroundProcess).toHaveBeenCalledTimes(3) + await vi.waitFor(() => + expect(write).toHaveBeenCalledWith( + 'pty-codex-woken', + expect.stringContaining( + `${AGENT_PROMPT_BRACKETED_PASTE_START}\nYou have 1 orchestration message` + ) + ) + ) + + expect(write).toHaveBeenCalledWith('pty-codex-woken', '\r') + expect(message.delivered_at).toEqual(expect.any(String)) + db.close() + } finally { + vi.useRealTimers() + } + }) })