From d625a16151bae30ebeb16f5c8c9325b59dbdd1ce Mon Sep 17 00:00:00 2001 From: Yuhan Lei Date: Tue, 26 May 2026 18:42:29 +0800 Subject: [PATCH 1/4] refactor(session): route safe recovery through retry policy --- packages/opencode/src/session/processor.ts | 38 ++++++---- packages/opencode/src/session/retry.ts | 29 +++++++ packages/opencode/test/session/retry.test.ts | 79 +++++++++++++++++++- 3 files changed, 129 insertions(+), 17 deletions(-) diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index e0477e0b9..7c3432f3b 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -1,4 +1,4 @@ -import { Cause, Deferred, Effect, Layer, Context, Scope } from "effect" +import { Cause, Deferred, Effect, Layer, Context, Scope, Schedule } from "effect" import * as Stream from "effect/Stream" import { Bus } from "@/bus" import { Config } from "@/config" @@ -31,7 +31,6 @@ import { currentLifecycleCloseAction, lifecycleCloseActionMeta } from "./lifecyc const log = Log.create({ service: "session.processor" }) const TOOL_CLEANUP_TIMEOUT_MS = 1_000 -const SAFE_RECOVERY_AUTO_RETRY_BACKOFF_MS = 1_000 export const REASONING_FIRST_ATTEMPT_CONNECT_TIMEOUT_MS = 60_000 export const REASONING_SAFE_RETRY_CONNECT_TIMEOUT_MS = 120_000 const LOCAL_LIFECYCLE_CLOSE_INTERRUPTION_MESSAGE = "The run was interrupted by a local lifecycle close." @@ -1301,6 +1300,20 @@ export const layer: Layer.Layer< safeRetryNoticeWritten = true }) + const safeRecoveryStep = yield* Schedule.toStepWithMetadata( + SessionRetry.safeRecoveryPolicy({ + set: (info) => + status.set(ctx.sessionID, { + type: "retry", + attempt: info.attempt, + message: info.message, + next: info.next, + presentation: info.presentation, + reason: info.reason, + }), + }), + ) + const runAttempt = Effect.fn("SessionProcessor.runAttempt")(function* () { ctx.currentText = undefined ctx.reasoningMap = {} @@ -1398,22 +1411,15 @@ export const layer: Layer.Layer< if (beforeRetry.allowed) { automaticStreamRetriesUsed += 1 yield* removeReasoningForAttempt(attemptID) - const next = Date.now() + SAFE_RECOVERY_AUTO_RETRY_BACKOFF_MS - yield* status.set(ctx.sessionID, { - type: "retry", - attempt: ctx.attemptCount, - message: - retryDecision.presentation === "safe_recovery" - ? "" - : (retrySignal.message ?? "Retrying interrupted stream"), - next, - ...(retryDecision.presentation === "safe_recovery" - ? { presentation: "safe_recovery" as const, reason: "network_connection_dropped" as const } - : {}), - }) - yield* Effect.sleep(`${SAFE_RECOVERY_AUTO_RETRY_BACKOFF_MS} millis`).pipe( + const safeRecoveryScheduled = yield* safeRecoveryStep(undefined).pipe( + Effect.as(true), + Effect.catchCause(() => Effect.succeed(false)), Effect.onInterrupt(() => recordProcessInterrupt(attemptID)), ) + if (!safeRecoveryScheduled) { + yield* writeSafeRetryFailedNotice(attemptID) + break + } const afterRetry = yield* retryStillAllowed("after_backoff") if (afterRetry.allowed) { ctx.runTrace.recordAutoRetryAttempted({ diff --git a/packages/opencode/src/session/retry.ts b/packages/opencode/src/session/retry.ts index 90e76e4b8..8b929f5f7 100644 --- a/packages/opencode/src/session/retry.ts +++ b/packages/opencode/src/session/retry.ts @@ -15,6 +15,8 @@ export const RETRY_BACKOFF_FACTOR = 2 export const RETRY_MAX_DELAY_NO_HEADERS = 30_000 // 30 seconds export const RETRY_MAX_DELAY = 2_147_483_647 // max 32-bit signed integer for setTimeout export const RETRY_MAX_ATTEMPTS = 10 +export const SAFE_RECOVERY_REPLAY_DELAY = 1_000 +export const SAFE_RECOVERY_MAX_ATTEMPTS = 1 function cap(ms: number) { return Math.min(ms, RETRY_MAX_DELAY) @@ -175,4 +177,31 @@ export function policy(opts: { ) } +export function safeRecoveryPolicy(opts: { + set: (input: { + attempt: number + message: string + next: number + presentation: "safe_recovery" + reason: "network_connection_dropped" + }) => Effect.Effect +}) { + return Schedule.fromStepWithMetadata( + Effect.succeed((meta: Schedule.InputMetadata) => { + if (meta.attempt > SAFE_RECOVERY_MAX_ATTEMPTS) return Cause.done(meta.attempt) + return Effect.gen(function* () { + const now = yield* Clock.currentTimeMillis + yield* opts.set({ + attempt: meta.attempt, + message: "", + next: now + SAFE_RECOVERY_REPLAY_DELAY, + presentation: "safe_recovery", + reason: "network_connection_dropped", + }) + return [meta.attempt, Duration.millis(SAFE_RECOVERY_REPLAY_DELAY)] as [number, Duration.Duration] + }) + }), + ) +} + export * as SessionRetry from "./retry" diff --git a/packages/opencode/test/session/retry.test.ts b/packages/opencode/test/session/retry.test.ts index b9696946e..6eea0f3dc 100644 --- a/packages/opencode/test/session/retry.test.ts +++ b/packages/opencode/test/session/retry.test.ts @@ -2,7 +2,7 @@ import { describe, expect, test } from "bun:test" import type { NamedError } from "@opencode-ai/util/error" import { APICallError } from "ai" import { setTimeout as sleep } from "node:timers/promises" -import { Effect, Exit, Schedule } from "effect" +import { Effect, Exit, Pull, Schedule } from "effect" import { SessionRetry } from "../../src/session/retry" import { MessageV2 } from "../../src/session/message-v2" import { ProviderID } from "../../src/provider/schema" @@ -128,6 +128,83 @@ describe("session.retry.delay", () => { }) }) + test("safe recovery policy emits lightweight retry presentation with separate attempt metadata", async () => { + await using tmp = await tmpdir() + await Instance.provide({ + directory: tmp.path, + fn: async () => { + const sessionID = SessionID.make("session-safe-recovery-retry-test") + + await Effect.runPromise( + Effect.gen(function* () { + const step = yield* Schedule.toStepWithMetadata( + SessionRetry.safeRecoveryPolicy({ + set: (info) => + Effect.promise(() => + AppRuntime.runPromise( + SessionStatus.Service.use((svc) => + svc.set(sessionID, { + type: "retry", + attempt: info.attempt, + message: info.message, + next: info.next, + presentation: info.presentation, + reason: info.reason, + }), + ), + ), + ), + }), + ) + yield* step(undefined) + }), + ) + + expect(await AppRuntime.runPromise(SessionStatus.Service.use((svc) => svc.get(sessionID)))).toMatchObject({ + type: "retry", + attempt: 1, + message: "", + presentation: "safe_recovery", + reason: "network_connection_dropped", + }) + }, + }) + }) + + test("safe recovery policy stops after the one replay budget is exhausted", async () => { + const statuses: Array<{ + attempt: number + message: string + next: number + presentation: "safe_recovery" + reason: "network_connection_dropped" + }> = [] + + const exit = await Effect.runPromise( + Effect.gen(function* () { + const step = yield* Schedule.toStepWithMetadata( + SessionRetry.safeRecoveryPolicy({ + set: (info) => Effect.sync(() => statuses.push(info)), + }), + ) + yield* step(undefined) + return yield* Effect.exit(step(undefined)) + }), + ) + + expect(Exit.isFailure(exit)).toBe(true) + if (Exit.isFailure(exit)) { + expect(Pull.isDoneCause(exit.cause)).toBe(true) + } + expect(statuses).toHaveLength(1) + expect(statuses[0]).toMatchObject({ + attempt: 1, + message: "", + presentation: "safe_recovery", + reason: "network_connection_dropped", + }) + }) + test("policy stops retrying after the configured max attempts", async () => { const attempts: number[] = [] let runs = 0 From fc54cbf4e6adbc19526b24a86a491c21f8204eed Mon Sep 17 00:00:00 2001 From: Yuhan Lei Date: Tue, 26 May 2026 19:33:43 +0800 Subject: [PATCH 2/4] chore: retrigger ci for PR #931 From 42459ee3895f2127c5087af0265de327097714bb Mon Sep 17 00:00:00 2001 From: Yuhan Lei Date: Tue, 26 May 2026 20:23:55 +0800 Subject: [PATCH 3/4] chore: retrigger ci for PR #931 --- .github/.ci-retrigger.tmp | 1 + 1 file changed, 1 insertion(+) create mode 100644 .github/.ci-retrigger.tmp diff --git a/.github/.ci-retrigger.tmp b/.github/.ci-retrigger.tmp new file mode 100644 index 000000000..8b1378917 --- /dev/null +++ b/.github/.ci-retrigger.tmp @@ -0,0 +1 @@ + From b26d661c83a06646d5e167d460139fe0ded15260 Mon Sep 17 00:00:00 2001 From: Yuhan Lei Date: Tue, 26 May 2026 20:23:55 +0800 Subject: [PATCH 4/4] chore: remove ci retrigger marker --- .github/.ci-retrigger.tmp | 1 - 1 file changed, 1 deletion(-) delete mode 100644 .github/.ci-retrigger.tmp diff --git a/.github/.ci-retrigger.tmp b/.github/.ci-retrigger.tmp deleted file mode 100644 index 8b1378917..000000000 --- a/.github/.ci-retrigger.tmp +++ /dev/null @@ -1 +0,0 @@ -