mirror of
https://github.com/stablyai/orca.git
synced 2026-10-07 16:02:29 +00:00
fix(native-chat): a resent send id gets its recorded answer, never an early refusal or a made-up record
The host now looks a resent send id up before preparing the session. A row that settled refused answers with its refusal before the chat is opened. A resend whose chat cannot be opened or made ready answers unknown instead of a refusal. A /clear in flight refuses only ids the ledger does not hold. A send row now records the journal epoch it was admitted into. An unsettled row with nothing written in that same epoch runs for the first time; under a later epoch the host answers unknown instead of reconstructing a submission it never had. The host advertises agent-session.send-answers-proof.v1.
This commit is contained in:
+17
-3
@@ -6,7 +6,9 @@
|
||||
// Admission is two-phase for a call that brings a `prepareSession`. The ledger's
|
||||
// answer comes first and places nothing; a call it will admit may then give the
|
||||
// session an owner, and only after that are the row placed and the lease
|
||||
// checked — against the lease as it stands once the owner is there.
|
||||
// checked — against the lease as it stands once the owner is there. An id the
|
||||
// ledger already refused for good is answered from its row before any of that,
|
||||
// so a resend never meets a refusal its first run did not make.
|
||||
|
||||
import {
|
||||
admitAgentSessionMutation,
|
||||
@@ -103,6 +105,17 @@ export async function admitAndRunAgentSessionMutation<TValue>(
|
||||
if (!ledger) {
|
||||
return refuseAgentSessionMutation(AGENT_SESSION_NOT_ATTACHED)
|
||||
}
|
||||
const recorded = ledger.decision.decision === 'replay' ? ledger.decision.row : null
|
||||
if (recorded?.outcome.status === 'failed') {
|
||||
const replay = resolveAgentSessionReplayOutcome({
|
||||
operationId: envelope.clientOperationId,
|
||||
outcome: recorded.outcome,
|
||||
reconstruct: () => null
|
||||
})
|
||||
if (replay.decision === 'refuse') {
|
||||
return refuseAgentSessionMutation(replay.refusal)
|
||||
}
|
||||
}
|
||||
if (ledger.decision.decision !== 'refused') {
|
||||
const prepared = await request.prepareSession(ledger.decision.decision, ledger.record)
|
||||
if (!prepared.ok) {
|
||||
@@ -120,7 +133,8 @@ export async function admitAndRunAgentSessionMutation<TValue>(
|
||||
hostFingerprint,
|
||||
now: request.now(),
|
||||
...(plan.operationIdScope ? { operationIdScope: plan.operationIdScope } : {}),
|
||||
...(plan.conversationWrite ? { conversationWrite: true } : {})
|
||||
...(plan.conversationWrite ? { conversationWrite: true } : {}),
|
||||
journalEpoch: journal.cursor().epoch
|
||||
}
|
||||
let admitted: AgentSessionMutationOperationDecision
|
||||
let ledgerRowWritten = true
|
||||
@@ -155,7 +169,7 @@ export async function admitAndRunAgentSessionMutation<TValue>(
|
||||
operationId: envelope.clientOperationId,
|
||||
outcome: admission.row.outcome,
|
||||
reconstruct: () => plan.replay(context, admission.row.outcome),
|
||||
rerunWhenReplayMissing: plan.rerunWhenReplayMissing?.(context),
|
||||
rerunWhenReplayMissing: plan.rerunWhenReplayMissing?.(context, admission.row),
|
||||
recoverUnknownFromDurableState: plan.recoverUnknownFromDurableState
|
||||
})
|
||||
if (replay.decision === 'refuse') {
|
||||
|
||||
@@ -3,11 +3,15 @@
|
||||
//
|
||||
// The replay half matters more than it looks. The ledger records only that an
|
||||
// operation happened, so the durable answer usually comes back out of the
|
||||
// journal. Send is fail-closed: admission alone cannot prove non-delivery.
|
||||
// journal. Send is fail-closed: admission alone cannot prove non-delivery, so a
|
||||
// send runs again only when the journal it wrote to proves it wrote nothing.
|
||||
|
||||
import type { AgentJournalMessageItem } from '../../../shared/agent-session-journal-types'
|
||||
import type { AgentChildWorkView } from '../../../shared/agent-status-child-work-view'
|
||||
import type { AgentSessionOperationOutcome } from '../../../shared/agent-session-operation-ledger'
|
||||
import type {
|
||||
AgentSessionOperationOutcome,
|
||||
AgentSessionOperationRow
|
||||
} from '../../../shared/agent-session-operation-ledger'
|
||||
import type {
|
||||
AgentSessionCancelResult,
|
||||
AgentSessionMutationEnvelope,
|
||||
@@ -54,7 +58,7 @@ export type MutationPlan<TValue> = {
|
||||
markUnknownBeforeRun?: boolean
|
||||
run: (ctx: AgentSessionTurnContext) => Promise<TurnOutcome<TValue>>
|
||||
replay: (ctx: AgentSessionTurnContext, outcome: AgentSessionOperationOutcome) => TValue | null
|
||||
rerunWhenReplayMissing?: (ctx: AgentSessionTurnContext) => boolean
|
||||
rerunWhenReplayMissing?: (ctx: AgentSessionTurnContext, row: AgentSessionOperationRow) => boolean
|
||||
recoverUnknownFromDurableState?: boolean
|
||||
settledOutcome?: (value: TValue) => AgentSessionOperationOutcome
|
||||
}
|
||||
@@ -104,7 +108,9 @@ export function sendPlan(params: {
|
||||
if (submission) {
|
||||
return { clientMessageId, submission }
|
||||
}
|
||||
if (outcome.status === 'failed') {
|
||||
// Only an accepted send whose row a later epoch dropped is answered without one; an
|
||||
// unsettled one with nothing written is decided by `rerunWhenReplayMissing`.
|
||||
if (outcome.status !== 'succeeded') {
|
||||
return null
|
||||
}
|
||||
const resolvedAt = ctx.now()
|
||||
@@ -122,7 +128,13 @@ export function sendPlan(params: {
|
||||
recovered: true
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
// A send writes its submission, draft or hand-off before anything can deliver it, and only a
|
||||
// new epoch removes one. So in the epoch it was admitted into, finding none proves it never
|
||||
// wrote: running it now is its first run. Under any other epoch, or a row with none recorded,
|
||||
// that proof is gone and the answer is unknown.
|
||||
rerunWhenReplayMissing: (ctx, row) =>
|
||||
row.journalEpoch !== undefined && row.journalEpoch === ctx.journal.cursor().epoch
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+25
-17
@@ -13,10 +13,22 @@ export type AgentSessionReplayOutcomeDecision<TValue> =
|
||||
| { decision: 'rerun' }
|
||||
| { decision: 'refuse'; refusal: AgentSessionWireRefusal }
|
||||
|
||||
/** A recorded id this host cannot answer from what it holds, so it neither ran it again nor
|
||||
* refused it: the caller resends under the same id. */
|
||||
export function agentSessionOperationOutcomeUnknown(operationId: string): AgentSessionWireRefusal {
|
||||
return refuse(
|
||||
'agent_session_operation_unknown',
|
||||
{ reason: 'outcomeUnknown' },
|
||||
`The outcome of operation ${operationId} is unknown; it was not run again.`
|
||||
)
|
||||
}
|
||||
|
||||
export function resolveAgentSessionReplayOutcome<TValue>(input: {
|
||||
operationId: string
|
||||
outcome: AgentSessionOperationOutcome
|
||||
reconstruct: () => TValue | null
|
||||
/** Whether a row with nothing to reconstruct may run for the first time. Undefined leaves an
|
||||
* unsettled row to the default (rerun); false refuses it as unknown. */
|
||||
rerunWhenReplayMissing?: boolean
|
||||
recoverUnknownFromDurableState?: boolean
|
||||
}): AgentSessionReplayOutcomeDecision<TValue> {
|
||||
@@ -41,14 +53,7 @@ export function resolveAgentSessionReplayOutcome<TValue>(input: {
|
||||
if (input.rerunWhenReplayMissing) {
|
||||
return { decision: 'rerun' }
|
||||
}
|
||||
return {
|
||||
decision: 'refuse',
|
||||
refusal: refuse(
|
||||
'agent_session_operation_unknown',
|
||||
{ reason: 'outcomeUnknown' },
|
||||
`The outcome of operation ${operationId} is unknown; it was not run again.`
|
||||
)
|
||||
}
|
||||
return { decision: 'refuse', refusal: agentSessionOperationOutcomeUnknown(operationId) }
|
||||
}
|
||||
const recorded = input.reconstruct()
|
||||
if (recorded) {
|
||||
@@ -57,15 +62,18 @@ export function resolveAgentSessionReplayOutcome<TValue>(input: {
|
||||
if (input.rerunWhenReplayMissing) {
|
||||
return { decision: 'rerun' }
|
||||
}
|
||||
return outcome.status === 'succeeded'
|
||||
? {
|
||||
decision: 'refuse',
|
||||
refusal: refuse(
|
||||
'agent_session_operation_unknown',
|
||||
{ reason: 'resultLost' },
|
||||
`Operation ${operationId} succeeded, but its result is no longer reconstructable.`
|
||||
)
|
||||
}
|
||||
if (outcome.status === 'succeeded') {
|
||||
return {
|
||||
decision: 'refuse',
|
||||
refusal: refuse(
|
||||
'agent_session_operation_unknown',
|
||||
{ reason: 'resultLost' },
|
||||
`Operation ${operationId} succeeded, but its result is no longer reconstructable.`
|
||||
)
|
||||
}
|
||||
}
|
||||
return input.rerunWhenReplayMissing === false
|
||||
? { decision: 'refuse', refusal: agentSessionOperationOutcomeUnknown(operationId) }
|
||||
: { decision: 'rerun' }
|
||||
}
|
||||
|
||||
|
||||
+168
@@ -0,0 +1,168 @@
|
||||
// What a resend of a send id gets: the answer its record holds, never a refusal made before the
|
||||
// host looked the id up, and never a made-up record.
|
||||
|
||||
import { afterEach, beforeEach, describe, expect, it, vi, type Mock } from 'vitest'
|
||||
import type { AgentSessionRecordStore } from '../../runtime/agent-session-record-store'
|
||||
import type { AgentSessionJournal } from '../agent-session-journal/journal-store'
|
||||
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
|
||||
import { DISPATCH_DOUBT_SUBMISSION_MISSING } from '../agent-session-journal/journal-dispatch-doubt-reasons'
|
||||
import type { StructuredAgentSessionAdapter } from './structured-agent-session-adapter'
|
||||
import type { StructuredAgentSessionHost } from './structured-agent-session-host'
|
||||
import {
|
||||
attach,
|
||||
CALLER,
|
||||
envelope,
|
||||
hostTestState
|
||||
} from './structured-agent-session-host-test-harness'
|
||||
import {
|
||||
HOST_TEST_SESSION as SESSION,
|
||||
hostTestMessage
|
||||
} from './structured-agent-session-host-test-data'
|
||||
|
||||
let root: string
|
||||
let store: AgentSessionRecordStore
|
||||
let host: StructuredAgentSessionHost
|
||||
let dispatch: Mock<StructuredAgentSessionAdapter['dispatch']>
|
||||
|
||||
beforeEach(() => {
|
||||
;({ root, store, host, dispatch } = hostTestState())
|
||||
})
|
||||
|
||||
afterEach(() => vi.restoreAllMocks())
|
||||
|
||||
function hostJournal(): AgentSessionJournal {
|
||||
return (
|
||||
host as unknown as { sessions: Map<string, { journal: AgentSessionJournal }> }
|
||||
).sessions.get(SESSION)!.journal
|
||||
}
|
||||
|
||||
function sendParams(text: string) {
|
||||
const body = hostTestMessage(text)
|
||||
return { envelope: envelope('agentSession.send', { body }), body }
|
||||
}
|
||||
|
||||
async function deliveredOnce(): Promise<void> {
|
||||
await vi.waitFor(() => expect(dispatch).toHaveBeenCalledTimes(1))
|
||||
}
|
||||
|
||||
describe('a resent send id', () => {
|
||||
it('is answered from a refused row without opening the chat again', async () => {
|
||||
await attach()
|
||||
vi.spyOn(hostJournal(), 'appendSubmission').mockRejectedValueOnce(new Error('disk full'))
|
||||
const params = sendParams('refused once')
|
||||
const first = await host.send(CALLER, params)
|
||||
expect(first).toMatchObject({
|
||||
ok: false,
|
||||
refusal: { code: 'agent_session_operation_invalid' }
|
||||
})
|
||||
await host.close(SESSION, 'evict')
|
||||
expect(host.hasSession(SESSION)).toBe(false)
|
||||
|
||||
const resent = await host.send(CALLER, params)
|
||||
|
||||
expect(resent).toMatchObject({
|
||||
ok: false,
|
||||
refusal: {
|
||||
code: 'agent_session_operation_invalid',
|
||||
details: { reason: 'journalWriteFailed' }
|
||||
}
|
||||
})
|
||||
// The answer came from the ledger alone: the closed chat was not opened to give it.
|
||||
expect(host.hasSession(SESSION)).toBe(false)
|
||||
expect(dispatch).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('answers unknown, never a refusal, when the chat holding its answer cannot be opened', async () => {
|
||||
await attach()
|
||||
const params = sendParams('recorded, then the chat would not open')
|
||||
await host.send(CALLER, params)
|
||||
await deliveredOnce()
|
||||
await host.close(SESSION, 'evict')
|
||||
const connection = openTestJournalHostDatabase(root).db
|
||||
const prepare = connection.prepare.bind(connection)
|
||||
vi.spyOn(connection, 'prepare').mockImplementation((sql: string) => {
|
||||
if (sql.includes('journal_')) {
|
||||
throw Object.assign(new Error('database disk image is malformed'), {
|
||||
code: 'ERR_SQLITE_ERROR',
|
||||
errcode: 11
|
||||
})
|
||||
}
|
||||
return prepare(sql)
|
||||
})
|
||||
vi.spyOn(console, 'warn').mockImplementation(() => undefined)
|
||||
|
||||
await expect(host.send(CALLER, params)).resolves.toMatchObject({
|
||||
ok: false,
|
||||
refusal: { code: 'agent_session_operation_unknown', details: { reason: 'outcomeUnknown' } }
|
||||
})
|
||||
// A new id meets the same chat as a first run, and is refused for what it is.
|
||||
await expect(host.send(CALLER, sendParams('a new message'))).resolves.toMatchObject({
|
||||
ok: false,
|
||||
refusal: { code: 'agent_session_journal_unreadable' }
|
||||
})
|
||||
expect(dispatch).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('is answered from its record while a /clear is in flight; a new id is refused', async () => {
|
||||
await attach()
|
||||
const params = sendParams('sent before the clear')
|
||||
expect(await host.send(CALLER, params)).toMatchObject({ ok: true, replayed: false })
|
||||
|
||||
const clearing = host.conversationCommand(CALLER, {
|
||||
command: 'clear',
|
||||
envelope: envelope('agentSession.conversationCommand', { command: 'clear' })
|
||||
})
|
||||
const resent = host.send(CALLER, params)
|
||||
const fresh = host.send(CALLER, sendParams('typed during the clear'))
|
||||
await clearing
|
||||
|
||||
await expect(resent).resolves.toMatchObject({
|
||||
ok: true,
|
||||
replayed: true,
|
||||
value: { submission: { clientMessageId: params.envelope.clientOperationId } }
|
||||
})
|
||||
await expect(fresh).resolves.toMatchObject({
|
||||
ok: false,
|
||||
refusal: { details: { reason: 'conversationCommandInFlight' } }
|
||||
})
|
||||
})
|
||||
|
||||
it('waits for an original still being accepted and answers with its submission', async () => {
|
||||
await attach()
|
||||
const journal = hostJournal()
|
||||
const append = journal.appendSubmission.bind(journal)
|
||||
let release!: () => void
|
||||
const held = new Promise<void>((resolve) => {
|
||||
release = resolve
|
||||
})
|
||||
vi.spyOn(journal, 'appendSubmission').mockImplementationOnce(async (...args) => {
|
||||
await held
|
||||
return append(...args)
|
||||
})
|
||||
const params = sendParams('resent while the first is mid-write')
|
||||
|
||||
const original = host.send(CALLER, params)
|
||||
const resent = host.send(CALLER, params)
|
||||
await vi.waitFor(() => expect(journal.appendSubmission).toHaveBeenCalledTimes(1))
|
||||
release()
|
||||
|
||||
const [first, second] = await Promise.all([original, resent])
|
||||
expect(first).toMatchObject({ ok: true, replayed: false })
|
||||
expect(second).toMatchObject({
|
||||
ok: true,
|
||||
replayed: true,
|
||||
value: { submission: { clientMessageId: params.envelope.clientOperationId } }
|
||||
})
|
||||
if (!second.ok || !('submission' in second.value)) {
|
||||
throw new Error('expected the submission arm')
|
||||
}
|
||||
expect(second.value.submission.reason).not.toBe(DISPATCH_DOUBT_SUBMISSION_MISSING)
|
||||
await deliveredOnce()
|
||||
expect(journal.submissions()).toHaveLength(1)
|
||||
expect(
|
||||
store
|
||||
.listOperationRows()
|
||||
.filter((row) => row.operationId === params.envelope.clientOperationId)
|
||||
).toHaveLength(1)
|
||||
})
|
||||
})
|
||||
+27
-12
@@ -274,16 +274,16 @@ describe('a send with no live owner', () => {
|
||||
expect(acquire).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('restarts nothing for a send the ledger holds but the journal never saw', async () => {
|
||||
const params = sendParams('claimed, then the host died')
|
||||
// The row was claimed and the host went down before the journal write: on replay, admission
|
||||
// reconstructs an unknown-outcome submission and never needs an owner.
|
||||
/** A row claimed for this send, then the host died before the journal write; `journalEpoch`
|
||||
* is what this build stamps, and an older build's row carries none. */
|
||||
async function claimedThenHostDied(params: ReturnType<typeof sendParams>, journalEpoch?: string) {
|
||||
await store.admitMutationOperation({
|
||||
callerKey: CALLER.callerKey,
|
||||
envelope: params.envelope,
|
||||
hostFingerprint: params.envelope.payloadFingerprint,
|
||||
now: NOW,
|
||||
operationIdScope: 'global'
|
||||
operationIdScope: 'global',
|
||||
...(journalEpoch ? { journalEpoch } : {})
|
||||
})
|
||||
await store.recordOperationOutcome({
|
||||
callerKey: CALLER.callerKey,
|
||||
@@ -300,24 +300,39 @@ describe('a send with no live owner', () => {
|
||||
})
|
||||
expect(store.getRecord(SESSION)?.lease.claimStatus).toBe('released')
|
||||
acquire.mockClear()
|
||||
|
||||
const result = await host.send(CALLER, {
|
||||
return {
|
||||
...params,
|
||||
envelope: {
|
||||
...params.envelope,
|
||||
expectedRuntimeFence: store.getRecord(SESSION)?.lease.runtimeFence ?? 0
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
expect(result).toMatchObject({
|
||||
ok: true,
|
||||
replayed: true,
|
||||
value: { submission: { dispatchState: 'unknown', recovered: true } }
|
||||
it('restarts nothing, and answers unknown, for a row with no epoch the journal never saw', async () => {
|
||||
const resent = await claimedThenHostDied(sendParams('claimed, then the host died'))
|
||||
|
||||
await expect(host.send(CALLER, resent)).resolves.toMatchObject({
|
||||
ok: false,
|
||||
refusal: { code: 'agent_session_operation_unknown', details: { reason: 'outcomeUnknown' } }
|
||||
})
|
||||
expect(acquire).not.toHaveBeenCalled()
|
||||
expect(dispatch).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('runs a send its own epoch never saw for the first time, restarting the owner once', async () => {
|
||||
const epoch = (await host.journalSnapshot(SESSION)).cursor.epoch
|
||||
const resent = await claimedThenHostDied(sendParams('claimed in this epoch'), epoch)
|
||||
|
||||
await expect(host.send(CALLER, resent)).resolves.toMatchObject({
|
||||
ok: true,
|
||||
replayed: false,
|
||||
value: { submission: { dispatchState: 'pending' } }
|
||||
})
|
||||
await eventually(async () => expect(dispatch).toHaveBeenCalledOnce())
|
||||
expect(acquire).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('accepts a send that arrives while a restart holds the queue, and hands both over in order', async () => {
|
||||
await loseOwner()
|
||||
const lostFence = store.getRecord(SESSION)?.lease.runtimeFence ?? 0
|
||||
|
||||
+13
-6
@@ -22,6 +22,7 @@ import {
|
||||
AGENT_SESSION_NOT_ATTACHED,
|
||||
type AgentSessionMutationSessionPreparation
|
||||
} from './structured-agent-session-mutation-admission'
|
||||
import { agentSessionOperationOutcomeUnknown } from './structured-agent-session-replay-outcome'
|
||||
import { rewindRefusal } from './structured-rewind-refusal'
|
||||
import type { StructuredAgentSessionMutationContext } from './structured-agent-session-host-mutations'
|
||||
import type { StructuredAgentSessionLogger } from './structured-agent-session-logger'
|
||||
@@ -132,21 +133,27 @@ export function openWithAgent(
|
||||
}
|
||||
|
||||
/** A rewind still in doubt once the conversation is open is one only its provider can settle —
|
||||
* the open settles every other — so a send starts the agent, whose attach recovers it. */
|
||||
* the open settles every other — so a send starts the agent, whose attach recovers it. For a
|
||||
* resend of a recorded id the answer is in the conversation: one that cannot be made ready leaves
|
||||
* that answer unknown, never refused. */
|
||||
export function sendPreparation(
|
||||
context: Pick<StructuredAgentSessionMutationContext, 'openConversation' | 'ensureAgent' | 'deps'>,
|
||||
envelope: AgentSessionMutationEnvelope
|
||||
): () => Promise<AgentSessionMutationSessionPreparation> {
|
||||
return async () => {
|
||||
): (ledger: 'admit' | 'replay') => Promise<AgentSessionMutationSessionPreparation> {
|
||||
return async (ledger) => {
|
||||
const opened = await openConversationForWrite(
|
||||
context.openConversation,
|
||||
envelope,
|
||||
context.deps.logger
|
||||
)
|
||||
const phase = context.deps.store.getRecord(envelope.sessionId)?.rewind?.phase
|
||||
return opened.ok && (phase === 'prepared' || phase === 'provider-succeeded')
|
||||
? context.ensureAgent(envelope.sessionId)
|
||||
: opened
|
||||
const prepared =
|
||||
opened.ok && (phase === 'prepared' || phase === 'provider-succeeded')
|
||||
? await context.ensureAgent(envelope.sessionId)
|
||||
: opened
|
||||
return ledger === 'replay' && !prepared.ok
|
||||
? { ok: false, refusal: agentSessionOperationOutcomeUnknown(envelope.clientOperationId) }
|
||||
: prepared
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -328,32 +328,32 @@ describe('send', () => {
|
||||
expect(dispatch).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('never reruns an admission-only send after the caller changes', async () => {
|
||||
it('runs an admission-only send for the first time when it is resent, once', async () => {
|
||||
await attach()
|
||||
const settlement = vi
|
||||
.spyOn(store, 'recordOperationOutcome')
|
||||
.mockRejectedValue(new Error('operation settlement failed'))
|
||||
const body = hostTestMessage('first delivery after caller recovery')
|
||||
const params = { envelope: envelope('agentSession.send', { body }), body }
|
||||
const clientMessageId = params.envelope.clientOperationId
|
||||
|
||||
await expect(host.send(CALLER, params)).rejects.toThrow('operation settlement failed')
|
||||
expect(dispatch).not.toHaveBeenCalled()
|
||||
expect(hostJournal().submissions()).toHaveLength(0)
|
||||
settlement.mockRestore()
|
||||
|
||||
// The row was placed and nothing was written in the epoch it names: this is its first run.
|
||||
await expect(host.send({ callerKey: 'client-after-recovery' }, params)).resolves.toMatchObject({
|
||||
ok: true,
|
||||
replayed: true,
|
||||
value: {
|
||||
submission: {
|
||||
dispatchState: 'unknown',
|
||||
reason: DISPATCH_DOUBT_SUBMISSION_MISSING
|
||||
}
|
||||
}
|
||||
replayed: false,
|
||||
value: { submission: { clientMessageId, dispatchState: 'pending' } }
|
||||
})
|
||||
expect(dispatch).not.toHaveBeenCalled()
|
||||
await delivered(clientMessageId)
|
||||
expect(
|
||||
store.listOperationRows().find((row) => row.operationId === params.envelope.clientOperationId)
|
||||
).toMatchObject({ callerKey: CALLER.callerKey, outcome: { status: 'pending' } })
|
||||
store.listOperationRows().find((row) => row.operationId === clientMessageId)
|
||||
).toMatchObject({ callerKey: CALLER.callerKey, outcome: { status: 'succeeded' } })
|
||||
await expect(host.send(CALLER, params)).resolves.toMatchObject({ ok: true, replayed: true })
|
||||
expect(dispatch).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('never redelivers after admission survives without its journal submission', async () => {
|
||||
@@ -377,23 +377,38 @@ describe('send', () => {
|
||||
await journal.rollEpoch('schema_unreadable', store.getRecord(SESSION)?.lease.runtimeFence ?? 1)
|
||||
expect(journal.submissions()).toHaveLength(0)
|
||||
|
||||
// A new epoch may have dropped what the send wrote, so nothing proves it never ran.
|
||||
await expect(
|
||||
host.send({ callerKey: 'client-after-recovery' }, { ...params, retryUnknown: true })
|
||||
).resolves.toMatchObject({
|
||||
ok: true,
|
||||
replayed: true,
|
||||
value: {
|
||||
submission: {
|
||||
dispatchState: 'unknown',
|
||||
reason: DISPATCH_DOUBT_SUBMISSION_MISSING,
|
||||
recovered: true
|
||||
}
|
||||
}
|
||||
ok: false,
|
||||
refusal: { code: 'agent_session_operation_unknown', details: { reason: 'outcomeUnknown' } }
|
||||
})
|
||||
expect(dispatch).toHaveBeenCalledTimes(1)
|
||||
expect(journal.submissions()).toHaveLength(0)
|
||||
})
|
||||
|
||||
it('answers an accepted send a new epoch dropped as recorded, never by running it again', async () => {
|
||||
await attach()
|
||||
const body = hostTestMessage('accepted, then the epoch was replaced')
|
||||
const params = { envelope: envelope('agentSession.send', { body }), body }
|
||||
await host.send(CALLER, params)
|
||||
await delivered(params.envelope.clientOperationId)
|
||||
await hostJournal().rollEpoch(
|
||||
'schema_unreadable',
|
||||
store.getRecord(SESSION)?.lease.runtimeFence ?? 1
|
||||
)
|
||||
|
||||
await expect(host.send(CALLER, params)).resolves.toMatchObject({
|
||||
ok: true,
|
||||
replayed: true,
|
||||
value: {
|
||||
submission: { dispatchState: 'unknown', reason: DISPATCH_DOUBT_SUBMISSION_MISSING }
|
||||
}
|
||||
})
|
||||
expect(dispatch).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('fails closed when a legacy pending row survives without its submission', async () => {
|
||||
await attach()
|
||||
const body = hostTestMessage('legacy pending send after caller recovery')
|
||||
@@ -414,14 +429,8 @@ describe('send', () => {
|
||||
await journal.rollEpoch('schema_unreadable', store.getRecord(SESSION)?.lease.runtimeFence ?? 1)
|
||||
|
||||
await expect(host.send({ callerKey: 'client-after-recovery' }, params)).resolves.toMatchObject({
|
||||
ok: true,
|
||||
replayed: true,
|
||||
value: {
|
||||
submission: {
|
||||
dispatchState: 'unknown',
|
||||
reason: DISPATCH_DOUBT_SUBMISSION_MISSING
|
||||
}
|
||||
}
|
||||
ok: false,
|
||||
refusal: { code: 'agent_session_operation_unknown', details: { reason: 'outcomeUnknown' } }
|
||||
})
|
||||
expect(dispatch).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
+8
-3
@@ -17,13 +17,18 @@ export class StructuredConversationCommandController {
|
||||
private readonly context: () => StructuredAgentSessionMutationContext,
|
||||
private readonly host: Pick<StructuredAgentSessionHost, 'waitForSendSettlement'>
|
||||
) {}
|
||||
/** A clear in flight refuses only a send it would be the first run of: a resent id the ledger
|
||||
* holds is answered from its record, behind the clear. */
|
||||
send = (
|
||||
caller: StructuredAgentSessionCaller,
|
||||
params: Parameters<typeof sendStructuredAgentSessionTurn>[2]
|
||||
): ReturnType<typeof sendStructuredAgentSessionTurn> =>
|
||||
this.pending.has(params.envelope.sessionId)
|
||||
): ReturnType<typeof sendStructuredAgentSessionTurn> => {
|
||||
const context = this.context()
|
||||
return this.pending.has(params.envelope.sessionId) &&
|
||||
!context.deps.store.holdsGlobalOperation(params.envelope.clientOperationId, context.now())
|
||||
? Promise.resolve({ ok: false, refusal: conversationCommandInFlight() })
|
||||
: sendStructuredAgentSessionTurn(this.context(), caller, params)
|
||||
: sendStructuredAgentSessionTurn(context, caller, params)
|
||||
}
|
||||
|
||||
run = (caller: StructuredAgentSessionCaller, params: ConversationCommandParams) => {
|
||||
if (params.command === 'compact') {
|
||||
|
||||
@@ -26,6 +26,8 @@ export type AgentSessionOperationAdmission = {
|
||||
operationId: string
|
||||
fingerprint: string
|
||||
now: number
|
||||
/** Stamped on a row this admission places: see `AgentSessionOperationRow.journalEpoch`. */
|
||||
journalEpoch?: string
|
||||
}
|
||||
|
||||
type OperationRows = Map<string, AgentSessionOperationRow>
|
||||
@@ -37,6 +39,8 @@ export type AgentSessionMutationOperationAdmission = {
|
||||
now: number
|
||||
operationIdScope?: 'global'
|
||||
conversationWrite?: true
|
||||
/** The journal epoch the mutation writes to. */
|
||||
journalEpoch?: string
|
||||
}
|
||||
|
||||
export type AgentSessionMutationOperationDecision = {
|
||||
@@ -129,7 +133,8 @@ function mutationOperation(
|
||||
callerKey: args.callerKey,
|
||||
operationId: args.envelope.clientOperationId,
|
||||
fingerprint: args.hostFingerprint,
|
||||
now: args.now
|
||||
now: args.now,
|
||||
...(args.journalEpoch !== undefined ? { journalEpoch: args.journalEpoch } : {})
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -11,6 +11,7 @@ import { setAgentSessionRecordConversationName } from './agent-session-record-co
|
||||
|
||||
import {
|
||||
agentSessionOperationKey,
|
||||
findAgentSessionGlobalOperationRow,
|
||||
type AgentSessionOperationClaim,
|
||||
type AgentSessionOperationDecision,
|
||||
type AgentSessionOperationOutcome,
|
||||
@@ -181,6 +182,10 @@ export class AgentSessionRecordStore {
|
||||
getOperationRow = (callerKey: string, operationId: string): AgentSessionOperationRow | null =>
|
||||
this.state.operations.get(agentSessionOperationKey(callerKey, operationId)) ?? null
|
||||
|
||||
/** Whether any caller's unexpired row holds this id, as a send's global admission reads it. */
|
||||
holdsGlobalOperation = (operationId: string, now: number): boolean =>
|
||||
findAgentSessionGlobalOperationRow(this.state.operations, operationId, now) !== undefined
|
||||
|
||||
isClaimKeyVerifiable = (keyId: string, now: number): boolean =>
|
||||
isAgentSessionClaimKeyVerifiable(this.state, keyId, now)
|
||||
|
||||
|
||||
@@ -76,6 +76,14 @@ export type AgentSessionOperationRow = {
|
||||
recordedAt: number
|
||||
expiresAt: number
|
||||
outcome: AgentSessionOperationOutcome
|
||||
/**
|
||||
* The chat journal epoch live when the row was placed, for an operation that writes to that
|
||||
* journal. While the same epoch is live, a row the operation never wrote is one it never wrote;
|
||||
* an epoch replaced since (a rewind, an import, a rebuild) may have dropped it. Absent on rows
|
||||
* no journal backs and on rows older builds wrote. Not checked by `isAgentSessionOperationRow`,
|
||||
* for the reason `launch` is not: a load drops a row it rejects.
|
||||
*/
|
||||
journalEpoch?: string
|
||||
}
|
||||
|
||||
export type AgentSessionOperationRefusalCode =
|
||||
@@ -235,6 +243,7 @@ export function evaluateAgentSessionOperation(args: {
|
||||
now: number
|
||||
perClientLimit?: number
|
||||
globalLimit?: number
|
||||
journalEpoch?: string
|
||||
}): AgentSessionOperationDecision {
|
||||
const { rows, callerKey, operationId, fingerprint, now } = args
|
||||
const operationTimestamp = parseAgentSessionOperationTimestamp(operationId)
|
||||
@@ -288,7 +297,13 @@ export function evaluateAgentSessionOperation(args: {
|
||||
}
|
||||
return {
|
||||
decision: 'admit',
|
||||
row: pendingAgentSessionOperationRow({ callerKey, operationId, fingerprint, now })
|
||||
row: pendingAgentSessionOperationRow({
|
||||
callerKey,
|
||||
operationId,
|
||||
fingerprint,
|
||||
now,
|
||||
...(args.journalEpoch !== undefined ? { journalEpoch: args.journalEpoch } : {})
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -298,6 +313,7 @@ export function pendingAgentSessionOperationRow(args: {
|
||||
operationId: string
|
||||
fingerprint: string
|
||||
now: number
|
||||
journalEpoch?: string
|
||||
}): AgentSessionOperationRow {
|
||||
const operationTimestamp = parseAgentSessionOperationTimestamp(args.operationId)
|
||||
if (operationTimestamp === null) {
|
||||
@@ -310,7 +326,8 @@ export function pendingAgentSessionOperationRow(args: {
|
||||
operationTimestamp,
|
||||
recordedAt: args.now,
|
||||
expiresAt: agentSessionOperationExpiry(operationTimestamp, args.now),
|
||||
outcome: { status: 'pending' }
|
||||
outcome: { status: 'pending' },
|
||||
...(args.journalEpoch !== undefined ? { journalEpoch: args.journalEpoch } : {})
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -196,6 +196,14 @@ export const AGENT_SESSION_PENDING_SEND_RESULT_RUNTIME_CAPABILITY =
|
||||
// mobile client lacks the capability; mobile must first show a rejected message in place.
|
||||
export const AGENT_SESSION_ACCEPTED_SEND_RUNTIME_CAPABILITY =
|
||||
'agent-session.accepted-send.v1' as const
|
||||
// Why: a host advertising this answers a resent send id from its ledger before anything else may
|
||||
// refuse it, so its answers to `agentSession.send` are proof. `agent_session_operation_unknown`
|
||||
// means it cannot tell yet: resend under the same id. `agent_session_operation_expired` means the
|
||||
// id outlived the window the host answers for: only the transcript can say. Any other refusal
|
||||
// means nothing under that id was recorded. An older host may refuse an id it recorded, so a
|
||||
// client must not read that host's refusal of a resend as proof.
|
||||
export const AGENT_SESSION_SEND_ANSWERS_PROOF_RUNTIME_CAPABILITY =
|
||||
'agent-session.send-answers-proof.v1' as const
|
||||
// Why: `agentSession.send`'s params are strict, so an older host rejects `delivery`; and only a
|
||||
// capable client can render the `queued` result arm, the draft list, and returned cards. DARK ON
|
||||
// PURPOSE — not in RUNTIME_CAPABILITIES: advertising still requires the integrated Codex steer
|
||||
@@ -400,6 +408,7 @@ export const RUNTIME_CAPABILITIES = [
|
||||
// The host side: it accepts a send before any agent has it, and a Stop with no writer before a
|
||||
// turn starts, so a client may gate on either.
|
||||
AGENT_SESSION_ACCEPTED_SEND_RUNTIME_CAPABILITY,
|
||||
AGENT_SESSION_SEND_ANSWERS_PROOF_RUNTIME_CAPABILITY,
|
||||
STRUCTURED_AGENT_SESSION_HOLD_RUNTIME_CAPABILITY,
|
||||
STRUCTURED_AGENT_SESSION_REVEAL_RUNTIME_CAPABILITY,
|
||||
STRUCTURED_AGENT_SESSION_RESUME_HISTORY_RUNTIME_CAPABILITY,
|
||||
|
||||
Reference in New Issue
Block a user