diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts index 398337b22e2c..fa1a28685db4 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts @@ -4975,16 +4975,18 @@ export const makeCodexAdapterV2 = Effect.fn("makeCodexAdapterV2")(function* ( driver: CODEX_PROVIDER, node: artifacts.node, }); + // The request is answerable once it is recorded, so its node and item + // land first and an answer always finds them to settle. yield* emitProviderEvent({ - type: "runtime_request.updated", + type: "turn_item.updated", driver: CODEX_PROVIDER, - threadId: artifacts.node.threadId, - runtimeRequest: artifacts.request, + turnItem: artifacts.turnItem, }); yield* emitProviderEvent({ - type: "turn_item.updated", + type: "runtime_request.updated", driver: CODEX_PROVIDER, - turnItem: artifacts.turnItem, + threadId: artifacts.node.threadId, + runtimeRequest: artifacts.request, }); yield* Ref.update(finalAnswerItemIdsByTurn, (current) => { const ids = current.get(payload.turnId); @@ -5082,16 +5084,18 @@ export const makeCodexAdapterV2 = Effect.fn("makeCodexAdapterV2")(function* ( driver: CODEX_PROVIDER, node: artifacts.node, }); + // The request is answerable once it is recorded, so its node and item + // land first and an answer always finds them to settle. yield* emitProviderEvent({ - type: "runtime_request.updated", + type: "turn_item.updated", driver: CODEX_PROVIDER, - threadId: artifacts.node.threadId, - runtimeRequest: artifacts.request, + turnItem: artifacts.turnItem, }); yield* emitProviderEvent({ - type: "turn_item.updated", + type: "runtime_request.updated", driver: CODEX_PROVIDER, - turnItem: artifacts.turnItem, + threadId: artifacts.node.threadId, + runtimeRequest: artifacts.request, }); const resolved = yield* Deferred.await(decision).pipe( @@ -5142,16 +5146,18 @@ export const makeCodexAdapterV2 = Effect.fn("makeCodexAdapterV2")(function* ( driver: CODEX_PROVIDER, node: artifacts.node, }); + // The request is answerable once it is recorded, so its node and item + // land first and an answer always finds them to settle. yield* emitProviderEvent({ - type: "runtime_request.updated", + type: "turn_item.updated", driver: CODEX_PROVIDER, - threadId: artifacts.node.threadId, - runtimeRequest: artifacts.request, + turnItem: artifacts.turnItem, }); yield* emitProviderEvent({ - type: "turn_item.updated", + type: "runtime_request.updated", driver: CODEX_PROVIDER, - turnItem: artifacts.turnItem, + threadId: artifacts.node.threadId, + runtimeRequest: artifacts.request, }); const resolved = yield* Deferred.await(decision).pipe( @@ -5203,16 +5209,18 @@ export const makeCodexAdapterV2 = Effect.fn("makeCodexAdapterV2")(function* ( driver: CODEX_PROVIDER, node: artifacts.node, }); + // The request is answerable once it is recorded, so its node and item + // land first and an answer always finds them to settle. yield* emitProviderEvent({ - type: "runtime_request.updated", + type: "turn_item.updated", driver: CODEX_PROVIDER, - threadId: artifacts.node.threadId, - runtimeRequest: artifacts.request, + turnItem: artifacts.turnItem, }); yield* emitProviderEvent({ - type: "turn_item.updated", + type: "runtime_request.updated", driver: CODEX_PROVIDER, - turnItem: artifacts.turnItem, + threadId: artifacts.node.threadId, + runtimeRequest: artifacts.request, }); const resolved = yield* Deferred.await(decision).pipe( @@ -5288,16 +5296,18 @@ export const makeCodexAdapterV2 = Effect.fn("makeCodexAdapterV2")(function* ( driver: CODEX_PROVIDER, node: artifacts.node, }); + // The request is answerable once it is recorded, so its node and item + // land first and an answer always finds them to settle. yield* emitProviderEvent({ - type: "runtime_request.updated", + type: "turn_item.updated", driver: CODEX_PROVIDER, - threadId: artifacts.node.threadId, - runtimeRequest: artifacts.request, + turnItem: artifacts.turnItem, }); yield* emitProviderEvent({ - type: "turn_item.updated", + type: "runtime_request.updated", driver: CODEX_PROVIDER, - turnItem: artifacts.turnItem, + threadId: artifacts.node.threadId, + runtimeRequest: artifacts.request, }); const resolved = yield* Deferred.await(decision).pipe( @@ -5347,16 +5357,18 @@ export const makeCodexAdapterV2 = Effect.fn("makeCodexAdapterV2")(function* ( driver: CODEX_PROVIDER, node: artifacts.node, }); + // The request is answerable once it is recorded, so its node and item + // land first and an answer always finds them to settle. yield* emitProviderEvent({ - type: "runtime_request.updated", + type: "turn_item.updated", driver: CODEX_PROVIDER, - threadId: artifacts.node.threadId, - runtimeRequest: artifacts.request, + turnItem: artifacts.turnItem, }); yield* emitProviderEvent({ - type: "turn_item.updated", + type: "runtime_request.updated", driver: CODEX_PROVIDER, - turnItem: artifacts.turnItem, + threadId: artifacts.node.threadId, + runtimeRequest: artifacts.request, }); const resolved = yield* Deferred.await(decision).pipe( @@ -5407,16 +5419,18 @@ export const makeCodexAdapterV2 = Effect.fn("makeCodexAdapterV2")(function* ( driver: CODEX_PROVIDER, node: artifacts.node, }); + // The request is answerable once it is recorded, so its node and item + // land first and an answer always finds them to settle. yield* emitProviderEvent({ - type: "runtime_request.updated", + type: "turn_item.updated", driver: CODEX_PROVIDER, - threadId: artifacts.node.threadId, - runtimeRequest: artifacts.request, + turnItem: artifacts.turnItem, }); yield* emitProviderEvent({ - type: "turn_item.updated", + type: "runtime_request.updated", driver: CODEX_PROVIDER, - turnItem: artifacts.turnItem, + threadId: artifacts.node.threadId, + runtimeRequest: artifacts.request, }); const resolved = yield* Deferred.await(decision).pipe( @@ -5465,16 +5479,18 @@ export const makeCodexAdapterV2 = Effect.fn("makeCodexAdapterV2")(function* ( driver: CODEX_PROVIDER, node: artifacts.node, }); + // The request is answerable once it is recorded, so its node and item + // land first and an answer always finds them to settle. yield* emitProviderEvent({ - type: "runtime_request.updated", + type: "turn_item.updated", driver: CODEX_PROVIDER, - threadId: artifacts.node.threadId, - runtimeRequest: artifacts.request, + turnItem: artifacts.turnItem, }); yield* emitProviderEvent({ - type: "turn_item.updated", + type: "runtime_request.updated", driver: CODEX_PROVIDER, - turnItem: artifacts.turnItem, + threadId: artifacts.node.threadId, + runtimeRequest: artifacts.request, }); const resolved = yield* Deferred.await(answers).pipe( diff --git a/apps/server/src/orchestration-v2/EventSink.ts b/apps/server/src/orchestration-v2/EventSink.ts index 039d1872da6f..a36711999968 100644 --- a/apps/server/src/orchestration-v2/EventSink.ts +++ b/apps/server/src/orchestration-v2/EventSink.ts @@ -84,6 +84,18 @@ export interface EventSinkV2Shape { }) => Effect.Effect, EventSinkV2Error>; readonly writeIfRunCurrent: (input: { readonly guardPendingUserInputCancellations?: boolean; + /** + * Guards approvals and every other request kind too. For a writer that + * cancels requests an earlier projection read showed pending; a + * provider's own approval cancellation stays authoritative. + */ + readonly guardPendingRequestCancellations?: boolean; + /** + * Drops subagent, node and turn-item updates for rows that ended after an + * earlier projection read showed them open. For a writer that settles + * work it no longer receives events for. + */ + readonly guardSettledWork?: boolean; readonly commandId?: CommandId; readonly threadId: ThreadId; readonly runId: RunId; @@ -263,14 +275,65 @@ const layerBase: Layer.Layer< // A user can answer after terminal normalization reads the pending request. // Recheck inside the write transaction so stale cleanup cannot erase answers. - const guardUserInputCancellations = (events: ReadonlyArray) => + // `allKinds` extends this from questions to approvals and every other kind. + const guardRequestCancellations = ( + events: ReadonlyArray, + allKinds: boolean, + ) => Effect.gen(function* () { const staleRequests = new Set(); const staleNodes = new Set(); + // A run's cleanup settles a request's node and item with the run's own + // status, which is not always `cancelled`. + const settles = (status: string) => + status === "cancelled" || (allKinds && (status === "failed" || status === "interrupted")); + // Whether a request was answered after the cleanup read it. A missing + // request leaves nothing to protect. + const answered = new Map(); + const wasAnswered = (threadId: ThreadId, requestId: RuntimeRequestId) => + Effect.gen(function* () { + const known = answered.get(requestId); + if (known !== undefined) return known; + const current = yield* projectionStore.getRuntimeRequest(threadId, requestId); + const result = + current !== undefined && + (current.status !== "pending" || current.responseCapability.type === "message"); + answered.set(requestId, result); + return result; + }); for (const event of events) { + // The cleanup reads nodes, requests and items at different moments, + // so it can settle a request's node or item after the request was + // answered and left its pending read. Each one is checked against + // its own request, with or without a cancellation beside it. + if (allKinds && event.type === "node.updated") { + const requestId = event.payload.runtimeRequestId; + if ( + requestId !== null && + settles(event.payload.status) && + (yield* wasAnswered(event.payload.threadId, requestId)) + ) { + staleNodes.add(event.payload.id); + } + continue; + } + if ( + allKinds && + event.type === "turn-item.updated" && + (event.payload.type === "user_input_request" || + event.payload.type === "approval_request") + ) { + if ( + settles(event.payload.status) && + (yield* wasAnswered(event.payload.threadId, event.payload.requestId)) + ) { + staleRequests.add(event.payload.requestId); + } + continue; + } if ( event.type !== "runtime-request.updated" || - event.payload.kind !== "user_input" || + (!allKinds && event.payload.kind !== "user_input") || event.payload.status !== "cancelled" ) continue; @@ -280,7 +343,7 @@ const layerBase: Layer.Layer< ); if ( current?.status !== "pending" || - current.kind !== "user_input" || + current.kind !== event.payload.kind || current.providerTurnId !== event.payload.providerTurnId || current.responseCapability.type === "message" ) { @@ -293,11 +356,12 @@ const layerBase: Layer.Layer< case "runtime-request.updated": return event.payload.status !== "cancelled" || !staleRequests.has(event.payload.id); case "node.updated": - return event.payload.status !== "cancelled" || !staleNodes.has(event.payload.id); + return !settles(event.payload.status) || !staleNodes.has(event.payload.id); case "turn-item.updated": return ( - event.payload.type !== "user_input_request" || - event.payload.status !== "cancelled" || + (event.payload.type !== "user_input_request" && + event.payload.type !== "approval_request") || + !settles(event.payload.status) || !staleRequests.has(event.payload.requestId) ); default: @@ -305,6 +369,68 @@ const layerBase: Layer.Layer< } }); }); + // A cleanup settles subagents, nodes and items it read as open. One that + // ended on its own before the commit keeps its own outcome and result. + const guardSettledWork = (events: ReadonlyArray) => + Effect.gen(function* () { + const ended = (status: string) => + status === "completed" || + status === "failed" || + status === "cancelled" || + status === "interrupted"; + const threads = new Map< + ThreadId, + ProjectionStore.ProjectionRecords<"nodes" | "subagents"> | undefined + >(); + const records = (threadId: ThreadId) => + Effect.gen(function* () { + if (threads.has(threadId)) return threads.get(threadId); + const read = yield* projectionStore + .getThreadRecords(threadId, ["nodes", "subagents"]) + .pipe( + Effect.catchTags({ + ProjectionStoreThreadNotFoundError: () => Effect.succeed(undefined), + }), + ); + threads.set(threadId, read); + return read; + }); + const kept: Array = []; + for (const event of events) { + const current = + event.type === "subagent.updated" + ? (yield* records(event.payload.threadId))?.subagents.find( + (row) => row.id === event.payload.id, + ) + : event.type === "node.updated" + ? (yield* records(event.payload.threadId))?.nodes.find( + (row) => row.id === event.payload.id, + ) + : event.type === "turn-item.updated" + ? yield* projectionStore.getTurnItem({ + threadId: event.payload.threadId, + itemId: event.payload.id, + }) + : undefined; + if (current == null || !ended(current.status)) kept.push(event); + } + return kept; + }); + const guardCancellations = (input: { + readonly guardPendingUserInputCancellations?: boolean; + readonly guardPendingRequestCancellations?: boolean; + readonly guardSettledWork?: boolean; + readonly events: ReadonlyArray; + }) => + Effect.gen(function* () { + const events: ReadonlyArray = + input.guardPendingRequestCancellations === true + ? yield* guardRequestCancellations(input.events, true) + : input.guardPendingUserInputCancellations === true + ? yield* guardRequestCancellations(input.events, false) + : input.events; + return input.guardSettledWork === true ? yield* guardSettledWork(events) : events; + }); const normalizeEvents = (events: ReadonlyArray) => { const runOrdinals = new Map( @@ -370,11 +496,7 @@ const layerBase: Layer.Layer< return yield* commitThenPublish( Effect.gen(function* () { - const normalized = yield* normalizeEvents( - input.guardPendingUserInputCancellations === true - ? yield* guardUserInputCancellations(input.events) - : input.events, - ); + const normalized = yield* normalizeEvents(yield* guardCancellations(input)); const committed = yield* eventStore.append({ ...(input.commandId === undefined ? {} : { commandId: input.commandId }), events: normalized, @@ -428,11 +550,7 @@ const layerBase: Layer.Layer< }; } - const normalized = yield* normalizeEvents( - input.guardPendingUserInputCancellations === true - ? yield* guardUserInputCancellations(input.events) - : input.events, - ); + const normalized = yield* normalizeEvents(yield* guardCancellations(input)); const storedEvents = yield* eventStore.append({ ...(input.commandId === undefined ? {} : { commandId: input.commandId }), events: normalized, @@ -493,11 +611,7 @@ const layerBase: Layer.Layer< }; } - const normalized = yield* normalizeEvents( - input.guardPendingUserInputCancellations === true - ? yield* guardUserInputCancellations(input.events) - : input.events, - ); + const normalized = yield* normalizeEvents(yield* guardCancellations(input)); const storedEvents = yield* eventStore.append({ ...(input.commandId === undefined ? {} : { commandId: input.commandId }), events: normalized, diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts index d6c12fab7ee1..2f3678f944f7 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts @@ -2,18 +2,28 @@ import { expect, it, vi } from "vite-plus/test"; import { it as effectIt } from "@effect/vitest"; import { CheckpointScopeId, + EventId, MessageId, NodeId, ProviderSessionId, ProviderThreadId, + ProviderTurnId, ProviderDriverKind, ProviderInstanceId, ProviderSetupError, RunAttemptId, RunId, + RuntimeRequestId, ThreadId, + TurnItemId, ProjectId, + type OrchestrationV2AppThread, + type OrchestrationV2ExecutionNode, + type OrchestrationV2Run, + type OrchestrationV2RuntimeRequest, + type OrchestrationV2Subagent, type OrchestrationV2ThreadProjection, + type OrchestrationV2TurnItem, OrchestrationV2DomainEvent, } from "@t3tools/contracts"; import * as Cause from "effect/Cause"; @@ -26,9 +36,11 @@ import * as Schema from "effect/Schema"; import * as GitWorkflow from "../git/GitWorkflowService.ts"; import * as ProjectService from "../project/ProjectService.ts"; +import * as SqlitePersistence from "../persistence/Sqlite.ts"; import * as ProviderAuthService from "../provider/ProviderAuthService.ts"; import * as ContextHandoffService from "./ContextHandoffService.ts"; import * as EventSink from "./EventSink.ts"; +import * as EventStore from "./EventStore.ts"; import * as IdAllocator from "@t3tools/provider-core/server/IdAllocator"; import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; import * as ProjectionStore from "./ProjectionStore.ts"; @@ -513,6 +525,7 @@ function makeLocalCommandHarness(input: { (m.text.trim().toLowerCase() !== "/compact" || m.attachments.length > 0), ), }), + getThreadRecords: () => Effect.succeed(projection as never), getTurnStartHistory: () => Effect.fail( new ProjectionStore.ProjectionStoreReadError({ @@ -855,3 +868,557 @@ for (const previousMessages of [[], ["/compact", " /COMPACT "]]) { }), ); } + +// Runs the start path against the real event sink and projection store, so the +// failure write goes through the same commit guard production uses. +function makePersistedStartFailureHarness() { + const now = DateTime.makeUnsafe("2026-10-06T12:00:00Z"); + const driver = ProviderDriverKind.make("codex"); + const instanceId = ProviderInstanceId.make("codex-start-failure"); + const threadId = ThreadId.make("thread-start-failure"); + const runId = RunId.make("run-start-failure"); + const attemptId = RunAttemptId.make("attempt-start-failure"); + const rootNodeId = NodeId.make("root-start-failure"); + const providerThreadId = ProviderThreadId.make("provider-thread-start-failure"); + const providerSessionId = ProviderSessionId.make("provider-session-start-failure"); + const providerTurnId = ProviderTurnId.make("provider-turn-start-failure"); + const checkpointScopeId = CheckpointScopeId.make("scope-start-failure"); + const messageId = MessageId.make("message-start-failure"); + const approvalNodeId = NodeId.make("node-approval-start-failure"); + const approvalRequestId = RuntimeRequestId.make("request-approval-start-failure"); + const approvalItemId = TurnItemId.make("item-approval-start-failure"); + const thread: OrchestrationV2AppThread = { + createdBy: "user", + creationSource: "web", + id: threadId, + projectId: ProjectId.make("project-start-failure"), + title: "Start failure", + providerInstanceId: instanceId, + modelSelection: { instanceId, model: "gpt-5.4" }, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + activeProviderThreadId: providerThreadId, + lineage: { parentThreadId: null, relationshipToParent: null, rootThreadId: threadId }, + forkedFrom: null, + createdAt: now, + updatedAt: now, + archivedAt: null, + settledOverride: null, + settledAt: null, + lastVisitedAt: null, + deletedAt: null, + }; + const run: OrchestrationV2Run = { + id: runId, + threadId, + ordinal: 1, + providerInstanceId: instanceId, + modelSelection: thread.modelSelection, + providerThreadId, + userMessageId: messageId, + rootNodeId, + activeAttemptId: attemptId, + status: "starting", + requestedAt: now, + startedAt: null, + completedAt: null, + checkpointId: null, + contextHandoffId: null, + }; + const approvalRequest: OrchestrationV2RuntimeRequest = { + id: approvalRequestId, + nodeId: approvalNodeId, + providerTurnId, + nativeRequestRef: null, + kind: "command", + status: "pending", + responseCapability: { type: "live", providerSessionId }, + createdAt: now, + resolvedAt: null, + }; + const approvalNode: OrchestrationV2ExecutionNode = { + id: approvalNodeId, + threadId, + runId, + parentNodeId: rootNodeId, + rootNodeId, + kind: "approval_request", + status: "waiting", + countsForRun: false, + providerThreadId, + providerTurnId, + nativeItemRef: null, + runtimeRequestId: approvalRequestId, + checkpointScopeId: null, + startedAt: now, + completedAt: null, + }; + const approvalItem: OrchestrationV2TurnItem = { + id: approvalItemId, + threadId, + runId, + nodeId: approvalNodeId, + providerThreadId, + providerTurnId, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + status: "waiting", + title: "Run npm test?", + startedAt: now, + completedAt: null, + updatedAt: now, + type: "approval_request", + requestId: approvalRequestId, + requestKind: "command", + }; + // The run as a restarted attempt inherits it: the earlier attempt left an + // approval waiting under the run. + const seedPayloads: ReadonlyArray< + Pick & { readonly runId?: RunId } + > = [ + { type: "thread.created", payload: thread }, + { + type: "message.updated", + payload: { + id: messageId, + threadId, + runId, + nodeId: rootNodeId, + role: "user", + createdBy: "user", + creationSource: "web", + text: "Continue", + attachments: [], + streaming: false, + createdAt: now, + updatedAt: now, + }, + }, + { type: "run.created", payload: run }, + { + type: "run-attempt.created", + payload: { + id: attemptId, + runId, + attemptOrdinal: 2, + rootNodeId, + providerInstanceId: instanceId, + providerThreadId, + providerTurnId: null, + reason: "steering_restart", + status: "pending", + startedAt: null, + completedAt: null, + }, + }, + { + type: "checkpoint-scope.created", + payload: { + id: checkpointScopeId, + threadId, + runId, + nodeId: rootNodeId, + parentScopeId: null, + providerThreadId, + kind: "root_run", + ordinalWithinParent: 0, + advancesAppRunCount: true, + cwd: "/tmp/start-failure", + createdAt: now, + }, + }, + { + type: "node.updated", + payload: { + id: rootNodeId, + threadId, + runId, + parentNodeId: null, + rootNodeId, + kind: "root_turn", + status: "running", + countsForRun: true, + providerThreadId, + providerTurnId: null, + nativeItemRef: null, + runtimeRequestId: null, + checkpointScopeId, + startedAt: now, + completedAt: null, + }, + }, + { + type: "provider-thread.updated", + payload: { + id: providerThreadId, + driver, + providerInstanceId: instanceId, + providerSessionId, + appThreadId: threadId, + ownerNodeId: null, + nativeThreadRef: null, + nativeConversationHeadRef: null, + status: "not_loaded", + firstRunOrdinal: 1, + lastRunOrdinal: 1, + handoffIds: [], + forkedFrom: null, + createdAt: now, + updatedAt: now, + }, + }, + { type: "node.updated", payload: approvalNode }, + { type: "runtime-request.updated", payload: approvalRequest }, + { type: "turn-item.updated", payload: approvalItem }, + ]; + const layerDatabase = SqlitePersistence.layerMemory; + const layerStores = Layer.mergeAll(EventStore.layer, ProjectionStore.layer).pipe( + Layer.provideMerge(layerDatabase), + ); + const layerSink = EventSink.layer.pipe(Layer.provide(layerStores)); + const layerPersistence = Layer.mergeAll(layerStores, layerSink, IdAllocator.layer); + // `beforeFailureWrite` runs once the start has read its projection and + // given up on the provider, right before the failure is committed. + const run_ = (options: { + /** Rows the run inherited besides the approval. */ + readonly seed?: ReadonlyArray; + readonly beforeFailureWrite?: Effect.Effect< + void, + EventSink.EventSinkV2Error, + EventSink.EventSinkV2 + >; + readonly projectionStore?: ( + store: ProjectionStore.ProjectionStoreV2Shape, + eventSink: EventSink.EventSinkV2Shape, + ) => ProjectionStore.ProjectionStoreV2Shape; + }) => + Effect.gen(function* () { + const eventSink = yield* EventSink.EventSinkV2; + const store = yield* ProjectionStore.ProjectionStoreV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + yield* eventSink.write({ + events: yield* Effect.forEach(seedPayloads, (event) => + Effect.gen(function* () { + return { + ...event, + id: yield* idAllocator.allocate.event({ threadId }), + threadId, + occurredAt: now, + } as OrchestrationV2DomainEvent; + }), + ), + }); + if (options.seed !== undefined) yield* eventSink.write({ events: options.seed }); + const gatedSink = EventSink.EventSinkV2.of({ + ...eventSink, + writeIfRunCurrent: (input) => + (options.beforeFailureWrite ?? Effect.void).pipe( + Effect.provideService(EventSink.EventSinkV2, eventSink), + Effect.orDie, + Effect.andThen(eventSink.writeIfRunCurrent(input)), + ), + }); + const service = yield* ProviderTurnStart.ProviderTurnStartServiceV2.pipe( + Effect.provide( + ProviderTurnStart.layer.pipe( + Layer.provide( + Layer.mergeAll( + Layer.succeed(EventSink.EventSinkV2, gatedSink), + Layer.succeed( + ProjectionStore.ProjectionStoreV2, + options.projectionStore?.(store, eventSink) ?? store, + ), + IdAllocator.layer, + Layer.mock(ContextHandoffService.ContextHandoffServiceV2)({}), + FileSystem.layerNoop({}), + Layer.mock(GitWorkflow.GitWorkflowService)({}), + Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(ProviderSessionManager.ProviderSessionManagerV2)({ + open: () => + Effect.fail( + new ProviderSessionManager.ProviderSessionOpenError({ + instanceId, + providerSessionId, + cause: "replacement session failed to open", + }), + ), + }), + Layer.mock(ProviderAuthService.ProviderAuthService)({ + tryHandlePromptCommand: () => Effect.succeed(false), + }), + Layer.mock(RunExecutionService.RunExecutionServiceV2)({ + startRootRun: () => Effect.die("a failed open must not start the run"), + }), + Layer.mock(RuntimePolicy.RuntimePolicyV2)({ + resolve: () => Effect.succeed({} as never), + }), + ), + ), + ), + ), + ); + yield* service.start({ threadId, runId }); + return yield* store.getThreadProjection(threadId); + }).pipe(Effect.provide(Layer.fresh(layerPersistence))); + return { + now, + driver, + instanceId, + threadId, + runId, + rootNodeId, + providerThreadId, + approvalRequest, + approvalNode, + approvalItem, + run: run_, + }; +} + +// The user accepts the approval the run inherited. +const acceptApproval = ( + harness: ReturnType, + eventSink: EventSink.EventSinkV2Shape, +) => + eventSink.write({ + events: [ + { + id: EventId.make("event-approval-accepted"), + type: "runtime-request.updated", + threadId: harness.threadId, + nodeId: harness.approvalNode.id, + occurredAt: harness.now, + payload: { + ...harness.approvalRequest, + status: "resolved", + decision: "accept", + resolvedAt: harness.now, + }, + }, + { + id: EventId.make("event-approval-node-completed"), + type: "node.updated", + threadId: harness.threadId, + nodeId: harness.approvalNode.id, + occurredAt: harness.now, + payload: { ...harness.approvalNode, status: "completed", completedAt: harness.now }, + }, + { + id: EventId.make("event-approval-item-completed"), + type: "turn-item.updated", + threadId: harness.threadId, + nodeId: harness.approvalNode.id, + occurredAt: harness.now, + payload: { + ...harness.approvalItem, + status: "completed", + completedAt: harness.now, + updatedAt: harness.now, + }, + }, + ], + }); + +const expectApprovalAccepted = (projection: OrchestrationV2ThreadProjection) => { + expect(projection.runs.map((run) => run.status)).toEqual(["failed"]); + expect(projection.runtimeRequests).toMatchObject([{ status: "resolved", decision: "accept" }]); + const approvalNode = projection.nodes.find((node) => node.kind === "approval_request"); + const approvalItem = projection.turnItems.find((item) => item.type === "approval_request"); + expect(approvalNode?.status).toBe("completed"); + expect(approvalItem?.status).toBe("completed"); +}; + +effectIt.effect("keeps an approval accepted while a failed start is written", () => + Effect.gen(function* () { + const harness = makePersistedStartFailureHarness(); + const projection = yield* harness.run({ + // The user approves after the failure read the approval as pending. + beforeFailureWrite: Effect.gen(function* () { + yield* acceptApproval(harness, yield* EventSink.EventSinkV2); + }), + }); + + expectApprovalAccepted(projection); + }), +); + +effectIt.effect("keeps an approval accepted while the failure reads what it inherited", () => + Effect.gen(function* () { + const harness = makePersistedStartFailureHarness(); + let settling = false; + let accepted = false; + const projection = yield* harness.run({ + // The user approves while the failure's recovery read runs: it loaded + // the approval's node and item as waiting, but the request had left the + // pending-requests read, so no cancellation accompanies their settlement. + projectionStore: (store, eventSink) => ({ + ...store, + // Only the inherited-work read asks for subagent links. + getThreadRecords: (threadId, fields, filter) => + Effect.sync(() => { + if (fields.some((field) => field === "subagents")) settling = true; + }).pipe(Effect.andThen(store.getThreadRecords(threadId, fields, filter))), + getRuntimeRecoveryProjection: (threadId) => + store.getRuntimeRecoveryProjection(threadId).pipe( + Effect.flatMap((recovery) => + !settling || accepted + ? Effect.succeed(recovery) + : Effect.sync(() => { + accepted = true; + }).pipe( + Effect.andThen(acceptApproval(harness, eventSink)), + Effect.orDie, + Effect.as({ + ...recovery, + runtimeRequests: recovery.runtimeRequests.filter( + (request) => request.id !== harness.approvalRequest.id, + ), + }), + ), + ), + ), + }), + }); + + expect(accepted).toBe(true); + + expectApprovalAccepted(projection); + }), +); + +effectIt.effect("fails the run when the work it inherited cannot be read", () => + Effect.gen(function* () { + const harness = makePersistedStartFailureHarness(); + const projection = yield* harness.run({ + // The start's own reads succeed; only the inherited-work read fails. + projectionStore: (store) => ({ + ...store, + getThreadRecords: (threadId, fields, filter) => + fields.some((field) => field === "subagents") + ? Effect.fail( + new ProjectionStore.ProjectionStoreReadError({ + threadId, + cause: "database unavailable", + }), + ) + : store.getThreadRecords(threadId, fields, filter), + }), + }); + + expect(projection.runs.map((run) => run.status)).toEqual(["failed"]); + expect(projection.attempts.map((attempt) => attempt.status)).toEqual(["failed"]); + expect(projection.turnItems).toMatchObject([ + { type: "approval_request", status: "waiting" }, + { type: "error", title: "Provider session failed to open" }, + ]); + }), +); + +effectIt.effect("keeps a native subagent that completed while a failed start is written", () => + Effect.gen(function* () { + const harness = makePersistedStartFailureHarness(); + const subagentNode: OrchestrationV2ExecutionNode = { + ...harness.approvalNode, + id: NodeId.make("node-subagent-start-failure"), + kind: "subagent", + status: "running", + runtimeRequestId: null, + }; + const subagent: OrchestrationV2Subagent = { + id: subagentNode.id, + threadId: harness.threadId, + runId: harness.runId, + parentNodeId: harness.rootNodeId, + origin: "provider_native", + createdBy: "agent", + driver: harness.driver, + providerInstanceId: harness.instanceId, + providerThreadId: harness.providerThreadId, + childThreadId: null, + nativeTaskRef: null, + prompt: "Explore the repo", + title: "Explorer", + model: null, + status: "running", + result: null, + startedAt: harness.now, + completedAt: null, + updatedAt: harness.now, + }; + const subagentItem: OrchestrationV2TurnItem = { + id: TurnItemId.make("item-subagent-start-failure"), + threadId: harness.threadId, + runId: harness.runId, + nodeId: subagentNode.id, + providerThreadId: harness.providerThreadId, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 2, + status: "running", + title: "Explorer", + startedAt: harness.now, + completedAt: null, + updatedAt: harness.now, + type: "subagent", + subagentId: subagent.id, + origin: "provider_native", + driver: harness.driver, + providerInstanceId: harness.instanceId, + childThreadId: null, + prompt: subagent.prompt, + result: null, + }; + const base = { threadId: harness.threadId, runId: harness.runId, occurredAt: harness.now }; + const subagentEvents = ( + label: string, + update: { readonly status: "running" | "completed"; readonly result: string | null }, + ): ReadonlyArray => { + const completedAt = update.status === "completed" ? harness.now : null; + return [ + { + ...base, + id: EventId.make(`event-subagent-node-${label}`), + type: "node.updated", + nodeId: subagentNode.id, + payload: { ...subagentNode, status: update.status, completedAt }, + }, + { + ...base, + id: EventId.make(`event-subagent-${label}`), + type: "subagent.updated", + nodeId: subagent.id, + payload: { ...subagent, ...update, completedAt }, + }, + { + ...base, + id: EventId.make(`event-subagent-item-${label}`), + type: "turn-item.updated", + nodeId: subagentNode.id, + payload: { ...subagentItem, ...update, completedAt }, + }, + ]; + }; + const projection = yield* harness.run({ + seed: subagentEvents("running", { status: "running", result: null }), + // The earlier attempt's session still reports on its native subagent; + // it finishes after the failure read it as running. + beforeFailureWrite: Effect.gen(function* () { + yield* (yield* EventSink.EventSinkV2).write({ + events: subagentEvents("completed", { status: "completed", result: "Found it" }), + }); + }), + }); + + expect(projection.runs.map((run) => run.status)).toEqual(["failed"]); + expect(projection.subagents).toMatchObject([{ status: "completed", result: "Found it" }]); + expect(projection.nodes.find((node) => node.id === subagentNode.id)?.status).toBe("completed"); + expect(projection.turnItems.find((item) => item.id === subagentItem.id)).toMatchObject({ + status: "completed", + result: "Found it", + }); + }), +); diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index 684bde2c0669..791556d6a5fb 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -2,6 +2,7 @@ import { modelSelectionsEqual } from "@t3tools/shared/model"; import { projectComposerContextForProvider } from "@t3tools/shared/composerContextReferences"; import { CommandId, + isProviderNativeSubagentThread, latestProviderTurnForAttempt, type OrchestrationV2DomainEvent, type OrchestrationV2ExecutionNode, @@ -178,6 +179,103 @@ export const layer: Layer.Layer< }; }; + /** + * Ends the work a run that never reached its provider still shows. A + * restart supersedes only the root turn, so requests, streaming replies and + * background work an earlier attempt left open stay with the run, and once + * it fails nothing will report on them. The run's projection rows feed + * RunExecutionService's terminal cascade, the same one a started run uses. + */ + const inheritedWorkSettlement = Effect.fn( + "orchestrationV2.providerTurnStart.inheritedWorkSettlement", + )(function* (input: { readonly run: OrchestrationV2Run; readonly now: DateTime.Utc }) { + const { run } = input; + // Lifetime links, settled rows included: a subagent row can settle before + // its child thread does. App-owned tasks (delegate_task) run on their own. + const nativeLinks = (threadId: ThreadId) => + projectionStore + .getThreadRecords(threadId, ["subagents", "turnItems"], { + turnItemTypes: ["subagent"], + ...(threadId === run.threadId ? { turnItemRunIds: [run.id] } : {}), + }) + .pipe( + Effect.map((records) => + [...records.subagents, ...records.turnItems].flatMap((row) => + (threadId !== run.threadId || row.runId === run.id) && + "childThreadId" in row && + row.childThreadId !== null && + row.origin === "provider_native" + ? [row.childThreadId] + : [], + ), + ), + ); + const pending = [...(yield* nativeLinks(run.threadId))]; + const linkedChildThreadIds = new Set(); + const threads = [yield* projectionStore.getRuntimeRecoveryProjection(run.threadId)]; + for (let threadId = pending.shift(); threadId !== undefined; threadId = pending.shift()) { + if (threadId === run.threadId || linkedChildThreadIds.has(threadId)) continue; + const child = yield* projectionStore.getRuntimeRecoveryProjection(threadId).pipe( + Effect.map(Option.some), + Effect.catchTags({ + ProjectionStoreThreadNotFoundError: () => Effect.succeed(Option.none()), + }), + ); + if (Option.isNone(child) || !isProviderNativeSubagentThread(child.value.thread)) continue; + linkedChildThreadIds.add(threadId); + // A child thread has no runs, so the recovery read leaves out its + // streaming replies, provider turns, subagents whose item already + // settled, and its pending requests' nodes and items; read them by + // thread instead. + const records = yield* projectionStore.getThreadRecords( + threadId, + ["messages", "nodes", "turnItems", "providerTurns", "subagents"], + { + messageRoles: ["assistant"], + turnItemTypes: ["approval_request", "user_input_request"], + turnItemStatuses: ["pending", "running", "waiting"], + }, + ); + const requestNodeIds = new Set( + child.value.runtimeRequests.flatMap((request) => + request.status === "pending" ? [request.nodeId] : [], + ), + ); + const knownNodeIds = new Set(child.value.nodes.map((node) => node.id)); + const knownItemIds = new Set(child.value.turnItems.map((item) => item.id)); + threads.push({ + ...child.value, + providerTurns: records.providerTurns, + subagents: records.subagents, + nodes: [ + ...child.value.nodes, + ...records.nodes.filter( + (node) => requestNodeIds.has(node.id) && !knownNodeIds.has(node.id), + ), + ], + turnItems: [ + ...child.value.turnItems, + ...records.turnItems.filter((item) => !knownItemIds.has(item.id)), + ], + messages: records.messages.filter((message) => message.streaming), + }); + pending.push(...(yield* nativeLinks(threadId))); + } + const linked = RunExecutionService.openRunOwnedWorkFromProjection({ + run, + threads, + linkedChildThreadIds, + }); + return yield* RunExecutionService.cascadeTerminalizeRunOwnedSubagents({ + run, + open: linked, + // The same status a started run's failure gives the work it owned. + status: "failed", + completedAt: input.now, + allocateEventId: () => idAllocator.allocate.event({ threadId: run.threadId }), + }); + }); + const makeDeliverySession = ( session: ProviderAdapter.ProviderAdapterV2SessionRuntime, startWithHandoffs: ( @@ -366,7 +464,31 @@ export const layer: Layer.Layer< runId, activeAttemptId: attempt.id, expectedStatus: "starting", - events, + // An answer or approval that lands while the failure is written wins + // over cancelling the request it resolved. + guardPendingRequestCancellations: true, + // A provider-native child the earlier attempt still reports on can + // finish between the inherited-work read and this commit. + guardSettledWork: true, + events: + status === "failed" + ? [ + // Failing the run matters more than settling what it + // inherited: a read that fails here must not leave it + // `starting`. + ...(yield* inheritedWorkSettlement({ run, now }).pipe( + Effect.catchCause((cause) => + Cause.hasInterruptsOnly(cause) + ? Effect.failCause(cause) + : Effect.logWarning( + "provider turn start could not settle the work a failed run inherited", + { threadId: projection.thread.id, runId, cause: Cause.pretty(cause) }, + ).pipe(Effect.as([])), + ), + )), + ...events, + ] + : events, }); }, ); diff --git a/apps/server/src/orchestration-v2/RunExecutionService.ts b/apps/server/src/orchestration-v2/RunExecutionService.ts index fa6d0fbd7f12..eb88155b259c 100644 --- a/apps/server/src/orchestration-v2/RunExecutionService.ts +++ b/apps/server/src/orchestration-v2/RunExecutionService.ts @@ -8,6 +8,7 @@ import { type NodeId, type OrchestrationV2AppThread, type OrchestrationV2CheckpointScope, + type OrchestrationV2ConversationMessage, type OrchestrationV2DomainEvent, type OrchestrationV2ExecutionNode, type OrchestrationV2ProviderFailure, @@ -15,6 +16,7 @@ import { type OrchestrationV2ProviderTurn, type OrchestrationV2Run, type OrchestrationV2RunAttempt, + type OrchestrationV2RuntimeRequest, type OrchestrationV2Subagent, type OrchestrationV2TurnItem, type ProviderSessionId, @@ -42,6 +44,7 @@ import * as CheckpointService from "./CheckpointService.ts"; import * as EventSink from "./EventSink.ts"; import * as IdAllocator from "@t3tools/provider-core/server/IdAllocator"; import * as ProviderAdapter from "@t3tools/provider-core/server/ProviderAdapter"; +import type * as ProjectionStore from "./ProjectionStore.ts"; import * as ProviderEventIngestor from "./ProviderEventIngestor.ts"; import type { ProjectionStoreV2Error } from "./ProjectionStore.ts"; import { @@ -147,13 +150,27 @@ export function selectInheritedBackgroundTurnItems(input: { type SubagentTurnItem = Extract; -type OpenRunOwnedSubagentProjection = { +export type OpenRunOwnedSubagentProjection = { readonly subagents: ReadonlyMap; readonly turnItems: ReadonlyMap; + /** Open items on linked child threads, and on the run's own thread for the run. */ readonly childTurnItems: ReadonlyMap; readonly nodes: ReadonlyMap; /** Child threads once linked by a root-run subagent row; kept for cascade. */ readonly linkedChildThreadIds: ReadonlySet; + // A live run's provider reports these itself. A run that never reached its + // provider has nothing to report them, so its read-model rows are passed in. + /** Pending requests; each settles with the node it waits on. */ + readonly runtimeRequests?: ReadonlyArray; + readonly streamingMessages?: ReadonlyArray; + /** + * Open provider turns of the run's own attempts and of linked child threads, + * with the thread each one runs on (the turn row does not carry it). + */ + readonly providerTurns?: ReadonlyArray<{ + readonly threadId: ThreadId; + readonly providerTurn: OrchestrationV2ProviderTurn; + }>; }; type RunOwnedSubagentTerminalStatus = Extract< @@ -203,6 +220,10 @@ function withLinkedChildThreadId( return { ...current, linkedChildThreadIds }; } +/** + * What a run that ends interrupted, failed or cancelled settles besides its + * root: the run's open work and the provider-native child threads it linked. + */ export function cascadeTerminalizeRunOwnedSubagents(input: { readonly run: OrchestrationV2Run; readonly open: OpenRunOwnedSubagentProjection; @@ -223,6 +244,7 @@ export function cascadeTerminalizeRunOwnedSubagents(input: { childThreadIds.add(item.childThreadId); } } + const settledNodeIds = new Set(); const keys = new Set([ ...input.open.subagents.keys(), ...input.open.turnItems.keys(), @@ -255,6 +277,7 @@ export function cascadeTerminalizeRunOwnedSubagents(input: { childThreadIds.has(node.threadId)) && isOpenExecutionNodeStatus(node.status) ) { + settledNodeIds.add(node.id); events.push({ id: yield* input.allocateEventId(), type: "node.updated", @@ -293,8 +316,11 @@ export function cascadeTerminalizeRunOwnedSubagents(input: { }); } } + const ownsThreadRow = (row: { readonly threadId: ThreadId; readonly runId: string | null }) => + childThreadIds.has(row.threadId) || + (row.threadId === input.run.threadId && row.runId === input.run.id); for (const turnItem of input.open.childTurnItems.values()) { - if (!childThreadIds.has(turnItem.threadId) || isSettledTurnItemStatus(turnItem.status)) { + if (!ownsThreadRow(turnItem) || isSettledTurnItemStatus(turnItem.status)) { continue; } events.push({ @@ -314,10 +340,144 @@ export function cascadeTerminalizeRunOwnedSubagents(input: { }, }); } + for (const request of input.open.runtimeRequests ?? []) { + if (request.status !== "pending" || !settledNodeIds.has(request.nodeId)) continue; + const node = input.open.nodes.get(request.nodeId); + events.push({ + id: yield* input.allocateEventId(), + type: "runtime-request.updated", + threadId: node?.threadId ?? input.run.threadId, + nodeId: request.nodeId, + providerInstanceId: input.run.providerInstanceId, + occurredAt: input.completedAt, + payload: { + ...request, + status: "cancelled", + responseCapability: { + type: "not_resumable", + reason: "The run ended before this request was resolved.", + }, + resolvedAt: input.completedAt, + }, + }); + } + for (const message of input.open.streamingMessages ?? []) { + if (!message.streaming || !ownsThreadRow(message)) continue; + events.push({ + id: yield* input.allocateEventId(), + type: "message.updated", + threadId: message.threadId, + runId: message.runId ?? input.run.id, + ...(message.nodeId === null ? {} : { nodeId: message.nodeId }), + providerInstanceId: input.run.providerInstanceId, + occurredAt: input.completedAt, + payload: { ...message, streaming: false, updatedAt: input.completedAt }, + }); + } + for (const { threadId, providerTurn } of input.open.providerTurns ?? []) { + if (isTerminalProviderTurnStatus(providerTurn.status)) continue; + events.push({ + id: yield* input.allocateEventId(), + type: "provider-turn.updated", + threadId, + runId: input.run.id, + nodeId: providerTurn.nodeId, + providerInstanceId: input.run.providerInstanceId, + occurredAt: input.completedAt, + payload: { ...providerTurn, status: input.status, completedAt: input.completedAt }, + }); + } return events; }); } +/** + * The cascade's input for a run that never reached its provider, read from the + * projection instead of tracked from a live event stream. `threads` holds the + * run's own thread and the provider-native child threads it linked. + */ +export function openRunOwnedWorkFromProjection(input: { + readonly run: OrchestrationV2Run; + readonly threads: ReadonlyArray; + readonly linkedChildThreadIds: ReadonlySet; +}): OpenRunOwnedSubagentProjection { + const { run } = input; + const owns = (row: { readonly threadId: ThreadId; readonly runId: string | null }) => + row.threadId === run.threadId + ? row.runId === run.id + : input.linkedChildThreadIds.has(row.threadId); + const rows = ( + select: (thread: ProjectionStore.ProjectionRuntimeRecoveryState) => ReadonlyArray, + ) => input.threads.flatMap(select); + const runAttemptIds = new Set( + rows((thread) => thread.attempts) + .filter((attempt) => attempt.runId === run.id) + .map((attempt) => attempt.id), + ); + // App-owned tasks (delegate_task) run in their own threads and outlive it. + const appOwnedTaskIds = new Set( + rows((thread) => thread.subagents) + .filter((subagent) => subagent.origin === "app_owned") + .map((subagent) => subagent.id), + ); + const ownedTurnItems = rows((thread) => thread.turnItems).filter( + (item) => + owns(item) && + !isSettledTurnItemStatus(item.status) && + !(item.type === "subagent" && item.origin === "app_owned"), + ); + return { + // A linked child's own subagents carry no run id, so they are taken by + // the thread they run on. + subagents: new Map( + rows((thread) => thread.subagents) + .filter( + (subagent) => + owns(subagent) && + subagent.origin === "provider_native" && + !isSettledSubagentStatus(subagent.status), + ) + .map((subagent) => [subagent.id, subagent]), + ), + turnItems: new Map( + ownedTurnItems.flatMap((item) => + item.type === "subagent" && item.runId === run.id ? [[item.subagentId, item] as const] : [], + ), + ), + childTurnItems: new Map( + ownedTurnItems + .filter((item) => !(item.type === "subagent" && item.runId === run.id)) + .map((item) => [item.id, item]), + ), + nodes: new Map( + rows((thread) => thread.nodes) + .filter( + (node) => + owns(node) && + // The never-started path writes the root with the run's own status. + node.id !== run.rootNodeId && + !appOwnedTaskIds.has(node.id) && + isOpenExecutionNodeStatus(node.status), + ) + .map((node) => [node.id, node]), + ), + linkedChildThreadIds: input.linkedChildThreadIds, + runtimeRequests: rows((thread) => thread.runtimeRequests), + streamingMessages: rows((thread) => thread.messages).filter((message) => message.streaming), + // A linked child's provider turns belong to no run attempt. + providerTurns: input.threads.flatMap((thread) => + thread.providerTurns.flatMap((providerTurn) => + (thread.thread.id === run.threadId + ? providerTurn.runAttemptId !== null && runAttemptIds.has(providerTurn.runAttemptId) + : input.linkedChildThreadIds.has(thread.thread.id)) && + !isTerminalProviderTurnStatus(providerTurn.status) + ? [{ threadId: thread.thread.id, providerTurn }] + : [], + ), + ), + }; +} + export function finalProviderThreadStatus( disposition: ProviderTerminalEvent["threadDisposition"], ): OrchestrationV2ProviderThread["status"] { diff --git a/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts b/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts index 958fb89cc8e4..ae033dfbfdc1 100644 --- a/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts +++ b/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts @@ -4,9 +4,12 @@ import { EventId, MessageId, type ModelSelection, + NodeId, + type OrchestrationV2DomainEvent, type OrchestrationV2ProviderCapabilities, type OrchestrationV2ProviderSession, type OrchestrationV2ProviderThread, + type OrchestrationV2ThreadProjection, ProjectId, ProviderDriverKind, ProviderInstanceId, @@ -14,7 +17,9 @@ import { ProviderThreadId, ProviderTurnId, type RunId, + RuntimeRequestId, ThreadId, + TurnItemId, } from "@t3tools/contracts"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; @@ -81,6 +86,431 @@ interface RestartAdapterState { }>; readonly closedSessionCount: number; readonly failedReplacementOpen: boolean; + /** Every replacement open fails, not just the first. */ + readonly replacementOpenAlwaysFails?: boolean; + /** The first turn leaves an approval, a streaming reply, a subagent and a command open. */ + readonly firstTurnLeavesWorkOpen?: boolean; +} + +const providerNativeChildThreadId = (threadId: ThreadId) => + ThreadId.make(`${threadId}:native-subagent`); + +// Work a turn leaves open when it is superseded: an approval it is waiting on, +// a reply still streaming, a provider-native subagent whose own thread is still +// running a command, and a background command. +function openTurnWork( + providerSessionId: ProviderSessionId, + active: ActiveTurn, + now: DateTime.Utc, +): ReadonlyArray { + const { input, providerTurnId } = active; + const base = { + threadId: input.threadId, + runId: input.runId, + providerThreadId: input.providerThread.id, + providerTurnId, + nativeItemRef: null, + parentItemId: null, + title: null, + startedAt: now, + completedAt: null, + updatedAt: now, + }; + const approvalNodeId = NodeId.make(`node:approval:${input.attemptId}`); + const subagentNodeId = NodeId.make(`node:subagent:${input.attemptId}`); + const childThreadId = providerNativeChildThreadId(input.threadId); + const grandchildThreadId = providerNativeChildThreadId(childThreadId); + const childApprovalNodeId = NodeId.make(`node:child-approval:${input.attemptId}`); + const childRequestId = RuntimeRequestId.make(`request:child-approval:${input.attemptId}`); + const nestedRunningSubagentId = NodeId.make(`node:nested-running-subagent:${input.attemptId}`); + const nestedSettledItemSubagentId = NodeId.make( + `node:nested-settled-item-subagent:${input.attemptId}`, + ); + const requestId = RuntimeRequestId.make(`request:approval:${input.attemptId}`); + return [ + { + type: "node.updated", + driver, + node: { + id: approvalNodeId, + threadId: input.threadId, + runId: input.runId, + parentNodeId: input.rootNodeId, + rootNodeId: input.rootNodeId, + kind: "approval_request", + status: "waiting", + countsForRun: false, + providerThreadId: input.providerThread.id, + providerTurnId, + nativeItemRef: null, + runtimeRequestId: requestId, + checkpointScopeId: null, + startedAt: now, + completedAt: null, + }, + }, + { + type: "runtime_request.updated", + driver, + threadId: input.threadId, + runtimeRequest: { + id: requestId, + nodeId: approvalNodeId, + providerTurnId, + nativeRequestRef: null, + kind: "command", + status: "pending", + responseCapability: { type: "live", providerSessionId }, + createdAt: now, + resolvedAt: null, + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:approval:${input.attemptId}`), + nodeId: approvalNodeId, + ordinal: 1, + status: "waiting", + type: "approval_request", + requestId, + requestKind: "command", + }, + }, + { + type: "message.updated", + driver, + message: { + id: MessageId.make(`message:reply:${input.attemptId}`), + threadId: input.threadId, + runId: input.runId, + nodeId: input.rootNodeId, + role: "assistant", + createdBy: "agent", + creationSource: "provider", + text: "Working on it", + attachments: [], + streaming: true, + createdAt: now, + updatedAt: now, + }, + }, + { + type: "app_thread.created", + driver, + appThread: { + ...input.appThread, + createdBy: "agent", + creationSource: "provider", + id: childThreadId, + title: "Explore", + activeProviderThreadId: null, + lineage: { + parentThreadId: input.threadId, + relationshipToParent: "subagent", + rootThreadId: input.threadId, + }, + forkedFrom: null, + createdAt: now, + updatedAt: now, + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:subagent:${input.attemptId}`), + nodeId: subagentNodeId, + ordinal: 2, + status: "running", + type: "subagent", + subagentId: subagentNodeId, + origin: "provider_native", + driver, + providerInstanceId: input.modelSelection.instanceId, + childThreadId, + prompt: "Explore the repo", + result: null, + }, + }, + // The subagent's own provider turn is running. It belongs to no run + // attempt. + { + type: "provider_turn.updated", + driver, + threadId: childThreadId, + providerTurn: { + id: ProviderTurnId.make(`provider-turn:child:${input.attemptId}`), + providerThreadId: ProviderThreadId.make(`provider-thread:${childThreadId}`), + nodeId: subagentNodeId, + runAttemptId: null, + nativeTurnRef: null, + ordinal: 1, + status: "running", + startedAt: now, + completedAt: null, + }, + }, + // The subagent is waiting on an approval in its own thread. Its rows carry + // no run id; the child thread has no runs. + { + type: "node.updated", + driver, + node: { + id: childApprovalNodeId, + threadId: childThreadId, + runId: null, + parentNodeId: null, + rootNodeId: childApprovalNodeId, + kind: "approval_request", + status: "waiting", + countsForRun: false, + providerThreadId: input.providerThread.id, + providerTurnId, + nativeItemRef: null, + runtimeRequestId: childRequestId, + checkpointScopeId: null, + startedAt: now, + completedAt: null, + }, + }, + { + type: "runtime_request.updated", + driver, + threadId: childThreadId, + runtimeRequest: { + id: childRequestId, + nodeId: childApprovalNodeId, + providerTurnId, + nativeRequestRef: null, + kind: "command", + status: "pending", + responseCapability: { type: "live", providerSessionId }, + createdAt: now, + resolvedAt: null, + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:child-approval:${input.attemptId}`), + threadId: childThreadId, + runId: null, + nodeId: childApprovalNodeId, + ordinal: 3, + status: "waiting", + type: "approval_request", + requestId: childRequestId, + requestKind: "command", + }, + }, + { + type: "message.updated", + driver, + message: { + id: MessageId.make(`message:child-reply:${input.attemptId}`), + threadId: childThreadId, + runId: null, + nodeId: null, + role: "assistant", + createdBy: "agent", + creationSource: "provider", + text: "Looking", + attachments: [], + streaming: true, + createdAt: now, + updatedAt: now, + }, + }, + // The subagent's own subagent already reported done, but its thread is + // still running a command. + { + type: "app_thread.created", + driver, + appThread: { + ...input.appThread, + createdBy: "agent", + creationSource: "provider", + id: grandchildThreadId, + title: "Search", + activeProviderThreadId: null, + lineage: { + parentThreadId: childThreadId, + relationshipToParent: "subagent", + rootThreadId: input.threadId, + }, + forkedFrom: null, + createdAt: now, + updatedAt: now, + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:nested-subagent:${input.attemptId}`), + threadId: childThreadId, + runId: null, + nodeId: null, + ordinal: 2, + status: "completed", + completedAt: now, + type: "subagent", + subagentId: NodeId.make(`node:nested-subagent:${input.attemptId}`), + origin: "provider_native", + driver, + providerInstanceId: input.modelSelection.instanceId, + childThreadId: grandchildThreadId, + prompt: "Search for TODOs", + result: "done", + }, + }, + // Another subagent it launched is still running. Like every row on the + // subagent's thread, it carries no run id. + { + type: "subagent.updated", + driver, + subagent: { + id: nestedRunningSubagentId, + threadId: childThreadId, + runId: null, + parentNodeId: subagentNodeId, + origin: "provider_native", + createdBy: "agent", + driver, + providerInstanceId: input.modelSelection.instanceId, + providerThreadId: null, + childThreadId: null, + nativeTaskRef: null, + prompt: "Read the docs", + title: null, + model: null, + status: "running", + result: null, + startedAt: now, + completedAt: null, + updatedAt: now, + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:nested-running-subagent:${input.attemptId}`), + threadId: childThreadId, + runId: null, + nodeId: null, + ordinal: 4, + status: "running", + type: "subagent", + subagentId: nestedRunningSubagentId, + origin: "provider_native", + driver, + providerInstanceId: input.modelSelection.instanceId, + childThreadId: null, + prompt: "Read the docs", + result: null, + }, + }, + // A third reported its item done, but its row is still running. + { + type: "subagent.updated", + driver, + subagent: { + id: nestedSettledItemSubagentId, + threadId: childThreadId, + runId: null, + parentNodeId: subagentNodeId, + origin: "provider_native", + createdBy: "agent", + driver, + providerInstanceId: input.modelSelection.instanceId, + providerThreadId: null, + childThreadId: null, + nativeTaskRef: null, + prompt: "Check the tests", + title: null, + model: null, + status: "running", + result: null, + startedAt: now, + completedAt: null, + updatedAt: now, + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:nested-settled-item-subagent:${input.attemptId}`), + threadId: childThreadId, + runId: null, + nodeId: null, + ordinal: 5, + status: "completed", + completedAt: now, + type: "subagent", + subagentId: nestedSettledItemSubagentId, + origin: "provider_native", + driver, + providerInstanceId: input.modelSelection.instanceId, + childThreadId: null, + prompt: "Check the tests", + result: "done", + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:grandchild-command:${input.attemptId}`), + threadId: grandchildThreadId, + runId: null, + nodeId: null, + ordinal: 1, + status: "running", + type: "command_execution", + input: "rg FIXME", + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:child-command:${input.attemptId}`), + threadId: childThreadId, + runId: null, + nodeId: null, + ordinal: 1, + status: "running", + type: "command_execution", + input: "rg TODO", + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:command:${input.attemptId}`), + nodeId: input.rootNodeId, + ordinal: 3, + status: "running", + type: "command_execution", + input: "npm run dev", + }, + }, + ]; } function makeRestartAdapter( @@ -103,7 +533,7 @@ function makeRestartAdapter( const failThisOpen = yield* Ref.modify(state, (current) => { const shouldFail = sessionInput.modelSelection.model === replacementSelection.model && - !current.failedReplacementOpen; + (current.replacementOpenAlwaysFails === true || !current.failedReplacementOpen); return [ shouldFail, { @@ -253,6 +683,15 @@ function makeRestartAdapter( completedAt: null, }, }); + if ((yield* Ref.get(state)).firstTurnLeavesWorkOpen === true) { + for (const event of openTurnWork( + sessionInput.providerSessionId, + active, + occurredAt, + )) { + yield* Queue.offer(events, event); + } + } return; } yield* publishTerminal(active, "completed"); @@ -534,6 +973,268 @@ it.live("restarts selection as a new attempt and retries after old-session clean ), ); +// A restart only replaces the root turn: the request, reply and command the +// superseded attempt left open become the restarted run's to settle. When the +// replacement never starts, failing the run must end them too, or the thread +// waits on work nothing is running. +it.live("settles the work a restarted run inherited when its replacement never opens", () => + Effect.scoped( + Effect.gen(function* () { + const name = "selection-restart-failed-start"; + const cwd = yield* checkpointWorkspace(name); + const threadId = ThreadId.make(`thread:${name}`); + const state = yield* Ref.make({ + activeTurn: null, + opened: [], + started: [], + closedSessionCount: 0, + failedReplacementOpen: false, + replacementOpenAlwaysFails: true, + firstTurnLeavesWorkOpen: true, + }); + const layerRegistry = ProviderAdapterRegistry.layerSingle(makeRestartAdapter(state)); + + const { + projection, + nativeChildItems, + nativeChildStreaming, + nativeChildRequests, + nativeChildSubagents, + nativeChildProviderTurns, + nativeGrandchildItems, + delegatedItems, + } = yield* Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const eventSink = yield* EventSink.EventSinkV2; + const awaitDomainEvent = (matches: (event: OrchestrationV2DomainEvent) => boolean) => + orchestrator.streamDomainEvents.pipe( + Stream.filter(matches), + Stream.take(1), + Stream.runDrain, + Effect.forkScoped({ startImmediately: true }), + ); + yield* orchestrator.dispatch({ + type: "thread.create", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make(`${name}:create`), + threadId, + projectId: ProjectId.make(`project:${name}`), + title: name, + modelSelection: initialSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: cwd, + }); + // The parent's command item is the last open work the first turn reports. + const workOpen = yield* awaitDomainEvent( + (event) => + event.type === "turn-item.updated" && + event.payload.type === "command_execution" && + event.payload.threadId === threadId, + ); + yield* orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make(`${name}:first`), + threadId, + messageId: MessageId.make(`${name}:first`), + text: "first", + attachments: [], + modelSelection: initialSelection, + dispatchMode: { type: "start_immediately" }, + }); + yield* Fiber.join(workOpen); + const started = yield* orchestrator.getThreadProjection(threadId); + const runId = started.runs[0]?.id; + if (runId === undefined) return yield* Effect.die("the first run is missing"); + // The run also delegated a task (delegate_task). Its thread runs on its + // own, so the failed restart must leave it running. + const now = yield* DateTime.now; + const delegatedThreadId = ThreadId.make(`${threadId}:delegated`); + const delegatedTaskId = NodeId.make(`node:delegated:${runId}`); + yield* eventSink.write({ + events: [ + { + id: EventId.make(`event:${name}:delegated-thread`), + type: "thread.created", + threadId: delegatedThreadId, + occurredAt: now, + payload: { + ...started.thread, + createdBy: "agent", + creationSource: "mcp", + id: delegatedThreadId, + activeProviderThreadId: null, + lineage: { + parentThreadId: threadId, + relationshipToParent: "subagent", + rootThreadId: threadId, + }, + }, + }, + { + id: EventId.make(`event:${name}:delegated-task`), + type: "subagent.updated", + threadId, + runId, + nodeId: delegatedTaskId, + driver, + providerInstanceId, + occurredAt: now, + payload: { + id: delegatedTaskId, + threadId, + runId, + parentNodeId: started.runs[0]!.rootNodeId!, + origin: "app_owned", + createdBy: "agent", + driver, + providerInstanceId, + providerThreadId: null, + childThreadId: delegatedThreadId, + nativeTaskRef: null, + prompt: "Write the tests", + title: null, + model: initialSelection.model, + status: "running", + result: null, + startedAt: now, + completedAt: null, + updatedAt: now, + }, + }, + { + id: EventId.make(`event:${name}:delegated-command`), + type: "turn-item.updated", + threadId: delegatedThreadId, + occurredAt: now, + payload: { + id: TurnItemId.make(`turn-item:${name}:delegated-command`), + threadId: delegatedThreadId, + runId: null, + nodeId: null, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + status: "running", + title: null, + startedAt: now, + completedAt: null, + updatedAt: now, + type: "command_execution", + input: "vp test run", + }, + }, + ], + }); + + const runFailed = yield* awaitDomainEvent( + (event) => event.type === "run.updated" && event.payload.status === "failed", + ); + yield* orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make(`${name}:second`), + threadId, + messageId: MessageId.make(`${name}:second`), + text: "second", + attachments: [], + modelSelection: replacementSelection, + dispatchMode: { type: "restart_active", targetRunId: runId }, + }); + yield* Fiber.join(runFailed); + const openItems = (projection: OrchestrationV2ThreadProjection) => + projection.turnItems.flatMap((item) => + item.type === "approval_request" || + item.type === "command_execution" || + item.type === "subagent" + ? [[item.type, item.status]] + : [], + ); + return { + projection: yield* orchestrator.getThreadProjection(threadId), + nativeChildItems: openItems( + yield* orchestrator.getThreadProjection(providerNativeChildThreadId(threadId)), + ), + nativeChildStreaming: (yield* orchestrator.getThreadProjection( + providerNativeChildThreadId(threadId), + )).messages.some((message) => message.streaming), + nativeChildRequests: (yield* orchestrator.getThreadProjection( + providerNativeChildThreadId(threadId), + )).runtimeRequests.map((request) => request.status), + nativeChildSubagents: (yield* orchestrator.getThreadProjection( + providerNativeChildThreadId(threadId), + )).subagents.map((subagent) => subagent.status), + nativeChildProviderTurns: (yield* orchestrator.getThreadProjection( + providerNativeChildThreadId(threadId), + )).providerTurns.map((turn) => turn.status), + nativeGrandchildItems: openItems( + yield* orchestrator.getThreadProjection( + providerNativeChildThreadId(providerNativeChildThreadId(threadId)), + ), + ), + delegatedItems: openItems(yield* orchestrator.getThreadProjection(delegatedThreadId)), + }; + }).pipe(Effect.provide(ProviderReplayHarness.layerWithRegistry({ name }, layerRegistry))); + + assert.deepEqual( + projection.runs.map((run) => run.status), + ["failed"], + ); + assert.deepEqual( + projection.runtimeRequests.map((request) => request.status), + ["cancelled"], + ); + assert.isFalse(projection.messages.some((message) => message.streaming)); + assert.deepEqual( + projection.turnItems.flatMap((item) => + item.type === "approval_request" || + item.type === "command_execution" || + item.type === "subagent" + ? [[item.type, item.status]] + : [], + ), + [ + ["approval_request", "failed"], + ["subagent", "failed"], + ["command_execution", "failed"], + ], + ); + // The provider-native subagent's own thread ends with the run, and so does + // the thread of a nested subagent whose row already settled. + assert.deepEqual(nativeChildItems, [ + ["approval_request", "failed"], + ["subagent", "completed"], + ["subagent", "failed"], + ["subagent", "completed"], + ["command_execution", "failed"], + ]); + assert.isFalse(nativeChildStreaming); + assert.deepEqual(nativeChildRequests, ["cancelled"]); + // Subagents the provider-native subagent launched end with it, including + // one whose item already reported done. + assert.deepEqual(nativeChildSubagents, ["failed", "failed"]); + // So does the subagent's own provider turn, which no run attempt owns. + assert.deepEqual(nativeChildProviderTurns, ["failed"]); + assert.deepEqual(nativeGrandchildItems, [["command_execution", "failed"]]); + // The delegated task and its thread keep running. + assert.deepEqual( + projection.subagents.flatMap((subagent) => + subagent.origin === "app_owned" ? [subagent.status] : [], + ), + ["running"], + ); + assert.deepEqual(delegatedItems, [["command_execution", "running"]]); + }), + ), +); + it.live.each(["stopped", "error"] as const)( "restarts the live session on a model change when a newer %s session record exists", (deadStatus) =>