From f73590b11405164d2ff83df3e0ee90b763a6044e Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Sat, 3 Oct 2026 06:28:56 +0200 Subject: [PATCH 1/6] fix(server): close orphaned runtime request transcripts --- .../ProviderRuntimeRecoveryService.test.ts | 111 ++++++++++++++++++ .../ProviderRuntimeRecoveryService.ts | 44 +++++++ .../ProviderSessionManager.test.ts | 109 +++++++++++++++++ .../ProviderSessionManager.ts | 13 +- 4 files changed, 276 insertions(+), 1 deletion(-) diff --git a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.test.ts b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.test.ts index 8df85818a761..0bac5f0dc03c 100644 --- a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.test.ts @@ -180,6 +180,117 @@ it.effect("expires orphaned runtime requests before command readiness", () => { }).pipe(Effect.provide(layer)); }); +it.effect.each([ + { trigger: "startup", requestType: "approval_request", runStatus: null }, + { trigger: "startup", requestType: "user_input_request", runStatus: null }, + { trigger: "shutdown", requestType: "approval_request", runStatus: null }, + { trigger: "shutdown", requestType: "user_input_request", runStatus: null }, + { trigger: "startup", requestType: "approval_request", runStatus: "running" }, + { trigger: "startup", requestType: "approval_request", runStatus: "completed" }, +] as const)( + "closes $requestType transcript entities on $trigger with run status $runStatus", + ({ trigger, requestType, runStatus }) => { + const threadId = ThreadId.make(`thread_recovery_${trigger}_${requestType}`); + const requestId = RuntimeRequestId.make("request_runless"); + const nodeId = NodeId.make("node_runless"); + const itemId = TurnItemId.make("item_runless"); + const terminalItemId = TurnItemId.make("item_already_terminal"); + const messageNodeId = NodeId.make("node_message"); + const runId = runStatus === null ? null : RunId.make("run_recovery_request"); + let committedInput: Parameters[0] | null = + null; + const projection = { + thread: { id: threadId }, + runtimeRequests: [ + { id: requestId, nodeId, status: "pending", responseCapability: { type: "live" } }, + { + id: RuntimeRequestId.make("request_message"), + nodeId: messageNodeId, + status: "pending", + responseCapability: { type: "message" }, + }, + ], + providerSessions: [], + providerThreads: [], + providerTurns: [], + runs: + runStatus === null + ? [] + : [ + { + id: runId, + status: runStatus, + providerInstanceId: ProviderInstanceId.make("codex"), + }, + ], + attempts: [], + subagents: [], + messages: [], + nodes: [ + { id: nodeId, runId, kind: requestType, status: "waiting" }, + { id: messageNodeId, runId, kind: "user_input_request", status: "waiting" }, + ], + turnItems: [ + { id: itemId, runId, nodeId, type: requestType, requestId, status: "waiting" }, + { + id: terminalItemId, + runId, + nodeId, + type: requestType, + requestId, + status: "completed", + }, + { + id: TurnItemId.make("item_message"), + runId, + nodeId: messageNodeId, + type: "user_input_request", + requestId: RuntimeRequestId.make("request_message"), + status: "waiting", + }, + ], + } as unknown as OrchestrationV2ThreadProjection; + const layer = ProviderRuntimeRecovery.layer.pipe( + Layer.provide(ServerSettings.layerTest()), + Layer.provide( + Layer.mergeAll( + Layer.mock(ProjectionStore.ProjectionStoreV2)({ + getRecoveryThreadIds: () => Effect.succeed([threadId]), + getRuntimeRecoveryProjection: () => Effect.succeed(projection), + }), + Layer.mock(EventSink.EventSinkV2)({ + commitCommand: (input) => { + committedInput = input; + return Effect.succeed({ committed: true, cancelledEffectCount: 0 } as never); + }, + }), + IdAllocator.layer, + Layer.mock(EffectOutbox.EffectOutboxV2)({ + reconcileAfterProcessLoss: Effect.succeed({ requeued: 0, cancelled: 0 }), + }), + ), + ), + ); + return Effect.gen(function* () { + const summary = + yield* (yield* ProviderRuntimeRecovery.ProviderRuntimeRecoveryService).reconcile(trigger); + assert.equal(summary.closedRequests, 1); + const events = committedInput?.events ?? []; + assert.equal(events.length, runStatus === "running" ? 4 : 3); + const requestEvent = events.find((event) => event.type === "runtime-request.updated"); + assert.equal(requestEvent?.payload.status, trigger === "startup" ? "expired" : "cancelled"); + const nodeEvent = events.find((event) => event.type === "node.updated"); + const itemEvent = events.find((event) => event.type === "turn-item.updated"); + assert.equal(nodeEvent?.payload.id, nodeId); + assert.equal(nodeEvent?.payload.status, "cancelled"); + assert.isNotNull(nodeEvent?.payload.completedAt); + assert.equal(itemEvent?.payload.id, itemId); + assert.equal(itemEvent?.payload.status, "cancelled"); + assert.isNotNull(itemEvent?.payload.completedAt); + }).pipe(Effect.provide(layer)); + }, +); + it.effect("preserves async questions across startup and shutdown", () => { const threadId = ThreadId.make("async-recovery-thread"); const nodeId = NodeId.make("async-recovery-node"); diff --git a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts index 67ae22f1883b..e91e67cdb84b 100644 --- a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts +++ b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts @@ -567,6 +567,50 @@ export const make = Effect.gen(function* () { }); } } + // Session-scoped requests may have no run. Close their transcript + // entities too, unless run/background recovery already closed them. + const terminalizedNodeIds = new Set( + events.flatMap((event) => (event.type === "node.updated" ? [event.payload.id] : [])), + ); + for (const request of requests) { + const node = projection.nodes.find((candidate) => candidate.id === request.nodeId); + if ( + node !== undefined && + isNonterminalNodeStatus(node.status) && + !terminalizedNodeIds.has(node.id) + ) { + terminalizedNodeIds.add(node.id); + events.push({ + id: yield* allocateEventId(), + type: "node.updated", + threadId: projection.thread.id, + ...(node.runId === null ? {} : { runId: node.runId }), + nodeId: node.id, + occurredAt: now, + payload: { ...node, status: "cancelled", completedAt: now }, + }); + } + for (const item of projection.turnItems ?? []) { + if ( + (item.type !== "approval_request" && item.type !== "user_input_request") || + item.requestId !== request.id || + !isNonterminalTurnItemStatus(item.status) || + cancelledStaleItemIds.has(item.id) + ) { + continue; + } + cancelledStaleItemIds.add(item.id); + events.push({ + id: yield* allocateEventId(), + type: "turn-item.updated", + threadId: projection.thread.id, + ...(item.runId === null ? {} : { runId: item.runId }), + ...(item.nodeId === null ? {} : { nodeId: item.nodeId }), + occurredAt: now, + payload: { ...item, status: "cancelled", completedAt: now, updatedAt: now }, + }); + } + } // All provider processes are gone on startup/shutdown: clear any // persisted Waiting roster (including idle threads from settled roots) // and idle active threads without resurrecting active status. diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts index 696a40d8d176..e100f69fa043 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts @@ -2643,6 +2643,115 @@ it.effect("ProviderSessionManagerV2 settles a request the event pump persists du }), ); +it.effect.each(["approval_request", "user_input_request"] as const)( + "ProviderSessionManagerV2 detaches only its thread's live %s", + (requestType) => + Effect.gen(function* () { + const state = yield* Ref.make(emptyState); + const mcpConfigs = yield* Ref.make< + ReadonlyArray + >([]); + const effect = Effect.gen(function* () { + const eventSink = yield* EventSink.EventSinkV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const manager = yield* ProviderSessionManager.ProviderSessionManagerV2; + const projectionStore = yield* ProjectionStore.ProjectionStoreV2; + const registry = yield* McpSessionRegistry.McpSessionRegistry; + const now = yield* DateTime.now; + const threadId = ThreadId.make(`detach_request_${requestType}`); + const siblingThreadId = ThreadId.make(`detach_request_sibling_${requestType}`); + const providerSessionId = idAllocator.derive.providerSession({ + providerInstanceId: modelSelection.instanceId, + }); + const otherSessionId = yield* idAllocator.allocate.providerSession({ + providerInstanceId: modelSelection.instanceId, + threadId, + }); + yield* eventSink.write({ + events: [ + yield* makeThreadCreatedEvent({ idAllocator, threadId, now }), + yield* makeThreadCreatedEvent({ idAllocator, threadId: siblingThreadId, now }), + ], + }); + const requests = yield* Effect.forEach( + [ + [threadId, providerSessionId], + [siblingThreadId, providerSessionId], + [threadId, otherSessionId], + ] as const, + Effect.fnUntraced(function* ([requestThreadId, requestSessionId]) { + const request = yield* makePendingRuntimeRequestEvents({ + idAllocator, + threadId: requestThreadId, + providerSessionId: requestSessionId, + providerThread: makeProviderThread({ + idAllocator, + threadId: requestThreadId, + providerSessionId: requestSessionId, + now, + }), + now, + }); + yield* eventSink.write({ + events: request.events.map((event) => + event.type === "turn-item.updated" && requestType === "user_input_request" + ? { ...event, payload: { ...event.payload, type: requestType, questions: [] } } + : event, + ), + }); + return request; + }), + ); + yield* manager.open({ threadId, providerSessionId, modelSelection, runtimePolicy }); + const token = (yield* Ref.get(mcpConfigs)) + .at(-1) + ?.authorizationHeader.replace(/^Bearer\s+/, ""); + assert.isDefined(token); + yield* manager.open({ + threadId: siblingThreadId, + providerSessionId, + modelSelection, + runtimePolicy, + }); + yield* manager.detach({ providerSessionId, threadId }); + // A duplicate detach must leave the sibling and replacement session alone. + yield* manager.detach({ providerSessionId, threadId }); + const projection = yield* projectionStore.getThreadProjection(threadId); + const detachedRequest = requests[0]!; + const closedRequest = projection.runtimeRequests.find( + (request) => request.id === detachedRequest.requestId, + ); + assert.equal(closedRequest?.status, "cancelled"); + assert.equal(closedRequest?.responseCapability.type, "not_resumable"); + assert.equal( + projection.nodes.find((node) => node.id === detachedRequest.nodeId)?.status, + "cancelled", + ); + assert.equal( + projection.turnItems.find( + (item) => item.type === requestType && item.requestId === detachedRequest.requestId, + )?.status, + "cancelled", + ); + const sibling = yield* projectionStore.getThreadProjection(siblingThreadId); + assert.equal(sibling.runtimeRequests[0]?.status, "pending"); + assert.equal(sibling.nodes[0]?.status, "waiting"); + assert.equal(sibling.turnItems[0]?.status, "waiting"); + assert.equal( + projection.runtimeRequests.find((request) => request.id === requests[2]!.requestId) + ?.status, + "pending", + ); + assert.isTrue(Option.isSome(yield* manager.get(providerSessionId))); + assert.equal((yield* Ref.get(state)).closeCount, 0); + assert.equal((yield* registry.resolve(token!))?.threadId, threadId); + }); + yield* effect.pipe( + Effect.provide(makeTestLayer({ state, idleTimeoutMs: 1_000, mcpConfigs })), + ); + }), +); + it.effect("ProviderSessionManagerV2 terminalizes a pending input transcript item on release", () => Effect.gen(function* () { const state = yield* Ref.make(emptyState); diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.ts index c1039c2a0207..6679c8a68477 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.ts @@ -615,6 +615,7 @@ export const layerWithOptions = ( readonly reason: ProviderSessionReleaseReason; /** Requests created later belong to a replacement session with the same id. */ readonly releasedAt: DateTime.Utc; + readonly threadIds?: ReadonlySet; }) => Effect.gen(function* () { const providerSessionId = input.entry.runtime.providerSessionId; @@ -626,7 +627,7 @@ export const layerWithOptions = ( : "Provider session was closed before this runtime request was resolved."; const events: Array = []; - for (const threadId of input.entry.attachedThreadIds) { + for (const threadId of input.threadIds ?? input.entry.attachedThreadIds) { const projection = yield* projectionStore.getThreadRecords( threadId, ["runtimeRequests", "nodes", "turnItems"], @@ -1934,6 +1935,16 @@ export const layerWithOptions = ( ); } } + // Persist request cleanup while the thread is still attached. + // Shared sessions stay live and cannot clean this thread on release. + if (currentEntry?.attachedThreadIds.has(input.threadId)) { + yield* writeReleasedRuntimeRequestEvents({ + entry: currentEntry, + reason: "manual_shutdown", + threadIds: new Set([input.threadId]), + releasedAt: yield* DateTime.now, + }); + } const detached = yield* Ref.modify(sessions, (current) => { const entry = current.get(key); if (entry === undefined || !entry.attachedThreadIds.has(input.threadId)) { From c603c29bbe5459e6495547177faf9992543da50a Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Sat, 3 Oct 2026 06:43:31 +0200 Subject: [PATCH 2/6] fix(server): serialize request cleanup with thread detach --- .../ProviderSessionManager.test.ts | 164 +++++++++++++++++- .../ProviderSessionManager.ts | 110 +++++++----- 2 files changed, 224 insertions(+), 50 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts index e100f69fa043..590896909c85 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts @@ -412,13 +412,26 @@ function makeTestLayer(input: { readonly beforeUnload?: Effect.Effect; readonly serverSettingsLayer?: ReturnType; readonly projectServiceLayer?: Layer.Layer; + readonly eventSinkLayer?: Layer.Layer< + EventSink.EventSinkV2, + Layer.Error + >; + readonly projectionStoreLayer?: Layer.Layer< + ProjectionStore.ProjectionStoreV2, + Layer.Error + >; }) { const configuredEventSinkLayer = - input.flakyReleaseWrites !== undefined + input.eventSinkLayer ?? + (input.flakyReleaseWrites !== undefined ? makeFlakyReleaseEventSinkLayer(input.flakyReleaseWrites) : input.failReleaseEventWrites ? FailingReleaseEventSinkLayer - : TestEventSinkLayer; + : TestEventSinkLayer); + const configuredStoresLayer = + input.projectionStoreLayer === undefined + ? TestStoresLayer + : Layer.merge(TestStoresLayer, input.projectionStoreLayer); const registryLayer = ProviderAdapterRegistry.makeSingleLayer( makeProviderAdapter(input.state, { failEventStream: input.failEventStream ?? false, @@ -439,13 +452,13 @@ function makeTestLayer(input: { Layer.mergeAll( configuredEventSinkLayer, IdAllocator.layer, - TestStoresLayer, + configuredStoresLayer, ThreadCommandExecutor.layer, ), ), ); return Layer.mergeAll( - TestStoresLayer, + configuredStoresLayer, configuredEventSinkLayer, IdAllocator.layer, TestMcpRegistryLayer, @@ -460,7 +473,7 @@ function makeTestLayer(input: { IdAllocator.layer, providerEventIngestorTestLayer, TestMcpRegistryLayer, - TestStoresLayer, + configuredStoresLayer, ...(input.serverSettingsLayer === undefined ? [] : [input.serverSettingsLayer]), ...(input.projectServiceLayer === undefined ? [] : [input.projectServiceLayer]), ), @@ -2752,6 +2765,147 @@ it.effect.each(["approval_request", "user_input_request"] as const)( }), ); +it.effect( + "ProviderSessionManagerV2 drains request writes before detach and rejects late artifacts", + () => + Effect.gen(function* () { + const state = yield* Ref.make(emptyState); + const incomingWriteStarted = yield* Deferred.make(); + const releaseIncomingWrite = yield* Deferred.make(); + const incomingCommitted = yield* Ref.make(false); + const requestReadStates = yield* Ref.make>([]); + const threadId = ThreadId.make("detach_request_race"); + const siblingThreadId = ThreadId.make("detach_request_race_sibling"); + const eventSinkLayer = Layer.effect( + EventSink.EventSinkV2, + Effect.gen(function* () { + const delegate = yield* EventSink.EventSinkV2; + return EventSink.EventSinkV2.of({ + ...delegate, + write: (input) => + Effect.gen(function* () { + const incomingRequest = input.events.some( + (event) => + event.threadId === threadId && + event.type === "runtime-request.updated" && + event.payload.status === "pending", + ); + if (incomingRequest && !(yield* Ref.get(incomingCommitted))) { + yield* Deferred.succeed(incomingWriteStarted, undefined); + yield* Deferred.await(releaseIncomingWrite); + } + const stored = yield* delegate.write(input); + if (incomingRequest) yield* Ref.set(incomingCommitted, true); + return stored; + }), + }); + }), + ).pipe(Layer.provide(TestEventSinkLayer)); + const projectionStoreLayer = Layer.effect( + ProjectionStore.ProjectionStoreV2, + Effect.gen(function* () { + const delegate = yield* ProjectionStore.ProjectionStoreV2; + return ProjectionStore.ProjectionStoreV2.of({ + ...delegate, + getThreadRecords: (requestedThreadId, fields, options) => { + // The fixture has no provider turns. Keep this read synchronous so + // fork's immediate start reaches the request permit before returning. + if (fields.some((field) => field === "providerThreads")) { + return Effect.succeed({ providerThreads: [], providerTurns: [] } as never); + } + if ( + requestedThreadId === threadId && + fields.some((field) => field === "runtimeRequests") + ) { + return Effect.gen(function* () { + const committed = yield* Ref.get(incomingCommitted); + yield* Ref.update(requestReadStates, (states) => [...states, committed]); + if (!committed) { + return { runtimeRequests: [], nodes: [], turnItems: [] } as never; + } + return yield* delegate.getThreadRecords(requestedThreadId, fields, options); + }); + } + return delegate.getThreadRecords(requestedThreadId, fields, options); + }, + }); + }), + ).pipe(Layer.provide(TestStoresLayer)); + const effect = Effect.gen(function* () { + const eventSink = yield* EventSink.EventSinkV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const manager = yield* ProviderSessionManager.ProviderSessionManagerV2; + const projectionStore = yield* ProjectionStore.ProjectionStoreV2; + const now = yield* DateTime.now; + const providerSessionId = idAllocator.derive.providerSession({ + providerInstanceId: modelSelection.instanceId, + }); + yield* eventSink.write({ + events: [ + yield* makeThreadCreatedEvent({ idAllocator, threadId, now }), + yield* makeThreadCreatedEvent({ idAllocator, threadId: siblingThreadId, now }), + ], + }); + const runtime = yield* manager.open({ + threadId, + providerSessionId, + modelSelection, + runtimePolicy, + }); + yield* manager.open({ + threadId: siblingThreadId, + providerSessionId, + modelSelection, + runtimePolicy, + }); + const pending = yield* makePendingRuntimeRequestEvents({ + idAllocator, + threadId, + providerSessionId, + providerThread: makeProviderThread({ idAllocator, threadId, providerSessionId, now }), + now, + }); + yield* eventSink.write({ + events: pending.events.filter((event) => event.type !== "runtime-request.updated"), + }); + const queue = (yield* Ref.get(state)).eventQueues.get(String(providerSessionId))!; + yield* Queue.offer(queue, pending.providerEvents[0]!); + yield* Deferred.await(incomingWriteStarted); + const detach = yield* manager + .detach({ providerSessionId, threadId }) + .pipe(Effect.forkScoped({ startImmediately: true })); + yield* Deferred.succeed(releaseIncomingWrite, undefined); + yield* Fiber.join(detach); + assert.deepEqual(yield* Ref.get(requestReadStates), [true]); + const subscribe = runtime.subscribeEvents; + assert.isDefined(subscribe); + if (subscribe === undefined) return; + const subscription = yield* subscribe; + // The pump processes this marker only after all late request artifacts. + yield* Queue.offerAll(queue, [ + ...pending.providerEvents, + { + type: "provider_session.updated", + driver: CODEX_DRIVER, + providerSession: runtime.providerSession, + }, + ]); + const marker = yield* subscription.events.pipe(Stream.runHead); + assert.isTrue(Option.isSome(marker)); + const projection = yield* projectionStore.getThreadProjection(threadId); + assert.equal(projection.runtimeRequests[0]?.status, "cancelled"); + assert.equal(projection.nodes[0]?.status, "cancelled"); + assert.equal(projection.turnItems[0]?.status, "cancelled"); + assert.equal((yield* Ref.get(state)).closeCount, 0); + }); + yield* effect.pipe( + Effect.provide( + makeTestLayer({ state, idleTimeoutMs: 1_000, eventSinkLayer, projectionStoreLayer }), + ), + ); + }), +); + it.effect("ProviderSessionManagerV2 terminalizes a pending input transcript item on release", () => Effect.gen(function* () { const state = yield* Ref.make(emptyState); diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.ts index 6679c8a68477..619c7462b6ed 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.ts @@ -1572,7 +1572,12 @@ export const layerWithOptions = ( const current = (yield* Ref.get(sessions)).get( sessionKey(entry.runtime.providerSessionId), ); - if (current?.runtime !== entry.runtime) return; + if ( + current?.runtime !== entry.runtime || + !current.attachedThreadIds.has(threadId) + ) { + return; + } yield* providerEventIngestor .ingestNormalized({ providerSessionId: entry.runtime.providerSessionId, @@ -1935,50 +1940,65 @@ export const layerWithOptions = ( ); } } - // Persist request cleanup while the thread is still attached. - // Shared sessions stay live and cannot clean this thread on release. - if (currentEntry?.attachedThreadIds.has(input.threadId)) { - yield* writeReleasedRuntimeRequestEvents({ - entry: currentEntry, - reason: "manual_shutdown", - threadIds: new Set([input.threadId]), - releasedAt: yield* DateTime.now, - }); - } - const detached = yield* Ref.modify(sessions, (current) => { - const entry = current.get(key); - if (entry === undefined || !entry.attachedThreadIds.has(input.threadId)) { - return [Option.none(), current] as const; - } - const attachedThreadIds = new Set(entry.attachedThreadIds); - attachedThreadIds.delete(input.threadId); - const loadedProviderThreadKeyByThread = new Map( - entry.loadedProviderThreadKeyByThread, - ); - loadedProviderThreadKeyByThread.delete(input.threadId); - // For a plain (workspace-change) detach, the credential id stays - // recorded: the thread may re-attach and reuse it, and - // releaseEntry revokes it when the provider process finally goes - // away. A terminal detach (archive/delete) prunes the record so - // nothing vetoes the revocation below. - const mcpCredentialIdByThread = - input.revokeMcpCredential === true - ? (() => { - const pruned = new Map(entry.mcpCredentialIdByThread); - pruned.delete(input.threadId); - return pruned; - })() - : entry.mcpCredentialIdByThread; - const updatedEntry = { - ...entry, - attachedThreadIds, - loadedProviderThreadKeyByThread, - mcpCredentialIdByThread, - }; - const updated = new Map(current); - updated.set(key, updatedEntry); - return [Option.some(updatedEntry), updated] as const; - }); + // Drain in-flight runless request writes, then close requests and + // remove the attachment before the pump can accept another write. + // Provider interrupts stay outside this permit: they may wait for + // an event that the pump needs to persist. + const detached = + currentEntry === undefined + ? Option.none() + : yield* Effect.gen(function* () { + const entry = (yield* Ref.get(sessions)).get(key); + if ( + entry?.runtime !== currentEntry.runtime || + !entry.attachedThreadIds.has(input.threadId) + ) { + return Option.none(); + } + yield* writeReleasedRuntimeRequestEvents({ + entry, + reason: "manual_shutdown", + threadIds: new Set([input.threadId]), + releasedAt: yield* DateTime.now, + }); + return yield* Ref.modify(sessions, (current) => { + const entry = current.get(key); + if ( + entry?.runtime !== currentEntry.runtime || + !entry.attachedThreadIds.has(input.threadId) + ) { + return [Option.none(), current] as const; + } + const attachedThreadIds = new Set(entry.attachedThreadIds); + attachedThreadIds.delete(input.threadId); + const loadedProviderThreadKeyByThread = new Map( + entry.loadedProviderThreadKeyByThread, + ); + loadedProviderThreadKeyByThread.delete(input.threadId); + // For a plain (workspace-change) detach, the credential id stays + // recorded: the thread may re-attach and reuse it, and + // releaseEntry revokes it when the provider process finally goes + // away. A terminal detach (archive/delete) prunes the record so + // nothing vetoes the revocation below. + const mcpCredentialIdByThread = + input.revokeMcpCredential === true + ? (() => { + const pruned = new Map(entry.mcpCredentialIdByThread); + pruned.delete(input.threadId); + return pruned; + })() + : entry.mcpCredentialIdByThread; + const updatedEntry = { + ...entry, + attachedThreadIds, + loadedProviderThreadKeyByThread, + mcpCredentialIdByThread, + }; + const updated = new Map(current); + updated.set(key, updatedEntry); + return [Option.some(updatedEntry), updated] as const; + }); + }).pipe(currentEntry.requestEventPermit.withPermits(1)); // Plain detaches deliberately do not revoke: a detached thread's // provider process may still be alive (shared multi-thread codex // session across a workspace handoff) and holds its MCP client's From 19c6c60c4cc826b12e9c3c7f8bdde1ed99ca4709 Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Sat, 3 Oct 2026 12:22:06 +0200 Subject: [PATCH 3/6] fix(server): select orphan request transcripts during recovery --- .../FoundationPersistence.test.ts | 197 ++++++++++++++++++ .../src/orchestration-v2/ProjectionStore.ts | 14 ++ 2 files changed, 211 insertions(+) diff --git a/apps/server/src/orchestration-v2/FoundationPersistence.test.ts b/apps/server/src/orchestration-v2/FoundationPersistence.test.ts index dc7eacc59ba5..2c4cd3e41176 100644 --- a/apps/server/src/orchestration-v2/FoundationPersistence.test.ts +++ b/apps/server/src/orchestration-v2/FoundationPersistence.test.ts @@ -22,6 +22,7 @@ import { ProviderThreadId, RunAttemptId, RunId, + RuntimeRequestId, ThreadId, TurnItemId, } from "@t3tools/contracts"; @@ -338,6 +339,202 @@ it.effect("keeps other database work runnable while discovering compaction candi ).pipe(Effect.provide(TestLayer)), ); +it.effect.each([ + { trigger: "startup", runStatus: null }, + { trigger: "shutdown", runStatus: null }, + { trigger: "startup", runStatus: "completed" }, + { trigger: "shutdown", runStatus: "completed" }, +] as const)( + "recovers request transcripts through SQLite on $trigger with run status $runStatus", + ({ trigger, runStatus }) => + Effect.gen(function* () { + const eventSink = yield* EventSink.EventSinkV2; + const projections = yield* ProjectionStore.ProjectionStoreV2; + const now = yield* DateTime.now; + const threadId = ThreadId.make(`thread:request-transcript:${trigger}:${runStatus}`); + const runId = runStatus === null ? null : RunId.make(`run:${threadId}`); + const events: Array = [ + threadCreatedEvent({ + id: `event:${threadId}:thread`, + thread: makeThread(threadId, now), + now, + }), + ]; + if (runId !== null) { + events.push({ + id: EventId.make(`event:${threadId}:run`), + type: "run.created", + threadId, + runId, + occurredAt: now, + payload: { + id: runId, + threadId, + ordinal: 1, + providerInstanceId, + modelSelection, + providerThreadId: null, + userMessageId: MessageId.make(`message:${threadId}`), + rootNodeId: null, + activeAttemptId: null, + status: "completed", + requestedAt: now, + startedAt: now, + completedAt: now, + checkpointId: null, + contextHandoffId: null, + }, + }); + } + for (const kind of ["approval", "question", "durable-question"] as const) { + const requestId = RuntimeRequestId.make(`request:${threadId}:${kind}`); + const nodeId = NodeId.make(`node:${threadId}:${kind}`); + const itemType = kind === "approval" ? "approval_request" : "user_input_request"; + const itemBase = { + id: TurnItemId.make(`item:${threadId}:${kind}`), + threadId, + runId, + nodeId, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: events.length, + status: "waiting" as const, + title: null, + startedAt: now, + completedAt: null, + updatedAt: now, + requestId, + }; + const item: OrchestrationV2TurnItem = + kind === "approval" + ? { + ...itemBase, + type: "approval_request", + requestKind: "command", + prompt: "Run command?", + } + : { ...itemBase, type: "user_input_request", questions: [] }; + events.push( + { + id: EventId.make(`event:${nodeId}`), + type: "node.updated", + threadId, + nodeId, + occurredAt: now, + payload: { + id: nodeId, + threadId, + runId, + parentNodeId: null, + rootNodeId: nodeId, + kind: itemType, + status: "waiting", + countsForRun: false, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + runtimeRequestId: requestId, + checkpointScopeId: null, + startedAt: now, + completedAt: null, + }, + }, + { + id: EventId.make(`event:${requestId}`), + type: "runtime-request.updated", + threadId, + occurredAt: now, + payload: { + id: requestId, + nodeId, + providerTurnId: null, + nativeRequestRef: null, + kind: kind === "approval" ? "command" : "user_input", + status: "pending", + responseCapability: + kind === "durable-question" + ? { type: "message" } + : { + type: "live", + providerSessionId: ProviderSessionId.make(`session:${threadId}`), + }, + createdAt: now, + resolvedAt: null, + }, + }, + { + id: EventId.make(`event:${item.id}`), + type: "turn-item.updated", + threadId, + occurredAt: now, + payload: item, + }, + ); + if (kind === "approval") { + events.push({ + id: EventId.make(`event:${item.id}:terminal`), + type: "turn-item.updated", + threadId, + occurredAt: now, + payload: { + ...item, + id: TurnItemId.make(`item:${threadId}:terminal`), + status: "completed", + completedAt: now, + }, + }); + } + } + yield* eventSink.commitCommand({ + commandId: CommandId.make(`command:${threadId}:seed`), + threadId, + commandType: "foundation.request-transcript", + acceptedAt: now, + events, + effects: [], + }); + const selected = yield* projections.getRuntimeRecoveryProjection(threadId); + assert.sameMembers( + selected.nodes.map((node) => node.id), + [NodeId.make(`node:${threadId}:approval`), NodeId.make(`node:${threadId}:question`)], + ); + assert.sameMembers( + selected.turnItems.map((item) => item.id), + [ + TurnItemId.make(`item:${threadId}:approval`), + TurnItemId.make(`item:${threadId}:question`), + ], + ); + const recovery = yield* ProviderRuntimeRecovery.make.pipe( + Effect.provide(ServerSettings.layerTest()), + ); + assert.equal((yield* recovery.reconcile(trigger)).closedRequests, 2); + const final = yield* projections.getThreadProjection(threadId); + for (const kind of ["approval", "question", "durable-question"] as const) { + const durable = kind === "durable-question"; + assert.equal( + final.runtimeRequests.find((request) => request.id === `request:${threadId}:${kind}`) + ?.status, + durable ? "pending" : trigger === "startup" ? "expired" : "cancelled", + ); + assert.equal( + final.nodes.find((node) => node.id === `node:${threadId}:${kind}`)?.status, + durable ? "waiting" : "cancelled", + ); + assert.equal( + final.turnItems.find((item) => item.id === `item:${threadId}:${kind}`)?.status, + durable ? "waiting" : "cancelled", + ); + } + assert.equal( + final.turnItems.find((item) => item.id === `item:${threadId}:terminal`)?.status, + "completed", + ); + }).pipe(Effect.provide(Layer.fresh(TestLayer))), +); + it.layer(TestLayer)("orchestration V2 foundation persistence", (it) => { it.effect("projects oversized tool bodies before both replay and live RPC retention", () => Effect.scoped( diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index ac0624b05806..4a138dc4f013 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -3843,6 +3843,11 @@ export const layer: Layer.Layer = WHERE item.thread_id = ${threadId} AND item.type = 'subagent' AND item.status IN ('pending', 'running', 'waiting') ) + OR node.node_id IN ( + SELECT request.node_id FROM orchestration_v2_projection_runtime_requests AS request + WHERE request.thread_id = ${threadId} AND request.status = 'pending' + AND json_extract(request.payload_json, '$.responseCapability.type') <> 'message' + ) ) ORDER BY COALESCE(node.started_at, ''), node.node_id ASC `, @@ -3963,6 +3968,15 @@ export const layer: Layer.Layer = AND status IN ('queued', 'preparing', 'starting', 'running', 'waiting') ) OR item.type IN ('command_execution', 'dynamic_tool', 'subagent') + OR ( + item.type IN ('approval_request', 'user_input_request') + AND json_extract(item.payload_json, '$.requestId') IN ( + SELECT request.runtime_request_id + FROM orchestration_v2_projection_runtime_requests AS request + WHERE request.thread_id = ${threadId} AND request.status = 'pending' + AND json_extract(request.payload_json, '$.responseCapability.type') <> 'message' + ) + ) OR ( item.run_id IS NULL AND item.node_id IN ( From af0f5c82d98c00390bbdfbf572e760677bf5d9c6 Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Sun, 4 Oct 2026 01:05:28 +0200 Subject: [PATCH 4/6] fix(server): explain request cancellation on thread detach --- .../ProviderSessionManager.test.ts | 8 +++- .../ProviderSessionManager.ts | 38 +++++++++++-------- 2 files changed, 28 insertions(+), 18 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts index 590896909c85..f0f4090b757b 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts @@ -2726,7 +2726,7 @@ it.effect.each(["approval_request", "user_input_request"] as const)( modelSelection, runtimePolicy, }); - yield* manager.detach({ providerSessionId, threadId }); + yield* manager.detach({ providerSessionId, threadId, detail: "Workspace changed." }); // A duplicate detach must leave the sibling and replacement session alone. yield* manager.detach({ providerSessionId, threadId }); const projection = yield* projectionStore.getThreadProjection(threadId); @@ -2735,7 +2735,11 @@ it.effect.each(["approval_request", "user_input_request"] as const)( (request) => request.id === detachedRequest.requestId, ); assert.equal(closedRequest?.status, "cancelled"); - assert.equal(closedRequest?.responseCapability.type, "not_resumable"); + assert.deepEqual(closedRequest?.responseCapability, { + type: "not_resumable", + reason: + "Thread detached from the provider session before this runtime request was resolved.", + }); assert.equal( projection.nodes.find((node) => node.id === detachedRequest.nodeId)?.status, "cancelled", diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.ts index 619c7462b6ed..29a04bfe7e07 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.ts @@ -4,7 +4,6 @@ import { OrchestrationV2DomainEvent, OrchestrationV2ProviderSession, type OrchestrationV2ProviderThread, - OrchestrationV2RuntimeRequest, ProviderInstanceId, ProviderSessionId, ThreadId, @@ -236,7 +235,7 @@ function releaseStatusFor( function releasedRuntimeRequestStatusFor( reason: ProviderSessionReleaseReason, -): OrchestrationV2RuntimeRequest["status"] { +): "cancelled" | "expired" { return reason === "manual_shutdown" || reason === "server_shutdown" ? "cancelled" : "expired"; } @@ -612,7 +611,9 @@ export const layerWithOptions = ( const writeReleasedRuntimeRequestEvents = (input: { readonly entry: LiveSessionEntry; - readonly reason: ProviderSessionReleaseReason; + readonly status: "cancelled" | "expired"; + readonly artifactStatus: "cancelled" | "failed"; + readonly reason: string; /** Requests created later belong to a replacement session with the same id. */ readonly releasedAt: DateTime.Utc; readonly threadIds?: ReadonlySet; @@ -620,11 +621,6 @@ export const layerWithOptions = ( Effect.gen(function* () { const providerSessionId = input.entry.runtime.providerSessionId; const now = yield* DateTime.now; - const status = releasedRuntimeRequestStatusFor(input.reason); - const reason = - input.reason === "runtime_error" - ? "Provider session failed before this runtime request was resolved." - : "Provider session was closed before this runtime request was resolved."; const events: Array = []; for (const threadId of input.threadIds ?? input.entry.attachedThreadIds) { @@ -654,10 +650,10 @@ export const layerWithOptions = ( occurredAt: now, payload: { ...request, - status, + status: input.status, responseCapability: { type: "not_resumable", - reason, + reason: input.reason, }, resolvedAt: now, }, @@ -678,7 +674,7 @@ export const layerWithOptions = ( occurredAt: now, payload: { ...requestNode, - status: input.reason === "runtime_error" ? "failed" : "cancelled", + status: input.artifactStatus, completedAt: now, }, }); @@ -703,7 +699,7 @@ export const layerWithOptions = ( occurredAt: now, payload: { ...turnItem, - status: input.reason === "runtime_error" ? "failed" : "cancelled", + status: input.artifactStatus, completedAt: now, updatedAt: now, }, @@ -734,9 +730,16 @@ export const layerWithOptions = ( ? Effect.succeed(Exit.void) : Effect.exit(writeReleasedSessionEvents(input)), Effect.exit( - writeReleasedRuntimeRequestEvents(input).pipe( - input.entry.requestEventPermit.withPermits(1), - ), + writeReleasedRuntimeRequestEvents({ + entry: input.entry, + releasedAt: input.releasedAt, + status: releasedRuntimeRequestStatusFor(input.reason), + artifactStatus: input.reason === "runtime_error" ? "failed" : "cancelled", + reason: + input.reason === "runtime_error" + ? "Provider session failed before this runtime request was resolved." + : "Provider session was closed before this runtime request was resolved.", + }).pipe(input.entry.requestEventPermit.withPermits(1)), ), ], { concurrency: 1 }, @@ -1957,7 +1960,10 @@ export const layerWithOptions = ( } yield* writeReleasedRuntimeRequestEvents({ entry, - reason: "manual_shutdown", + status: "cancelled", + artifactStatus: "cancelled", + reason: + "Thread detached from the provider session before this runtime request was resolved.", threadIds: new Set([input.threadId]), releasedAt: yield* DateTime.now, }); From b47359db8971ea9d17db117a67547017504d0273 Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Sun, 4 Oct 2026 13:38:35 +0200 Subject: [PATCH 5/6] fix(server): reject detached turn-scoped request artifacts --- .../ProviderSessionManager.test.ts | 24 ++++++++++++++++ .../ProviderSessionManager.ts | 28 +++++++++++++++++-- 2 files changed, 50 insertions(+), 2 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts index f0f4090b757b..2753050417ae 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts @@ -13,6 +13,8 @@ import { ProjectId, ProviderDriverKind, ProviderInstanceId, + ProviderTurnId, + RunId, type ProviderSessionId, ThreadId, } from "@t3tools/contracts"; @@ -2888,6 +2890,27 @@ it.effect( // The pump processes this marker only after all late request artifacts. yield* Queue.offerAll(queue, [ ...pending.providerEvents, + ...pending.providerEvents.map((event): ProviderAdapterV2Event => { + switch (event.type) { + case "runtime_request.updated": + return { + ...event, + runtimeRequest: { + ...event.runtimeRequest, + providerTurnId: ProviderTurnId.make("detached-turn"), + }, + }; + case "node.updated": + return { ...event, node: { ...event.node, runId: RunId.make("detached-run") } }; + case "turn_item.updated": + return { + ...event, + turnItem: { ...event.turnItem, runId: RunId.make("detached-run") }, + }; + default: + return event; + } + }), { type: "provider_session.updated", driver: CODEX_DRIVER, @@ -2896,6 +2919,7 @@ it.effect( ]); const marker = yield* subscription.events.pipe(Stream.runHead); assert.isTrue(Option.isSome(marker)); + if (Option.isSome(marker)) assert.equal(marker.value.type, "provider_session.updated"); const projection = yield* projectionStore.getThreadProjection(threadId); assert.equal(projection.runtimeRequests[0]?.status, "cancelled"); assert.equal(projection.nodes[0]?.status, "cancelled"); diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.ts index 29a04bfe7e07..52472b4e8acb 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.ts @@ -266,6 +266,22 @@ function sessionScopedRuntimeRequestThreadId(event: ProviderAdapterV2Event): Thr } } +function runtimeRequestArtifactThreadId(event: ProviderAdapterV2Event): ThreadId | undefined { + switch (event.type) { + case "runtime_request.updated": + return event.threadId; + case "node.updated": + return event.node.runtimeRequestId !== null ? event.node.threadId : undefined; + case "turn_item.updated": + return event.turnItem.type === "approval_request" || + event.turnItem.type === "user_input_request" + ? event.turnItem.threadId + : undefined; + default: + return undefined; + } +} + function providerThreadRuntimeKey( providerThread: Parameters[0]["providerThread"], ): string { @@ -1570,17 +1586,25 @@ export const layerWithOptions = ( // their runless request artifacts directly so the normal T3 // request UI can answer them and unblock session setup. const threadId = sessionScopedRuntimeRequestThreadId(event); - if (threadId !== undefined) { + const requestThreadId = runtimeRequestArtifactThreadId(event); + if (requestThreadId !== undefined) { yield* Effect.gen(function* () { const current = (yield* Ref.get(sessions)).get( sessionKey(entry.runtime.providerSessionId), ); if ( current?.runtime !== entry.runtime || - !current.attachedThreadIds.has(threadId) + !current.attachedThreadIds.has(requestThreadId) ) { return; } + if (threadId === undefined) { + yield* publishToSubscribers(entry.eventSubscribers, { + type: "event", + event, + }); + return; + } yield* providerEventIngestor .ingestNormalized({ providerSessionId: entry.runtime.providerSessionId, From 95fc18a898d9baf4b88bca06ae012e2e643a5a76 Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Sun, 4 Oct 2026 13:39:31 +0200 Subject: [PATCH 6/6] test(server): prove detach cleanup retries after commit failures --- .../ProviderSessionManager.test.ts | 20 ++++++++++++++++++- docs/user/permission-modes.md | 5 +++++ 2 files changed, 24 insertions(+), 1 deletion(-) diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts index 2753050417ae..6c5456f59671 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts @@ -2663,6 +2663,11 @@ it.effect.each(["approval_request", "user_input_request"] as const)( (requestType) => Effect.gen(function* () { const state = yield* Ref.make(emptyState); + const flaky: FlakyReleaseWrites = { + failing: yield* Ref.make<"none" | "session" | "session-and-requests">("none"), + failures: yield* Queue.unbounded(), + }; + const mcpConfigs = yield* Ref.make< ReadonlyArray >([]); @@ -2728,6 +2733,17 @@ it.effect.each(["approval_request", "user_input_request"] as const)( modelSelection, runtimePolicy, }); + yield* Ref.set(flaky.failing, "session-and-requests"); + assert.isTrue( + Exit.isFailure(yield* Effect.exit(manager.detach({ providerSessionId, threadId }))), + ); + yield* Queue.take(flaky.failures); + assert.equal( + (yield* projectionStore.getThreadProjection(threadId)).runtimeRequests[0]?.status, + "pending", + ); + assert.equal((yield* Ref.get(state)).closeCount, 0); + yield* Ref.set(flaky.failing, "none"); yield* manager.detach({ providerSessionId, threadId, detail: "Workspace changed." }); // A duplicate detach must leave the sibling and replacement session alone. yield* manager.detach({ providerSessionId, threadId }); @@ -2766,7 +2782,9 @@ it.effect.each(["approval_request", "user_input_request"] as const)( assert.equal((yield* registry.resolve(token!))?.threadId, threadId); }); yield* effect.pipe( - Effect.provide(makeTestLayer({ state, idleTimeoutMs: 1_000, mcpConfigs })), + Effect.provide( + makeTestLayer({ state, idleTimeoutMs: 1_000, mcpConfigs, flakyReleaseWrites: flaky }), + ), ); }), ); diff --git a/docs/user/permission-modes.md b/docs/user/permission-modes.md index c7e45519acdb..ccddda662ac0 100644 --- a/docs/user/permission-modes.md +++ b/docs/user/permission-modes.md @@ -18,6 +18,11 @@ and modes you choose in a draft keep their permissions. Approve or reject requests in the conversation to let the agent continue. Permission modes do not prevent the agent from asking questions about the task. +A provider approval or question tied to a live process is cancelled when its +thread detaches, such as after changing workspaces, settling the thread, or +switching its provider model or runtime mode. Start or resume a turn to receive +a new prompt. Questions answered through a message keep their resumable behavior. + ## Provider differences Providers enforce permissions differently. Some read-only actions can proceed in **Supervised**.