mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 08:03:20 +00:00
fix(native-chat): settle a structured send whose replay lands after the timeout
`dispatchClaudeTurn` waits a bounded window for Claude to echo the uuid it
sent. On timeout the send returns `unknown` and the host records that on
the submission. The late replay IS matched — `resolveClaudeReplayWaiter`
finds the retired waiter and `recoverLateIdentity` repairs turn identity —
but nothing ever writes the journal, and the only other writer moves
`pending -> unknown`. So `unknown` was terminal.
That is load-bearing for the client. The outbox only drops an entry on
`dispatchState === 'accepted'`, so the entry never left, and three
user-visible defects followed from one stuck row:
1. "Message delivery is unconfirmed." appears ~12s into a mid-turn send
and never clears, even though the message was delivered and answered.
2. The banner's Retry sets `retryUnknown`, which bypasses the host's
replay-an-existing-submission guard and delivers the message to the
agent a second time.
3. Retry targets `outbox.find(entry => entry.state === 'unconfirmed')` —
the OLDEST stuck entry, not the message the user just typed — so a
click re-sends a message from much earlier in the session.
A user replay is Claude echoing the message, which proves delivery. Carry
the `clientMessageId` on the dispatch waiter, and when a retired waiter is
matched by a user replay, settle its submission `accepted`.
The notification is deliberately not gated on `dispatchSequence`: turn
identity must stay with the newer dispatch, but the older message still
provably landed.
`settleLateDispatch` only moves a row that is still `unknown`, so an
accepted or rejected submission is never rewritten, and it never
re-dispatches.
No wire shape change: `accepted` is an existing `dispatchState` that every
client already handles. This only makes it get written when it is true.
This commit is contained in:
@@ -65,6 +65,49 @@ describe('Claude structured dispatch image limits', () => {
|
||||
expect(session.activeTurnSequence).toBe(session.dispatchSequence)
|
||||
})
|
||||
|
||||
it('reports the settled submission when a timed-out replay arrives late', async () => {
|
||||
const session = sessionFor()
|
||||
const accepted: { clientMessageId: string; uuid: string }[] = []
|
||||
session.onLateDispatchAccepted = (input) => accepted.push(input)
|
||||
const dispatched = dispatchClaudeTurn(
|
||||
session,
|
||||
{ clientMessageId: 'client-1', body: userMessage([{ type: 'text', text: 'one' }]) },
|
||||
500
|
||||
)
|
||||
await vi.waitFor(() => expect(session.dispatchWaiters).toHaveLength(1))
|
||||
const sentUuid = (session.dispatchWaiters[0] as { sentUuid?: string }).sentUuid
|
||||
// The send gave up waiting, so the host has already recorded this submission `unknown`.
|
||||
await expect(dispatched).resolves.toMatchObject({ state: 'unknown' })
|
||||
|
||||
resolveClaudeReplayWaiter(session, userReplayFrame(sentUuid!, 'one'))
|
||||
|
||||
expect(accepted).toEqual([{ clientMessageId: 'client-1', uuid: sentUuid }])
|
||||
})
|
||||
|
||||
it('reports a late replay even after a newer dispatch has started', async () => {
|
||||
const session = sessionFor()
|
||||
const accepted: { clientMessageId: string; uuid: string }[] = []
|
||||
session.onLateDispatchAccepted = (input) => accepted.push(input)
|
||||
const first = dispatchClaudeTurn(
|
||||
session,
|
||||
{ clientMessageId: 'client-1', body: userMessage([{ type: 'text', text: 'one' }]) },
|
||||
500
|
||||
)
|
||||
await vi.waitFor(() => expect(session.dispatchWaiters).toHaveLength(1))
|
||||
const firstUuid = (session.dispatchWaiters[0] as { sentUuid?: string }).sentUuid
|
||||
await expect(first).resolves.toMatchObject({ state: 'unknown' })
|
||||
dispatchClaudeTurn(
|
||||
session,
|
||||
{ clientMessageId: 'client-2', body: userMessage([{ type: 'text', text: 'two' }]) },
|
||||
100
|
||||
)
|
||||
await vi.waitFor(() => expect(session.dispatchWaiters).toHaveLength(1))
|
||||
|
||||
// Turn identity must stay with dispatch B, but message A still provably landed.
|
||||
expect(resolveClaudeReplayWaiter(session, userReplayFrame(firstUuid!, 'one'))).toBe(false)
|
||||
expect(accepted).toEqual([{ clientMessageId: 'client-1', uuid: firstUuid }])
|
||||
})
|
||||
|
||||
it('never lets a late replay for dispatch A resolve dispatch B', async () => {
|
||||
const session = sessionFor()
|
||||
const first = dispatchClaudeTurn(
|
||||
|
||||
@@ -147,6 +147,12 @@ function recoverLateIdentity(
|
||||
session.activeTurnId = uuid
|
||||
session.activeTurnSequence = waiter.dispatchSequence
|
||||
}
|
||||
// The waiter is retired, so its send already returned `unknown`. A user replay is Claude
|
||||
// echoing the message, which proves delivery — regardless of whether a newer dispatch has
|
||||
// since started, so this is deliberately not gated on `dispatchSequence`.
|
||||
if (isUserReplay) {
|
||||
session.onLateDispatchAccepted?.({ clientMessageId: waiter.clientMessageId, uuid })
|
||||
}
|
||||
return isUserReplay && waiter.dispatchSequence === session.dispatchSequence
|
||||
}
|
||||
|
||||
@@ -155,12 +161,14 @@ function waitForReplay(
|
||||
timeoutMs: number,
|
||||
acceptsResult: boolean,
|
||||
sentUuid: string,
|
||||
replayContentKey: string
|
||||
replayContentKey: string,
|
||||
clientMessageId: string
|
||||
): { waiter: ClaudeDispatchWaiter; promise: Promise<string | null> } {
|
||||
let waiter!: ClaudeDispatchWaiter
|
||||
const promise = new Promise<string | null>((resolve) => {
|
||||
waiter = {
|
||||
acceptsResult,
|
||||
clientMessageId,
|
||||
sentUuid,
|
||||
dispatchSequence: session.dispatchSequence,
|
||||
replayContentKey,
|
||||
@@ -220,7 +228,8 @@ export async function dispatchClaudeTurn(
|
||||
timeoutMs,
|
||||
acceptsResult,
|
||||
sentUuid,
|
||||
claudeDispatchContentKey(content)
|
||||
claudeDispatchContentKey(content),
|
||||
input.clientMessageId
|
||||
)
|
||||
const replayed = replay.promise
|
||||
try {
|
||||
|
||||
@@ -265,6 +265,12 @@ export async function acquireClaudeSession({
|
||||
})
|
||||
const acquired: AgentSessionAcquisition = publication.acquisition
|
||||
liveSession = publication.session
|
||||
if (deps.onLateDispatchAccepted) {
|
||||
const notify = deps.onLateDispatchAccepted
|
||||
const bound = publication.session
|
||||
bound.onLateDispatchAccepted = ({ clientMessageId, uuid }) =>
|
||||
notify({ sessionId, clientMessageId, uuid, providerSessionId: bound.providerSessionId })
|
||||
}
|
||||
await restoreClaudeStructuredSessionOptions(liveSession, deps.requestTimeoutMs)
|
||||
acquisitions.assertCurrent(sessionId, attempt)
|
||||
acquisitions.deleteIfCurrent(sessionId, attempt)
|
||||
|
||||
@@ -60,6 +60,15 @@ export type ClaudeStructuredSessionAdapterDeps = {
|
||||
sessionId: string,
|
||||
state: AgentSessionBackgroundTaskState | null
|
||||
) => void
|
||||
/** A replay that lands after its dispatch timed out proves the message reached Claude, so the
|
||||
* submission the host already recorded `unknown` can still be settled `accepted`. */
|
||||
onLateDispatchAccepted?: (input: {
|
||||
sessionId: string
|
||||
clientMessageId: string
|
||||
uuid: string
|
||||
/** Claude's own session id — the journal identity is provider-scoped, not Orca-scoped. */
|
||||
providerSessionId: string
|
||||
}) => void
|
||||
openConnection?: typeof openClaudeStreamJsonConnection
|
||||
readProcessStartTime?: (pid: number) => Promise<number | null>
|
||||
mintLinkId?: () => string
|
||||
@@ -87,6 +96,8 @@ export type ClaudeDispatchWaiter = {
|
||||
resolve: (uuid: string | null) => void
|
||||
timer: ReturnType<typeof setTimeout>
|
||||
acceptsResult: boolean
|
||||
/** The submission this dispatch belongs to, so a late replay can settle its journal row. */
|
||||
clientMessageId: string
|
||||
/** Client uuid echoed by Claude so a replay is tied to its own dispatch. */
|
||||
sentUuid: string
|
||||
/** Sequence used to fence a late identity from a newer dispatch. */
|
||||
@@ -100,6 +111,8 @@ export type ClaudeDispatchWaiter = {
|
||||
}
|
||||
|
||||
export type ClaudeSession = {
|
||||
/** Set by the adapter; see `onLateDispatchAccepted` on the adapter deps. */
|
||||
onLateDispatchAccepted?: (input: { clientMessageId: string; uuid: string }) => void
|
||||
connection: ClaudeStreamJsonConnection
|
||||
providerSessionId: string
|
||||
/** Durable transcript files live under this account's `projects` directory. */
|
||||
|
||||
@@ -338,6 +338,54 @@ describe('send', () => {
|
||||
expect(result).toMatchObject({ ok: true, value: { submission: { dispatchState: 'unknown' } } })
|
||||
})
|
||||
|
||||
it('settles an unknown submission accepted when its replay lands late', async () => {
|
||||
await attach()
|
||||
dispatch.mockResolvedValueOnce({ state: 'unknown', reason: 'no replay in time' })
|
||||
const body = hostTestMessage('add a retry')
|
||||
const result = await host.send(CALLER, {
|
||||
envelope: envelope('agentSession.send', { body }),
|
||||
body
|
||||
})
|
||||
if (!result.ok) {
|
||||
throw new Error(`expected a send, got ${result.refusal.code}`)
|
||||
}
|
||||
expect(result.value.submission.dispatchState).toBe('unknown')
|
||||
|
||||
await host.settleLateDispatch({
|
||||
sessionId: SESSION,
|
||||
clientMessageId: result.value.clientMessageId,
|
||||
providerIdentity: { provider: 'codex', threadId: THREAD, turnId: 'turn-late', ordinal: 0 }
|
||||
})
|
||||
|
||||
const page = host.history({ sessionId: SESSION, direction: 'tail' })
|
||||
const settled = (page.ok ? page.page.submissions : []).find(
|
||||
(entry) => entry.clientMessageId === result.value.clientMessageId
|
||||
)
|
||||
expect(settled?.dispatchState).toBe('accepted')
|
||||
// The message was already delivered; a late settlement must never re-dispatch it.
|
||||
expect(dispatch).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('leaves an already accepted submission alone when a late replay arrives', async () => {
|
||||
await attach()
|
||||
const body = hostTestMessage('add a retry')
|
||||
const result = await host.send(CALLER, {
|
||||
envelope: envelope('agentSession.send', { body }),
|
||||
body
|
||||
})
|
||||
if (!result.ok) {
|
||||
throw new Error(`expected a send, got ${result.refusal.code}`)
|
||||
}
|
||||
const before = host.history({ sessionId: SESSION, direction: 'tail' })
|
||||
await host.settleLateDispatch({
|
||||
sessionId: SESSION,
|
||||
clientMessageId: result.value.clientMessageId,
|
||||
providerIdentity: { provider: 'codex', threadId: THREAD, turnId: 'turn-late', ordinal: 0 }
|
||||
})
|
||||
const after = host.history({ sessionId: SESSION, direction: 'tail' })
|
||||
expect(after.ok && after.page.fence).toBe(before.ok && before.page.fence)
|
||||
})
|
||||
|
||||
it('replays a retried send from the journal without dispatching twice', async () => {
|
||||
await attach()
|
||||
const body = hostTestMessage('add a retry')
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
// Structured agent-session host: where the lease, journal, and provider adapter meet.
|
||||
// Mutations share one durable admission path and serialize per session.
|
||||
|
||||
import type { AgentJournalItemIdentity } from '../../../shared/agent-session-journal-types'
|
||||
import type { AgentSessionExecutionLocation } from '../../../shared/agent-session-record'
|
||||
import type * as SessionWire from '../../../shared/agent-session-wire'
|
||||
import type { AgentSessionAttachParams } from './structured-agent-session-attach'
|
||||
@@ -337,6 +338,37 @@ export class StructuredAgentSessionHost {
|
||||
) => this.backgroundTasks.publish(sessionId, state)
|
||||
unsubscribe = (sessionId: string, id: string): void => this.subscribers.close(sessionId, id)
|
||||
|
||||
/** A provider replay that arrived after its dispatch timed out. The send already recorded
|
||||
* `unknown`, which is terminal for the outbox — the client shows "delivery is unconfirmed"
|
||||
* forever and its Retry redelivers a message the agent already has. The replay proves
|
||||
* delivery, so settle the submission the host was unsure about. */
|
||||
settleLateDispatch = (input: {
|
||||
sessionId: string
|
||||
clientMessageId: string
|
||||
providerIdentity: AgentJournalItemIdentity
|
||||
}): Promise<void> =>
|
||||
this.serialize(input.sessionId, async () => {
|
||||
const session = this.sessions.get(input.sessionId)
|
||||
if (!session) {
|
||||
return
|
||||
}
|
||||
// Only an `unknown` row is ours to move; anything else is already the truth.
|
||||
const submission = session.journal
|
||||
.submissions()
|
||||
.find((entry) => entry.clientMessageId === input.clientMessageId)
|
||||
if (submission?.dispatchState !== 'unknown') {
|
||||
return
|
||||
}
|
||||
await session.journal.resolveDispatch({
|
||||
clientMessageId: input.clientMessageId,
|
||||
state: 'accepted',
|
||||
providerIdentity: input.providerIdentity,
|
||||
fence: session.fence,
|
||||
recovered: true
|
||||
})
|
||||
this.subscribers.publish(input.sessionId, session.journal)
|
||||
})
|
||||
|
||||
/** Every session's projected status for session lists; unlike `subscribe`, retains nothing. */
|
||||
subscribeStatus = (
|
||||
subscriber: Parameters<StructuredAgentSessionStatusFeed['subscribe']>[0]
|
||||
|
||||
@@ -260,6 +260,17 @@ async function install(deps: StructuredAgentSessionRuntimeDeps): Promise<Install
|
||||
},
|
||||
onBackgroundTasksChanged: (sessionId, state) =>
|
||||
host?.publishBackgroundTaskState(sessionId, state),
|
||||
onLateDispatchAccepted: ({ sessionId, clientMessageId, uuid, providerSessionId }) => {
|
||||
void host
|
||||
?.settleLateDispatch({
|
||||
sessionId,
|
||||
clientMessageId,
|
||||
providerIdentity: { provider: 'claude', sessionId: providerSessionId, uuid }
|
||||
})
|
||||
.catch((error) =>
|
||||
deps.onError?.({ scope: `structured-agent-session-late-dispatch:${sessionId}`, error })
|
||||
)
|
||||
},
|
||||
...(deps.openClaudeConnection ? { openClaudeConnection: deps.openClaudeConnection } : {}),
|
||||
...(deps.readProcessStartTime ? { readProcessStartTime: deps.readProcessStartTime } : {})
|
||||
})
|
||||
|
||||
@@ -34,6 +34,7 @@ export type StructuredClaudeRuntimeAdapterDeps = {
|
||||
sessionId: string,
|
||||
state: AgentSessionBackgroundTaskState | null
|
||||
) => void
|
||||
onLateDispatchAccepted?: ClaudeStructuredSessionAdapterDeps['onLateDispatchAccepted']
|
||||
}
|
||||
|
||||
export function createStructuredClaudeRuntimeAdapter(
|
||||
@@ -100,6 +101,7 @@ export function createStructuredClaudeRuntimeAdapter(
|
||||
...(deps.onBackgroundTasksChanged
|
||||
? { onBackgroundTasksChanged: deps.onBackgroundTasksChanged }
|
||||
: {}),
|
||||
...(deps.onLateDispatchAccepted ? { onLateDispatchAccepted: deps.onLateDispatchAccepted } : {}),
|
||||
...(deps.openClaudeConnection ? { openConnection: deps.openClaudeConnection } : {}),
|
||||
...(deps.readProcessStartTime ? { readProcessStartTime: deps.readProcessStartTime } : {})
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user