Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 12 additions & 3 deletions apps/server/src/orchestration-v2/EffectWorker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void, OrchestrationEffectExecutionError>;
}

Expand Down Expand Up @@ -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({
Expand Down Expand Up @@ -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) =>
Expand Down Expand Up @@ -297,6 +303,7 @@ export const executorLayer: Layer.Layer<
providerTurnStart.start({
threadId: effect.threadId,
runId: effect.request.runId,
willRetry,
}),
),
Effect.mapError(
Expand Down Expand Up @@ -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)),
);
Expand Down
22 changes: 22 additions & 0 deletions apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)),
};
}

Expand Down Expand Up @@ -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({
Expand Down
8 changes: 8 additions & 0 deletions apps/server/src/orchestration-v2/ProviderTurnStartService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -61,9 +61,15 @@ export class ProviderTurnStartError extends Schema.TaggedError<ProviderTurnStart
const isProviderTurnStartError = Schema.is(ProviderTurnStartError);

export interface ProviderTurnStartServiceV2Shape {
/**
* Starts the run's provider turn. When `willRetry` is true, a session open
* failure is returned so the caller can retry. Otherwise the run is settled
* as failed.
*/
readonly start: (input: {
readonly threadId: ThreadId;
readonly runId: RunId;
readonly willRetry?: boolean;
}) => Effect.Effect<void, ProviderTurnStartError>;
}

Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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;
Expand Down
Loading