diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-admission.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-admission.ts index 5bb072f2594..316df7a5ef9 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-admission.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-admission.ts @@ -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( 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( 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( 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') { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-plans.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-plans.ts index a81ea5270c8..2fe72e23a4a 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-plans.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-mutation-plans.ts @@ -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 = { markUnknownBeforeRun?: boolean run: (ctx: AgentSessionTurnContext) => Promise> 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 } } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-replay-outcome.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-replay-outcome.ts index 00fa33a30fb..902a40f872f 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-replay-outcome.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-replay-outcome.ts @@ -13,10 +13,22 @@ export type AgentSessionReplayOutcomeDecision = | { 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(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 { @@ -41,14 +53,7 @@ export function resolveAgentSessionReplayOutcome(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(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' } } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-resend-answer.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-resend-answer.test.ts new file mode 100644 index 00000000000..713889afa3c --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-resend-answer.test.ts @@ -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 + +beforeEach(() => { + ;({ root, store, host, dispatch } = hostTestState()) +}) + +afterEach(() => vi.restoreAllMocks()) + +function hostJournal(): AgentSessionJournal { + return ( + host as unknown as { sessions: Map } + ).sessions.get(SESSION)!.journal +} + +function sendParams(text: string) { + const body = hostTestMessage(text) + return { envelope: envelope('agentSession.send', { body }), body } +} + +async function deliveredOnce(): Promise { + 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((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) + }) +}) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-send-preparation.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-send-preparation.test.ts index cf746e29b74..647f0a1d86d 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-send-preparation.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-send-preparation.test.ts @@ -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, 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 diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-send-preparation.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-send-preparation.ts index 77413f45634..edbc4f809d4 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-send-preparation.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-send-preparation.ts @@ -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, envelope: AgentSessionMutationEnvelope -): () => Promise { - return async () => { +): (ledger: 'admit' | 'replay') => Promise { + 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 } } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-send.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-send.test.ts index f0163b22791..1eaaea9bf1f 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-send.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-send.test.ts @@ -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) }) diff --git a/src/main/native-chat/agent-session-wire/structured-conversation-command-controller.ts b/src/main/native-chat/agent-session-wire/structured-conversation-command-controller.ts index 3c7f71dbe1f..15e562fb04f 100644 --- a/src/main/native-chat/agent-session-wire/structured-conversation-command-controller.ts +++ b/src/main/native-chat/agent-session-wire/structured-conversation-command-controller.ts @@ -17,13 +17,18 @@ export class StructuredConversationCommandController { private readonly context: () => StructuredAgentSessionMutationContext, private readonly host: Pick ) {} + /** 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[2] - ): ReturnType => - this.pending.has(params.envelope.sessionId) + ): ReturnType => { + 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') { diff --git a/src/main/runtime/agent-session-operation-admission.ts b/src/main/runtime/agent-session-operation-admission.ts index f16d48efb9d..79fa72d4a31 100644 --- a/src/main/runtime/agent-session-operation-admission.ts +++ b/src/main/runtime/agent-session-operation-admission.ts @@ -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 @@ -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 } : {}) } } diff --git a/src/main/runtime/agent-session-record-store.ts b/src/main/runtime/agent-session-record-store.ts index cd87b024c0f..0ab692e047f 100644 --- a/src/main/runtime/agent-session-record-store.ts +++ b/src/main/runtime/agent-session-record-store.ts @@ -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) diff --git a/src/shared/agent-session-operation-ledger.ts b/src/shared/agent-session-operation-ledger.ts index 136b86af2b7..ea56287ad77 100644 --- a/src/shared/agent-session-operation-ledger.ts +++ b/src/shared/agent-session-operation-ledger.ts @@ -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 } : {}) } } diff --git a/src/shared/protocol-version.ts b/src/shared/protocol-version.ts index 47deda80d81..76a8a94e7a2 100644 --- a/src/shared/protocol-version.ts +++ b/src/shared/protocol-version.ts @@ -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,