fix: settle failed outbox mutations and retire launch callers

This commit is contained in:
Merge Sim
2026-09-07 23:58:12 -07:00
parent daf83559a1
commit e9c73bd572
7 changed files with 385 additions and 18 deletions
@@ -34,6 +34,6 @@ export function enqueueStructuredAgentSessionLaunchPrompt(
return staged
}
export function discardStructuredAgentSessionLaunchOutbox(sessionId: string): void {
transitionOutbox(sessionId, () => [])
export function discardStructuredAgentSessionLaunchOutbox(sessionId: string): boolean {
return transitionOutbox(sessionId, () => []).ok
}
@@ -27,7 +27,7 @@ export function transitionOutbox(
transitionRevision: (previous?.transitionRevision ?? 0) + 1
}
})
return writeOutbox(sessionId, stamped, () => {
const ok = writeOutbox(sessionId, stamped, () => {
const retained = new Set(stamped.map(claimKey))
for (const entry of current) {
if (!retained.has(claimKey(entry))) {
@@ -36,8 +36,15 @@ export function transitionOutbox(
}
publishOutboxSettlements(current, stamped, accepted)
})
? { ok: true, entries: stamped }
: { ok: false, entries: current }
if (!ok) {
const successors = new Map(next.map((entry) => [claimKey(entry), entry]))
for (const entry of current) {
if (successors.get(claimKey(entry)) !== entry) {
settleOutboxObservation(entry, 'unavailable')
}
}
}
return { ok, entries: ok ? stamped : current }
}
export function transitionOutboxEntry(
@@ -63,9 +70,6 @@ export function transitionOutboxEntry(
}),
accepted ? [expected] : []
)
if (!result.ok) {
settleOutboxObservation(expected, 'unavailable')
}
return {
ok: result.ok,
changed: changed && result.ok,
@@ -0,0 +1,247 @@
// @vitest-environment happy-dom
import { act, cleanup, renderHook } from '@testing-library/react'
import { afterEach, beforeEach, expect, it, vi } from 'vitest'
import type { AgentJournalSubmission } from '../../../../shared/agent-session-journal-types'
import {
enqueueStructuredAgentSessionLaunchPrompt as enqueue,
discardStructuredAgentSessionLaunchOutbox as discard
} from './structured-agent-session-launch-outbox'
import { readOutbox } from './structured-agent-session-outbox-storage'
import { observeOutboxSettlement } from './structured-agent-session-outbox-settlement'
import { settleStructuredAgentLaunchPrompt as settle } from '../../lib/structured-agent-session-launch-prompt'
import {
addStructuredLaunchCaller,
createStructuredLaunchCallerGroup,
settleStructuredLaunchCallersWithoutFallback,
claimStructuredLaunchCallerFallback
} from '../../lib/structured-agent-session-launch-callers'
const mocks = vi.hoisted(() => ({ call: vi.fn() }))
vi.mock('@/runtime/structured-agent-session-client', () => ({
callStructuredAgentSession: mocks.call
}))
import { useStructuredAgentSessionOutbox } from './use-structured-agent-session-outbox'
const sessionId = 'independent-restart-edge'
const target = { kind: 'local' } as const
const empty: AgentJournalSubmission[] = []
const receipt = { sessionId, fence: 1 }
const accepted = { ok: true, value: { submission: { dispatchState: 'accepted' } } }
function mount() {
return renderHook(
({ submissions }) =>
useStructuredAgentSessionOutbox({ sessionId, target, fence: 1, submissions }),
{ initialProps: { submissions: empty } }
)
}
function launch(stagedEntry: NonNullable<ReturnType<typeof enqueue>>, callback = vi.fn()) {
return settle({
stagedEntry,
options: { prompt: 'synthetic', onPromptDelivered: callback },
launchResult: Promise.resolve(receipt)
})!
}
beforeEach(() => {
localStorage.clear()
mocks.call.mockReset().mockImplementation(() => new Promise(() => {}))
vi.useFakeTimers()
})
afterEach(() => {
vi.restoreAllMocks()
cleanup()
discard(sessionId)
vi.useRealTimers()
})
it('control: journal acceptance settles the caller with a still-pending send RPC', async () => {
const staged = enqueue(sessionId, 'synthetic')!
const pane = mount()
const callback = vi.fn()
const delivery = launch(staged, callback)
await act(async () => {})
pane.rerender({
submissions: [
{
clientMessageId: staged.clientMessageId,
dispatchState: 'accepted'
} as AgentJournalSubmission
]
})
await act(async () => {})
expect(await delivery).toMatchObject({ delivered: true })
expect(callback).toHaveBeenCalledTimes(1)
expect(mocks.call).toHaveBeenCalledTimes(1)
})
it('journal acceptance write failure terminates observers as unavailable just like RPC completion write failure', async () => {
const staged = enqueue(sessionId, 'synthetic')!
const pane = mount()
let outcome = 'pending'
const callback = vi.fn()
void launch(staged, callback).then((result) => {
outcome = result.delivered ? 'accepted' : 'unavailable'
})
await act(async () => {})
const spy = vi.spyOn(localStorage, 'removeItem').mockImplementation(() => {
throw new Error('synthetic storage failure')
})
pane.rerender({
submissions: [
{
clientMessageId: staged.clientMessageId,
dispatchState: 'accepted'
} as AgentJournalSubmission
]
})
await act(async () => {})
spy.mockRestore()
await act(async () => {
await vi.advanceTimersByTimeAsync(160000)
})
expect(pane.result.current.error).toBe('Message could not be saved to the outbox')
expect(pane.result.current.recoveryPaused).toBe(false)
expect(readOutbox(sessionId, false)[0].state).toBe('dispatching')
expect(callback).not.toHaveBeenCalled()
console.log('journal-storage-failure', { outcome, calls: mocks.call.mock.calls.length })
expect(outcome).toBe('unavailable')
})
it('control: failed RPC acceptance persistence settles unavailable without a hanging caller', async () => {
const staged = enqueue(sessionId, 'synthetic')!
const response = Promise.withResolvers<unknown>()
mocks.call.mockReturnValue(response.promise)
mount()
const delivery = launch(staged)
await act(async () => {})
const spy = vi.spyOn(localStorage, 'removeItem').mockImplementation(() => {
throw new Error('synthetic storage failure')
})
await act(async () => {
response.resolve(accepted)
})
spy.mockRestore()
expect(await delivery).toMatchObject({ delivered: false })
expect(await observeOutboxSettlement(staged)).toBe('unavailable')
})
it('failed discard terminates waiting launch observers without claiming acceptance', async () => {
const staged = enqueue(sessionId, 'synthetic')!
let outcome = 'pending'
void launch(staged).then((result) => {
outcome = result.delivered ? 'accepted' : 'unavailable'
})
await act(async () => {})
const spy = vi.spyOn(localStorage, 'removeItem').mockImplementation(() => {
throw new Error('synthetic storage failure')
})
expect(discard(sessionId)).toBe(false)
spy.mockRestore()
await act(async () => {
await vi.advanceTimersByTimeAsync(160000)
})
expect(readOutbox(sessionId, false)).toHaveLength(1)
expect(outcome).toBe('unavailable')
})
it('completed caller registrations retire while a subsequent operation keeps the launch group live', async () => {
const group = createStructuredLaunchCallerGroup()
settleStructuredLaunchCallersWithoutFallback(group, 'published')
const responses: ReturnType<typeof Promise.withResolvers<unknown>>[] = []
mocks.call.mockImplementation(() => {
const response = Promise.withResolvers<unknown>()
responses.push(response)
return response.promise
})
function join() {
const stagedEntry = enqueue(sessionId, 'synthetic')!
const caller = addStructuredLaunchCaller({
group,
launchResult: Promise.resolve(receipt),
options: { prompt: 'synthetic' },
stagedEntry
})
void claimStructuredLaunchCallerFallback(group, caller, () => ({
delivered: false,
failureNotified: false
}))
return caller
}
const held: ReturnType<typeof join>[] = []
let current = join()
await act(async () => {})
for (let i = 0; i < 64; i++) {
const next = join()
await act(async () => {
responses[i].resolve(accepted)
})
expect(await current.promptDeliveryResult).toMatchObject({ delivered: true })
held.push(current)
current = next
}
mocks.call.mockRejectedValue(new Error('connection closed'))
await act(async () => {
responses[64].resolve({ ok: true, value: { submission: { dispatchState: 'unknown' } } })
})
const pane = mount()
for (let i = 0; i < 10; i++) {
await act(async () => {
await vi.advanceTimersByTimeAsync(16000)
})
}
expect(pane.result.current.recoveryPaused).toBe(true)
expect(group.promptDeliveryResults.size).toBe(1)
expect(readOutbox(sessionId, false)).toHaveLength(1)
const retainedCallbacks = [...group.entries].filter(
(caller) => caller.refusalFallback.callback !== null
).length
console.log('caller-retention', {
registrations: group.entries.size,
retainedCallbacks,
activeResults: group.promptDeliveryResults.size,
calls: mocks.call.mock.calls.length
})
expect(group.entries.size).toBe(1)
expect(retainedCallbacks).toBe(0)
for (const caller of held) {
expect(caller.refusalFallback.callback).toBe(null)
expect(await caller.promptDeliveryResult).toMatchObject({ delivered: true })
}
})
it('a finite unknown RPC after failed journal acceptance cannot strand an already-accepted operation', async () => {
const staged = enqueue(sessionId, 'synthetic')!
const response = Promise.withResolvers<unknown>()
mocks.call.mockReturnValue(response.promise)
const pane = mount()
let outcome = 'pending'
void launch(staged).then((result) => {
outcome = result.delivered ? 'accepted' : 'unavailable'
})
await act(async () => {})
const spy = vi.spyOn(localStorage, 'removeItem').mockImplementation(() => {
throw new Error('synthetic storage failure')
})
pane.rerender({
submissions: [
{
clientMessageId: staged.clientMessageId,
dispatchState: 'accepted'
} as AgentJournalSubmission
]
})
spy.mockRestore()
await act(async () => {
response.resolve({ ok: true, value: { submission: { dispatchState: 'unknown' } } })
})
await act(async () => {
await vi.advanceTimersByTimeAsync(160000)
})
expect(pane.result.current.error).toBe(null)
expect(pane.result.current.recoveryPaused).toBe(false)
expect(readOutbox(sessionId, false)[0].state).toBe('unconfirmed')
expect(mocks.call).toHaveBeenCalledTimes(1)
console.log('finite-journal-storage-failure', {
outcome,
error: pane.result.current.error,
recoveryPaused: pane.result.current.recoveryPaused
})
expect(outcome).not.toBe('pending')
})
@@ -0,0 +1,64 @@
// @vitest-environment happy-dom
import { expect, it, vi } from 'vitest'
import {
addStructuredLaunchCaller,
claimStructuredLaunchCallerFallback,
createStructuredLaunchCallerGroup,
settleStructuredLaunchCallersWithFallback,
settleStructuredLaunchCallersWithoutFallback
} from './structured-agent-session-launch-callers'
function join(group: ReturnType<typeof createStructuredLaunchCallerGroup>) {
return addStructuredLaunchCaller({
group,
launchResult: Promise.resolve({ sessionId: 'caller-retirement', fence: 1 }),
options: {},
stagedEntry: null
})
}
it('retires promptless callers and refuses to retain late unusable fallbacks', async () => {
const group = createStructuredLaunchCallerGroup()
const caller = join(group)
settleStructuredLaunchCallersWithoutFallback(group, 'published')
const fallback = vi.fn()
expect(await claimStructuredLaunchCallerFallback(group, caller, fallback)).toBe(false)
expect(group.entries.size).toBe(0)
expect(caller.refusalFallback.callback).toBe(null)
expect(fallback).not.toHaveBeenCalled()
})
it('keeps unresolved fallbacks and preserves successful receipts after membership retirement', async () => {
const group = createStructuredLaunchCallerGroup()
const first = join(group)
const second = join(group)
const pending = Promise.withResolvers<void>()
const ran = claimStructuredLaunchCallerFallback(group, first, () => pending.promise)
const absent = second.refusalFallback.promise
settleStructuredLaunchCallersWithFallback(group)
expect(group.entries.size).toBe(1)
expect(group.entries.has(first)).toBe(true)
expect(group.refusalSettlement.settled).toBe(false)
pending.resolve()
expect(await ran).toBe(true)
expect(await absent).toBe(false)
expect(await group.refusalSettlement.promise).toBe(true)
expect(group.entries.size).toBe(0)
expect(first.refusalFallback.callback).toBe(null)
expect(await first.refusalFallback.promise).toBe(true)
})
it('retires failed fallbacks without losing the failure receipt', async () => {
const group = createStructuredLaunchCallerGroup()
const caller = join(group)
const failure = new Error('synthetic fallback failed')
const result = claimStructuredLaunchCallerFallback(group, caller, () => {
throw failure
})
const callerFailure = expect(result).rejects.toBe(failure)
const groupFailure = expect(group.refusalSettlement.promise).rejects.toBe(failure)
settleStructuredLaunchCallersWithFallback(group)
await Promise.all([callerFailure, groupFailure])
expect(group.entries.size).toBe(0)
expect(caller.refusalFallback.callback).toBe(null)
})
@@ -44,6 +44,7 @@ export type StructuredLaunchCallerGroup = {
resolve: (ran: boolean) => void
reject: (error: unknown) => void
settled: boolean
ran: boolean
failure: { error: unknown } | null
}
onSettled: () => void
@@ -60,19 +61,37 @@ export function createStructuredLaunchCallerGroup(): StructuredLaunchCallerGroup
resolve: refusalSettlement.resolve,
reject: refusalSettlement.reject,
settled: false,
ran: false,
failure: null
},
onSettled: () => {}
}
}
function settleCallerWithoutFallback(caller: StructuredLaunchCaller): void {
function retireSettledCaller(
group: StructuredLaunchCallerGroup,
caller: StructuredLaunchCaller
): void {
if (
caller.refusalFallback.settled &&
(!caller.promptDeliveryResult || !group.promptDeliveryResults.has(caller.promptDeliveryResult))
) {
group.entries.delete(caller)
}
}
function settleCallerWithoutFallback(
group: StructuredLaunchCallerGroup,
caller: StructuredLaunchCaller
): void {
if (caller.refusalFallback.settled) {
return
}
caller.refusalFallback.settled = true
caller.refusalFallback.callback = null
caller.refusalFallback.resolve(false)
caller.refusalFallback.resolvePromptDelivery(null)
retireSettledCaller(group, caller)
}
function finalizeRefusalSettlement(group: StructuredLaunchCallerGroup): void {
@@ -87,7 +106,7 @@ function finalizeRefusalSettlement(group: StructuredLaunchCallerGroup): void {
if (group.refusalSettlement.failure) {
group.refusalSettlement.reject(group.refusalSettlement.failure.error)
} else {
group.refusalSettlement.resolve([...group.entries].some((caller) => caller.refusalFallback.ran))
group.refusalSettlement.resolve(group.refusalSettlement.ran)
}
group.onSettled()
}
@@ -102,7 +121,7 @@ function runCallerRefusalFallback(
caller.refusalFallback.started = true
const fallback = caller.refusalFallback.callback
if (!fallback) {
settleCallerWithoutFallback(caller)
settleCallerWithoutFallback(group, caller)
finalizeRefusalSettlement(group)
return
}
@@ -111,6 +130,7 @@ function runCallerRefusalFallback(
.then(
(result) => {
caller.refusalFallback.ran = true
group.refusalSettlement.ran = true
caller.refusalFallback.resolve(true)
caller.refusalFallback.resolvePromptDelivery(result ?? null)
},
@@ -122,17 +142,21 @@ function runCallerRefusalFallback(
)
.finally(() => {
caller.refusalFallback.settled = true
caller.refusalFallback.callback = null
retireSettledCaller(group, caller)
finalizeRefusalSettlement(group)
})
}
function trackPromptDelivery(
group: StructuredLaunchCallerGroup,
promptDeliveryResult: Promise<StructuredPromptDeliveryResult>
promptDeliveryResult: Promise<StructuredPromptDeliveryResult>,
caller: StructuredLaunchCaller
): void {
group.promptDeliveryResults.add(promptDeliveryResult)
const settled = (): void => {
group.promptDeliveryResults.delete(promptDeliveryResult)
retireSettledCaller(group, caller)
group.onSettled()
}
void promptDeliveryResult.then(settled, settled)
@@ -177,10 +201,10 @@ export function addStructuredLaunchCaller(args: {
return { delivered: false, failureNotified: true }
})
if (caller.promptDeliveryResult) {
trackPromptDelivery(args.group, caller.promptDeliveryResult)
trackPromptDelivery(args.group, caller.promptDeliveryResult, caller)
}
if (['published', 'failed', 'cancelled'].includes(args.group.outcome)) {
settleCallerWithoutFallback(caller)
settleCallerWithoutFallback(args.group, caller)
} else if (args.group.outcome === 'refused') {
queueMicrotask(() => runCallerRefusalFallback(args.group, caller))
}
@@ -193,7 +217,7 @@ export function settleStructuredLaunchCallersWithoutFallback(
): void {
group.outcome = outcome
for (const caller of group.entries) {
settleCallerWithoutFallback(caller)
settleCallerWithoutFallback(group, caller)
}
if (!group.refusalSettlement.settled) {
group.refusalSettlement.settled = true
@@ -220,7 +244,9 @@ export function claimStructuredLaunchCallerFallback(
caller: StructuredLaunchCaller,
fallback: StructuredRefusalFallback
): Promise<boolean> {
caller.refusalFallback.callback ??= fallback
if (!caller.refusalFallback.settled && !caller.refusalFallback.started) {
caller.refusalFallback.callback ??= fallback
}
if (group.outcome === 'refused') {
runCallerRefusalFallback(group, caller)
}
@@ -234,7 +260,7 @@ export function releaseStructuredLaunchCallerAfterUnknownOutcome(
if (group.outcome !== 'unknown' || !group.entries.delete(caller)) {
return false
}
settleCallerWithoutFallback(caller)
settleCallerWithoutFallback(group, caller)
group.onSettled()
return true
}
@@ -69,6 +69,7 @@ import {
import { refreshLocalStructuredSessionTabs } from '@/runtime/local-structured-session-tabs-sync'
import {
cancelStructuredAgentLaunch,
getStructuredAgentLaunchStatus,
startStructuredAgentLaunch
} from './structured-agent-session-launch'
import { readOutbox } from '@/components/native-chat/structured-agent-session-outbox-storage'
@@ -672,6 +673,29 @@ describe('startStructuredAgentLaunch', () => {
await flushLaunchSettlement()
})
it('keeps cancellation retryable when discarding the staged prompt cannot persist', async () => {
const worktreeId = 'wt-close-storage-failure'
const intent = launchIntent(worktreeId)
mocks.createIntent.mockReturnValueOnce(intent)
mocks.launch.mockImplementationOnce(() => new Promise(() => {}))
startStructuredAgentLaunch(worktreeId, 'codex', { prompt: 'retain on failed cancellation' })
const staged = readOutbox(intent.sessionId, false)
const spy = vi.spyOn(localStorage, 'removeItem').mockImplementation(() => {
throw new Error('synthetic storage failure')
})
try {
expect(cancelStructuredAgentLaunch(worktreeId, intent.sessionId)).toBe(false)
expect(readOutbox(intent.sessionId, false)).toEqual(staged)
expect(mocks.abandonIntent).not.toHaveBeenCalled()
expect(getStructuredAgentLaunchStatus(worktreeId, 'codex')).toBe('pending')
} finally {
spy.mockRestore()
}
expect(cancelStructuredAgentLaunch(worktreeId, intent.sessionId)).toBe(true)
expect(readOutbox(intent.sessionId, false)).toEqual([])
expect(mocks.abandonIntent).toHaveBeenCalledWith(intent)
})
it('suppresses a close that races the retry verification catch', async () => {
const worktreeId = 'wt-retry-close-race'
const intent = launchIntent(worktreeId, 'session-retry-close-race')
@@ -305,10 +305,12 @@ export function cancelStructuredAgentLaunch(worktreeId: string, sessionId: strin
if (!state) {
return false
}
if (!discardStructuredAgentSessionLaunchOutbox(state.intent.sessionId)) {
return false
}
state.cancelled = true
settleStructuredLaunchCallersWithoutFallback(state.callers, 'cancelled')
cleanupLaunchState(state)
discardStructuredAgentSessionLaunchOutbox(state.intent.sessionId)
abandonStructuredAgentSessionLaunchIntent(state.intent)
notifyStructuredLaunchListeners()
return true