mirror of
https://github.com/stablyai/orca.git
synced 2026-09-28 08:02:43 +00:00
Fix structured chat pending-work lifecycle and mobile cancellation
This commit is contained in:
@@ -71,6 +71,7 @@ export function MobileNativeChatOverlay({
|
||||
error={session.error}
|
||||
agent={controller.nativeChatAgent}
|
||||
agentWorking={controller.nativeChatAgentWorking}
|
||||
canStop={controller.nativeChatCanStop}
|
||||
structuredActivityUi={controller.nativeChatStructured}
|
||||
streaming={streaming}
|
||||
onStop={controller.handleNativeChatStop}
|
||||
|
||||
@@ -74,6 +74,7 @@ type Overrides = {
|
||||
pending?: Parameters<typeof MobileNativeChatView>[0]['pending']
|
||||
structuredActivityUi?: boolean
|
||||
agentWorking?: boolean
|
||||
canStop?: boolean
|
||||
sendSurfaceId?: string
|
||||
}
|
||||
|
||||
@@ -118,6 +119,18 @@ describe('MobileNativeChatView', () => {
|
||||
}
|
||||
|
||||
/** Ids of the rows the list is currently rendering. */
|
||||
it('keeps Stop hidden during a structured dispatch until a provider turn can be cancelled', async () => {
|
||||
const props = { structuredActivityUi: true, agentWorking: true, canStop: false }
|
||||
await render(props)
|
||||
const stops = () =>
|
||||
renderer!.root.findAll((node) => node.props.accessibilityLabel === 'Stop the agent')
|
||||
expect(stops()).toHaveLength(0)
|
||||
await update({ ...props, canStop: true })
|
||||
expect(stops()).toHaveLength(1)
|
||||
await update({ agentWorking: true })
|
||||
expect(stops()).toHaveLength(1)
|
||||
})
|
||||
|
||||
function listIds(): string[] {
|
||||
const list = renderer!.root.find((node) => node.type === 'FlatList')
|
||||
return (list.props.data as { id: string }[]).map((row) => row.id)
|
||||
|
||||
@@ -49,10 +49,11 @@ type Props = {
|
||||
/** Resolved agent for this chat; names the empty-state copy (desktop parity). */
|
||||
agent?: string | null
|
||||
agentWorking?: boolean
|
||||
canStop?: boolean
|
||||
/** Structured lane: per-turn "Working for N" status plus live tool progress,
|
||||
* replacing the bridge lane's static three-dot working row (desktop parity). */
|
||||
structuredActivityUi?: boolean
|
||||
/** Interrupt the agent mid-turn (shown as a Stop button on the working bar). */
|
||||
/** Interrupt a provider turn. */
|
||||
onStop?: () => void
|
||||
/** Live partial assistant text to show as an in-progress bubble, already gated
|
||||
* by the overlay against the transcript catching up. */
|
||||
@@ -129,6 +130,7 @@ export function MobileNativeChatView({
|
||||
error,
|
||||
agent,
|
||||
agentWorking,
|
||||
canStop = agentWorking,
|
||||
structuredActivityUi = false,
|
||||
onStop,
|
||||
streaming,
|
||||
@@ -396,8 +398,6 @@ export function MobileNativeChatView({
|
||||
question={question}
|
||||
onAnswerQuestion={onAnswerQuestion}
|
||||
/>
|
||||
{/* Chrome row above the composer: the working indicator and the global
|
||||
tool-calls expand/collapse toggle on the left, Stop in the far corner. */}
|
||||
<View style={styles.chromeRow}>
|
||||
<View style={styles.chromeLeft}>
|
||||
{agentWorking && !structuredActivityUi ? <MobileAgentWorkingIndicator /> : null}
|
||||
@@ -414,7 +414,7 @@ export function MobileNativeChatView({
|
||||
<Text style={styles.chromeToggleLabel}>{toolsExpanded ? 'Collapse' : 'Tools'}</Text>
|
||||
</Pressable>
|
||||
</View>
|
||||
{agentWorking ? (
|
||||
{canStop ? (
|
||||
<Pressable
|
||||
style={({ pressed }) => [styles.stopButton, pressed && styles.pressed]}
|
||||
onPress={onStop}
|
||||
|
||||
@@ -28,6 +28,7 @@ export type MobileNativeChatController = {
|
||||
/** Structured lane: drives the per-turn status row and live tool progress. */
|
||||
nativeChatStructured: boolean
|
||||
nativeChatAgentWorking: boolean
|
||||
nativeChatCanStop: boolean
|
||||
nativeChatStreamingText?: string
|
||||
/** Agent mid-turn, regardless of whether chat is the visible view. */
|
||||
nativeChatStreamLive: boolean
|
||||
|
||||
@@ -55,6 +55,7 @@ const structuredQuestion = {
|
||||
allowOther: true,
|
||||
optionTokens: ['choice-a', 'choice-b']
|
||||
}
|
||||
const structuredActivity = { isWorking: false, turnId: null as string | null }
|
||||
const structuredSessionState = {
|
||||
messages: [] as unknown[],
|
||||
status: 'ready',
|
||||
@@ -86,8 +87,7 @@ vi.mock('./use-mobile-native-chat-session', () => ({
|
||||
vi.mock('./use-mobile-structured-agent-session', () => ({
|
||||
useMobileStructuredAgentSession: () => ({
|
||||
session: structuredSessionState,
|
||||
isWorking: false,
|
||||
turnId: null,
|
||||
...structuredActivity,
|
||||
sendWithOutcome: structuredSendWithOutcome,
|
||||
cancel: structuredCancel,
|
||||
permission: structuredPermission,
|
||||
@@ -341,6 +341,37 @@ describe('useMobileNativeChatController handleNativeChatSend', () => {
|
||||
expect(clientStub.sendRequest).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('separates structured working status from provider cancellation availability', async () => {
|
||||
const props = {
|
||||
tab: {
|
||||
type: 'agent-session',
|
||||
id: 'agent-tab-1',
|
||||
title: 'Chat',
|
||||
sessionId: 'session-structured',
|
||||
agent: 'codex',
|
||||
isActive: true
|
||||
},
|
||||
activeHandle: null,
|
||||
inputLeaseReady: false
|
||||
}
|
||||
structuredActivity.isWorking = true
|
||||
try {
|
||||
await act(async () => {
|
||||
renderer?.update(createElement(Harness, props))
|
||||
})
|
||||
expect(controller?.nativeChatAgentWorking).toBe(true)
|
||||
expect(controller?.nativeChatCanStop).toBe(false)
|
||||
structuredActivity.turnId = 'provider-turn'
|
||||
await act(async () => {
|
||||
renderer?.update(createElement(Harness, props))
|
||||
})
|
||||
expect(controller?.nativeChatCanStop).toBe(true)
|
||||
} finally {
|
||||
structuredActivity.isWorking = false
|
||||
structuredActivity.turnId = null
|
||||
}
|
||||
})
|
||||
|
||||
it('exposes structured prompt cards and session options on structured tabs', async () => {
|
||||
await act(async () => {
|
||||
renderer?.update(
|
||||
|
||||
@@ -299,6 +299,9 @@ export function useMobileNativeChatController(args: {
|
||||
/** Structured lane: drives the per-turn status row and live tool progress. */
|
||||
nativeChatStructured: activeChatStructured,
|
||||
nativeChatAgentWorking,
|
||||
nativeChatCanStop: activeChatStructured
|
||||
? structuredNativeChat.turnId !== null
|
||||
: nativeChatAgentWorking,
|
||||
nativeChatStreamingText,
|
||||
nativeChatStreamLive,
|
||||
nativeChatStreamScopeKey: streamScopeKey,
|
||||
|
||||
@@ -297,7 +297,7 @@ export function useMobileStructuredAgentSession(args: {
|
||||
// A dispatch the provider has not answered yet is already work — see the desktop hook.
|
||||
isWorking:
|
||||
activeStructuredAgentSessionTurnId(state.items) !== null ||
|
||||
hasUnansweredStructuredAgentSessionDispatch(state.submissions),
|
||||
hasUnansweredStructuredAgentSessionDispatch(state.submissions, state.fence),
|
||||
turnId: activeStructuredAgentSessionTurnId(state.items),
|
||||
sendWithOutcome,
|
||||
cancel,
|
||||
|
||||
@@ -15,6 +15,7 @@ import type {
|
||||
AgentJournalMessageItem,
|
||||
AgentSessionJournalIdentity
|
||||
} from '../../../shared/agent-session-journal-types'
|
||||
import { hasUnansweredStructuredAgentSessionDispatch } from '../../../shared/structured-agent-session-projection'
|
||||
import { digestPayload } from './journal-payload-bounds'
|
||||
import {
|
||||
reconcileSubmissions,
|
||||
@@ -141,6 +142,30 @@ describe('crash between provider accept and journal commit', () => {
|
||||
expect(restarted.receiptFor('cm_1')?.providerItemId).toBe(agentJournalItemKey(outcome.identity))
|
||||
})
|
||||
|
||||
it('retires an ack timeout on restart without changing its delivery verdict', async () => {
|
||||
const journal = await open()
|
||||
await journal.appendSubmission({
|
||||
clientMessageId: 'cm_timeout',
|
||||
payloadFingerprint: digestPayload('slow'),
|
||||
body: userMessage('slow'),
|
||||
fence: 1
|
||||
})
|
||||
await journal.resolveDispatch({
|
||||
clientMessageId: 'cm_timeout',
|
||||
state: 'unknown',
|
||||
reason: 'ack timeout',
|
||||
fence: 1
|
||||
})
|
||||
expect(hasUnansweredStructuredAgentSessionDispatch(journal.submissions())).toBe(true)
|
||||
const restarted = await open()
|
||||
await restarted.markPendingSubmissionsUnknown(2)
|
||||
expect(restarted.submissions()[0]?.dispatchState).toBe('unknown')
|
||||
expect(hasUnansweredStructuredAgentSessionDispatch(restarted.submissions())).toBe(false)
|
||||
const cursor = restarted.cursor()
|
||||
await restarted.markPendingSubmissionsUnknown(2)
|
||||
expect(restarted.cursor()).toEqual(cursor)
|
||||
})
|
||||
|
||||
it('reports a rejected submission as never delivered, and never re-sends it', async () => {
|
||||
const journal = await open()
|
||||
await journal.appendSubmission({
|
||||
|
||||
@@ -2,14 +2,22 @@ import type { AgentSessionJournal } from './journal-store'
|
||||
|
||||
export async function markJournalPendingSubmissionsUnknown(
|
||||
journal: AgentSessionJournal,
|
||||
fence: number
|
||||
fence: number,
|
||||
reason = 'host_restarted_before_acknowledgement'
|
||||
): Promise<string[]> {
|
||||
const pending = journal.pendingSubmissions().map((entry) => entry.clientMessageId)
|
||||
const pending = journal
|
||||
.submissions()
|
||||
.filter(
|
||||
(entry) =>
|
||||
entry.dispatchState === 'pending' ||
|
||||
(entry.dispatchState === 'unknown' && entry.recovered !== true)
|
||||
)
|
||||
.map((entry) => entry.clientMessageId)
|
||||
for (const clientMessageId of pending) {
|
||||
await journal.resolveDispatch({
|
||||
clientMessageId,
|
||||
state: 'unknown',
|
||||
reason: 'host_restarted_before_acknowledgement',
|
||||
reason,
|
||||
fence,
|
||||
recovered: true
|
||||
})
|
||||
|
||||
@@ -256,12 +256,15 @@ function applyDispatch(
|
||||
if (submission.dispatchState === 'rejected' || submission.dispatchState === 'accepted') {
|
||||
return
|
||||
}
|
||||
submission.fence = row.fence
|
||||
submission.dispatchState = row.state
|
||||
submission.providerItemId = row.providerItemId
|
||||
submission.reason = row.reason
|
||||
submission.resolvedAt = row.ts
|
||||
if (row.recovered) {
|
||||
submission.recovered = row.recovered
|
||||
} else {
|
||||
delete submission.recovered
|
||||
}
|
||||
if (row.state !== 'accepted' || !row.providerItemId) {
|
||||
return
|
||||
|
||||
@@ -245,10 +245,9 @@ export class AgentSessionJournal {
|
||||
}))
|
||||
}
|
||||
|
||||
/** On restart every `pending` submission becomes `unknown` before the session
|
||||
* accepts a writer. Orca never re-sends on the user's behalf. */
|
||||
async markPendingSubmissionsUnknown(fence: number): Promise<string[]> {
|
||||
return markJournalPendingSubmissionsUnknown(this, fence)
|
||||
/** Retire unanswered sends after their execution owner ended, without assuming delivery. */
|
||||
async markPendingSubmissionsUnknown(fence: number, reason?: string): Promise<string[]> {
|
||||
return markJournalPendingSubmissionsUnknown(this, fence, reason)
|
||||
}
|
||||
|
||||
/** The escape hatch for corruption, an unreconcilable prefix, a forked handle,
|
||||
|
||||
@@ -86,6 +86,13 @@ export function createStructuredAgentSessionHostHandoff(
|
||||
host.publishStatus?.(sessionId)
|
||||
try {
|
||||
await host.flush(sessionId)
|
||||
const session = host.session(sessionId)
|
||||
await session.journal.markPendingSubmissionsUnknown(
|
||||
session.fence,
|
||||
'provider_exited_before_acknowledgement'
|
||||
)
|
||||
host.subscribers.publish(sessionId, session.journal)
|
||||
host.publishStatus?.(sessionId)
|
||||
host.eventSink(sessionId).unbind()
|
||||
return { state: 'stopped' }
|
||||
} catch (error) {
|
||||
|
||||
+34
@@ -3,6 +3,7 @@ import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { AgentJournalMessageItem } from '../../../shared/agent-session-journal-types'
|
||||
import { hasUnansweredStructuredAgentSessionDispatch } from '../../../shared/structured-agent-session-projection'
|
||||
import { structuredAgentSessionPayloadFingerprint } from '../../../shared/structured-agent-session-mutation'
|
||||
import { createTrackedJournalOpener } from '../agent-session-journal/journal-store-test-open'
|
||||
import type { AgentSessionJournal } from '../agent-session-journal/journal-store'
|
||||
@@ -34,6 +35,39 @@ afterEach(async () => {
|
||||
})
|
||||
|
||||
describe('structured send idempotency', () => {
|
||||
it('publishes a recovered retry as working before waiting for its provider', async () => {
|
||||
const body: AgentJournalMessageItem = {
|
||||
kind: 'message',
|
||||
role: 'user',
|
||||
blocks: [{ type: 'text', text: 'retry' }]
|
||||
}
|
||||
const input = { clientMessageId: 'retry-id', payloadFingerprint: 'fingerprint', body }
|
||||
await journal.appendSubmission({ ...input, fence: 1 })
|
||||
await journal.markPendingSubmissionsUnknown(2)
|
||||
const originalItem = journal.snapshot().items[0]
|
||||
const publish = vi.fn()
|
||||
const dispatch = vi.fn(async () => {
|
||||
expect(publish).toHaveBeenCalledOnce()
|
||||
expect(hasUnansweredStructuredAgentSessionDispatch(journal.submissions(), 2)).toBe(true)
|
||||
return { state: 'unknown' as const, reason: 'ack timeout' }
|
||||
})
|
||||
await performSend(
|
||||
{
|
||||
sessionId: 'session-1',
|
||||
journal,
|
||||
fence: 2,
|
||||
adapter: { dispatch } as unknown as StructuredAgentSessionAdapter,
|
||||
persistOptions: async () => undefined,
|
||||
resolvedBy: 'caller',
|
||||
publish,
|
||||
now: () => 1
|
||||
},
|
||||
{ ...input, retryUnknown: true }
|
||||
)
|
||||
expect(hasUnansweredStructuredAgentSessionDispatch(journal.submissions(), 2)).toBe(true)
|
||||
expect(journal.snapshot().items).toEqual([originalItem])
|
||||
})
|
||||
|
||||
it('does not redispatch one send id reused across caller ledgers', async () => {
|
||||
const body: AgentJournalMessageItem = {
|
||||
kind: 'message',
|
||||
|
||||
+26
-1
@@ -56,9 +56,11 @@ async function openJournal(sessionId = SESSION, now?: () => number) {
|
||||
function indexed(session: {
|
||||
journal: Awaited<ReturnType<typeof openJournal>>
|
||||
hasProviderChild?: boolean
|
||||
fence?: number
|
||||
}) {
|
||||
return {
|
||||
journal: session.journal,
|
||||
fence: session.fence ?? 1,
|
||||
...(session.hasProviderChild !== undefined
|
||||
? { hasProviderChild: session.hasProviderChild }
|
||||
: {}),
|
||||
@@ -69,7 +71,7 @@ function indexed(session: {
|
||||
function feedFor(
|
||||
sessions: Map<
|
||||
string,
|
||||
{ journal: Awaited<ReturnType<typeof openJournal>>; hasProviderChild?: boolean }
|
||||
{ journal: Awaited<ReturnType<typeof openJournal>>; hasProviderChild?: boolean; fence?: number }
|
||||
>,
|
||||
record: Partial<AgentSessionRecord> | null = null,
|
||||
onStatusChanged?: StructuredAgentSessionStatusFeedDeps['onStatusChanged']
|
||||
@@ -153,6 +155,29 @@ describe('StructuredAgentSessionStatusFeed', () => {
|
||||
])
|
||||
})
|
||||
|
||||
it('stops projecting an old-host unknown submission after the owner fence advances', async () => {
|
||||
const journal = await openJournal()
|
||||
const session = { journal, fence: 1 }
|
||||
const { feed, events } = feedFor(new Map([[SESSION, session]]))
|
||||
await journal.appendSubmission({
|
||||
clientMessageId: 'old-host',
|
||||
payloadFingerprint: 'fp',
|
||||
body: { kind: 'message', role: 'user', blocks: [{ type: 'text', text: 'slow' }] },
|
||||
fence: 1
|
||||
})
|
||||
await journal.resolveDispatch({
|
||||
clientMessageId: 'old-host',
|
||||
state: 'unknown',
|
||||
reason: 'ack timeout',
|
||||
fence: 1
|
||||
})
|
||||
feed.publish(SESSION)
|
||||
expect(events.at(-1)).toMatchObject({ session: { status: 'working' } })
|
||||
session.fence = 2
|
||||
feed.publish(SESSION)
|
||||
expect(events.at(-1)).toMatchObject({ session: { status: 'idle' } })
|
||||
})
|
||||
|
||||
it('publishes working from the pending submission, before the provider replays the turn', async () => {
|
||||
const journal = await openJournal()
|
||||
const { feed, events } = feedFor(new Map([[SESSION, { journal }]]))
|
||||
|
||||
@@ -31,6 +31,7 @@ type StatusFeedSession = {
|
||||
journal: AgentSessionJournal
|
||||
params: { location: { workspaceId: string }; provider: AgentSessionRecord['provider'] }
|
||||
hasProviderChild?: boolean
|
||||
fence?: number
|
||||
}
|
||||
|
||||
export type StructuredAgentSessionStatusFeedDeps = {
|
||||
@@ -166,7 +167,7 @@ export class StructuredAgentSessionStatusFeed {
|
||||
workspaceId: session.params.location.workspaceId,
|
||||
agent: session.params.provider,
|
||||
...(session.hasProviderChild ? { hostExecutionOwned: true as const } : {}),
|
||||
...projectStructuredAgentSessionStatusSummary(items, submissions),
|
||||
...projectStructuredAgentSessionStatusSummary(items, submissions, session.fence),
|
||||
...(record?.rewind?.phase === 'prepared' || record?.rewind?.phase === 'provider-succeeded'
|
||||
? { rewindBlockedReason: 'outcome-unknown' as const }
|
||||
: {}),
|
||||
|
||||
+6
@@ -8,6 +8,7 @@ import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { afterEach, beforeEach, describe, expect, it, vi, type Mock } from 'vitest'
|
||||
import type { AgentSessionOwnerProbe } from '../../../shared/agent-session-lease-adjudication'
|
||||
import { hasUnansweredStructuredAgentSessionDispatch } from '../../../shared/structured-agent-session-projection'
|
||||
import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope'
|
||||
import type {
|
||||
AgentSessionMutationEnvelope,
|
||||
@@ -334,6 +335,11 @@ describe('an unexpected provider exit', () => {
|
||||
acquisitionGeneration: 'generation-1'
|
||||
})
|
||||
|
||||
const recoveredHistory = host.history({ sessionId: SESSION, direction: 'tail' })
|
||||
expect(
|
||||
recoveredHistory.ok &&
|
||||
hasUnansweredStructuredAgentSessionDispatch(recoveredHistory.page.submissions)
|
||||
).toBe(false)
|
||||
expect(acquire).toHaveBeenCalledTimes(2)
|
||||
expect(dispatch).toHaveBeenCalledOnce()
|
||||
expect(store.getRecord(SESSION)?.lease).toMatchObject({
|
||||
|
||||
@@ -97,6 +97,15 @@ export async function performSend(
|
||||
if (!(input.retryUnknown && existing?.dispatchState === 'unknown')) {
|
||||
await ctx.journal.appendSubmission({ ...input, fence: ctx.fence })
|
||||
ctx.publish()
|
||||
} else {
|
||||
// Retry resumes work without moving or duplicating the original message.
|
||||
await ctx.journal.resolveDispatch({
|
||||
clientMessageId: input.clientMessageId,
|
||||
state: 'unknown',
|
||||
reason: 'dispatch_retry_in_progress',
|
||||
fence: ctx.fence
|
||||
})
|
||||
ctx.publish()
|
||||
}
|
||||
|
||||
const outcome = await dispatchSafely(ctx, input.clientMessageId, input.body)
|
||||
@@ -124,8 +133,7 @@ export async function performSend(
|
||||
clientMessageId: input.clientMessageId,
|
||||
state: 'unknown',
|
||||
reason: 'dispatch_result_persistence_failed',
|
||||
fence: ctx.fence,
|
||||
recovered: true
|
||||
fence: ctx.fence
|
||||
})
|
||||
} catch {
|
||||
// Nothing further to record; the pending row is settled on the next attach.
|
||||
|
||||
+10
-1
@@ -49,7 +49,11 @@ describe('provider-exit recovery tickets', () => {
|
||||
hasProviderChild: true,
|
||||
fence: 7,
|
||||
acquisitionGeneration: GENERATION,
|
||||
journal: { snapshot: () => ({ items: [] }), appendLifecycleBatch }
|
||||
journal: {
|
||||
snapshot: () => ({ items: [] }),
|
||||
appendLifecycleBatch,
|
||||
markPendingSubmissionsUnknown: vi.fn(async () => [])
|
||||
}
|
||||
} as unknown as StructuredAgentSessionHostSession
|
||||
const store = {
|
||||
getRecord: () => ({
|
||||
@@ -89,6 +93,10 @@ describe('provider-exit recovery tickets', () => {
|
||||
|
||||
expect(result).toMatchObject({ settlementRetryRequired: false, releasedFence: 8 })
|
||||
expect(appendLifecycleBatch).toHaveBeenCalledOnce()
|
||||
expect(session.journal.markPendingSubmissionsUnknown).toHaveBeenCalledWith(
|
||||
7,
|
||||
'provider_exited_before_acknowledgement'
|
||||
)
|
||||
expect(session.hasProviderChild).toBe(false)
|
||||
})
|
||||
|
||||
@@ -98,6 +106,7 @@ describe('provider-exit recovery tickets', () => {
|
||||
fence: 7,
|
||||
acquisitionGeneration: GENERATION,
|
||||
journal: {
|
||||
markPendingSubmissionsUnknown: vi.fn(async () => []),
|
||||
snapshot: () => ({ items: [] }),
|
||||
appendLifecycleBatch: vi.fn(async () => {
|
||||
throw new Error('journal still unavailable')
|
||||
|
||||
@@ -81,6 +81,15 @@ export async function settleUnexpectedStructuredAgentSessionExit(
|
||||
settlementRetryRequired = true
|
||||
context.onBarrierError?.(unexpectedEvent.sessionId, error)
|
||||
}
|
||||
try {
|
||||
await session.journal.markPendingSubmissionsUnknown(
|
||||
session.fence,
|
||||
'provider_exited_before_acknowledgement'
|
||||
)
|
||||
} catch (error) {
|
||||
settlementRetryRequired = true
|
||||
context.onBarrierError?.(unexpectedEvent.sessionId, error)
|
||||
}
|
||||
if (unexpectedEvent.settlementRetryRequired || settlementRetryRequired) {
|
||||
const retried = await retryUnexpectedExitSettlement({
|
||||
context,
|
||||
@@ -170,6 +179,10 @@ export async function retryUnexpectedExitSettlement(input: {
|
||||
stableSettlementId: string
|
||||
}): Promise<boolean> {
|
||||
try {
|
||||
await input.session.journal.markPendingSubmissionsUnknown(
|
||||
input.session.fence,
|
||||
'provider_exited_before_acknowledgement'
|
||||
)
|
||||
const mutations = unexpectedExitFallbackMutations(
|
||||
input.event,
|
||||
input.session,
|
||||
|
||||
@@ -86,8 +86,7 @@ export function useStructuredAgentSession(args: {
|
||||
// A dispatch the provider has not answered is already work; Claude's running row trails the
|
||||
// send by seconds, and only a provider-minted turn is cancellable, so the two stay separate.
|
||||
const isWorking =
|
||||
turnId !== null ||
|
||||
hasUnansweredStructuredAgentSessionDispatch(state.submissions)
|
||||
turnId !== null || hasUnansweredStructuredAgentSessionDispatch(state.submissions, state.fence)
|
||||
const turnActivity = useMemo(
|
||||
() => selectStructuredAgentTurnActivity(state.items, turnId, state.activity),
|
||||
[state.activity, state.items, turnId]
|
||||
|
||||
@@ -193,6 +193,7 @@ export type AgentJournalDispatchState = (typeof AGENT_JOURNAL_DISPATCH_STATES)[n
|
||||
* the turn reads as delivery unconfirmed, never as sent and never as failed. */
|
||||
export type AgentJournalSubmission = {
|
||||
clientMessageId: string
|
||||
/** Execution fence of the latest dispatch attempt or recovery. */
|
||||
fence: number
|
||||
payloadFingerprint: string
|
||||
dispatchState: AgentJournalDispatchState
|
||||
|
||||
@@ -1,9 +1,6 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { AGENT_STATUS_MAX_FIELD_LENGTH } from './agent-status-field-normalization'
|
||||
import type {
|
||||
AgentJournalRenderItem,
|
||||
AgentJournalSubmission
|
||||
} from './agent-session-journal-types'
|
||||
import type { AgentJournalRenderItem, AgentJournalSubmission } from './agent-session-journal-types'
|
||||
import { parsePaneKey } from './stable-pane-id'
|
||||
import {
|
||||
activeStructuredAgentSessionTurnId,
|
||||
@@ -180,6 +177,20 @@ describe('structured agent session status projection', () => {
|
||||
})
|
||||
})
|
||||
|
||||
it('does not resurrect old-host unknown work after its execution fence advances', () => {
|
||||
const oldHostSubmission = { ...submission('m1', 'unknown'), fence: 2 }
|
||||
expect(hasUnansweredStructuredAgentSessionDispatch([oldHostSubmission], 2)).toBe(true)
|
||||
expect(hasUnansweredStructuredAgentSessionDispatch([oldHostSubmission], 3)).toBe(false)
|
||||
})
|
||||
|
||||
it('recognizes recovery from an older host without the optional marker', () => {
|
||||
expect(
|
||||
hasUnansweredStructuredAgentSessionDispatch([
|
||||
{ ...submission('m1', 'unknown'), reason: 'host_restarted_before_acknowledgement' }
|
||||
])
|
||||
).toBe(false)
|
||||
})
|
||||
|
||||
it('stops reading a resolved dispatch as work, and lets a pending prompt outrank it', () => {
|
||||
const asked = item('asked', 1, {
|
||||
kind: 'message',
|
||||
@@ -205,9 +216,9 @@ describe('structured agent session status projection', () => {
|
||||
{ ...submission('m1', 'unknown'), recovered: true }
|
||||
])
|
||||
).toBe(false)
|
||||
expect(projectStructuredAgentSessionStatus([asked, prompt], [submission('m1', 'pending')])).toBe(
|
||||
'attention'
|
||||
)
|
||||
expect(
|
||||
projectStructuredAgentSessionStatus([asked, prompt], [submission('m1', 'pending')])
|
||||
).toBe('attention')
|
||||
expect(projectStructuredAgentSessionStatusSummary([], [])).toEqual({
|
||||
status: null,
|
||||
latestPrompt: ''
|
||||
|
||||
@@ -190,12 +190,17 @@ export function hasPersistedStructuredAgentSessionTurn(
|
||||
* that sent it, so there is nothing still running to report.
|
||||
*/
|
||||
export function hasUnansweredStructuredAgentSessionDispatch(
|
||||
submissions: readonly AgentJournalSubmission[]
|
||||
submissions: readonly AgentJournalSubmission[],
|
||||
currentFence?: number | null
|
||||
): boolean {
|
||||
return submissions.some(
|
||||
(submission) =>
|
||||
submission.dispatchState === 'pending' ||
|
||||
(submission.dispatchState === 'unknown' && submission.recovered !== true)
|
||||
(currentFence == null || submission.fence >= currentFence) &&
|
||||
(submission.dispatchState === 'pending' ||
|
||||
(submission.dispatchState === 'unknown' &&
|
||||
submission.recovered !== true &&
|
||||
// Older hosts publish the recovery reason but omit the optional marker.
|
||||
submission.reason !== 'host_restarted_before_acknowledgement'))
|
||||
)
|
||||
}
|
||||
|
||||
@@ -207,7 +212,8 @@ export function structuredAgentSessionTabId(sessionId: string): string {
|
||||
|
||||
export function projectStructuredAgentSessionStatus(
|
||||
items: readonly AgentJournalRenderItem[],
|
||||
submissions: readonly AgentJournalSubmission[] = []
|
||||
submissions: readonly AgentJournalSubmission[] = [],
|
||||
currentFence?: number | null
|
||||
): StructuredAgentSessionProjectedStatus {
|
||||
if (
|
||||
items.some(
|
||||
@@ -219,7 +225,7 @@ export function projectStructuredAgentSessionStatus(
|
||||
return 'attention'
|
||||
}
|
||||
return activeStructuredAgentSessionTurnId(items) ||
|
||||
hasUnansweredStructuredAgentSessionDispatch(submissions)
|
||||
hasUnansweredStructuredAgentSessionDispatch(submissions, currentFence)
|
||||
? 'working'
|
||||
: 'idle'
|
||||
}
|
||||
@@ -299,17 +305,18 @@ export type StructuredAgentSessionStatusProjection = {
|
||||
* has to stay small even though the row only ever renders one line of it. */
|
||||
export function projectStructuredAgentSessionStatusSummary(
|
||||
items: readonly AgentJournalRenderItem[],
|
||||
submissions: readonly AgentJournalSubmission[] = []
|
||||
submissions: readonly AgentJournalSubmission[] = [],
|
||||
currentFence?: number | null
|
||||
): StructuredAgentSessionStatusProjection {
|
||||
// A first send has no journalled message until the provider replays it, so the pending
|
||||
// dispatch is also what makes a brand-new session listable at all.
|
||||
if (
|
||||
!hasPersistedStructuredAgentSessionTurn(items) &&
|
||||
!hasUnansweredStructuredAgentSessionDispatch(submissions)
|
||||
!hasUnansweredStructuredAgentSessionDispatch(submissions, currentFence)
|
||||
) {
|
||||
return { status: null, latestPrompt: '' }
|
||||
}
|
||||
const status = projectStructuredAgentSessionStatus(items, submissions)
|
||||
const status = projectStructuredAgentSessionStatus(items, submissions, currentFence)
|
||||
const activeToolCall = status === 'working' ? activeStructuredAgentSessionToolCall(items) : null
|
||||
const toolName = activeToolCall
|
||||
? normalizeOptionalField(activeToolCall.name, AGENT_STATUS_TOOL_NAME_MAX_LENGTH)
|
||||
|
||||
Reference in New Issue
Block a user