diff --git a/apps/desktop/src/main/__tests__/runtime-host-session-observer.test.ts b/apps/desktop/src/main/__tests__/runtime-host-session-observer.test.ts index d0550e32af..eea95a4dd0 100644 --- a/apps/desktop/src/main/__tests__/runtime-host-session-observer.test.ts +++ b/apps/desktop/src/main/__tests__/runtime-host-session-observer.test.ts @@ -892,7 +892,7 @@ test('broadcasts durable admission and transcript changes from the same message' await observer.close(); }); -for (const resolution of ['owned', 'cancelled', 'pending', 'unavailable'] as const) { +for (const resolution of ['owned', 'cancelled', 'not_admitted', 'pending', 'unavailable'] as const) { test(`proves a removed follow-up is ${resolution} before successor content`, async (t) => { const events = new AsyncFrameQueue(); const queries: string[][] = []; @@ -944,7 +944,7 @@ for (const resolution of ['owned', 'cancelled', 'pending', 'unavailable'] as con const admissions = target.events.filter((event) => event.type === 'message_admission'); assert.deepEqual(admissions.map((event) => ({ messageId: event.messageId, turnId: event.turnId, outcome: event.outcome, - })), resolution === 'owned' || resolution === 'cancelled' ? [{ + })), resolution === 'owned' || resolution === 'cancelled' || resolution === 'not_admitted' ? [{ messageId: 'followup-1', turnId: 'turn-2', outcome: resolution === 'owned' ? 'admitted' : 'retracted', }] : []); diff --git a/apps/desktop/src/main/__tests__/session-local.test.ts b/apps/desktop/src/main/__tests__/session-local.test.ts index e49bf808a9..4ac568bd05 100644 --- a/apps/desktop/src/main/__tests__/session-local.test.ts +++ b/apps/desktop/src/main/__tests__/session-local.test.ts @@ -453,6 +453,166 @@ test('lost ACK recovery never changes epoch or ID and does not block another Ses assert.equal(calls.find((call) => call.messageId === 'message-2')?.originHostEpoch, 'epoch-2'); }); +test('a Message dispatched in a released epoch is settled from the Host, never replayed', async (t) => { + const { store, beforeClose } = await database(t); + const submissions: TurnMessageSubmitInput[] = []; + const queried: string[][] = []; + let quiescent = false; + const target: DesktopSessionLocalTarget = { + partition: 'authority', + profileId: 'profile', + scope: { hostId: 'root', targetEpoch: 'target-1' }, + client: { + ...client('epoch-2'), + async queryMessageExecutions(input) { + queried.push([...input.messageIds]); + return { + resolutions: [ + { messageId: 'message-1', state: 'owned', turnId: 'turn-recovered', runId: 'run-1' }, + ], + }; + }, + }, + submit: async (input) => { + submissions.push(input); + return accepted; + }, + }; + const service = new DesktopSessionLocalService(store, { + targets: () => [target], + changed() {}, + onError: (error) => assert.fail(String(error)), + }); + beforeClose.push(() => service.close()); + // Dispatch once in the previous epoch, then lose the answer the way a Host + // restart does: the local copy is left `unknown` with a dead dispatch epoch. + const dispatched = store.enqueue('authority', intent()); + store.update({ + ...dispatched, + state: 'unknown', + intent: { ...dispatched.intent, originHostEpoch: 'epoch-1' }, + }); + service.wake(); + await waitFor(() => store.get('authority', 'message-1')?.state === 'accepted'); + assert.deepEqual(queried, [['message-1']]); + // The old epoch's submit is never replayed: that answer can only be + // `outcome_unknown`, which would freeze the copy forever. + assert.deepEqual(submissions, []); + assert.equal(store.get('authority', 'message-1')?.result?.disposition, 'turn_started'); +}); + +test('a stale-epoch Message the Host never admitted is retired so later sends proceed', async (t) => { + const { store, beforeClose } = await database(t); + const submissions: string[] = []; + const target: DesktopSessionLocalTarget = { + partition: 'authority', + profileId: 'profile', + scope: { hostId: 'root', targetEpoch: 'target-1' }, + client: { + ...client('epoch-2'), + async queryMessageExecutions(input) { + // The Host holds no receipt, steering proof, tombstone, or admission + // for this identity, so it reports that absence positively. + return { + resolutions: input.messageIds.map( + (messageId) => ({ messageId, state: 'not_admitted' as const }), + ), + }; + }, + }, + submit: async (input) => { + submissions.push(input.messageId); + return accepted; + }, + }; + const service = new DesktopSessionLocalService(store, { + targets: () => [target], + changed() {}, + onError: (error) => assert.fail(String(error)), + }); + beforeClose.push(() => service.close()); + const dispatched = store.enqueue('authority', intent()); + store.update({ + ...dispatched, + state: 'unknown', + intent: { ...dispatched.intent, originHostEpoch: 'epoch-1' }, + }); + service.wake(); + await waitFor(() => store.get('authority', 'message-1')?.state === 'failed'); + // A settled non-delivery releases the Session's ordering. + store.enqueue('authority', intent('message-2')); + service.wake(); + await waitFor(() => store.get('authority', 'message-2')?.state === 'accepted'); + assert.deepEqual(submissions, ['message-2']); +}); + +test('a running Host that cannot yet resolve a stale Message leaves the copy unresolved', async (t) => { + const { store, beforeClose } = await database(t); + const target: DesktopSessionLocalTarget = { + partition: 'authority', + profileId: 'profile', + scope: { hostId: 'root', targetEpoch: 'target-1' }, + client: { + ...client('epoch-2'), + async queryMessageExecutions() { + throw new RuntimeHostOperationError('turn.message.execution.query', 'host_not_ready', 'busy'); + }, + }, + submit: async () => accepted, + }; + const service = new DesktopSessionLocalService(store, { + targets: () => [target], + changed() {}, + onError: (error) => assert.fail(String(error)), + }); + beforeClose.push(() => service.close()); + const dispatched = store.enqueue('authority', intent()); + store.update({ + ...dispatched, + state: 'unknown', + intent: { ...dispatched.intent, originHostEpoch: 'epoch-1' }, + }); + service.wake(); + await nextTurn(); + // Guessing "never delivered" here would let the user resend a Message the + // Host may already own, so the copy stays unresolved instead. + assert.equal(store.get('authority', 'message-1')?.state, 'unknown'); +}); + +test('an omitted resolution leaves the stale copy unresolved rather than failed', async (t) => { + const { store, beforeClose } = await database(t); + const target: DesktopSessionLocalTarget = { + partition: 'authority', + profileId: 'profile', + scope: { hostId: 'root', targetEpoch: 'target-1' }, + client: { + ...client('epoch-2'), + // The Host answered successfully but could not assert anything about the + // identity, so it omitted it. That is "cannot say yet", not "not admitted". + async queryMessageExecutions() { + return { resolutions: [] }; + }, + }, + submit: async () => accepted, + }; + const service = new DesktopSessionLocalService(store, { + targets: () => [target], + changed() {}, + onError: (error) => assert.fail(String(error)), + }); + beforeClose.push(() => service.close()); + const dispatched = store.enqueue('authority', intent()); + store.update({ + ...dispatched, + state: 'unknown', + intent: { ...dispatched.intent, originHostEpoch: 'epoch-1' }, + }); + service.wake(); + await waitFor(() => store.get('authority', 'message-1')?.state === 'unknown'); + // Only a positive `not_admitted` may retire the copy; silence must not. + assert.equal(store.get('authority', 'message-1')?.state, 'unknown'); +}); + test('a removed authority cannot be repopulated by an in-flight admission', async (t) => { const { store, beforeClose } = await database(t); const response = deferred(); diff --git a/apps/desktop/src/main/runtime-host-session-observer.ts b/apps/desktop/src/main/runtime-host-session-observer.ts index 660eb7b7b0..51ab35ee2d 100644 --- a/apps/desktop/src/main/runtime-host-session-observer.ts +++ b/apps/desktop/src/main/runtime-host-session-observer.ts @@ -857,13 +857,23 @@ export class RuntimeHostSessionObserver { const turnId = resolution.state === 'owned' ? resolution.turnId : (next.rootTurn ?? previous.rootTurn)?.turnId; if (!turnId) continue; + // `owned` admits; `cancelled` and the positive `not_admitted` — proof + // the Message can never execute — both retract it. Naming the two + // retracting states keeps a future addition from silently inheriting + // this outcome through the `else`. + const outcome = resolution.state === 'owned' + ? 'admitted' as const + : resolution.state === 'cancelled' || resolution.state === 'not_admitted' + ? 'retracted' as const + : undefined; + if (!outcome) continue; this.#broadcast(state.sessionId, { type: 'message_admission', id: `host-message-resolution:${next.queue.hostEpoch}:${next.queue.queueRevision}:${resolution.messageId}`, turnId, ts: this.#now(), messageId: resolution.messageId, - outcome: resolution.state === 'owned' ? 'admitted' : 'retracted', + outcome, }); } } catch { diff --git a/apps/desktop/src/main/session-local-service.ts b/apps/desktop/src/main/session-local-service.ts index 7470afe20f..7620d220db 100644 --- a/apps/desktop/src/main/session-local-service.ts +++ b/apps/desktop/src/main/session-local-service.ts @@ -27,6 +27,7 @@ import { RuntimeHostRequestInterruptedError, } from '@maka/runtime-host/client'; import type { + TurnMessageExecutionResolution, TurnMessageSubmitInput, TurnMessageSubmitResult, WorkspaceTarget, @@ -57,6 +58,9 @@ import { } from './session-local-store.js'; import type { DesktopTranscriptReplicaSnapshot } from './desktop-transcript-replica.js'; +/** The Host resolved no Skill for a Message that is settled from durable facts. */ +const EMPTY_SKILL_INVOCATION = { loaded: [], failed: [], receipts: [] } as const; + export interface DesktopSessionLocalTarget { readonly partition: string; readonly scope: DesktopTargetScope; @@ -64,7 +68,10 @@ export interface DesktopSessionLocalTarget { readonly client?: Pick< DesktopRuntimeHostClient, 'hostEpoch' | 'createSession' | 'getSession' | 'listSessions' | 'ingestAttachment' - >; + > & { + /** Absent only in narrow test doubles; production clients always provide it. */ + readonly queryMessageExecutions?: DesktopRuntimeHostClient['queryMessageExecutions']; + }; readonly submit?: (input: TurnMessageSubmitInput) => Promise; } @@ -383,6 +390,19 @@ export class DesktopSessionLocalService { const stillOwned = () => this.#current(target) && this.store.get(record.partition, record.messageId) !== undefined; try { + // A Message whose immutable dispatch epoch is gone cannot be replayed: + // the running Host has no in-memory submit for it and answers + // `outcome_unknown` for any identity lacking durable proof, so retrying + // the same submit loops forever and blocks the Session's queue. Settle it + // from the Host's durable facts instead. + if ( + record.intent.originHostEpoch !== undefined && + record.intent.originHostEpoch !== client.hostEpoch && + client.queryMessageExecutions !== undefined + ) { + await this.#settleDispatchedEpoch(target, record); + return; + } const creation = this.store.creation(target.partition, record.sessionId); if (creation) { // session.create already has a durable request fingerprint. Replaying @@ -476,6 +496,106 @@ export class DesktopSessionLocalService { } this.deps.changed(target.scope, record.sessionId); } + + /** + * Settles a Message whose immutable dispatch epoch is no longer the running + * Host's. The running Host released the previous epoch's in-memory submits, + * so replaying this submit can never be proven and answers `outcome_unknown` + * forever, which held the Session's queue behind it. The identity's real fate + * is still on record: the Host resolves it from its durable receipts, + * steering proofs, cancellation tombstones, and pending admissions. Read that + * instead of replaying the doomed submit. + */ + async #settleDispatchedEpoch( + target: DesktopSessionLocalTarget, + record: LocalOutboxRecord, + ): Promise { + const client = target.client!; + const key = `${target.partition}:${record.messageId}`; + const stillOwned = () => this.#current(target) && this.store.get(target.partition, record.messageId) !== undefined; + let resolution: TurnMessageExecutionResolution | undefined; + try { + const resolved = await client.queryMessageExecutions!({ + sessionId: record.sessionId, + messageIds: [record.messageId], + }); + resolution = resolved.resolutions.find((entry) => entry.messageId === record.messageId); + } catch { + // The running Host cannot answer yet. Keep the copy unresolved rather + // than guess: a wrong "never delivered" would let the user resend a + // Message the Host may already own. + if (!stillOwned()) return; + this.store.update({ + ...record, + state: 'unknown', + error: 'Host outcome is unknown; the running Host has not confirmed the original message.', + }); + if (!this.#closed && !this.#retries.has(key)) this.#scheduleRetry(key); + this.deps.changed(target.scope, record.sessionId); + return; + } + if (!stillOwned()) return; + const current = this.store.get(target.partition, record.messageId)!; + if (resolution?.state === 'owned') { + // A durable receipt or steering proof names this identity; the Turn it + // opened is this local copy's settlement. + this.store.update({ + ...current, + state: 'accepted', + result: { + disposition: 'turn_started', + turnId: resolution.turnId, + skillInvocation: EMPTY_SKILL_INVOCATION, + }, + error: undefined, + }); + } else if (resolution?.state === 'pending') { + // The Host durably holds a queued admission for it. Delivery succeeded; + // the Host queue owns the ordering, so this local copy yields. + this.store.update({ + ...current, + state: 'accepted', + result: { disposition: 'followup', skillInvocation: EMPTY_SKILL_INVOCATION }, + error: undefined, + }); + } else if (resolution?.state === 'cancelled' || resolution?.state === 'not_admitted') { + // The Host answered positively that this identity is settled: either it + // carries a cancellation tombstone, or it has no durable record at all + // and no in-flight submit, so no epoch ever admitted it and it can never + // execute. Record that as an explicit non-delivery: it releases the + // Session's ordering and gives the user a removable local copy. + this.store.update({ + ...current, + state: 'failed', + error: + resolution.state === 'cancelled' + ? 'The Host cancelled this message; the local copy is retained.' + : 'The Host never admitted this message; the local copy is retained.', + result: undefined, + }); + } else { + // The Host omitted the identity: it cannot yet assert anything about it + // (`recovering`). Keep the copy unresolved rather than guess — a wrong + // "never delivered" would let the user resend a Message the Host may + // already own. + this.store.update({ + ...current, + state: 'unknown', + error: 'Host outcome is unknown; the running Host has not confirmed the original message.', + }); + if (!this.#closed && !this.#retries.has(key)) this.#scheduleRetry(key); + this.deps.changed(target.scope, record.sessionId); + return; + } + const timer = this.#retries.get(key); + if (timer) { + clearTimeout(timer); + this.#retries.delete(key); + } + this.#probed.delete(key); + this.#catalogFresh.delete(target.partition); + this.deps.changed(target.scope, record.sessionId); + } } export function registerDesktopSessionLocalIpc(deps: { diff --git a/apps/desktop/src/renderer/features/workbar/tools/side-chat/use-quote-companion.ts b/apps/desktop/src/renderer/features/workbar/tools/side-chat/use-quote-companion.ts index eb42cc70a5..1c80fc59c5 100644 --- a/apps/desktop/src/renderer/features/workbar/tools/side-chat/use-quote-companion.ts +++ b/apps/desktop/src/renderer/features/workbar/tools/side-chat/use-quote-companion.ts @@ -663,13 +663,17 @@ export function useQuoteCompanion(input: UseQuoteCompanionInput): UseQuoteCompan try { const { resolutions } = await sideChat.queryMessageExecutions(forkId, [...messageIds]); if (!mountedRef.current || companionIdRef.current !== forkId) return; - const cancelled = new Set(); + const retired = new Set(); const unprovenOwnedTurnIds = new Set(); let ownershipChanged = false; for (const resolution of resolutions) { const pending = pendingAdmissionRef.current; - if (resolution.state === 'cancelled') { - cancelled.add(resolution.messageId); + // `cancelled` and `not_admitted` are both positive proof that this + // Message will never execute, so each retires the transient and frees + // the Composer's admission slot. Only an omitted identity means the + // Host cannot say yet, and that keeps its slot. + if (resolution.state === 'cancelled' || resolution.state === 'not_admitted') { + retired.add(resolution.messageId); if (pending?.messageId === resolution.messageId) releaseAdmission(pending); } else if (resolution.state === 'owned') { bindPendingMessageTurn(resolution.messageId, resolution.turnId); @@ -689,7 +693,7 @@ export function useQuoteCompanion(input: UseQuoteCompanionInput): UseQuoteCompan setHasContent(true); setOwnTurnTick((tick) => tick + 1); } - for (const messageId of cancelled) pendingUserMessagesRef.current.delete(messageId); + for (const messageId of retired) pendingUserMessagesRef.current.delete(messageId); if (unprovenOwnedTurnIds.size > 0) { const recovered = await Promise.allSettled( [...unprovenOwnedTurnIds].map((turnId) => @@ -703,10 +707,10 @@ export function useQuoteCompanion(input: UseQuoteCompanionInput): UseQuoteCompan } } reconcilePendingUserMessages(); - if (cancelled.size > 0) { + if (retired.size > 0) { setMessageQueue((current) => ({ ...current, - entries: current.entries.filter((entry) => !cancelled.has(entry.messageId)), + entries: current.entries.filter((entry) => !retired.has(entry.messageId)), })); } } catch { diff --git a/apps/desktop/src/renderer/features/workhub/model/delegation-feedback.ts b/apps/desktop/src/renderer/features/workhub/model/delegation-feedback.ts index e11c19ce90..0a0af754ab 100644 --- a/apps/desktop/src/renderer/features/workhub/model/delegation-feedback.ts +++ b/apps/desktop/src/renderer/features/workhub/model/delegation-feedback.ts @@ -30,6 +30,9 @@ export function projectWorkHubDelegationState(input: { }): WorkHubDelegationState { if (input.executionReadFailed || !input.resolution) return 'recovering'; if (input.resolution.state === 'cancelled') return 'aborted'; + // The Host positively reported this identity was never admitted: the + // delegation's Message never became work, so it can only be reported failed. + if (input.resolution.state === 'not_admitted') return 'failed'; if (input.resolution.state === 'pending') return 'accepted'; if (input.turn?.statusSource === 'recorded' && input.turn.status !== 'running') { return input.turn.status; diff --git a/packages/runtime-host/src/__tests__/message-coordinator.test.ts b/packages/runtime-host/src/__tests__/message-coordinator.test.ts index 19a97e5ffb..ecef7c4772 100644 --- a/packages/runtime-host/src/__tests__/message-coordinator.test.ts +++ b/packages/runtime-host/src/__tests__/message-coordinator.test.ts @@ -660,11 +660,54 @@ test('message execution query reports the Turn that durably owns each Message', turnId: ROOT.turnId, runId: ROOT.runId, }, + { + // No receipt, steering proof, tombstone, or admission names it, so + // the Host reports the absence positively rather than omitting it. + messageId: 'unknown-message', + state: 'not_admitted', + }, ], }, }); }); +test('message execution query never observes a submit mid-admission', async () => { + const fixture = createFixture(); + fixture.coordinator.reserveRootTurn(ROOT); + // Hold admission open inside the submit so the query must queue behind it. + const preparing = deferred(); + const release = deferred(); + fixture.setMessagePreparation(async () => { + preparing.resolve(undefined); + await release.promise; + return { kind: 'ready', content: { text: 'raced' }, skillInvocation: EMPTY_SKILL_INVOCATION }; + }); + const submit = fixture.coordinator.handlers['turn.message.submit']( + { + originHostEpoch: 'epoch-1', + sessionId: ROOT.sessionId, + messageId: 'in-flight-message', + content: { text: 'raced' }, + placement: 'next_turn', + } as const, + operationContext(), + ); + await preparing.promise; + // Issued while the submit holds the Session admission. If the read were not + // gated it would see no admission row and answer `not_admitted`, handing the + // user a resend for a Message this Host is in the middle of admitting. + const query = fixture.coordinator.handlers['turn.message.execution.query']( + { sessionId: ROOT.sessionId, messageIds: ['in-flight-message'] }, + operationContext(), + ); + release.resolve(undefined); + await submit; + assert.deepEqual(await query, { + ok: true, + result: { resolutions: [{ messageId: 'in-flight-message', state: 'pending' }] }, + }); +}); + test('message execution disposition reuses a held Session admission', async () => { const fixture = createFixture(); const content = { text: 'pending delegation' }; diff --git a/packages/runtime-host/src/protocol/index.ts b/packages/runtime-host/src/protocol/index.ts index f4ed0e72a3..eac5206acb 100644 --- a/packages/runtime-host/src/protocol/index.ts +++ b/packages/runtime-host/src/protocol/index.ts @@ -101,7 +101,11 @@ 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 = 168 as const; +export const RUNTIME_HOST_COMPATIBILITY_EPOCH = 169 as const; +// 169: The message execution query reports an identity the Host can prove was +// never admitted as a positive `not_admitted` resolution instead of omitting +// it, so silence stops meaning both "not admitted" and "cannot say yet". +// Epoch-168 peers reject the new state as an invalid frame. // 167: Removed the turn.regenerate operation. Older peers can no longer // safely interoperate because they may submit or advertise that operation. // 166: Connection usage reads add an operation, an accepted availability reason diff --git a/packages/runtime-host/src/protocol/message.ts b/packages/runtime-host/src/protocol/message.ts index 89806dc726..21819e6417 100644 --- a/packages/runtime-host/src/protocol/message.ts +++ b/packages/runtime-host/src/protocol/message.ts @@ -138,6 +138,7 @@ export interface TurnMessageExecutionQueryResult { export type TurnMessageExecutionResolution = | { readonly messageId: string; readonly state: 'pending' } | { readonly messageId: string; readonly state: 'cancelled' } + | { readonly messageId: string; readonly state: 'not_admitted' } | { readonly messageId: string; readonly state: 'owned'; @@ -401,6 +402,16 @@ function decodeTurnMessageExecutionQueryResult(value: unknown): TurnMessageExecu state: 'cancelled', }; } + if (resolution.state === 'not_admitted') { + assertExactKeys(resolution, 'turn.message.execution.query not_admitted resolution', [ + 'messageId', + 'state', + ]); + return { + messageId: requireEntityId(resolution.messageId, 'messageId'), + state: 'not_admitted', + }; + } if (resolution.state === 'owned') { assertExactKeys(resolution, 'turn.message.execution.query owned resolution', [ 'messageId', diff --git a/packages/runtime-host/src/server/message-coordinator.ts b/packages/runtime-host/src/server/message-coordinator.ts index a5aded5881..687a8274cb 100644 --- a/packages/runtime-host/src/server/message-coordinator.ts +++ b/packages/runtime-host/src/server/message-coordinator.ts @@ -195,6 +195,22 @@ export type HostMessageExecutionDisposition = | HostMessageResolvedDisposition | { readonly kind: 'pending' }; +/** + * The wire shape one resolved identity reports. `not_admitted` is a positive + * statement that no durable record names the identity and no admission write + * is in flight — distinct from omission, which means the Host cannot say yet. + */ +type MessageExecutionResolutionOutcome = + | { readonly messageId: string; readonly state: 'pending' } + | { readonly messageId: string; readonly state: 'cancelled' } + | { readonly messageId: string; readonly state: 'not_admitted' } + | { + readonly messageId: string; + readonly state: 'owned'; + readonly turnId: string; + readonly runId: string; + }; + /** Root execution operations that must share the message coordinator's Session gate. */ export interface HostMessageRootPort { readLatestRootTurnLineage(identity: { @@ -463,50 +479,55 @@ export class HostMessageCoordinator implements RuntimeMessageAuthority { return success({ cancelledMessageIds }); } - async queryMessageExecutions(input: { + queryMessageExecutions(input: { sessionId: string; messageIds: readonly string[]; - }): Promise< - MessageOutcome<{ - resolutions: Array< - | { messageId: string; state: 'pending' } - | { messageId: string; state: 'cancelled' } - | { messageId: string; state: 'owned'; turnId: string; runId: string } - >; - }> - > { - const resolutions: Array< - | { messageId: string; state: 'pending' } - | { messageId: string; state: 'cancelled' } - | { messageId: string; state: 'owned'; turnId: string; runId: string } - > = []; - for (const messageId of input.messageIds) { - const disposition = await this.#resolveMessageExecution(input.sessionId, messageId); - if (disposition.kind === 'owned_root' || disposition.kind === 'shared_turn') { - // This read projects current execution, including safe-boundary - // continuations. The Message's durable admission ownership is unchanged. - const latest = await this.#root.readLatestRootTurnLineage({ - sessionId: input.sessionId, - turnId: disposition.turnId, - runId: disposition.runId, - }); - resolutions.push({ - messageId, - state: 'owned', - turnId: latest.turnId, - runId: latest.runId, - }); - continue; - } - if (disposition.kind === 'cancelled') { - resolutions.push({ messageId, state: 'cancelled' }); - continue; - } - if (disposition.kind === 'pending') { - resolutions.push({ messageId, state: 'pending' }); + }): Promise }>> { + // Enter the Session admission the way every other admission reader does. + // `not_admitted` claims that no epoch ever admitted this identity, and it + // is only sound while no admission write can be in flight. WorkHub writes + // admission rows without an in-memory submit to observe, so the gate — + // not `#pendingSubmits` alone — is what makes the read atomic. + return this.#sessionAdmission.runOrJoin(input.sessionId, async () => { + const resolutions: Array = []; + for (const messageId of input.messageIds) { + const disposition = await this.#resolveMessageExecution(input.sessionId, messageId); + if (disposition.kind === 'owned_root' || disposition.kind === 'shared_turn') { + // This read projects current execution, including safe-boundary + // continuations. The Message's durable admission ownership is unchanged. + const latest = await this.#root.readLatestRootTurnLineage({ + sessionId: input.sessionId, + turnId: disposition.turnId, + runId: disposition.runId, + }); + resolutions.push({ + messageId, + state: 'owned', + turnId: latest.turnId, + runId: latest.runId, + }); + continue; + } + if (disposition.kind === 'cancelled') { + resolutions.push({ messageId, state: 'cancelled' }); + continue; + } + if (disposition.kind === 'pending') { + resolutions.push({ messageId, state: 'pending' }); + continue; + } + // `recovering` means no durable receipt, steering proof, cancellation + // tombstone, or pending admission names this identity. Under the + // Session gate no admission write is in flight, so that silence is + // itself the proof: nothing in this epoch — or any prior one, since a + // stale epoch's submit can never commit here — ever admitted it. + // Report that fact positively instead of omitting the identity, so a + // missing entry stops meaning both "not admitted" and "cannot say yet". + if (this.#pendingSubmits.has(operationKey(input.sessionId, messageId))) continue; + resolutions.push({ messageId, state: 'not_admitted' }); } - } - return success({ resolutions }); + return success({ resolutions }); + }); } /**