fix(orchestration): retry resumed Codex mail delivery

This commit is contained in:
Merge Sim
2026-09-05 00:30:12 -07:00
parent 70ea5388b7
commit 76feffa1ec
7 changed files with 269 additions and 15 deletions
@@ -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
}
}
@@ -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<TWaiter extends OrchestrationMessageWaiter> = {
export class OrchestrationMailboxPointerDelivery<TWaiter extends OrchestrationMessageWaiter> {
private readonly state = new OrchestrationMailboxPointerState()
private readonly statuslessCodexProofs: OrchestrationMailboxStatuslessCodexProofCoordinator
private readonly statuslessCodexRedrives: OrchestrationMailboxStatuslessCodexRedrive
constructor(private readonly deps: PointerDeliveryDependencies<TWaiter>) {
this.statuslessCodexProofs = new OrchestrationMailboxStatuslessCodexProofCoordinator(deps)
this.statuslessCodexRedrives = new OrchestrationMailboxStatuslessCodexRedrive((mailboxHandle) =>
this.redrive(mailboxHandle, true)
)
}
deliverForHandle(handle: string, reservedTypes?: ReadonlySet<string>): void {
@@ -122,7 +127,7 @@ export class OrchestrationMailboxPointerDelivery<TWaiter extends OrchestrationMe
if (!db || !mailboxHandle.startsWith('run:')) {
return
}
if (!this.deps.getTerminalHandleForLeafKey(this.leafKey(leaf))) {
if (!this.deps.getTerminalHandleForLeafKey(this.deps.getLeafKey(leaf.tabId, leaf.leafId))) {
return
}
if (
@@ -171,7 +176,7 @@ export class OrchestrationMailboxPointerDelivery<TWaiter extends OrchestrationMe
mailboxHandle,
newestSequence,
leaf.ptyId,
this.leafKey(leaf)
this.deps.getLeafKey(leaf.tabId, leaf.leafId)
)
) {
return
@@ -202,7 +207,19 @@ export class OrchestrationMailboxPointerDelivery<TWaiter extends OrchestrationMe
? { requestSleepingRecipientWake: this.deps.requestSleepingRecipientWake }
: {}),
...(this.deps.submitStatuslessCodexPointer
? { submitStatuslessCodexPointer: this.deps.submitStatuslessCodexPointer }
? {
submitStatuslessCodexPointer: this.deps.submitStatuslessCodexPointer,
deferRedriveUntilPtyOutput: (
ptyId: string,
mailboxHandle: string,
sequence: number
) => 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<TWaiter extends OrchestrationMe
retirePty(ptyId: string): void {
this.statuslessCodexProofs.retirePty(ptyId)
this.statuslessCodexRedrives.retirePty(ptyId)
const { flight, releasedMailboxes } = this.state.retirePty(ptyId)
if (flight?.enterTimer != null) {
clearTimeout(flight.enterTimer)
@@ -236,13 +254,17 @@ export class OrchestrationMailboxPointerDelivery<TWaiter extends OrchestrationMe
}
}
redriveAfterPtyOutput(ptyId: string): void {
this.statuslessCodexRedrives.handlePtyOutput(ptyId)
}
private redeliverAfterProbe(
leaf: OrchestrationMailboxLeaf,
ptyId: string,
mailboxHandle: string,
statuslessIdleProof?: OrchestrationStatuslessIdleProof
): void {
const currentLeaf = this.deps.getLeaf(this.leafKey(leaf))
const currentLeaf = this.deps.getLeaf(this.deps.getLeafKey(leaf.tabId, leaf.leafId))
if (
currentLeaf?.ptyId === ptyId &&
((currentLeaf.lastAgentStatus === 'idle' && currentLeaf.lastAgentStatusObservedLive) ||
@@ -267,7 +289,9 @@ export class OrchestrationMailboxPointerDelivery<TWaiter extends OrchestrationMe
return
}
for (const [mailboxHandle, delivery] of parked) {
const currentLeaf = this.deps.getLeaf(this.leafKey(delivery.leaf))
const currentLeaf = this.deps.getLeaf(
this.deps.getLeafKey(delivery.leaf.tabId, delivery.leaf.leafId)
)
if (
currentLeaf?.ptyId !== ptyId ||
this.deps.mailboxOwner.resolve(currentLeaf, mailboxHandle) !== mailboxHandle
@@ -297,8 +321,4 @@ export class OrchestrationMailboxPointerDelivery<TWaiter extends OrchestrationMe
}
})
}
private leafKey(leaf: OrchestrationMailboxLeaf): string {
return this.deps.getLeafKey(leaf.tabId, leaf.leafId)
}
}
@@ -33,6 +33,8 @@ type PointerStageDependencies<TWaiter extends OrchestrationMessageWaiter> = {
isLeafPtyProvenAbsent: (ptyId: string) => Promise<boolean>
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<boolean>
settle: (ptyId: string, flight: OrchestrationMailboxDeliveryFlight) => void
redrive: (mailboxHandle: string, force?: boolean) => void
@@ -63,7 +65,12 @@ export function stageOrchestrationMailboxPointer<TWaiter extends OrchestrationMe
) {
return
}
if (input.statuslessIdleProof && deps.submitStatuslessCodexPointer) {
if (
input.statuslessIdleProof &&
deps.submitStatuslessCodexPointer &&
deps.deferRedriveUntilPtyOutput &&
deps.clearDeferredOutputRedrive
) {
submitStatuslessCodexMailboxPointer(
{
mailboxOwner: deps.mailboxOwner,
@@ -74,6 +81,8 @@ export function stageOrchestrationMailboxPointer<TWaiter extends OrchestrationMe
getMessageWaiters: deps.getMessageWaiters,
getTerminalProcessIncarnation: deps.getTerminalProcessIncarnation,
submitStatuslessCodexPointer: deps.submitStatuslessCodexPointer,
deferRedriveUntilPtyOutput: deps.deferRedriveUntilPtyOutput,
clearDeferredOutputRedrive: deps.clearDeferredOutputRedrive,
settle: deps.settle,
redrive: deps.redrive
},
@@ -35,7 +35,11 @@ function makeHarness(
db.insertMessage({ from: 'term-sender', to: MAILBOX, subject: 'wake work' })
}
const writePty = vi.fn().mockReturnValue(true)
const delivery = new OrchestrationMailboxPointerDelivery({
let delivery: OrchestrationMailboxPointerDelivery<never>
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<void>()
const submitStatuslessCodexPointer = vi
.fn<SubmitStatuslessCodexPointer>()
.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<SubmitStatuslessCodexPointer>()
.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'))
@@ -0,0 +1,95 @@
const REDRIVE_FALLBACK_MS = 1_000
type DeferredRedrive = {
sequence: number
armed: boolean
timer: ReturnType<typeof setTimeout> | null
}
export class OrchestrationMailboxStatuslessCodexRedrive {
private readonly redrivesByPtyId = new Map<string, Map<string, DeferredRedrive>>()
private readonly armedPtyIds = new Set<string>()
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<typeof setTimeout> | 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)
}
}
@@ -28,6 +28,8 @@ type StatuslessCodexSubmitDependencies<TWaiter extends OrchestrationMessageWaite
getMessageWaiters: (mailboxHandle: string) => ReadonlySet<TWaiter> | 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<TWaiter extends Orchestratio
let submitted = false
let clearAndRedrive = false
let releaseWithoutRedrive = false
let deferredUntilOutput = false
let finalizeReservation = true
void Promise.resolve()
.then(() =>
@@ -100,15 +103,21 @@ export function submitStatuslessCodexMailboxPointer<TWaiter extends Orchestratio
const state = targetState(deps, input, ptyId, flight)
clearAndRedrive = state === 'invalid'
releaseWithoutRedrive = state === 'released'
deferredUntilOutput =
state === 'current' &&
deps.deferRedriveUntilPtyOutput(ptyId, input.mailboxHandle, input.newestSequence)
})
.finally(() => {
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)
+74 -2
View File
@@ -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<void> {
for (let index = 0; index < 10; index += 1) {
await Promise.resolve()
}
}
function sleepingRecord(
overrides: Partial<SleepingAgentSessionRecord> = {}
): SleepingAgentSessionRecord {
@@ -61,6 +68,7 @@ async function sleptPaneRuntime(record: SleepingAgentSessionRecord): Promise<{
connected: boolean
tabMountSends: unknown[][]
write: ReturnType<typeof vi.fn>
confirmForegroundProcess: ReturnType<typeof vi.fn>
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<string | null>()
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()
}
})
})