diff --git a/apps/server/src/orchestration-v2/RunExecutionService.test.ts b/apps/server/src/orchestration-v2/RunExecutionService.test.ts index 4e433dc7bb1f..e1cae4f85514 100644 --- a/apps/server/src/orchestration-v2/RunExecutionService.test.ts +++ b/apps/server/src/orchestration-v2/RunExecutionService.test.ts @@ -23,9 +23,12 @@ import { ProviderTurnId, RunAttemptId, RunId, + ServerSettingsError, ThreadId, TurnItemId, } from "@t3tools/contracts"; +import * as Cause from "effect/Cause"; +import * as Exit from "effect/Exit"; import * as DateTime from "effect/DateTime"; import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; @@ -780,6 +783,175 @@ it.effect("starts the provider when checkpoint baseline capture fails", () => }), ); +for (const scenario of ["failure", "interruption", "stale-attempt", "start-guard"] as const) { + it.effect(`handles ${scenario} before the provider turn starts`, () => + Effect.gen(function* () { + const threadId = ThreadId.make("thread:run-execution-settings-failure"); + const runId = RunId.make("run:run-execution-settings-failure"); + const attemptId = RunAttemptId.make("attempt:run-execution-settings-failure"); + const providerInstanceId = ProviderInstanceId.make("codex"); + const providerSessionId = ProviderSessionId.make("session:run-execution-settings-failure"); + const providerThreadId = ProviderThreadId.make( + "provider-thread:run-execution-settings-failure", + ); + const rootNodeId = NodeId.make("node:run-execution-settings-failure"); + const checkpointScope = { + id: CheckpointScopeId.make("checkpoint-scope:run-execution-settings-failure"), + } as OrchestrationV2CheckpointScope; + const providerStarts = yield* Ref.make(0); + const refreshes = yield* Ref.make(0); + const guardedWrites = yield* Ref.make(0); + const writes = yield* Ref.make>>([]); + const testLayer = runExecutionServiceLayer.pipe( + Layer.provide( + Layer.mergeAll( + Layer.mock(CheckpointServiceV2)({ + captureBaseline: () => + scenario === "start-guard" ? Effect.void : Effect.die("not reached"), + }), + Layer.mock(EventSinkV2)({ + writeIfRunCurrent: (input) => + Effect.gen(function* () { + assert.equal(input.threadId, threadId); + assert.equal(input.runId, runId); + assert.equal(input.activeAttemptId, attemptId); + assert.equal(input.expectedStatus, "running"); + yield* Ref.update(guardedWrites, (count) => count + 1); + if (scenario === "stale-attempt") { + return { committed: false, storedEvents: [] }; + } + yield* Ref.update(writes, (current) => [...current, input.events]); + return { committed: true, storedEvents: [] }; + }), + }), + idAllocatorLayer, + Layer.mock(ProviderEventIngestorV2)({ ingestNormalized: () => Effect.succeed([]) }), + scenario === "start-guard" + ? ServerSettingsService.layerTest() + : Layer.mock(ServerSettingsService)({ + getSettings: + scenario === "interruption" + ? Effect.interrupt + : Effect.fail( + new ServerSettingsError({ + settingsPath: "", + operation: "read-file", + cause: new Error("settings read failed"), + }), + ), + }), + Layer.succeed(RunFinalizationObserver, { + refresh: () => Effect.void, + refreshAfterTurn: () => Ref.update(refreshes, (count) => count + 1), + }), + ), + ), + ); + + const result = yield* Effect.gen(function* () { + const runExecution = yield* RunExecutionServiceV2; + yield* runExecution.startRootRun({ + commandId: CommandId.make("command:run-execution-settings-failure"), + appThread: { id: threadId } as OrchestrationV2AppThread, + providerSessionId, + session: { + events: Stream.never, + startTurn: () => Ref.update(providerStarts, (count) => count + 1), + } as unknown as ProviderAdapterV2SessionRuntime, + run: { + id: runId, + threadId, + ordinal: 1, + providerInstanceId, + status: "running", + } as OrchestrationV2Run, + rootNode: { id: rootNodeId, status: "running" } as OrchestrationV2ExecutionNode, + checkpointScope, + providerThread: { + id: providerThreadId, + driver, + } as OrchestrationV2ProviderThread, + attempt: { + id: attemptId, + providerTurnId: null, + status: "running", + } as OrchestrationV2RunAttempt, + attemptId, + providerTurnOrdinal: 1, + // A declined start is a normal exit, not a preparation failure. + ...(scenario === "start-guard" + ? { shouldStartProviderTurn: () => Effect.succeed(false) } + : {}), + message: { + messageId: MessageId.make("message:run-execution-settings-failure"), + text: "Start after settings fail.", + attachments: [], + createdBy: "user", + creationSource: "web", + }, + modelSelection: { instanceId: providerInstanceId, model: "gpt-5.4" }, + runtimePolicy: { + runtimeMode: "full-access", + interactionMode: "default", + cwd: process.cwd(), + approvalPolicy: "never", + sandboxPolicy: { + type: "readOnly", + access: { type: "fullAccess" }, + networkAccess: false, + }, + }, + }); + }).pipe(Effect.provide(testLayer), Effect.exit); + + assert.equal(yield* Ref.get(providerStarts), 0); + const events = (yield* Ref.get(writes)).flat(); + if (scenario === "interruption") { + assert.isTrue(Exit.isFailure(result)); + if (Exit.isFailure(result)) assert.isTrue(Cause.hasInterruptsOnly(result.cause)); + assert.equal(yield* Ref.get(guardedWrites), 0); + assert.equal(yield* Ref.get(refreshes), 0); + assert.isEmpty(events); + return; + } + assert.isTrue(Exit.isSuccess(result)); + if (scenario === "start-guard") { + assert.equal(yield* Ref.get(guardedWrites), 0); + assert.equal(yield* Ref.get(refreshes), 0); + assert.isEmpty(events); + return; + } + assert.equal(yield* Ref.get(guardedWrites), 1); + if (scenario === "stale-attempt") { + assert.equal(yield* Ref.get(refreshes), 0); + assert.isEmpty(events); + return; + } + assert.equal(yield* Ref.get(refreshes), 1); + assert.deepEqual( + events + .filter( + (event) => + event.type === "run.updated" || + event.type === "run-attempt.updated" || + event.type === "node.updated", + ) + .map((event) => event.payload.status), + ["failed", "failed", "failed"], + ); + const errorItem = events.find( + (event) => event.type === "turn-item.updated" && event.payload.type === "error", + ); + assert.isDefined(errorItem); + if (errorItem?.type === "turn-item.updated" && errorItem.payload.type === "error") { + // The persisted item carries a bounded curated message; the exact + // underlying text stays in the logged cause. + assert.equal(errorItem.payload.failure.message, "Run preparation failed."); + } + }), + ); +} + it.effect("keeps ingesting owned child events after the root turn terminalizes", () => Effect.gen(function* () { const threadId = ThreadId.make("thread:run-execution-late-child"); diff --git a/apps/server/src/orchestration-v2/RunExecutionService.ts b/apps/server/src/orchestration-v2/RunExecutionService.ts index 98c3b27d8719..040644eff674 100644 --- a/apps/server/src/orchestration-v2/RunExecutionService.ts +++ b/apps/server/src/orchestration-v2/RunExecutionService.ts @@ -560,6 +560,10 @@ export const layer: Layer.Layer< readonly terminal: ProviderTerminalEvent; readonly failureItemPersisted: boolean; readonly refreshAfterTurn: Effect.Effect; + readonly writeIfRunCurrent?: { + readonly activeAttemptId: RunAttemptId; + readonly expectedStatus: OrchestrationV2Run["status"]; + }; }) => Effect.gen(function* () { const completedAt = yield* DateTime.now; @@ -652,7 +656,7 @@ export const layer: Layer.Layer< const checkpointCaptureCommandId = CommandId.make( `command:effect:checkpoint.capture:${input.run.id}`, ); - yield* eventSink.writeWithEffects({ + const finalization = { effects: input.terminal.status === "completed" ? [ @@ -760,57 +764,27 @@ export const layer: Layer.Layer< payload: finalizedProviderThread, }, ], - }); + } satisfies Parameters[0]; + if (input.writeIfRunCurrent !== undefined) { + const result = yield* eventSink.writeIfRunCurrent({ + threadId: input.run.threadId, + runId: input.run.id, + activeAttemptId: input.writeIfRunCurrent.activeAttemptId, + expectedStatus: input.writeIfRunCurrent.expectedStatus, + events: finalization.events, + }); + if (!result.committed) { + return; + } + } else { + yield* eventSink.writeWithEffects(finalization); + } yield* input.refreshAfterTurn; }); return RunExecutionServiceV2.of({ startRootRun: (input) => Effect.gen(function* () { - const responseStreamingMode = yield* serverSettings.getSettings.pipe( - Effect.map( - (settings) => - resolveProjectSettings(settings, input.appThread.projectId).settings - .responseStreamingMode, - ), - Effect.mapError( - (cause) => - new RunExecutionStartError({ - commandId: input.commandId, - runId: input.run.id, - cause, - }), - ), - ); - yield* checkpointService - .captureBaseline({ - scope: input.checkpointScope, - ordinalWithinScope: Math.max(0, input.run.ordinal - 1), - }) - .pipe( - Effect.catchCause((cause) => - Cause.hasInterruptsOnly(cause) - ? Effect.failCause(cause) - : Effect.logWarning( - "orchestration V2 checkpoint baseline capture failed; starting provider without a baseline", - { runId: input.run.id }, - ), - ), - Effect.mapError( - (cause) => - new RunExecutionStartError({ - commandId: input.commandId, - runId: input.run.id, - cause, - }), - ), - ); - if ( - input.shouldStartProviderTurn !== undefined && - !(yield* input.shouldStartProviderTurn()) - ) { - return; - } // Startup failure and stream shutdown can report the same attempt. const refreshAfterTurn = yield* Effect.cached( finalizationObserver.refreshAfterTurn(input.appThread.projectId).pipe( @@ -823,7 +797,6 @@ export const layer: Layer.Layer< ), ), ); - const terminalEvent = yield* Ref.make(null); const makeFailedTerminalEvent = ( failure: OrchestrationV2ProviderFailure, failureItemOrdinal: number, @@ -843,6 +816,85 @@ export const layer: Layer.Layer< failure, threadDisposition: "reusable", }); + const responseStreamingMode = yield* Effect.gen(function* () { + const responseStreamingMode = yield* serverSettings.getSettings.pipe( + Effect.map( + (settings) => + resolveProjectSettings(settings, input.appThread.projectId).settings + .responseStreamingMode, + ), + ); + yield* checkpointService + .captureBaseline({ + scope: input.checkpointScope, + ordinalWithinScope: Math.max(0, input.run.ordinal - 1), + }) + .pipe( + Effect.catchCause((cause) => + Cause.hasInterruptsOnly(cause) + ? Effect.failCause(cause) + : Effect.logWarning( + "orchestration V2 checkpoint baseline capture failed; starting provider without a baseline", + { runId: input.run.id }, + ), + ), + ); + if ( + input.shouldStartProviderTurn !== undefined && + !(yield* input.shouldStartProviderTurn()) + ) { + return null; + } + return responseStreamingMode; + }).pipe( + Effect.catchCause((cause) => + Effect.gen(function* () { + if (Cause.hasInterruptsOnly(cause)) { + return yield* Effect.failCause(cause); + } + yield* Effect.logError("orchestration V2 run preparation failed", { + runId: input.run.id, + cause, + }); + yield* writeFinalRunEvents({ + run: input.run, + rootNode: input.rootNode, + checkpointScope: input.checkpointScope, + providerThread: input.providerThread, + attempt: input.attempt, + terminal: makeFailedTerminalEvent( + makeProviderFailure({ + cause: Cause.squash(cause), + // Keep exact underlying text in the logged cause only; + // the persisted turn item gets a bounded curated message. + message: "Run preparation failed.", + class: "unknown", + }), + input.providerTurnOrdinal * 100 + 1, + ), + failureItemPersisted: false, + refreshAfterTurn, + writeIfRunCurrent: { + activeAttemptId: input.attemptId, + expectedStatus: "running", + }, + }); + return null; + }), + ), + Effect.mapError( + (cause) => + new RunExecutionStartError({ + commandId: input.commandId, + runId: input.run.id, + cause, + }), + ), + ); + if (responseStreamingMode === null) { + return; + } + const terminalEvent = yield* Ref.make(null); const latestTurnItemOrdinal = yield* Ref.make(input.providerTurnOrdinal * 100); const latestProviderThread = yield* Ref.make(input.providerThread); const routeIdentity: ProviderEventRouteIdentity = {