mirror of
https://github.com/stablyai/orca.git
synced 2026-10-03 00:02:19 +00:00
fix(native-chat): persist late dispatch receipts before session close
This commit is contained in:
@@ -122,19 +122,7 @@ export function readStructuredAgentSessionOptions(
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
* Settle a send whose ack window expired but which the provider later proved it
|
||||
* had received.
|
||||
*
|
||||
* Without this the submission stays `unknown` for the life of the session: the
|
||||
* client renders an unconfirmed bubble whose Retry redispatches, so the user is
|
||||
* invited to deliver the same message to the agent a second time. Every send made
|
||||
* while a turn is already running takes this path, because the provider does not
|
||||
* echo the new message until the running turn ends.
|
||||
*
|
||||
* Serialized with the session's other journal writes, and a no-op once the row is
|
||||
* `accepted` or `rejected` — the reducer treats both as terminal.
|
||||
*/
|
||||
/** Settle provider-proven delivery independently of an in-flight client mutation. */
|
||||
export async function settleStructuredAgentSessionLateDispatch(
|
||||
context: StructuredAgentSessionMutationContext,
|
||||
input: {
|
||||
@@ -147,13 +135,12 @@ export async function settleStructuredAgentSessionLateDispatch(
|
||||
if (!session) {
|
||||
return
|
||||
}
|
||||
await context.serialize(input.sessionId, async () => {
|
||||
await session.journal.resolveDispatch({
|
||||
clientMessageId: input.clientMessageId,
|
||||
state: 'accepted',
|
||||
providerIdentity: input.providerIdentity,
|
||||
fence: session.fence
|
||||
})
|
||||
context.publish(input.sessionId, session.journal)
|
||||
// The journal queue drains before close; the host queue would defer this past teardown.
|
||||
await session.journal.resolveDispatch({
|
||||
clientMessageId: input.clientMessageId,
|
||||
state: 'accepted',
|
||||
providerIdentity: input.providerIdentity,
|
||||
fence: session.fence
|
||||
})
|
||||
context.publish(input.sessionId, session.journal)
|
||||
}
|
||||
|
||||
+76
-4
@@ -3,9 +3,11 @@ import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { afterEach, beforeEach, describe, expect, it, vi, type Mock } from 'vitest'
|
||||
import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope'
|
||||
import type { AgentSessionMutationEnvelope } from '../../../shared/agent-session-wire'
|
||||
import type {
|
||||
AgentSessionMutationEnvelope,
|
||||
AgentSessionSubscribeEvent
|
||||
} from '../../../shared/agent-session-wire'
|
||||
import { AgentSessionRecordStore } from '../../runtime/agent-session-record-store'
|
||||
import { createTrackedJournalOpener } from '../agent-session-journal/journal-store-test-open'
|
||||
import type {
|
||||
AgentSessionDispatchOutcome,
|
||||
StructuredAgentSessionAdapter
|
||||
@@ -21,13 +23,13 @@ import {
|
||||
resetHostTestOperationIds
|
||||
} from './structured-agent-session-host-test-data'
|
||||
|
||||
const journals = createTrackedJournalOpener()
|
||||
const CALLER = { callerKey: 'client-1' }
|
||||
|
||||
let root: string
|
||||
let store: AgentSessionRecordStore
|
||||
let host: StructuredAgentSessionHost
|
||||
let dispatch: Mock<StructuredAgentSessionAdapter['dispatch']>
|
||||
let closeSession: Mock<NonNullable<StructuredAgentSessionAdapter['closeSession']>>
|
||||
|
||||
function accepted(): AgentSessionDispatchOutcome {
|
||||
return {
|
||||
@@ -65,6 +67,7 @@ beforeEach(async () => {
|
||||
root = await mkdtemp(join(tmpdir(), 'orca-wire-late-settle-'))
|
||||
resetHostTestOperationIds()
|
||||
dispatch = vi.fn(async () => accepted())
|
||||
closeSession = vi.fn(async () => true)
|
||||
store = await AgentSessionRecordStore.open({ directory: join(root, 'store'), hostId: 'local' })
|
||||
host = new StructuredAgentSessionHost({
|
||||
store,
|
||||
@@ -86,6 +89,7 @@ beforeEach(async () => {
|
||||
})),
|
||||
releaseAcquisition: vi.fn(async () => true),
|
||||
dispatch,
|
||||
closeSession,
|
||||
cancelTurn: vi.fn(async () => ({ cancelled: true })),
|
||||
answerPrompt: vi.fn(async () => undefined),
|
||||
setOption: vi.fn(async () => undefined)
|
||||
@@ -99,12 +103,80 @@ beforeEach(async () => {
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
await journals.closeAll()
|
||||
await host.flushAllStreamedEvents()
|
||||
await host.close(SESSION)
|
||||
await rm(root, { recursive: true, force: true })
|
||||
})
|
||||
|
||||
describe('settling a send the provider proves it received after the ack window', () => {
|
||||
it('publishes acceptance during a pending send and never reopens it for retry', async () => {
|
||||
let finishDispatch!: (outcome: AgentSessionDispatchOutcome) => void
|
||||
dispatch.mockImplementationOnce(
|
||||
() =>
|
||||
new Promise((resolve) => {
|
||||
finishDispatch = resolve
|
||||
})
|
||||
)
|
||||
const events: AgentSessionSubscribeEvent[] = []
|
||||
const unsubscribe = host.subscribe({
|
||||
id: 'late-receipt',
|
||||
sessionId: SESSION,
|
||||
emit: (event) => events.push(event)
|
||||
})
|
||||
const params = sendParams('echo before send completes')
|
||||
const pending = host.send(CALLER, params)
|
||||
await vi.waitFor(() => expect(dispatch).toHaveBeenCalledTimes(1))
|
||||
try {
|
||||
await host.settleLateDispatch({
|
||||
sessionId: SESSION,
|
||||
clientMessageId: params.envelope.clientOperationId,
|
||||
providerIdentity: { provider: 'claude', sessionId: THREAD, uuid: 'early-echo' }
|
||||
})
|
||||
expect(events.at(-1)).toMatchObject({
|
||||
type: 'batch',
|
||||
batch: {
|
||||
submissions: [
|
||||
{ clientMessageId: params.envelope.clientOperationId, dispatchState: 'accepted' }
|
||||
]
|
||||
}
|
||||
})
|
||||
} finally {
|
||||
finishDispatch({ state: 'unknown', reason: 'ack timeout' })
|
||||
unsubscribe()
|
||||
}
|
||||
await expect(pending).resolves.toMatchObject({
|
||||
ok: true,
|
||||
value: { submission: { dispatchState: 'accepted' } }
|
||||
})
|
||||
await expect(host.send(CALLER, { ...params, retryUnknown: true })).resolves.toMatchObject({
|
||||
ok: true,
|
||||
value: { submission: { dispatchState: 'accepted' } }
|
||||
})
|
||||
expect(dispatch).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('persists an echo received while the provider is closing', async () => {
|
||||
dispatch.mockResolvedValueOnce({ state: 'unknown', reason: 'ack timeout' })
|
||||
const params = sendParams('received just before shutdown')
|
||||
await host.send(CALLER, params)
|
||||
let settlement: Promise<void> | undefined
|
||||
closeSession.mockImplementationOnce(async () => {
|
||||
settlement = host.settleLateDispatch({
|
||||
sessionId: SESSION,
|
||||
clientMessageId: params.envelope.clientOperationId,
|
||||
providerIdentity: { provider: 'claude', sessionId: THREAD, uuid: 'closing-echo' }
|
||||
})
|
||||
void settlement.catch(() => undefined)
|
||||
return true
|
||||
})
|
||||
|
||||
await host.close(SESSION)
|
||||
await expect(settlement).resolves.toBeUndefined()
|
||||
await host.revealSession(SESSION)
|
||||
expect(submissions()).toMatchObject([{ dispatchState: 'accepted' }])
|
||||
expect(dispatch).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('moves a durable unknown to accepted so nothing offers to send it again', async () => {
|
||||
dispatch.mockRejectedValueOnce(new Error('socket closed'))
|
||||
const params = sendParams('sent while a turn was running')
|
||||
|
||||
Reference in New Issue
Block a user