diff --git a/apps/server/src/orchestration-v2/EffectWorker.ts b/apps/server/src/orchestration-v2/EffectWorker.ts index e78c44cedc7f..ee3aec77665a 100644 --- a/apps/server/src/orchestration-v2/EffectWorker.ts +++ b/apps/server/src/orchestration-v2/EffectWorker.ts @@ -67,8 +67,13 @@ export function isNonRetryableProviderTurnControlFailure( } export interface OrchestrationEffectExecutorV2Shape { + /** + * Runs one claimed effect. `willRetry` is true when the worker will retry a + * failure, so a step can fail and try again instead of settling the run. + */ readonly execute: ( effect: OrchestrationEffectV2, + options?: { readonly willRetry: boolean }, ) => Effect.Effect; } @@ -103,7 +108,8 @@ export const executorLayer: Layer.Layer< const threads = yield* ThreadManagementService; const settings = yield* ServerSettingsService; return OrchestrationEffectExecutorV2.of({ - execute: (effect) => { + execute: (effect, options) => { + const willRetry = options?.willRetry ?? false; switch (effect.request.type) { case "provider-runtime.continue": return continueRestartedRun({ @@ -143,7 +149,7 @@ export const executorLayer: Layer.Layer< ); case "provider-turn.start": return providerTurnStart - .start({ threadId: effect.threadId, runId: effect.request.runId }) + .start({ threadId: effect.threadId, runId: effect.request.runId, willRetry }) .pipe( Effect.mapError( (cause) => @@ -297,6 +303,7 @@ export const executorLayer: Layer.Layer< providerTurnStart.start({ threadId: effect.threadId, runId: effect.request.runId, + willRetry, }), ), Effect.mapError( @@ -596,7 +603,9 @@ export const layerWithOptions = ( }).pipe(Effect.onError((cause) => requeueClaim(effect, cause))); if (cancelledBeforeExecution) return true; - const execution = executor.execute(effect).pipe(Effect.as("executed" as const)); + const execution = executor + .execute(effect, { willRetry: effect.attemptCount < maxAttempts }) + .pipe(Effect.as("executed" as const)); const exit = yield* Effect.exit(Effect.raceFirst(execution, cancellation)).pipe( Effect.ensuring(outbox.clearCancellation(effect.id)), ); diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts index f364edd5a49f..5761689fcedb 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts @@ -445,6 +445,13 @@ function makeLocalCommandHarness(input: { start: Effect.gen(function* () { yield* (yield* ProviderTurnStart.ProviderTurnStartServiceV2).start({ threadId, runId }); }).pipe(Effect.provide(layer)), + startWithRetry: Effect.gen(function* () { + yield* (yield* ProviderTurnStart.ProviderTurnStartServiceV2).start({ + threadId, + runId, + willRetry: true, + }); + }).pipe(Effect.provide(layer)), }; } @@ -482,6 +489,21 @@ effectIt.effect("terminalizes a starting run when its provider session cannot op }), ); +effectIt.effect("leaves the run starting when a session-open failure will be retried", () => + Effect.gen(function* () { + const harness = makeLocalCommandHarness({ + text: "Continue", + openFailure: new Error("provider session rejected"), + }); + + const error = yield* harness.startWithRetry.pipe(Effect.flip); + + expect(error._tag).toBe("ProviderTurnStartError"); + expect(harness.writeIfRunCurrent).not.toHaveBeenCalled(); + expect(harness.projection().runs.at(-1)?.status).toBe("starting"); + }), +); + effectIt.effect("keeps a session-open failure retryable when terminal persistence fails", () => Effect.gen(function* () { const harness = makeLocalCommandHarness({ diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index 4d24dcb43372..9ea71cb4c9d5 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -61,9 +61,15 @@ export class ProviderTurnStartError extends Schema.TaggedError Effect.Effect; } @@ -205,6 +211,7 @@ export const layer: Layer.Layer< const start = Effect.fn("orchestrationV2.providerTurnStart.start")(function* (input: { readonly threadId: ThreadId; readonly runId: RunId; + readonly willRetry?: boolean; }) { const { runId } = input; const projection = yield* projectionStore.getTurnStartContext(input.threadId, runId); @@ -531,6 +538,7 @@ export const layer: Layer.Layer< }), ); if (sessionResult._tag === "Failure") { + if (input.willRetry === true) return yield* sessionResult.failure; const failedAt = yield* DateTime.now; const openError = sessionResult.failure; const nestedCause = "cause" in openError ? openError.cause : undefined;