diff --git a/apps/desktop/src/main/__tests__/app-shell-session-ui-state.test.ts b/apps/desktop/src/main/__tests__/app-shell-session-ui-state.test.ts index 3b7f2b6b0f..a465dfc847 100644 --- a/apps/desktop/src/main/__tests__/app-shell-session-ui-state.test.ts +++ b/apps/desktop/src/main/__tests__/app-shell-session-ui-state.test.ts @@ -373,6 +373,85 @@ describe('shellSessionRowEqual', () => { assert.equal(shellSessionRowEqual(future, row), false); }); + it('orders same-revision rows by the live run epoch (#5713)', async () => { + const { root } = installReactRenderer(); + try { + const catalog = createSessionCatalogController(); + catalog.commitSessions([{ + ...row, + revision: 5, + runningTurnIds: ['turn-1'], + runHostGeneration: 'host-1', + runEpoch: 2, + }]); + + // A read taken before the turn started lands after the running patch: + // same revision, same host generation, older epoch — it must not flip + // the row back to idle. + await act(async () => { + catalog.commitPatch(row.id, { + ...row, + revision: 5, + runningTurnIds: [], + runHostGeneration: 'host-1', + runEpoch: 1, + }); + }); + assert.deepEqual( + selectSessionById(catalog.getState(), row.id)?.runningTurnIds, + ['turn-1'], + 'the older live state must not overwrite the newer', + ); + + // A genuinely newer epoch updates the row even at the same revision. + await act(async () => { + catalog.commitPatch(row.id, { + ...row, + revision: 5, + runningTurnIds: [], + runHostGeneration: 'host-1', + runEpoch: 3, + }); + }); + assert.deepEqual(selectSessionById(catalog.getState(), row.id)?.runningTurnIds, []); + + // A Host restart is a new generation: the previous host is gone, so + // its row cannot out-rank the restarted host's first read, whatever + // each side's epoch counter reads — a wall clock is not monotonic + // across processes (#5713 review round two). + await act(async () => { + catalog.commitPatch(row.id, { + ...row, + revision: 5, + runningTurnIds: ['turn-2'], + runHostGeneration: 'host-2', + runEpoch: 1, + }); + }); + assert.deepEqual( + selectSessionById(catalog.getState(), row.id)?.runningTurnIds, + ['turn-2'], + 'the restarted host must take over the row', + ); + + // Within the restarted generation the counter orders reads again. + await act(async () => { + catalog.commitPatch(row.id, { + ...row, + revision: 5, + runningTurnIds: ['turn-2'], + runHostGeneration: 'host-2', + runEpoch: 0, + }); + }); + assert.deepEqual( + selectSessionById(catalog.getState(), row.id)?.runningTurnIds, + ['turn-2'], + 'the older read of the restarted generation must not win', + ); + } finally { cleanupFakeDom(); } + }); + it('keeps a catalog row subscriber mounted through rail-only patches', async () => { const { root } = installReactRenderer(); try { diff --git a/apps/desktop/src/renderer/application/contracts/session-catalog/session-catalog-state.ts b/apps/desktop/src/renderer/application/contracts/session-catalog/session-catalog-state.ts index 7a55a9f64b..0ef32a1e32 100644 --- a/apps/desktop/src/renderer/application/contracts/session-catalog/session-catalog-state.ts +++ b/apps/desktop/src/renderer/application/contracts/session-catalog/session-catalog-state.ts @@ -188,9 +188,39 @@ export function waitForCatalogSession( }); } -/** A committed row at a newer revision is authoritative over an older snapshot of it. */ +/** + * A committed row at a newer revision is authoritative over an older snapshot + * of it. Equal revisions tie on the live run state's own order: a turn + * starting or ending does not move `revision`, so two same-revision reads can + * disagree about `runningTurnIds` — the run epoch says which observation is + * older, and the stale one must not overwrite the fresher (#5713). + * + * The epoch counter only orders observations of one Host generation. + * Generations themselves are not ordered, so a read from a different + * generation is never stale: a restarted Host must take the row over from its + * predecessor whatever the two counters read (#5713 review). A successful + * cross-generation response cannot exist on the wire, either: closing a + * connection rejects every in-flight request with `connection_lost` + * (client/connection.ts), so a lagging predecessor read never delivers after + * the successor's row has landed. + */ function isStaleSummary(prior: DesktopSessionSummary, next: DesktopSessionSummary): boolean { - return prior.revision > next.revision; + if (prior.revision !== next.revision) return prior.revision > next.revision; + const priorGeneration = prior.runHostGeneration; + const nextGeneration = next.runHostGeneration; + if ( + priorGeneration !== undefined && + nextGeneration !== undefined && + priorGeneration !== nextGeneration + ) { + return false; + } + const priorEpoch = prior.runEpoch; + const nextEpoch = next.runEpoch; + if (priorEpoch === undefined || nextEpoch === undefined || priorEpoch === nextEpoch) { + return false; + } + return priorEpoch > nextEpoch; } export const selectSessions = (state: SessionCatalogState): readonly DesktopSessionSummary[] => diff --git a/packages/core/src/session.ts b/packages/core/src/session.ts index 9ae0ac39bd..993f02924e 100644 --- a/packages/core/src/session.ts +++ b/packages/core/src/session.ts @@ -403,6 +403,26 @@ export interface SessionSummary { * the header alone and omits it. */ runningTurnIds?: string[]; + /** + * Bumped by the runtime each time a turn of this session starts or ends. + * `revision` does not move for those transitions, so two same-revision + * summaries can disagree about `runningTurnIds` — the epoch orders them: + * the higher epoch is the newer observation (#5713). Present alongside + * `runningTurnIds` under the same population rules. + * + * The counter restarts at zero with a fresh Host process, so it only orders + * observations of one host generation: summaries whose `runHostGeneration` + * differs are not comparable by epoch, and the newer generation's host owns + * the row outright. + */ + runEpoch?: number; + /** + * Identifies the Host process generation that produced this live-run + * observation. Summaries from different generations are not ordered by + * `runEpoch` — a restarted Host supersedes every observation its + * predecessor published, whatever the epoch counters read (#5713). + */ + runHostGeneration?: string; parentSessionId?: string; branchOfTurnId?: string; subagent?: SessionSubagentProjection; diff --git a/packages/runtime-host/src/__tests__/authenticated-websocket.test.ts b/packages/runtime-host/src/__tests__/authenticated-websocket.test.ts index b3a96a6a4f..c58e887866 100644 --- a/packages/runtime-host/src/__tests__/authenticated-websocket.test.ts +++ b/packages/runtime-host/src/__tests__/authenticated-websocket.test.ts @@ -58,6 +58,7 @@ const PROTOCOL = { const KNOWN_EMPTY_LIVE_RUN_STATE = { schemaVersion: SESSION_CATALOG_LIVE_RUN_STATE_SCHEMA_VERSION, runningTurnIds: [], + runEpoch: 0, } as const; async function configureTestModel(local: RuntimeHostConnection): Promise { @@ -309,15 +310,25 @@ test('one Local IPC owner and one authenticated WebSocket Client control the sam if (closed.value?.kind === 'subscription.closed') { assert.equal(closed.value.reason, 'access_revoked'); } + const sharedRead = await remote.request('session.catalog.query', { + kind: 'get', + sessionId: 'shared-session', + }); + assert.equal(sharedRead.kind, 'session'); + const sharedSession = sharedRead.kind === 'session' ? sharedRead.session : null; + assert.ok(sharedSession && !('kind' in sharedSession)); + if ('kind' in sharedSession) assert.fail('the shared row must be a current projection'); + // The host generation is unique per Host process, so assert its shape and + // compare the known-empty remainder. + const { hostGeneration: sharedGeneration, ...sharedKnownEmpty } = + sharedSession.liveRunState ?? {}; + if (typeof sharedGeneration !== 'string' || sharedGeneration.length === 0) { + assert.fail('the host generation must be a non-empty string'); + } + assert.deepEqual(sharedKnownEmpty, KNOWN_EMPTY_LIVE_RUN_STATE); assert.deepEqual( - await remote.request('session.catalog.query', { - kind: 'get', - sessionId: 'shared-session', - }), - { - kind: 'session', - session: { ...created, liveRunState: KNOWN_EMPTY_LIVE_RUN_STATE }, - }, + { ...sharedSession, liveRunState: KNOWN_EMPTY_LIVE_RUN_STATE }, + { ...created, liveRunState: KNOWN_EMPTY_LIVE_RUN_STATE }, ); const catalogChanged = new Promise((resolve) => { @@ -330,16 +341,23 @@ test('one Local IPC owner and one authenticated WebSocket Client control the sam }); assert.equal(renamed.kind, 'committed'); assert.equal(await catalogChanged, 'shared-session'); + const localRead = await local.request('session.catalog.query', { + kind: 'get', + sessionId: 'shared-session', + }); + assert.equal(localRead.kind, 'session'); + const localSession = localRead.kind === 'session' ? localRead.session : null; + assert.ok(localSession && !('kind' in localSession)); + if ('kind' in localSession) assert.fail('the local row must be a current projection'); + const { hostGeneration: localGeneration, ...localKnownEmpty } = localSession.liveRunState ?? {}; + if (typeof localGeneration !== 'string' || localGeneration.length === 0) { + assert.fail('the host generation must be a non-empty string'); + } + assert.deepEqual(localKnownEmpty, KNOWN_EMPTY_LIVE_RUN_STATE); assert.deepEqual( - await local.request('session.catalog.query', { - kind: 'get', - sessionId: 'shared-session', - }), + { ...localSession, liveRunState: KNOWN_EMPTY_LIVE_RUN_STATE }, renamed.kind === 'committed' - ? { - kind: 'session', - session: { ...renamed.session, liveRunState: KNOWN_EMPTY_LIVE_RUN_STATE }, - } + ? { ...renamed.session, liveRunState: KNOWN_EMPTY_LIVE_RUN_STATE } : assert.fail('Remote Session rename did not commit'), ); @@ -616,7 +634,29 @@ test('an authenticated WebSocket Client reconnects after service restart to cano websocket: { host: '127.0.0.1', port }, }); - assert.deepEqual(await recovered, expected); + const canonical = await recovered; + assert.equal(canonical.kind, 'session'); + assert.equal(expected.kind, 'session'); + const beforeSession = expected.kind === 'session' ? expected.session : null; + const afterSession = canonical.kind === 'session' ? canonical.session : null; + assert.ok(beforeSession && afterSession); + if ('kind' in beforeSession || 'kind' in afterSession) { + assert.fail('the canonical reads must be current projections'); + } + const beforeLive = beforeSession.liveRunState; + const afterLive = afterSession.liveRunState; + assert.ok(beforeLive && afterLive); + // The restart must change the host generation — the wire signal that lets + // clients tell the restarted Host's reads apart from its predecessor's + // (#5713). + assert.notEqual(afterLive.hostGeneration, beforeLive.hostGeneration); + assert.deepEqual( + { + ...afterSession, + liveRunState: { ...afterLive, hostGeneration: beforeLive.hostGeneration }, + }, + beforeSession, + ); assert.notEqual(remote.hostEpoch, firstHostEpoch); } finally { await Promise.allSettled([remote?.close(), local?.close()]); diff --git a/packages/runtime-host/src/__tests__/session-catalog-coordinator.test.ts b/packages/runtime-host/src/__tests__/session-catalog-coordinator.test.ts index 67ee6fb4e3..48516f4ea2 100644 --- a/packages/runtime-host/src/__tests__/session-catalog-coordinator.test.ts +++ b/packages/runtime-host/src/__tests__/session-catalog-coordinator.test.ts @@ -314,6 +314,8 @@ test('catalog queries project known-empty and running state from Runtime authori assert.deepEqual(emptyOutcome.result.session.liveRunState, { schemaVersion: 1, runningTurnIds: [], + runEpoch: 0, + hostGeneration: 'test-host-generation', }); runningTurnIds = ['turn-live']; @@ -332,6 +334,8 @@ test('catalog queries project known-empty and running state from Runtime authori assert.deepEqual(session.liveRunState, { schemaVersion: 1, runningTurnIds: ['turn-live'], + runEpoch: 0, + hostGeneration: 'test-host-generation', }); }); @@ -379,6 +383,8 @@ test('catalog queries de-duplicate Runtime live turn ids in stable order', async assert.deepEqual(outcome.result.session.liveRunState, { schemaVersion: 1, runningTurnIds: ['turn-a', 'turn-b'], + runEpoch: 0, + hostGeneration: 'test-host-generation', }); }); @@ -2202,6 +2208,8 @@ function createFixture( const runtimePolicy = options.runtimePolicy ?? runtimePolicyFixture(options.connection ?? {}); const manager: ConfigurationAuthority = { runningTurnIds: () => [], + sessionRunEpoch: () => 0, + sessionHostGeneration: () => 'test-host-generation', transitionSessionConfiguration: async (_sessionId, input) => { header = { ...header, diff --git a/packages/runtime-host/src/__tests__/session-catalog-protocol.test.ts b/packages/runtime-host/src/__tests__/session-catalog-protocol.test.ts index e2748485c6..1df2737117 100644 --- a/packages/runtime-host/src/__tests__/session-catalog-protocol.test.ts +++ b/packages/runtime-host/src/__tests__/session-catalog-protocol.test.ts @@ -91,6 +91,96 @@ describe('Session catalog protocol', () => { ); }); + test('decodes a Session attention payload on a catalog change', () => { + assert.deepEqual( + decodeHostFrame({ + kind: 'session.catalog.changed', + revision: 4, + sessionId: 'session-1', + attention: { + kind: 'errored', + eventId: 'terminal-1', + body: 'Provider request failed', + }, + }), + { + kind: 'session.catalog.changed', + revision: 4, + sessionId: 'session-1', + attention: { + kind: 'errored', + eventId: 'terminal-1', + body: 'Provider request failed', + }, + }, + ); + }); + + test('accepts an optional run epoch in the live run state and rejects a bad one', () => { + const withEpoch = { + ...projection(), + liveRunState: { schemaVersion: 1, runningTurnIds: ['turn-1'], runEpoch: 7 }, + }; + assert.deepEqual(decodeSessionCatalogItem(withEpoch), withEpoch); + + // A two-field live run state from a host that does not track the epoch + // still decodes, and stays two fields. + const withoutEpoch = { + ...projection(), + liveRunState: { schemaVersion: 1, runningTurnIds: ['turn-1'] }, + }; + assert.deepEqual(decodeSessionCatalogItem(withoutEpoch), withoutEpoch); + + assert.throws( + () => + decodeSessionCatalogItem({ + ...projection(), + liveRunState: { schemaVersion: 1, runningTurnIds: ['turn-1'], runEpoch: -1 }, + }), + isProtocolError, + ); + assert.throws( + () => + decodeSessionCatalogItem({ + ...projection(), + liveRunState: { schemaVersion: 1, runningTurnIds: ['turn-1'], runEpoch: 1.5 }, + }), + isProtocolError, + ); + }); + + test('accepts an optional host generation in the live run state and rejects a bad one', () => { + const withGeneration = { + ...projection(), + liveRunState: { + schemaVersion: 1, + runningTurnIds: ['turn-1'], + runEpoch: 3, + hostGeneration: 'host-gen-1', + }, + }; + assert.deepEqual(decodeSessionCatalogItem(withGeneration), withGeneration); + + // Hosts that do not track the generation keep decoding. + const withoutGeneration = { + ...projection(), + liveRunState: { schemaVersion: 1, runningTurnIds: ['turn-1'] }, + }; + assert.deepEqual(decodeSessionCatalogItem(withoutGeneration), withoutGeneration); + + for (const hostGeneration of [42, '', `x`.repeat(129), 'gen\u0000-1']) { + assert.throws( + () => + decodeSessionCatalogItem({ + ...projection(), + liveRunState: { schemaVersion: 1, runningTurnIds: ['turn-1'], hostGeneration }, + }), + isProtocolError, + `hostGeneration ${JSON.stringify(hostGeneration)} must be rejected`, + ); + } + }); + test('bounds the live running-turn collection explicitly', () => { const atLimit = Array.from( { length: SESSION_CATALOG_RUNNING_TURN_MAX_ITEMS }, @@ -751,28 +841,3 @@ test('executor configuration rejects ambiguous routes and malformed values', () isProtocolError, ); }); - -test('decodes a Session attention payload on a catalog change', () => { - assert.deepEqual( - decodeHostFrame({ - kind: 'session.catalog.changed', - revision: 4, - sessionId: 'session-1', - attention: { - kind: 'errored', - eventId: 'terminal-1', - body: 'Provider request failed', - }, - }), - { - kind: 'session.catalog.changed', - revision: 4, - sessionId: 'session-1', - attention: { - kind: 'errored', - eventId: 'terminal-1', - body: 'Provider request failed', - }, - }, - ); -}); diff --git a/packages/runtime-host/src/__tests__/session-catalog-two-client-uds.test.ts b/packages/runtime-host/src/__tests__/session-catalog-two-client-uds.test.ts index d53f65f79d..ee15a3b816 100644 --- a/packages/runtime-host/src/__tests__/session-catalog-two-client-uds.test.ts +++ b/packages/runtime-host/src/__tests__/session-catalog-two-client-uds.test.ts @@ -53,6 +53,7 @@ import { SESSION_CATALOG_LIVE_RUN_STATE_SCHEMA_VERSION, type ClientFrame, type SessionCatalogItem, + type SessionCatalogLiveRunState, type SessionCatalogProjection, type SessionCreateInput, type SubscriptionFrame, @@ -69,8 +70,25 @@ const WIRE_OVERSIZED_MODEL_ID = '😀'.repeat(256); const KNOWN_EMPTY_LIVE_RUN_STATE = { schemaVersion: SESSION_CATALOG_LIVE_RUN_STATE_SCHEMA_VERSION, runningTurnIds: [], + runEpoch: 0, } as const; +// The host generation is unique per Host process, so it cannot be spelled in +// advance: assert its shape, then compare the known-empty remainder. +function expectKnownEmptyLiveRunState(actual: SessionCatalogLiveRunState | undefined): void { + assert.ok(actual !== undefined, 'the live run state must be present'); + const { hostGeneration, ...knownEmpty } = actual; + if (typeof hostGeneration !== 'string' || hostGeneration.length === 0) { + assert.fail('the host generation must be a non-empty string'); + } + assert.deepEqual(knownEmpty, KNOWN_EMPTY_LIVE_RUN_STATE); +} + +function querySessionReconciled(summary: SessionCatalogProjection): SessionCatalogProjection { + expectKnownEmptyLiveRunState(summary.liveRunState); + return { ...summary, liveRunState: KNOWN_EMPTY_LIVE_RUN_STATE }; +} + test('two Clients share stable Session creation, CAS configuration, and catalog continuity', { skip: process.platform === 'win32' ? 'Windows SQLite shutdown lifecycle' : false, timeout: 120_000, @@ -303,7 +321,7 @@ test('two Clients share stable Session creation, CAS configuration, and catalog assert.fail('One Session configuration must commit'); } const configuredSession = requireSessionProjection(committedConfiguration.session); - assert.deepEqual(await querySession(desktop, created.id), { + assert.deepEqual(querySessionReconciled(await querySession(desktop, created.id)), { ...configuredSession, liveRunState: KNOWN_EMPTY_LIVE_RUN_STATE, }); @@ -359,7 +377,7 @@ test('two Clients share stable Session creation, CAS configuration, and catalog relocatedSession.workspace.hostCwd === (await realpath(firstCwd)) || relocatedSession.workspace.hostCwd === (await realpath(secondCwd)), ); - assert.deepEqual(await querySession(tui, narrowedSession.id), { + assert.deepEqual(querySessionReconciled(await querySession(tui, narrowedSession.id)), { ...relocatedSession, liveRunState: KNOWN_EMPTY_LIVE_RUN_STATE, }); @@ -397,7 +415,7 @@ test('two Clients share stable Session creation, CAS configuration, and catalog (error: unknown) => error instanceof RuntimeHostProtocolError && error.code === 'invalid_frame', ); - assert.deepEqual(await querySession(desktop, relocatedSession.id), { + assert.deepEqual(querySessionReconciled(await querySession(desktop, relocatedSession.id)), { ...relocatedSession, liveRunState: KNOWN_EMPTY_LIVE_RUN_STATE, }); @@ -786,7 +804,7 @@ test('stable Session creation survives response loss and Host restart', { try { const retried = requireSessionProjection(await retrying.request('session.create', input)); const { liveRunState, ...persistedCommitted } = committed; - assert.deepEqual(liveRunState, KNOWN_EMPTY_LIVE_RUN_STATE); + expectKnownEmptyLiveRunState(liveRunState); assert.deepEqual(retried, persistedCommitted); } finally { await retrying.close(); diff --git a/packages/runtime-host/src/client/session-catalog-summary.ts b/packages/runtime-host/src/client/session-catalog-summary.ts index 8eafc68e7f..13ba8cf4fb 100644 --- a/packages/runtime-host/src/client/session-catalog-summary.ts +++ b/packages/runtime-host/src/client/session-catalog-summary.ts @@ -44,7 +44,15 @@ export function projectSessionCatalogSummary( ...(session.statusUpdatedAt === undefined ? {} : { statusUpdatedAt: session.statusUpdatedAt }), ...(session.liveRunState === undefined ? {} - : { runningTurnIds: [...session.liveRunState.runningTurnIds] }), + : { + runningTurnIds: [...session.liveRunState.runningTurnIds], + ...(session.liveRunState.runEpoch === undefined + ? {} + : { runEpoch: session.liveRunState.runEpoch }), + ...(session.liveRunState.hostGeneration === undefined + ? {} + : { runHostGeneration: session.liveRunState.hostGeneration }), + }), ...(session.parentSessionId === undefined ? {} : { parentSessionId: session.parentSessionId }), ...(session.branchOfTurnId === undefined ? {} : { branchOfTurnId: session.branchOfTurnId }), ...(session.subagent === undefined ? {} : { subagent: session.subagent }), diff --git a/packages/runtime-host/src/protocol/index.ts b/packages/runtime-host/src/protocol/index.ts index 571152f867..96266f0a2e 100644 --- a/packages/runtime-host/src/protocol/index.ts +++ b/packages/runtime-host/src/protocol/index.ts @@ -103,7 +103,14 @@ export const RUNTIME_HOST_REGISTRATION_SCHEMA_VERSION = 1 as const; export const RUNTIME_HOST_PROTOCOL_VERSION = 0 as const; // Increment when the same protocol version no longer guarantees safe Client-Host // interoperability. Mismatches are rejected before domain commands are admitted. -export const RUNTIME_HOST_COMPATIBILITY_EPOCH = 192 as const; +export const RUNTIME_HOST_COMPATIBILITY_EPOCH = 195 as const; +// 195: Session catalog live run state may carry the host generation that +// produced it, so clients can tell same-revision reads of a restarted Host +// apart from its predecessor's instead of ordering them by a per-process +// epoch. Epoch-194 peers reject the unknown key. +// 194: Session catalog live run state may carry the runtime's run epoch, which +// lets clients order same-revision reads. Epoch-193 peers reject the unknown +// key, so a newer Desktop against an older Host loses session catalog reads. // 192: Message quotes carry an optional annotation written by the user, which // the model reads beside the excerpt. An epoch-191 peer rejects the field. // 191: Session catalog change frames can carry attention events. diff --git a/packages/runtime-host/src/protocol/session-catalog.ts b/packages/runtime-host/src/protocol/session-catalog.ts index 38987f3cc7..4d30ab35c5 100644 --- a/packages/runtime-host/src/protocol/session-catalog.ts +++ b/packages/runtime-host/src/protocol/session-catalog.ts @@ -71,6 +71,7 @@ export const SESSION_CATALOG_MODEL_MAX_BYTES = SESSION_MODEL_ID_MAX_BYTES; export const SESSION_CATALOG_CONNECTION_SLUG_MAX_BYTES = 256; export const SESSION_CATALOG_LIVE_RUN_STATE_SCHEMA_VERSION = 1 as const; export const SESSION_CATALOG_RUNNING_TURN_MAX_ITEMS = 64; +export const SESSION_CATALOG_HOST_GENERATION_MAX_CHARS = 128; const QUERY_ERRORS = [ 'host_not_ready', @@ -225,6 +226,23 @@ export interface SessionExecutionBoundaryQueryInput { export interface SessionCatalogLiveRunState { readonly schemaVersion: typeof SESSION_CATALOG_LIVE_RUN_STATE_SCHEMA_VERSION; readonly runningTurnIds: readonly string[]; + /** + * The runtime's own order for this live state, bumped on every turn start + * and end. `revision` does not move for those transitions, so two + * same-revision reads can disagree about `runningTurnIds` — the epoch says + * which read is older (#5713). Absent from hosts that do not track it. The + * counter is per-process; across a Host restart only `hostGeneration` + * orders observations, never the epoch. + */ + readonly runEpoch?: number; + /** + * Identifies the Host process generation that produced this live state. + * Rows survive a Host restart while the epoch counter restarts at zero, so + * clients must not order same-revision reads across generations by epoch — + * a restarted Host supersedes every observation its predecessor published. + * Absent from hosts that do not track it. + */ + readonly hostGeneration?: string; } export interface SessionCatalogProjection { @@ -1003,30 +1021,57 @@ function optionalLiveRunState( record: Record, ): Pick | Record { if (record.liveRunState === undefined) return {}; - const state = requireExactRecord(record.liveRunState, 'Session catalog live run state', [ + const liveRunState = requireRecord(record.liveRunState, 'Session catalog live run state'); + assertAllowedKeys(liveRunState, 'Session catalog live run state', [ 'schemaVersion', 'runningTurnIds', + 'runEpoch', + 'hostGeneration', ]); - if (state.schemaVersion !== SESSION_CATALOG_LIVE_RUN_STATE_SCHEMA_VERSION) { + if (liveRunState.schemaVersion !== SESSION_CATALOG_LIVE_RUN_STATE_SCHEMA_VERSION) { throw invalidProtocolFrame('Unsupported Session catalog live run state schema version'); } if ( - !Array.isArray(state.runningTurnIds) || - state.runningTurnIds.length > SESSION_CATALOG_RUNNING_TURN_MAX_ITEMS + !Object.hasOwn(liveRunState, 'runningTurnIds') || + !Array.isArray(liveRunState.runningTurnIds) || + liveRunState.runningTurnIds.length > SESSION_CATALOG_RUNNING_TURN_MAX_ITEMS ) { throw invalidProtocolFrame('Invalid Session catalog running turn ids'); } const runningTurnIds: string[] = []; - for (let index = 0; index < state.runningTurnIds.length; index += 1) { - runningTurnIds.push(requireEntityId(state.runningTurnIds[index], 'Session running turn id')); + for (let index = 0; index < liveRunState.runningTurnIds.length; index += 1) { + runningTurnIds.push( + requireEntityId(liveRunState.runningTurnIds[index], 'Session running turn id'), + ); } if (new Set(runningTurnIds).size !== runningTurnIds.length) { throw invalidProtocolFrame('Duplicate Session catalog running turn id'); } + if ( + liveRunState.runEpoch !== undefined && + (typeof liveRunState.runEpoch !== 'number' || + !Number.isSafeInteger(liveRunState.runEpoch) || + liveRunState.runEpoch < 0) + ) { + throw invalidProtocolFrame('Invalid Session catalog run epoch'); + } + if ( + liveRunState.hostGeneration !== undefined && + (typeof liveRunState.hostGeneration !== 'string' || + liveRunState.hostGeneration.length === 0 || + liveRunState.hostGeneration.length > SESSION_CATALOG_HOST_GENERATION_MAX_CHARS || + /[\u0000-\u001f\u007f]/.test(liveRunState.hostGeneration)) + ) { + throw invalidProtocolFrame('Invalid Session catalog host generation'); + } return { liveRunState: { schemaVersion: SESSION_CATALOG_LIVE_RUN_STATE_SCHEMA_VERSION, runningTurnIds, + ...(liveRunState.runEpoch === undefined ? {} : { runEpoch: liveRunState.runEpoch }), + ...(liveRunState.hostGeneration === undefined + ? {} + : { hostGeneration: liveRunState.hostGeneration }), }, }; } diff --git a/packages/runtime-host/src/server/session-catalog-coordinator.ts b/packages/runtime-host/src/server/session-catalog-coordinator.ts index 7287904034..e66a6bb6f4 100644 --- a/packages/runtime-host/src/server/session-catalog-coordinator.ts +++ b/packages/runtime-host/src/server/session-catalog-coordinator.ts @@ -138,7 +138,11 @@ type SessionRuntimePolicyStores = { type SessionConfigurationAuthority = Pick< SessionManager, - 'transitionSessionConfiguration' | 'relocateSessionWorkspace' | 'runningTurnIds' + | 'transitionSessionConfiguration' + | 'relocateSessionWorkspace' + | 'runningTurnIds' + | 'sessionRunEpoch' + | 'sessionHostGeneration' >; type SessionContinuity = Pick; @@ -516,7 +520,11 @@ export class HostSessionCatalogCoordinator { session: record ? projectSharedSessionCatalogRecord( record, - projectCatalogLiveRunState(this.#manager.runningTurnIds(record.header.id)), + projectCatalogLiveRunState( + this.#manager.runningTurnIds(record.header.id), + this.#manager.sessionRunEpoch(record.header.id), + this.#manager.sessionHostGeneration(), + ), ) : null, }, @@ -532,7 +540,11 @@ export class HostSessionCatalogCoordinator { #projectCatalogQueryRecord(record: SessionCatalogRecord): SessionCatalogItem { return projectSessionCatalogRecord( record, - projectCatalogLiveRunState(this.#manager.runningTurnIds(record.header.id)), + projectCatalogLiveRunState( + this.#manager.runningTurnIds(record.header.id), + this.#manager.sessionRunEpoch(record.header.id), + this.#manager.sessionHostGeneration(), + ), ); } @@ -1737,12 +1749,16 @@ function projectSharedSessionCatalogRecord( function projectCatalogLiveRunState( runningTurnIds: readonly string[], + runEpoch?: number, + hostGeneration?: string, ): SessionCatalogLiveRunState | undefined { const uniqueRunningTurnIds = [...new Set(runningTurnIds)]; if (uniqueRunningTurnIds.length > SESSION_CATALOG_RUNNING_TURN_MAX_ITEMS) return undefined; return { schemaVersion: SESSION_CATALOG_LIVE_RUN_STATE_SCHEMA_VERSION, runningTurnIds: uniqueRunningTurnIds, + runEpoch, + ...(hostGeneration === undefined ? {} : { hostGeneration }), }; } diff --git a/packages/runtime/src/__tests__/runtime-kernel-interaction.test.ts b/packages/runtime/src/__tests__/runtime-kernel-interaction.test.ts index 6c5e99557e..acb01b366a 100644 --- a/packages/runtime/src/__tests__/runtime-kernel-interaction.test.ts +++ b/packages/runtime/src/__tests__/runtime-kernel-interaction.test.ts @@ -261,6 +261,62 @@ describe('RuntimeKernel Interaction close cleanup', () => { await drainIterator(iterator); }); + test('the session run epoch bumps on every turn start and end', async () => { + const store = memoryStore(); + const backends = new BackendRegistry(); + const backend = new BlockingBackend(SESSION_ID, {}); + backends.register('ai-sdk', () => backend); + let id = 0; + const kernel = new RuntimeKernel({ + store, + backends, + newId: () => `epoch-id-${++id}`, + now: () => id, + }); + + assert.equal(kernel.sessionRunEpoch(SESSION_ID), 0); + + const iterator = kernel + .startTurn(SESSION_ID, { turnId: 'turn-epoch-1', text: 'go' }) + [Symbol.asyncIterator](); + await iterator.next(); + const afterStart = kernel.sessionRunEpoch(SESSION_ID); + assert.ok(afterStart >= 1, 'a turn entering the active set bumps the epoch'); + + backend.releaseBlockedSend(); + await drainIterator(iterator).catch(() => undefined); + assert.ok( + kernel.sessionRunEpoch(SESSION_ID) > afterStart, + 'the turn leaving the active set bumps the epoch again', + ); + }); + + test('each kernel gets its own host generation; epochs restart per process (#5713 review)', () => { + const backends = new BackendRegistry(); + backends.register('ai-sdk', () => new BlockingBackend(SESSION_ID, {})); + const first = new RuntimeKernel({ + store: memoryStore(), + backends, + newId: () => 'unused', + now: () => 0, + hostGeneration: () => 'host-generation-1', + }); + // A second kernel stands in for the restarted Host process: a fresh + // generation and epoch counters back to zero, even though the previous + // process may have left a higher epoch on a client's catalog row. + const second = new RuntimeKernel({ + store: memoryStore(), + backends, + newId: () => 'unused', + now: () => 0, + hostGeneration: () => 'host-generation-2', + }); + assert.equal(first.sessionHostGeneration(), 'host-generation-1'); + assert.equal(second.sessionHostGeneration(), 'host-generation-2'); + assert.equal(first.sessionRunEpoch(SESSION_ID), 0); + assert.equal(second.sessionRunEpoch(SESSION_ID), 0); + }); + test('a generation stopped after Run reservation cannot send on the stale backend', async () => { const store = memoryStore(); const backends = new BackendRegistry(); diff --git a/packages/runtime/src/runtime-kernel.ts b/packages/runtime/src/runtime-kernel.ts index a5de4ddedf..35045cc54a 100644 --- a/packages/runtime/src/runtime-kernel.ts +++ b/packages/runtime/src/runtime-kernel.ts @@ -20,6 +20,7 @@ import type { WorkHubActionReceipt } from '@maka/core/workhub-action-result'; import type { AgentRunStore } from '@maka/core/agent-run'; import { agentRunCompositionFromEvents } from '@maka/core/agent-run'; +import { randomUUID } from 'node:crypto'; import type { RuntimeInvocationRecord } from '@maka/core/runtime-invocation'; import { decodeRuntimeBoundaryCursor, @@ -197,6 +198,21 @@ export interface RuntimeKernelLike { * from an arbitrary one of them. */ runningTurnIds?(sessionId: string): string[]; + /** + * Bumped every time a run enters or leaves the session's active set, so two + * same-revision catalog reads can be ordered by their live run state even + * when a turn started or ended between them (#5713). The counter belongs to + * this process alone; pair it with `sessionHostGeneration` to tell which + * process an observation came from. + */ + sessionRunEpoch?(sessionId: string): number; + /** + * Identifies this kernel process's run-epoch generation. Catalog rows + * survive a Host restart while the per-process epoch counters restart at + * zero, so clients must not order rows across generations by epoch — a + * restarted Host supersedes every observation its predecessor published. + */ + sessionHostGeneration?(): string; hasActiveRun?(sessionId: string, runId: string, turnId?: string): boolean; requestRunHandoff?( sessionId: string, @@ -278,6 +294,12 @@ export interface RuntimeKernelDeps { toolBoundaryProtocol?: ToolBoundaryProtocol; backends: BackendRegistry; newId: () => string; + /** + * The run-epoch generation of this kernel process. Omit to draw a random + * UUID; injecting one lets tests pin `sessionHostGeneration()` without + * touching the `newId` sequence. + */ + hostGeneration?: () => string; now: () => number; childTools?: readonly MakaTool[]; resolveChildTools?: (sessionId: string) => Promise; @@ -407,6 +429,13 @@ export class RuntimeKernel implements RuntimeKernelLike { if (deps.runStore && !deps.runtimeEventStore) { throw new Error('RuntimeEventStore is required when AgentRunStore is configured'); } + // One identity per kernel process: catalog rows survive a Host restart + // while the per-session epoch counters restart at zero, so clients pair + // the generation with the epoch instead of comparing epochs across + // processes (#5713 review). Drawn outside the `newId` sequence so a + // restart identity never shifts the ids callers observe, and defaulted + // for deps that omit `newId` entirely. + this.#hostGeneration = deps.hostGeneration?.() ?? randomUUID(); this.historyCompactCoordinator = new HistoryCompactCheckpointCoordinator(deps); } @@ -567,6 +596,11 @@ export class RuntimeKernel implements RuntimeKernelLike { } execution.run = run; execution.phase = 'attached'; + // A host-operation claim's run is visible through the claim itself before + // any backend reservation: the epoch must follow that visibility change + // too, or two same-revision reads can disagree about runningTurnIds + // (#5713 review). + if (execution.hostOperation) this.#bumpSessionRunEpoch(execution.sessionId); if (execution.stopIntent) { run.stop(execution.stopIntent.input.source, execution.stopIntent.input.workHubActionId); } @@ -626,6 +660,11 @@ export class RuntimeKernel implements RuntimeKernelLike { execution: PendingExecutionClaim, outcome: ExecutionClaimOutcome, ): void { + // The claim removal takes a host-operation run out of the visible set: + // order that transition like the attach and the unregister (#5713 review). + if (execution.hostOperation && execution.run) { + this.#bumpSessionRunEpoch(execution.sessionId); + } const claims = this.executionClaims.get(execution.sessionId); claims?.delete(execution); if (claims?.size === 0) this.executionClaims.delete(execution.sessionId); @@ -2183,6 +2222,21 @@ export class RuntimeKernel implements RuntimeKernelLike { return [...new Set(this.activeRunsFor(sessionId).map((run) => run.turnId))]; } + readonly #sessionRunEpochs = new Map(); + readonly #hostGeneration: string; + + sessionRunEpoch(sessionId: string): number { + return this.#sessionRunEpochs.get(sessionId) ?? 0; + } + + sessionHostGeneration(): string { + return this.#hostGeneration; + } + + #bumpSessionRunEpoch(sessionId: string): void { + this.#sessionRunEpochs.set(sessionId, (this.#sessionRunEpochs.get(sessionId) ?? 0) + 1); + } + hasActiveRun(sessionId: string, runId: string, turnId?: string): boolean { return this.activeRunsFor(sessionId).some( (run) => run.runId === runId && (turnId === undefined || run.turnId === turnId), @@ -2721,6 +2775,7 @@ export class RuntimeKernel implements RuntimeKernelLike { } active.activeRuns.set(run.runId, run); active.turnToRunId.set(run.turnId, run.runId); + this.#bumpSessionRunEpoch(active.sessionId); } private assertRunCanDispatch(run: AgentRun, backend: AgentBackend): void { @@ -2748,6 +2803,7 @@ export class RuntimeKernel implements RuntimeKernelLike { if (active.turnToRunId.get(run.turnId) === run.runId) { active.turnToRunId.delete(run.turnId); } + this.#bumpSessionRunEpoch(active.sessionId); } private async unregisterParentRun(active: AgentRunActiveSession, run: AgentRun): Promise { diff --git a/packages/runtime/src/session-manager.ts b/packages/runtime/src/session-manager.ts index 40b020fc83..523ebade47 100644 --- a/packages/runtime/src/session-manager.ts +++ b/packages/runtime/src/session-manager.ts @@ -35,7 +35,7 @@ import { listRecallCandidateSessions, type RecallCandidateStores, } from './recall-candidates.js'; -import { createHash } from 'node:crypto'; +import { createHash, randomUUID } from 'node:crypto'; import { isDeepStrictEqual } from 'node:util'; import { setTimeout as delay } from 'node:timers/promises'; import type { @@ -967,12 +967,42 @@ export class SessionManager { return this.runtimeKernel.runningTurnIds?.(sessionId) ?? []; } + /** + * The live run state's own order, bumped on every turn start and end. Two + * same-revision catalog reads can disagree about `runningTurnIds`; the epoch + * says which one is older (#5713). The counter is per-process — pair it with + * `sessionHostGeneration` to tell which process an observation came from. + */ + sessionRunEpoch(sessionId: string): number { + return this.runtimeKernel.sessionRunEpoch?.(sessionId) ?? 0; + } + + /** + * Identifies this process's run-epoch generation. Catalog rows survive a + * Host restart while the per-process epoch counters restart at zero, so + * clients order same-revision reads by generation first, never by epoch + * across restarts (#5713). Falls back to a random identity — fixed once, + * and drawn outside the `newId` sequence — when the kernel does not expose + * one. + */ + readonly #hostGenerationFallback = randomUUID(); + + sessionHostGeneration(): string { + return this.runtimeKernel.sessionHostGeneration?.() ?? this.#hostGenerationFallback; + } + #projectLiveRunState(sessions: SessionSummary[]): SessionSummary[] { const runningTurnIds = this.runtimeKernel.runningTurnIds?.bind(this.runtimeKernel); if (!runningTurnIds) return sessions; + const sessionRunEpoch = this.runtimeKernel.sessionRunEpoch?.bind(this.runtimeKernel); + const sessionHostGeneration = this.runtimeKernel.sessionHostGeneration?.bind( + this.runtimeKernel, + ); return sessions.map((session) => ({ ...session, runningTurnIds: runningTurnIds(session.id), + ...(sessionRunEpoch ? { runEpoch: sessionRunEpoch(session.id) } : {}), + ...(sessionHostGeneration ? { runHostGeneration: sessionHostGeneration() } : {}), })); }