mirror of
https://github.com/stablyai/orca.git
synced 2026-10-01 16:02:03 +00:00
fix(native-chat): preserve restart refusal and teardown evidence
This commit is contained in:
@@ -67,11 +67,13 @@ export function structuredAgentSessionHostTeardownPhases(collaborators: {
|
||||
{
|
||||
name: 'record-resume-markers',
|
||||
run: () =>
|
||||
withPhaseTimeout(collaborators.recordResumeMarkers, RESUME_MARKER_RECORD_TIMEOUT_MS).catch(
|
||||
(error: unknown) => {
|
||||
console.warn('[structured-agent-session] recording recovery capsule failed', error)
|
||||
}
|
||||
)
|
||||
withPhaseTimeout(async () => {
|
||||
// Capture accepted prompts before eviction cancels their waiting-on-user evidence.
|
||||
await collaborators.runtimeState.flushAllEventSinks()
|
||||
await collaborators.recordResumeMarkers()
|
||||
}, RESUME_MARKER_RECORD_TIMEOUT_MS).catch((error: unknown) => {
|
||||
console.warn('[structured-agent-session] recording recovery capsule failed', error)
|
||||
})
|
||||
},
|
||||
{ name: 'dispose-holds', run: () => collaborators.holds.dispose() },
|
||||
{ name: 'stop-lease-renewal', run: () => collaborators.runtimeState.stopLeaseRenewal() },
|
||||
|
||||
+9
-1
@@ -42,7 +42,15 @@ export async function runSettledAgentSessionMutation<TValue>(input: {
|
||||
)
|
||||
return outcome
|
||||
} catch (error) {
|
||||
await settle({ status: 'unknown' })
|
||||
try {
|
||||
await settle({ status: 'unknown' })
|
||||
} catch (settlementError) {
|
||||
// Bookkeeping must not replace the operation's proof of whether dispatch began.
|
||||
console.warn(
|
||||
'[structured-agent-session] operation uncertainty persistence failed',
|
||||
settlementError
|
||||
)
|
||||
}
|
||||
throw error
|
||||
}
|
||||
}
|
||||
|
||||
+57
-24
@@ -216,9 +216,14 @@ it('refuses a completed turn even when its original delivery acknowledgement is
|
||||
expect(host.isHeld(SESSION)).toBe(false)
|
||||
})
|
||||
|
||||
it.each(['completion', 'message'] as const)(
|
||||
'includes a provider %s queued at send admission in the restart decision',
|
||||
async (newer) => {
|
||||
it.each([
|
||||
{ newer: 'completion', settlementFails: false },
|
||||
{ newer: 'message', settlementFails: false },
|
||||
{ newer: 'completion', settlementFails: true },
|
||||
{ newer: 'message', settlementFails: true }
|
||||
])(
|
||||
'refuses queued $newer before dispatch even if later settlement fails: $settlementFails',
|
||||
async ({ newer, settlementFails }) => {
|
||||
const { host, store, acquire, dispatch } = await interruptedRestart('submission', false)
|
||||
expect(await host.restartResume.list()).toHaveLength(1)
|
||||
await host.hold(SESSION, 'pane')
|
||||
@@ -226,9 +231,18 @@ it.each(['completion', 'message'] as const)(
|
||||
if (!events) {
|
||||
throw new Error('missing resumed provider event sink')
|
||||
}
|
||||
const warning = vi.spyOn(console, 'warn').mockImplementation(() => {})
|
||||
const settle = store.recordOperationOutcome.bind(store)
|
||||
vi.spyOn(store, 'recordOperationOutcome').mockImplementationOnce(async (input) => {
|
||||
let admitted = false
|
||||
vi.spyOn(store, 'recordOperationOutcome').mockImplementation(async (input) => {
|
||||
if (admitted) {
|
||||
if (settlementFails) {
|
||||
throw new Error('operation outcome could not be persisted')
|
||||
}
|
||||
return settle(input)
|
||||
}
|
||||
await settle(input)
|
||||
admitted = true
|
||||
events.appendItem(
|
||||
{ provider: 'codex', threadId: THREAD, turnId: 'original-turn', ordinal: 2 },
|
||||
{ kind: 'status', text: 'Provider finished its last action' }
|
||||
@@ -251,6 +265,18 @@ it.each(['completion', 'message'] as const)(
|
||||
expect(host.journalSnapshot(SESSION).submissions).toHaveLength(1)
|
||||
host.release(SESSION, 'pane')
|
||||
expect(host.isHeld(SESSION)).toBe(false)
|
||||
expect(await host.restartResume.continueAfterRestart([SESSION], 'retry')).toEqual({
|
||||
resumed: [],
|
||||
continued: []
|
||||
})
|
||||
expect(dispatch).not.toHaveBeenCalled()
|
||||
if (settlementFails) {
|
||||
expect(warning).toHaveBeenCalledWith(
|
||||
'[structured-agent-session] operation uncertainty persistence failed',
|
||||
expect.any(Error)
|
||||
)
|
||||
}
|
||||
warning.mockRestore()
|
||||
}
|
||||
)
|
||||
|
||||
@@ -273,28 +299,35 @@ it('replays the same logical continuation through the durable send ledger', asyn
|
||||
expect(dispatch).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('keeps delivery unconfirmed when operation settlement fails after dispatch', async () => {
|
||||
const { host, store, dispatch } = await interruptedRestart()
|
||||
expect(await host.restartResume.list()).toHaveLength(1)
|
||||
await host.hold(SESSION, 'pane')
|
||||
const settle = store.recordOperationOutcome.bind(store)
|
||||
vi.spyOn(store, 'recordOperationOutcome').mockImplementation(async (input) => {
|
||||
if (input.outcome.status === 'succeeded') {
|
||||
throw new Error('operation outcome could not be persisted')
|
||||
}
|
||||
return settle(input)
|
||||
})
|
||||
it.each([false, true])(
|
||||
'keeps dispatched delivery unconfirmed when settlement fails (uncertainty write fails: %s)',
|
||||
async (uncertaintyFails) => {
|
||||
const { host, store, dispatch } = await interruptedRestart()
|
||||
expect(await host.restartResume.list()).toHaveLength(1)
|
||||
await host.hold(SESSION, 'pane')
|
||||
const settle = store.recordOperationOutcome.bind(store)
|
||||
let dispatched = false
|
||||
vi.spyOn(store, 'recordOperationOutcome').mockImplementation(async (input) => {
|
||||
if (input.outcome.status === 'succeeded') {
|
||||
dispatched = true
|
||||
}
|
||||
if (dispatched && (input.outcome.status === 'succeeded' || uncertaintyFails)) {
|
||||
throw new Error('operation outcome could not be persisted')
|
||||
}
|
||||
return settle(input)
|
||||
})
|
||||
|
||||
const result = await host.restartResume.continueAfterRestart([SESSION], 'modal')
|
||||
const result = await host.restartResume.continueAfterRestart([SESSION], 'modal')
|
||||
|
||||
expect(dispatch).toHaveBeenCalledTimes(1)
|
||||
expect(host.journalSnapshot(SESSION).submissions[0]?.dispatchState).toBe('accepted')
|
||||
expect(result.continued).toMatchObject([{ sessionId: SESSION, outcome: 'unknown' }])
|
||||
host.release(SESSION, 'pane')
|
||||
expect(host.isHeld(SESSION)).toBe(false)
|
||||
await host.restartResume.continueAfterRestart([SESSION], 'retry')
|
||||
expect(dispatch).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
expect(dispatch).toHaveBeenCalledTimes(1)
|
||||
expect(host.journalSnapshot(SESSION).submissions[0]?.dispatchState).toBe('accepted')
|
||||
expect(result.continued).toMatchObject([{ sessionId: SESSION, outcome: 'unknown' }])
|
||||
host.release(SESSION, 'pane')
|
||||
expect(host.isHeld(SESSION)).toBe(false)
|
||||
await host.restartResume.continueAfterRestart([SESSION], 'retry')
|
||||
expect(dispatch).toHaveBeenCalledTimes(1)
|
||||
}
|
||||
)
|
||||
|
||||
it('releases a failed acquisition without retrying the spent offer', async () => {
|
||||
const { host, acquire, dispatch, closeSession } = await interruptedRestart()
|
||||
|
||||
+39
@@ -0,0 +1,39 @@
|
||||
import { expect, it } from 'vitest'
|
||||
import { AgentSessionRecoveryCapsule } from '../../runtime/agent-session-recovery-capsule'
|
||||
import { attach, hostTestState } from './structured-agent-session-host-test-harness'
|
||||
import { pendingApproval } from './structured-agent-session-restart-resume-test-harness'
|
||||
import {
|
||||
HOST_TEST_NOW as NOW,
|
||||
HOST_TEST_SESSION as SESSION,
|
||||
HOST_TEST_THREAD as THREAD
|
||||
} from './structured-agent-session-host-test-data'
|
||||
|
||||
it.each(['approval', 'completed'])(
|
||||
'does not offer work with an accepted %s event queued at quit',
|
||||
async (event) => {
|
||||
await attach()
|
||||
const { host, root, acquire } = hostTestState()
|
||||
const events = acquire.mock.calls[0]?.[0].events
|
||||
if (!events) {
|
||||
throw new Error('missing provider event sink')
|
||||
}
|
||||
events.appendItem(
|
||||
{ provider: 'codex', threadId: THREAD, turnId: 'working', ordinal: 1 },
|
||||
{ kind: 'turn', turnId: 'working', state: 'running' }
|
||||
)
|
||||
await host.flushStreamedEvents(SESSION)
|
||||
events.appendItem(
|
||||
{ provider: 'codex', threadId: THREAD, turnId: 'working', ordinal: 2 },
|
||||
{ kind: 'status', text: 'Provider is requesting approval' }
|
||||
)
|
||||
events.appendItem(
|
||||
{ provider: 'codex', threadId: THREAD, turnId: 'working', ordinal: 3 },
|
||||
event === 'approval'
|
||||
? pendingApproval().body
|
||||
: { kind: 'turn', turnId: 'working', state: 'completed' },
|
||||
{ lifecycle: true }
|
||||
)
|
||||
await host.flushAllStreamedEvents()
|
||||
expect(await new AgentSessionRecoveryCapsule(root).take(NOW)).toEqual([])
|
||||
}
|
||||
)
|
||||
+37
-31
@@ -171,38 +171,44 @@ describe('structured agent-session host teardown', () => {
|
||||
])
|
||||
})
|
||||
|
||||
it('bounds and diagnoses a stalled capsule write without preventing later cleanup', async () => {
|
||||
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
|
||||
const warning = vi.spyOn(console, 'warn').mockImplementation(() => {})
|
||||
const pending = Promise.withResolvers<void>()
|
||||
const cleaned = vi.fn(async () => {})
|
||||
const phases = structuredAgentSessionHostTeardownPhases({
|
||||
holds: { dispose: cleaned },
|
||||
runtimeState: { stopLeaseRenewal: () => {}, flushAllEventSinks: cleaned },
|
||||
handoffs: { stopTuiHistoryCatchup: () => {}, drain: cleaned },
|
||||
tasks: { drainAttaches: cleaned },
|
||||
evictOwnedSessions: cleaned,
|
||||
recordResumeMarkers: () => pending.promise
|
||||
})
|
||||
try {
|
||||
const teardown = (async () => {
|
||||
for (const phase of phases) {
|
||||
await phase.run()
|
||||
}
|
||||
})()
|
||||
await vi.advanceTimersByTimeAsync(2000)
|
||||
await teardown
|
||||
expect(cleaned).toHaveBeenCalledTimes(5)
|
||||
expect(warning).toHaveBeenCalledWith(
|
||||
'[structured-agent-session] recording recovery capsule failed',
|
||||
expect.objectContaining({ message: expect.stringContaining('2000ms') })
|
||||
)
|
||||
} finally {
|
||||
pending.resolve()
|
||||
warning.mockRestore()
|
||||
vi.useRealTimers()
|
||||
it.each(['capsule', 'events'])(
|
||||
'bounds stalled recovery %s without preventing later cleanup',
|
||||
async (stalled) => {
|
||||
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
|
||||
const warning = vi.spyOn(console, 'warn').mockImplementation(() => {})
|
||||
const pending = Promise.withResolvers<void>()
|
||||
const cleaned = vi.fn(async () => {})
|
||||
const flush = vi
|
||||
.fn(async () => cleaned())
|
||||
.mockImplementationOnce(() => (stalled === 'events' ? pending.promise : Promise.resolve()))
|
||||
const phases = structuredAgentSessionHostTeardownPhases({
|
||||
holds: { dispose: cleaned },
|
||||
runtimeState: { stopLeaseRenewal: () => {}, flushAllEventSinks: flush },
|
||||
handoffs: { stopTuiHistoryCatchup: () => {}, drain: cleaned },
|
||||
tasks: { drainAttaches: cleaned },
|
||||
evictOwnedSessions: cleaned,
|
||||
recordResumeMarkers: () => (stalled === 'capsule' ? pending.promise : Promise.resolve())
|
||||
})
|
||||
try {
|
||||
const teardown = (async () => {
|
||||
for (const phase of phases) {
|
||||
await phase.run()
|
||||
}
|
||||
})()
|
||||
await vi.advanceTimersByTimeAsync(2000)
|
||||
await teardown
|
||||
expect(cleaned).toHaveBeenCalledTimes(5)
|
||||
expect(warning).toHaveBeenCalledWith(
|
||||
'[structured-agent-session] recording recovery capsule failed',
|
||||
expect.objectContaining({ message: expect.stringContaining('2000ms') })
|
||||
)
|
||||
} finally {
|
||||
pending.resolve()
|
||||
warning.mockRestore()
|
||||
vi.useRealTimers()
|
||||
}
|
||||
}
|
||||
})
|
||||
)
|
||||
|
||||
it('gives up on a wedged handoff instead of holding the quit open', async () => {
|
||||
const request = requests.request('to-tui', 'now', { operationId: hostTestOperationId() })
|
||||
|
||||
Reference in New Issue
Block a user