diff --git a/docs/reference/remote-wire-compatibility.md b/docs/reference/remote-wire-compatibility.md index d72f93770c5..b2bcb3c40f4 100644 --- a/docs/reference/remote-wire-compatibility.md +++ b/docs/reference/remote-wire-compatibility.md @@ -181,8 +181,10 @@ hand-written. It covers the three skews that surface can fail on: - a new client against the old dispatcher always gets an answer rather than silence, and `method_not_found` for every method that release does not register, so the absence is visible during negotiation instead of by calling; -- a cursor survives a host restart: the client's fence is refused as stale with the live - one attached, and resuming from the held cursor replays only what it missed. +- a cursor survives a host restart: a reattach at the client's fence is refused as stale + with the live one attached, a write still carrying that fence is delivered (writes are + named by their target, and every released client still sends a fence), and resuming from + the held cursor replays only what it missed. Run it with: diff --git a/src/main/native-chat/agent-session-wire/agent-session-history-forward-read-budget.test.ts b/src/main/native-chat/agent-session-wire/agent-session-history-forward-read-budget.test.ts index 5a9a6f20801..f240d65273b 100644 --- a/src/main/native-chat/agent-session-wire/agent-session-history-forward-read-budget.test.ts +++ b/src/main/native-chat/agent-session-wire/agent-session-history-forward-read-budget.test.ts @@ -105,11 +105,10 @@ describe('forward history SQL read budget', () => { const journal = await seedJournal(count) const { returnedRows, parse } = observeForwardReads() const events: AgentSessionSubscribeEvent[] = [] - new AgentSessionSubscribers().open({ + new AgentSessionSubscribers({ readFence: () => 1 }).open({ id: 'reader', sessionId: identity.sessionId, journal, - fence: 1, cursor: { epoch: journal.epoch, sequence: 1 }, emit: (event) => events.push(event) }) @@ -126,11 +125,10 @@ describe('forward history SQL read budget', () => { const journal = await seedJournal(2_000) const render = vi.spyOn(journalReducer, 'renderJournalState') const events: AgentSessionSubscribeEvent[] = [] - new AgentSessionSubscribers().open({ + new AgentSessionSubscribers({ readFence: () => 1 }).open({ id: 'reader', sessionId: identity.sessionId, journal, - fence: 1, cursor: { epoch: journal.epoch, sequence: 1 }, emit: (event) => events.push(event) }) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-attach-context.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-attach-context.ts index 4b9b46716cb..1464ae5ae6c 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-attach-context.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-attach-context.ts @@ -21,13 +21,8 @@ export type StructuredAgentSessionAttachContext = { runtimeState: StructuredAgentSessionHostRuntimeState sessions: Map subscribers: { - reset: ( - sessionId: string, - journal: AgentSessionJournal, - reset: AgentJournalResetReason, - fence: number - ) => void - snapshot: (sessionId: string, journal: AgentSessionJournal, fence: number) => void + reset: (sessionId: string, journal: AgentSessionJournal, reset: AgentJournalResetReason) => void + snapshot: (sessionId: string, journal: AgentSessionJournal) => void publish: ( sessionId: string, journal: AgentSessionJournal, diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-attach-orchestration.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-attach-orchestration.ts index 5482826d072..aba50a33ff4 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-attach-orchestration.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-attach-orchestration.ts @@ -21,7 +21,6 @@ import { pinnedAgentSessionLaunchEnv } from './structured-agent-session-launch-env' import { refuseAgentSessionMutation } from './structured-agent-session-mutation-admission' -import { isResumableStructuredAgentSessionRecord } from './structured-agent-session-resume-eligibility' import { retryPendingStructuredAgentSessionSettlement } from './structured-agent-session-settlement-retry' import { settleStaleSessionStateOnAcquire } from './structured-agent-session-stale-turn-verdict' import type { StructuredAgentSessionAttachContext } from './structured-agent-session-attach-context' @@ -135,13 +134,6 @@ async function runAttach( // attach makes it the session's; any other exit closes it with whatever the child queued. const attemptSink = context.runtimeState.mintEventSink(sessionId) let attemptSinkAdopted = false - // A lease handed back cleanly is what a resume replaces. A writer current as of that owner is - // rebased onto the fence this attach publishes, since the restart is the only thing that moved it. - const released = context.deps.store.getRecord(sessionId) - const resumedFromFence = - released && isResumableStructuredAgentSessionRecord(released) - ? released.lease.runtimeFence - : undefined const attached = stampFailedCreateOwnerVerdict( context.deps.store, callerKey, @@ -227,8 +219,7 @@ async function runAttach( providerChildPhase: acquiredOwner ? providerChildPhase : (previous?.providerChildPhase ?? 'ready'), - acquisitionGeneration: acquisitionGeneration ?? previous?.acquisitionGeneration ?? null, - resumedFromFence: acquiredOwner ? resumedFromFence : previous?.resumedFromFence + acquisitionGeneration: acquisitionGeneration ?? previous?.acquisitionGeneration ?? null }) await recoverStructuredRewind( context.deps.store, @@ -240,9 +231,9 @@ async function runAttach( ) await recoverInterruptedCompaction(context.deps.store, sessionId, attached.journal, fence) if (attached.recovery) { - context.subscribers.reset(sessionId, attached.journal, attached.recovery.reset, fence) + context.subscribers.reset(sessionId, attached.journal, attached.recovery.reset) } else if (previousFence !== undefined && previousFence !== fence) { - context.subscribers.snapshot(sessionId, attached.journal, fence) + context.subscribers.snapshot(sessionId, attached.journal) } else { context.subscribers.publish(sessionId, attached.journal) } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-background-task-channel.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-background-task-channel.ts index 5d922d72350..1fbc73e4afa 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-background-task-channel.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-background-task-channel.ts @@ -48,7 +48,6 @@ export class StructuredAgentSessionBackgroundTaskChannel { return this.subscribers.open({ ...input, journal: session.journal, - fence: this.deps.store.getRecord(input.sessionId)?.lease.runtimeFence ?? 0, ...(backgroundTasks !== undefined ? { backgroundTasks } : {}) }) } @@ -57,7 +56,7 @@ export class StructuredAgentSessionBackgroundTaskChannel { const session = this.sessions.get(sessionId) const state = publishedState !== undefined ? publishedState : this.state(sessionId) if (session && state !== undefined) { - this.subscribers.backgroundTasks(sessionId, state, session.fence) + this.subscribers.backgroundTasks(sessionId, state) this.onPublished(sessionId) } } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-client-delivery.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-client-delivery.ts index 7681e0f3436..34dc46f23b4 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-client-delivery.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-client-delivery.ts @@ -36,6 +36,7 @@ export class StructuredAgentSessionClientDelivery { ) this.waitForSendSettlement = this.sendSettlement.wait this.subscribers = new AgentSessionSubscribers({ + readFence: (sessionId) => deps().store.getRecord(sessionId)?.lease.runtimeFence ?? 0, readCommands: (sessionId) => deps().adapter.readCommands?.(sessionId), onJournalPublished: (sessionId, journal) => this.publishJournal(sessionId, journal) }) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-command-publication.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-command-publication.test.ts index b9b7021e854..80c89f48ce2 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-command-publication.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-command-publication.test.ts @@ -65,23 +65,25 @@ it('delivers catalog changes through existing frames without resending them on o let commands: AgentSessionSlashCommand[] | undefined = [ { name: 'loaded', kind: 'command', kindUnspecified: true } ] - const subscribers = new AgentSessionSubscribers({ readCommands: () => commands }) + const subscribers = new AgentSessionSubscribers({ + readFence: () => 7, + readCommands: () => commands + }) const close = subscribers.open({ id: 'one', sessionId, journal, - fence: 7, emit: coalescer.push }) expect(state.commands).toEqual(commands) for (let i = 0; i < 25; i++) { - subscribers.backgroundTasks(sessionId, null, 7) + subscribers.backgroundTasks(sessionId, null) } coalescer.flush() expect(events.filter((event) => 'commands' in event)).toHaveLength(1) commands = [] subscribers.publish(sessionId, journal) - subscribers.backgroundTasks(sessionId, null, 7) + subscribers.backgroundTasks(sessionId, null) coalescer.flush() expect(state.commands).toEqual([]) expect(events.filter((event) => 'commands' in event)).toHaveLength(2) @@ -96,12 +98,11 @@ it('delivers catalog changes through existing frames without resending them on o sessionId, journal, cursor: journal.cursor(), - fence: 8, emit: coalescer.push }) coalescer.flush() expect(state.commands).toEqual(commands) - subscribers.reset(sessionId, journal, 'epoch_changed', 8) + subscribers.reset(sessionId, journal, 'epoch_changed') expect(state.commands).toEqual(commands) commands = undefined subscribers.open({ @@ -109,7 +110,6 @@ it('delivers catalog changes through existing frames without resending them on o sessionId, journal, cursor: journal.cursor(), - fence: 9, emit: coalescer.push }) coalescer.flush() diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host-types.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host-types.ts index 574ced3cc38..5db65f320ac 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host-types.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host-types.ts @@ -44,10 +44,6 @@ export type StructuredAgentSessionHostSession = { owesProviderChildWindDown?: boolean /** Exact adapter acquisition behind `hasProviderChild`; retained after exit to fence recovery. */ acquisitionGeneration: string | null - /** The fence of the released owner this child replaced, when it was resumed into a lease handed - * back cleanly. A writer current as of that owner is admitted at `fence`: the restart is the - * only thing that moved it. Absent for a create or a journal restored for reading. */ - resumedFromFence?: number } export type StructuredAgentSessionHostDeps = { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts index ba508b67bcf..b13a042a85f 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-host.ts @@ -116,8 +116,7 @@ export class StructuredAgentSessionHost { store: deps.store, sessions: this.sessions, flushLifecycle: (sessionId) => this.runtimeState.lifecycleBarrier(sessionId), - publishFence: (sessionId, session) => - this.subscribers.snapshot(sessionId, session.journal, session.fence), + publishFence: (sessionId, session) => this.subscribers.snapshot(sessionId, session.journal), publishStatus: this.clientDelivery.publishStatusAndSettlement, hasResumeCapableHolder: (sessionId) => this.holds.hasResumeCapableHolder(sessionId), restartReleaseGrace: (sessionId) => this.holds.renew(sessionId), 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 04001450621..62be8900159 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 @@ -5,8 +5,8 @@ // // 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 and -// fence checked — against the lease as it stands once the owner is there. +// 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. import { admitAgentSessionMutation, @@ -44,7 +44,7 @@ export function refuseAgentSessionMutation(refusal: AgentSessionWireRefusal): { } export type AgentSessionMutationSessionPreparation = - | { ok: true; envelope: AgentSessionMutationEnvelope } + | { ok: true } | { ok: false; refusal: AgentSessionWireRefusal } export type AgentSessionMutationRequest = { @@ -56,8 +56,7 @@ export type AgentSessionMutationRequest = { /** Journal of the attached session, read after `prepareSession`; absent when this host holds none. */ journal: () => AgentSessionJournal | undefined /** Between the ledger's answer and the lease check, for a call that may first have to make the - * session ready for itself. Answers with the envelope to admit — the caller's, or one moved - * onto a fence the preparation itself published — or with the refusal that ends the call. */ + * session ready for itself. Answers with the refusal that ends the call, if any. */ prepareSession?: ( ledger: Exclude, record: AgentSessionRecord @@ -71,8 +70,7 @@ export type AgentSessionMutationRequest = { export async function admitAndRunAgentSessionMutation( request: AgentSessionMutationRequest ): Promise> { - const { plan } = request - let { envelope } = request + const { plan, envelope } = request const hostFingerprint = computeAgentSessionPayloadFingerprint({ method: plan.method, sessionId: envelope.sessionId, @@ -98,7 +96,6 @@ export async function admitAndRunAgentSessionMutation( if (!prepared.ok) { return prepared } - envelope = prepared.envelope } } const journal = request.journal() @@ -137,9 +134,9 @@ export async function admitAndRunAgentSessionMutation( return { ok: true, replayed: true, fence, cursor: journal.cursor(), value: replay.value } } // Nothing durable landed, so this id is about to run for the first time. A - // refused call leaves its ledger row behind, and replaying past the lease and - // the fence would let a resend act under an owner that has since changed — so - // a first run pays the full admission price either way. + // refused call leaves its ledger row behind, and replaying past the lease + // would let a resend act with no live owner — so a first run pays the full + // admission price either way. const rerun = admitAgentSessionMutation({ envelope, hostFingerprint, diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-refusal-retry.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-refusal-retry.test.ts index 566156e096f..87f4f9e3d65 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-refusal-retry.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-refusal-retry.test.ts @@ -231,7 +231,10 @@ const UNREACHABLE = new Set([ // Send reconstructs doubt from its global tombstone instead of refusing it. 'agentSession.send:agent_session_operation_unknown', // Only a send restarts a lost owner. - 'agentSession.setOption:agent_session_owner_restart_failed' + 'agentSession.setOption:agent_session_owner_restart_failed', + // A write names its target, not an owner generation; only an attach compares fences. + 'agentSession.setOption:agent_session_checkpoint_stale', + 'agentSession.send:agent_session_checkpoint_stale' ]) describe('agentSessionRefusalOperationState host oracle', () => { @@ -242,20 +245,8 @@ describe('agentSessionRefusalOperationState host oracle', () => { const stale = await createHarness() for (const method of METHODS) { - const spec = { - method, - operationId: operationId(), - expectedRuntimeFence: 99 - } - record( - await assertHostAgreement(stale, spec, 'agent_session_checkpoint_stale', async () => ({ - harness: stale, - spec: { - ...spec, - expectedRuntimeFence: stale.store.getRecord(SESSION)?.lease.runtimeFence ?? 1 - } - })) - ) + const spec = { method, operationId: operationId(), expectedRuntimeFence: 99 } + await expect(invoke(stale, spec), method).resolves.toMatchObject({ ok: true }) } expect(stale.setOption).toHaveBeenCalledTimes(1) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-rewind.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-rewind.test.ts index 19ae18dd5f4..b82248cb06f 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-rewind.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-rewind.test.ts @@ -229,14 +229,8 @@ describe('host rewind', () => { expect(recoverRewind).toHaveBeenCalledTimes(2) expect(rewind).toHaveBeenCalledTimes(1) }) - it('fences stale owners and the second of two concurrent rewinds', async () => { + it('refuses the second of two concurrent rewinds by the epoch it targets', async () => { const target = await seed() - const stale = params(target) - stale.envelope.expectedRuntimeFence++ - expect(await host.rewind(caller, stale)).toMatchObject({ - ok: false, - refusal: { code: 'agent_session_checkpoint_stale' } - }) let finish!: () => void rewind.mockImplementation( () => 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 f48af3b4d4b..e5866241105 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 @@ -6,7 +6,10 @@ import { join } from 'node:path' import { afterEach, beforeEach, describe, expect, it, vi, type Mock } from 'vitest' import { computeAgentSessionPayloadFingerprint } from '../../../shared/agent-session-mutation-envelope' import { agentSessionRefusalOperationState } from '../../../shared/agent-session-refusal-retry' -import type { AgentSessionMutationEnvelope } from '../../../shared/agent-session-wire' +import type { + AgentSessionMutationEnvelope, + AgentSessionSubscribeEvent +} from '../../../shared/agent-session-wire' import { AgentSessionRecordStore } from '../../runtime/agent-session-record-store' import type { StructuredAgentSessionAdapter } from './structured-agent-session-adapter' import { StructuredAgentSessionHost } from './structured-agent-session-host' @@ -676,3 +679,70 @@ describe('a send with no live owner', () => { expect(acquire).not.toHaveBeenCalled() }) }) + +// A pane keeps the fence of the last frame it read. An idle release and the restart after it +// each move the lease, so that fence can be several generations behind the one a write lands on. +describe('a write fenced to an owner the pane has not seen replaced', () => { + it('delivers a send fenced to the owner an idle release retired', async () => { + const seenFence = store.getRecord(SESSION)?.lease.runtimeFence ?? 0 + const frames: AgentSessionSubscribeEvent[] = [] + host.subscribe({ id: 'pane', sessionId: SESSION, emit: (event) => frames.push(event) }) + await loseOwner() + const params = sendParams('after the release') + params.envelope.expectedRuntimeFence = seenFence + + expect(await host.send(CALLER, params)).toMatchObject({ + ok: true, + replayed: false, + value: { submission: { dispatchState: 'accepted' } } + }) + + expect(acquire).toHaveBeenCalledOnce() + expect(dispatch).toHaveBeenCalledOnce() + const fence = store.getRecord(SESSION)?.lease.runtimeFence ?? 0 + expect(fence).toBeGreaterThan(seenFence + 1) + // The pane learns the fence the host is at, not the one it subscribed under. + expect(frames.at(-1)).toMatchObject({ fence }) + }) + + it('admits a Stop and a send queued behind the cold start that replaced their owner', async () => { + await loseOwner() + const lostFence = store.getRecord(SESSION)?.lease.runtimeFence ?? 0 + let claimed = () => {} + let release = () => {} + const claim = new Promise((resolve) => (claimed = resolve)) + const spawn = new Promise((resolve) => (release = resolve)) + acquire.mockImplementationOnce(async (input) => { + claimed() + await spawn + return spawnChild(input) + }) + + const first = host.send(CALLER, sendParams('starts the agent')) + await claim + const cancelFields = { turnId: 'turn-1' } + const stop = host.cancel(CALLER, { + envelope: { + sessionId: SESSION, + clientOperationId: hostTestOperationId(), + expectedRuntimeFence: lostFence, + payloadFingerprint: computeAgentSessionPayloadFingerprint({ + method: 'agentSession.cancel', + sessionId: SESSION, + fields: cancelFields + }) + }, + ...cancelFields + }) + const late = sendParams('typed during the start') + late.envelope.expectedRuntimeFence = lostFence - 1 + const second = host.send(CALLER, late) + release() + + expect(await first).toMatchObject({ ok: true }) + expect(await stop).toMatchObject({ ok: true, replayed: false }) + expect(await second).toMatchObject({ ok: true, replayed: false }) + expect(acquire).toHaveBeenCalledOnce() + expect(dispatch).toHaveBeenCalledTimes(2) + }) +}) 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 e4a332f7acf..fa57156d05a 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 @@ -12,13 +12,8 @@ // // The ledger's answer comes first, so a send it already holds a row for restarts nothing: // admission replays or refuses it whoever owns the session now, and a closed session is made -// readable for that, never given a child. Otherwise a child that dies at startup moves the fence, -// the client resends the same message against the new fence, and each replay spawns another -// child that dies the same way. -// -// Running inside the send's serialize is what makes "no child" exact and the fence bookkeeping -// simple: the owner this send (or a hold just ahead of it) replaced is the one the client was -// current as of, so the send is admitted at the fence the restart published. +// readable for that, never given a child. Otherwise each resend of a message whose child died at +// startup would spawn another child that dies the same way. // // A child that has not proven its start is still the owner: the send is admitted against it and // the adapter holds the message until startup lands, or rejects it with the child's own reason @@ -128,7 +123,7 @@ export async function prepareStructuredAgentSessionSend( if (!context.sessions.has(sessionId)) { await context.restoreReadable(sessionId) } - return { ok: true, envelope } + return { ok: true } } if (structuredAgentSessionSendNeedsOwner(context.sessions.get(sessionId), record)) { const refusal = await restartOwnerForSend(context, envelope, record) @@ -136,7 +131,7 @@ export async function prepareStructuredAgentSessionSend( return { ok: false, refusal } } } - return { ok: true, envelope: admitAtResumedFence(context.sessions.get(sessionId), envelope) } + return { ok: true } } /** One restart attempt. Answers with the refusal that ends the send, or null when the send goes @@ -229,16 +224,3 @@ async function recordFailedRestart( context.deps.onEventSinkError?.({ sessionId, error }) } } - -/** A writer current as of the owner this child replaced is current now: the restart was the only - * thing that moved the fence, whether this send ran it or one just ahead of it did. */ -function admitAtResumedFence( - session: StructuredAgentSessionHostSession | undefined, - envelope: AgentSessionMutationEnvelope -): AgentSessionMutationEnvelope { - return session?.hasProviderChild && - session.resumedFromFence !== undefined && - envelope.expectedRuntimeFence === session.resumedFromFence - ? { ...envelope, expectedRuntimeFence: session.fence } - : envelope -} 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 f2d1ab12230..638d8a8679f 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 @@ -446,24 +446,7 @@ describe('send', () => { expect(journal.submissions()).toHaveLength(1) }) - it('refuses a stale fence and hands back the current one', async () => { - const record = await attach() - const body = hostTestMessage('add a retry') - const result = await host.send(CALLER, { - envelope: envelope( - 'agentSession.send', - { body }, - { expectedRuntimeFence: (record?.lease.runtimeFence ?? 1) + 5 } - ), - body - }) - expect(result).toMatchObject({ - ok: false, - refusal: { code: 'agent_session_checkpoint_stale', currentFence: record?.lease.runtimeFence } - }) - }) - - it('reuses a pending send admission after the client refreshes its fence', async () => { + it('admits a send fenced to another generation, and replays it by id once the fence catches up', async () => { const record = await attach() const body = hostTestMessage('add a retry') const params = { @@ -475,26 +458,14 @@ describe('send', () => { body } expect(await host.send(CALLER, params)).toMatchObject({ - ok: false, - refusal: { code: 'agent_session_checkpoint_stale' } - }) - expect( - store - .listOperationRows() - .filter((row) => row.operationId === params.envelope.clientOperationId) - ).toEqual([]) - const retry = { - ...params, - envelope: { - ...params.envelope, - expectedRuntimeFence: record?.lease.runtimeFence ?? 1 - } - } - expect(await host.send(CALLER, retry)).toMatchObject({ ok: true, replayed: false, value: { submission: { dispatchState: 'accepted' } } }) + const retry = { + ...params, + envelope: { ...params.envelope, expectedRuntimeFence: record?.lease.runtimeFence ?? 1 } + } expect(await host.send(CALLER, retry)).toMatchObject({ ok: true, replayed: true }) expect(dispatch).toHaveBeenCalledTimes(1) }) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.test.ts index 82f8c551642..a5a588559cd 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.test.ts @@ -46,7 +46,7 @@ describe('AgentSessionSubscribers', () => { }, journalDir: join(root, 'activity-churn-journal') }) - const subscribers = new AgentSessionSubscribers() + const subscribers = new AgentSessionSubscribers({ readFence: () => 1 }) for (let index = 0; index < MAX_RETAINED_SESSION_ACTIVITIES + 4; index += 1) { subscribers.publish(`session-${index}`, journal, { turnId: `turn-${index}`, text: 'working' }) @@ -68,11 +68,10 @@ describe('AgentSessionSubscribers', () => { }) const events: AgentSessionSubscribeEvent[] = [] - new AgentSessionSubscribers().open({ + new AgentSessionSubscribers({ readFence: () => 7 }).open({ id: 'subscriber-1', sessionId: SESSION, journal, - fence: 7, cursor: journal.cursor(), emit: (event) => events.push(event) }) @@ -107,20 +106,20 @@ describe('AgentSessionSubscribers', () => { }) let now = 1_000 const events: AgentSessionSubscribeEvent[] = [] - const subscribers = new AgentSessionSubscribers({ now: () => (now += 1) }) + const subscribers = new AgentSessionSubscribers({ readFence: () => 1, now: () => (now += 1) }) const emit = (event: AgentSessionSubscribeEvent): void => { events.push(event) } - subscribers.open({ id: 'one', sessionId: SESSION, journal, fence: 1, emit }) - subscribers.open({ id: 'two', sessionId: SESSION, journal, fence: 1, emit }) + subscribers.open({ id: 'one', sessionId: SESSION, journal, emit }) + subscribers.open({ id: 'two', sessionId: SESSION, journal, emit }) await journal.appendItem( { provider: 'orca', clientMessageId: 'clocked' }, { kind: 'status', text: 'Clocked' }, { fence: 1 } ) subscribers.publish(SESSION, journal) - subscribers.backgroundTasks(SESSION, null, 1) - subscribers.reset(SESSION, journal, 'epoch_changed', 1) + subscribers.backgroundTasks(SESSION, null) + subscribers.reset(SESSION, journal, 'epoch_changed') expect(events.map((event) => ('hostNow' in event ? event.hostNow : null))).toEqual([ 1_001, 1_002, @@ -152,12 +151,14 @@ describe('AgentSessionSubscribers', () => { }) let commands = [{ name: 'first', kind: 'skill' as const }] const events: AgentSessionSubscribeEvent[] = [] - const subscribers = new AgentSessionSubscribers({ readCommands: () => commands }) + const subscribers = new AgentSessionSubscribers({ + readFence: () => 7, + readCommands: () => commands + }) subscribers.open({ id: 'one', sessionId: SESSION, journal, - fence: 7, emit: (event) => events.push(event) }) expect(events[0]).toMatchObject({ type: 'snapshot', commands }) @@ -176,7 +177,6 @@ describe('AgentSessionSubscribers', () => { sessionId: SESSION, journal, cursor: journal.cursor(), - fence: 7, emit: (event) => events.push(event) }) expect(events[2]).toMatchObject({ type: 'batch', commands }) @@ -195,6 +195,7 @@ describe('AgentSessionSubscribers', () => { }) const published: string[] = [] const subscribers = new AgentSessionSubscribers({ + readFence: () => 1, onJournalPublished: (sessionId, published_journal) => { expect(published_journal).toBe(journal) published.push(sessionId) @@ -202,8 +203,8 @@ describe('AgentSessionSubscribers', () => { }) subscribers.publish(SESSION, journal) - subscribers.reset(SESSION, journal, 'epoch_changed', 1) - subscribers.snapshot(SESSION, journal, 1) + subscribers.reset(SESSION, journal, 'epoch_changed') + subscribers.snapshot(SESSION, journal) expect(published).toEqual([SESSION, SESSION, SESSION]) }) @@ -244,6 +245,7 @@ describe('AgentSessionSubscribers', () => { now: () => 1_000 }) const subscribers = new AgentSessionSubscribers({ + readFence: () => 1, onJournalPublished: (sessionId, published) => statusFeed.publish(sessionId, published) }) const statuses: AgentSessionStatusEvent[] = [] @@ -276,7 +278,7 @@ describe('AgentSessionSubscribers', () => { }) }) - it('publishes background lifecycle without advancing the journal and carries its fence forward', async () => { + it('publishes background lifecycle without advancing the journal', async () => { const journal = await journals.open({ identity: { sessionId: SESSION, @@ -287,13 +289,12 @@ describe('AgentSessionSubscribers', () => { }, journalDir: join(root, 'background-journal') }) - const subscribers = new AgentSessionSubscribers() + const subscribers = new AgentSessionSubscribers({ readFence: () => 1 }) const events: AgentSessionSubscribeEvent[] = [] subscribers.open({ id: 'subscriber-1', sessionId: SESSION, journal, - fence: 1, backgroundTasks: null, emit: (event) => events.push(event) }) @@ -303,26 +304,51 @@ describe('AgentSessionSubscribers', () => { state: 'monitoring' as const, tasks: [{ id: 'task-1', kind: 'command' as const, description: 'run the build' }] } - subscribers.backgroundTasks(SESSION, backgroundTasks, 2) + subscribers.backgroundTasks(SESSION, backgroundTasks) expect(journal.cursor()).toEqual(cursor) expect(events.at(-1)).toEqual({ type: 'batch', sessionId: SESSION, batch: { cursor, items: [], removedItemIds: [], submissions: [] }, - fence: 2, + fence: 1, hostNow: expect.any(Number), backgroundTasks }) + }) + it('stamps a fence that moved without a replay on the next frame', async () => { + const journal = await journals.open({ + identity: { + sessionId: SESSION, + workspaceId: 'workspace-1', + hostId: 'local', + agent: 'codex', + providerHandle: { kind: 'codex', threadId: 'thread-1' } + }, + journalDir: join(root, 'moved-fence-journal') + }) + let fence = 1 + const subscribers = new AgentSessionSubscribers({ readFence: () => fence }) + const events: AgentSessionSubscribeEvent[] = [] + subscribers.open({ + id: 'subscriber-1', + sessionId: SESSION, + journal, + emit: (event) => events.push(event) + }) + expect(events.at(-1)).toMatchObject({ type: 'snapshot', fence: 1 }) + + // An idle release and the resume after it move the lease with no snapshot in between. + fence = 3 await journal.appendItem( - { provider: 'orca', clientMessageId: 'after-background-fence' }, - { kind: 'status', text: 'After background state' }, - { fence: 2 } + { provider: 'orca', clientMessageId: 'after-release' }, + { kind: 'status', text: 'After the release' }, + { fence: 3 } ) subscribers.publish(SESSION, journal) - expect(events.at(-1)).toMatchObject({ type: 'batch', fence: 2 }) + expect(events.at(-1)).toMatchObject({ type: 'batch', fence: 3 }) }) it('publishes latest turn activity without advancing or adding journal rows', async () => { @@ -336,13 +362,12 @@ describe('AgentSessionSubscribers', () => { }, journalDir: join(root, 'activity-journal') }) - const subscribers = new AgentSessionSubscribers() + const subscribers = new AgentSessionSubscribers({ readFence: () => 1 }) const events: AgentSessionSubscribeEvent[] = [] subscribers.open({ id: 'subscriber-1', sessionId: SESSION, journal, - fence: 1, emit: (event) => events.push(event) }) const cursor = journal.cursor() @@ -368,7 +393,6 @@ describe('AgentSessionSubscribers', () => { id: 'reconnected', sessionId: SESSION, journal, - fence: 1, cursor, emit: (event) => events.push(event) }) @@ -440,13 +464,12 @@ describe('AgentSessionSubscribers', () => { journalDir }) - const subscribers = new AgentSessionSubscribers() + const subscribers = new AgentSessionSubscribers({ readFence: () => 1 }) const events: AgentSessionSubscribeEvent[] = [] subscribers.open({ id: 'subscriber-1', sessionId: SESSION, journal, - fence: 1, cursor: resumeCursor, emit: (event) => events.push(event) }) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.ts index af5e231730c..4c188d8c72d 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-subscribers.ts @@ -3,6 +3,9 @@ // Each subscriber advances independently: a client that connected two epochs // ago gets a reset while a caught-up one gets a batch from the same publish. // Nothing raw reaches a subscriber — every event carries reducer output. +// +// Every frame carries the fence read when it is sent. A per-subscriber copy went +// stale whenever the lease moved without a replay, such as an idle release. import type { AgentJournalCursor, @@ -36,11 +39,12 @@ type Subscriber = { sessionId: string emit: AgentSessionSubscriberEmit cursor: AgentJournalCursor - fence: number commands?: AgentSessionSlashCommand[] | null } export type AgentSessionSubscribersHooks = { + /** The session's current runtime fence, stamped on every frame for clients that still read it. */ + readFence: (sessionId: string) => number readCommands?: (sessionId: string) => AgentSessionSlashCommand[] | undefined /** Fires after publications that can change journal content. */ onJournalPublished?: (sessionId: string, journal: AgentSessionJournal) => void @@ -51,7 +55,7 @@ export class AgentSessionSubscribers { private readonly bySession = new Map>() private readonly activityBySession = new Map() - constructor(private readonly hooks: AgentSessionSubscribersHooks = {}) {} + constructor(private readonly hooks: AgentSessionSubscribersHooks) {} get retainedActivityCountForTests(): number { return this.activityBySession.size @@ -61,7 +65,6 @@ export class AgentSessionSubscribers { id: string sessionId: string journal: AgentSessionJournal - fence: number emit: AgentSessionSubscriberEmit cursor?: AgentJournalCursor backgroundTasks?: AgentSessionBackgroundTaskState | null @@ -71,8 +74,7 @@ export class AgentSessionSubscribers { id: input.id, sessionId: input.sessionId, emit: input.emit, - cursor: input.cursor ?? { epoch: liveCursor.epoch, sequence: 0 }, - fence: input.fence + cursor: input.cursor ?? { epoch: liveCursor.epoch, sequence: 0 } } const session = this.bySession.get(input.sessionId) ?? new Map() session.set(input.id, subscriber) @@ -82,12 +84,13 @@ export class AgentSessionSubscribers { if (input.cursor) { this.deliver(subscriber, input.journal, hostNow, true, input.backgroundTasks) } else { - const page = readAgentSessionHydrationPage(input.journal, input.fence) + const fence = this.hooks.readFence(input.sessionId) + const page = readAgentSessionHydrationPage(input.journal, fence) this.emit(subscriber, { type: 'snapshot', sessionId: input.sessionId, page, - fence: input.fence, + fence, hostNow, ...(input.backgroundTasks !== undefined ? { backgroundTasks: input.backgroundTasks } : {}), ...this.activityField(input.sessionId) @@ -138,28 +141,26 @@ export class AgentSessionSubscribers { sessionId: string, journal: AgentSessionJournal, reason: AgentJournalResetReason, - fence: number, backgroundTasks?: AgentSessionBackgroundTaskState | null ): void { - this.replay(sessionId, journal, fence, backgroundTasks, { type: 'reset', reset: reason }) + this.replay(sessionId, journal, backgroundTasks, { type: 'reset', reset: reason }) } snapshot( sessionId: string, journal: AgentSessionJournal, - fence: number, backgroundTasks?: AgentSessionBackgroundTaskState | null ): void { - this.replay(sessionId, journal, fence, backgroundTasks, { type: 'snapshot' }) + this.replay(sessionId, journal, backgroundTasks, { type: 'snapshot' }) } private replay( sessionId: string, journal: AgentSessionJournal, - fence: number, backgroundTasks: AgentSessionBackgroundTaskState | null | undefined, frame: { type: 'snapshot' } | { type: 'reset'; reset: AgentJournalResetReason } ): void { + const fence = this.hooks.readFence(sessionId) const page = readAgentSessionHydrationPage(journal, fence) const hostNow = this.now() for (const subscriber of this.subscribers(sessionId)) { @@ -173,17 +174,13 @@ export class AgentSessionSubscribers { ...this.activityField(sessionId) }) subscriber.cursor = page.liveCursor ?? page.window.nextCursor - subscriber.fence = fence } this.hooks.onJournalPublished?.(sessionId, journal) } - backgroundTasks( - sessionId: string, - state: AgentSessionBackgroundTaskState | null, - fence: number - ): void { + backgroundTasks(sessionId: string, state: AgentSessionBackgroundTaskState | null): void { const hostNow = this.now() + const fence = this.hooks.readFence(sessionId) for (const subscriber of this.subscribers(sessionId)) { this.emit(subscriber, { type: 'batch', @@ -193,7 +190,6 @@ export class AgentSessionSubscribers { backgroundTasks: state, hostNow }) - subscriber.fence = fence } } @@ -213,6 +209,7 @@ export class AgentSessionSubscribers { ? this.activityField(subscriber.sessionId).activity : undefined const publishedActivity = activity !== undefined ? activity : checkpointActivity + const fence = this.hooks.readFence(subscriber.sessionId) const readPage = createAgentSessionCatchUpReader(journal) while (true) { const result = readPage({ @@ -222,13 +219,13 @@ export class AgentSessionSubscribers { limit: AGENT_SESSION_HISTORY_MAX_LIMIT }) if (!result.ok) { - const page = { ...result.page, fence: subscriber.fence } + const page = { ...result.page, fence } this.emit(subscriber, { type: 'reset', sessionId: subscriber.sessionId, reset: result.reset, page, - fence: subscriber.fence, + fence, hostNow, ...(backgroundTasks !== undefined ? { backgroundTasks } : {}), ...(publishedActivity !== undefined ? { activity: publishedActivity } : {}) @@ -247,7 +244,7 @@ export class AgentSessionSubscribers { type: 'batch', sessionId: subscriber.sessionId, batch: emptyAgentSessionBatch(page.window.nextCursor), - fence: subscriber.fence, + fence, hostNow, ...(backgroundTasks !== undefined ? { backgroundTasks } : {}), ...(publishedActivity !== undefined ? { activity: publishedActivity } : {}) @@ -264,7 +261,7 @@ export class AgentSessionSubscribers { removedItemIds: page.removedItemIds, submissions: page.submissions }, - fence: subscriber.fence, + fence, hostNow, ...(backgroundTasks !== undefined ? { backgroundTasks } : {}), ...(publishedActivity !== undefined ? { activity: publishedActivity } : {}) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-wire-admission.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-wire-admission.test.ts index 585cbd4fbbd..82a7d683a18 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-wire-admission.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-wire-admission.test.ts @@ -54,21 +54,22 @@ describe('structured agent-session outbound admission', () => { REMOTE_RUNTIME_MAX_OUTBOUND_JSON_BYTES ) - const subscribers = new AgentSessionSubscribers() + let fence = 1 + const subscribers = new AgentSessionSubscribers({ readFence: () => fence }) const initial: AgentSessionSubscribeEvent[] = [] const dispose = subscribers.open({ id: 'initial', sessionId: SESSION, journal, - fence: 1, emit: (event) => initial.push(event) }) expect(initial).toHaveLength(1) expect(initial[0]).toMatchObject({ type: 'snapshot', page: { hasOlder: true } }) expectAdmitted(initial[0]) - subscribers.backgroundTasks(SESSION, null, 2) - subscribers.snapshot(SESSION, journal, 2) + fence = 2 + subscribers.backgroundTasks(SESSION, null) + subscribers.snapshot(SESSION, journal) expect(initial.slice(1)).toHaveLength(2) initial.slice(1).forEach(expectAdmitted) @@ -77,7 +78,6 @@ describe('structured agent-session outbound admission', () => { id: 'old-epoch', sessionId: SESSION, journal, - fence: 2, cursor: { epoch: 'retired-epoch', sequence: 1 }, emit: (event) => epochReset.push(event) }) @@ -109,7 +109,7 @@ describe('structured agent-session outbound admission', () => { }) it('splits valid-cursor catch-up into admitted batch frames', () => { - const subscribers = new AgentSessionSubscribers() + const subscribers = new AgentSessionSubscribers({ readFence: () => 2 }) const catchup: AgentSessionSubscribeEvent[] = [] const firstItemSequence = journal.snapshot().items[0]!.sequence @@ -117,7 +117,6 @@ describe('structured agent-session outbound admission', () => { id: 'catchup', sessionId: SESSION, journal, - fence: 2, cursor: { epoch: journal.epoch, sequence: firstItemSequence }, emit: (event) => catchup.push(event) }) diff --git a/src/main/native-chat/agent-session-wire/structured-conversation-command.test.ts b/src/main/native-chat/agent-session-wire/structured-conversation-command.test.ts index a261c6dd151..37d8d4ee83d 100644 --- a/src/main/native-chat/agent-session-wire/structured-conversation-command.test.ts +++ b/src/main/native-chat/agent-session-wire/structured-conversation-command.test.ts @@ -245,14 +245,11 @@ describe('host conversation commands', () => { }) }) - it('rejects stale fences before provider execution', async () => { + it('runs a command whose fence the client has not caught up to', async () => { const params = commandParams('compact') params.envelope.expectedRuntimeFence++ - expect(await host.conversationCommand(caller, params)).toMatchObject({ - ok: false, - refusal: { code: 'agent_session_checkpoint_stale' } - }) - expect(compact).not.toHaveBeenCalled() + expect(await host.conversationCommand(caller, params)).toMatchObject({ ok: true }) + expect(compact).toHaveBeenCalledTimes(1) }) it('allows cancellation while compaction is awaiting completion and refuses a second client', async () => { let finish!: (value: {}) => void diff --git a/src/shared/agent-session-mutation-envelope.test.ts b/src/shared/agent-session-mutation-envelope.test.ts index e929d628de5..306450f8a63 100644 --- a/src/shared/agent-session-mutation-envelope.test.ts +++ b/src/shared/agent-session-mutation-envelope.test.ts @@ -109,16 +109,15 @@ describe('admitAgentSessionMutation', () => { expect(admission.decision).toBe('replay') }) - it('refuses a stale fence and hands back the current one', () => { - const admission = admitAgentSessionMutation({ - ...base, - envelope: envelope({ expectedRuntimeFence: 3 }), - ledger: ADMIT('f'.repeat(64)) - }) - expect(admission).toMatchObject({ - decision: 'refused', - refusal: { code: 'agent_session_checkpoint_stale', currentFence: 4 } - }) + it('admits a writer whatever fence the client last saw', () => { + for (const expectedRuntimeFence of [3, 5, null]) { + const admission = admitAgentSessionMutation({ + ...base, + envelope: envelope({ expectedRuntimeFence }), + ledger: ADMIT('f'.repeat(64)) + }) + expect(admission, String(expectedRuntimeFence)).toMatchObject({ decision: 'admit' }) + } }) it('refuses a writer while the lease is unreconciled', () => { diff --git a/src/shared/agent-session-mutation-envelope.ts b/src/shared/agent-session-mutation-envelope.ts index 83fdb65e9ae..fe439c20b08 100644 --- a/src/shared/agent-session-mutation-envelope.ts +++ b/src/shared/agent-session-mutation-envelope.ts @@ -11,10 +11,7 @@ import type { AgentSessionOperationDecision, AgentSessionOperationRow } from './agent-session-operation-ledger' -import { - agentSessionLeaseAdmitsWriter, - isAgentSessionFenceCurrent -} from './agent-session-lease-adjudication' +import { agentSessionLeaseAdmitsWriter } from './agent-session-lease-adjudication' import type { AgentSessionLease } from './agent-session-record' import type { AgentSessionMutationEnvelope, AgentSessionWireRefusal } from './agent-session-wire' @@ -81,10 +78,11 @@ export type AgentSessionMutationAdmission = /** * Fixed order: fingerprint agreement, then the ledger (so a retry replays - * before anything else can refuse it), then the lease, then the fence. Putting - * the ledger ahead of the fence is deliberate — a retry that crossed an owner - * change must still return its recorded answer instead of a stale-checkpoint - * refusal the client would then resend as a second effect. + * before anything else can refuse it), then the lease. + * + * `expectedRuntimeFence` is not checked: each write names its own target (a + * turn, an item revision, an epoch) or is last-writer-wins, so an owner restart + * the client has not seen yet refuses nothing. Older hosts still check it. */ export function admitAgentSessionMutation(input: { envelope: AgentSessionMutationEnvelope @@ -115,19 +113,6 @@ export function admitAgentSessionMutation(input: { if (leaseRefusal) { return { decision: 'refused', refusal: leaseRefusal } } - if ( - envelope.expectedRuntimeFence === null || - !isAgentSessionFenceCurrent(lease, envelope.expectedRuntimeFence) - ) { - return { - decision: 'refused', - refusal: { - code: 'agent_session_checkpoint_stale', - message: `Expected runtime fence ${envelope.expectedRuntimeFence ?? 'none'}; the session is at ${lease.runtimeFence}.`, - currentFence: lease.runtimeFence - } - } - } return { decision: 'admit', row: ledger.row } } diff --git a/tests/e2e/cross-version-wire/cross-version-agent-session-wire.unit.test.ts b/tests/e2e/cross-version-wire/cross-version-agent-session-wire.unit.test.ts index d244fba1ba6..4d18e81ced6 100644 --- a/tests/e2e/cross-version-wire/cross-version-agent-session-wire.unit.test.ts +++ b/tests/e2e/cross-version-wire/cross-version-agent-session-wire.unit.test.ts @@ -726,15 +726,17 @@ describe('cross-version structured agent sessions', () => { expect(batch?.cursor.sequence).toBeGreaterThan(held.sequence) }) - it('refuses a write still fenced to the host generation that died', async () => { + // Every released client still sends the fence it last saw; this host names a write by its + // target and ignores that fence. Only the attach keeps comparing one, which `reattach` pins. + it('delivers a write still fenced to the host generation that died', async () => { const created = await answer('agentSession.create', createIntentParams()) await bootHost('b') const reattached = await reattach(created.fence) expect(reattached.fence).toBeGreaterThan(created.fence) expect(await answer('agentSession.send', sendParams('stale', created.fence))).toMatchObject({ - ok: false, - refusal: { code: 'agent_session_checkpoint_stale' } + ok: true, + fence: reattached.fence }) }) })