Skip to content
36 changes: 36 additions & 0 deletions packages/opencode/src/session/lifecycle-provenance.ts
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ const activeRunsByDirectory = new Map<string, number>()
const idleWaiters = new Set<{ directories: readonly string[]; resolve: () => void }>()
const closingByDirectory = new Map<string, number>()
const closeWaiters = new Set<{ directory: string; resolve: (release: () => void) => void }>()
const closingStartWaiters = new Set<{ directory: string; resolve: () => void }>()

export function directoryKey(directory: string): string {
const digest = createHash("sha256").update(directory).digest("hex").slice(0, 16)
Expand Down Expand Up @@ -103,6 +104,7 @@ export async function withLifecycleCloseAction<T>(
stack.push(action)
activeByDirectory.set(directory, stack)
}
notifyClosingStartWaiters(directories)
try {
return await fn()
} finally {
Expand Down Expand Up @@ -142,6 +144,39 @@ function hasLifecycleClose(directories: readonly string[]): boolean {
return directories.some((directory) => (closingByDirectory.get(directory) ?? 0) > 0)
}

export function isLifecycleClosing(directory: string): boolean {
return (closingByDirectory.get(directory) ?? 0) > 0
}

export function whenLifecycleCloseBegins(directory: string): { promise: Promise<void>; cancel: () => void } {
if (isLifecycleClosing(directory) || currentLifecycleCloseAction(directory)) {
return { promise: Promise.resolve(), cancel: () => {} }
}
const waiter = { directory, resolve: () => {} }
const promise = new Promise<void>((resolve) => {
waiter.resolve = resolve
})
closingStartWaiters.add(waiter)
return {
promise,
cancel: () => {
closingStartWaiters.delete(waiter)
},
}
}

export function closingStartWaiterCount(): number {
return closingStartWaiters.size
}

function notifyClosingStartWaiters(directories: readonly string[]) {
for (const waiter of [...closingStartWaiters]) {
if (!directories.includes(waiter.directory)) continue
closingStartWaiters.delete(waiter)
waiter.resolve()
}
}

function notifyCloseWaiters() {
for (const waiter of [...closeWaiters]) {
if (hasLifecycleClose([waiter.directory])) continue
Expand Down Expand Up @@ -225,6 +260,7 @@ export function beginLifecycleClose(directories: readonly string[]): () => void
for (const directory of uniqueDirectories) {
closingByDirectory.set(directory, (closingByDirectory.get(directory) ?? 0) + 1)
}
notifyClosingStartWaiters(uniqueDirectories)
let released = false
return () => {
if (released) return
Expand Down
75 changes: 58 additions & 17 deletions packages/opencode/src/session/processor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,12 @@ import { InstanceState } from "@/effect/instance-state"
import { TurnChange } from "./turn-change"
import { LLMTrace } from "./llm-trace"
import { RunObservability } from "./run-observability"
import { currentLifecycleCloseAction, lifecycleCloseActionMeta } from "./lifecycle-provenance"
import {
currentLifecycleCloseAction,
isLifecycleClosing,
lifecycleCloseActionMeta,
whenLifecycleCloseBegins,
} from "./lifecycle-provenance"

const log = Log.create({ service: "session.processor" })
const TOOL_CLEANUP_TIMEOUT_MS = 1_000
Expand Down Expand Up @@ -98,6 +103,7 @@ type Input = {
assistantMessage: MessageV2.Assistant
sessionID: SessionID
model: Provider.Model
safeRecoveryDelay?: (attempt: number) => number
}

export interface Interface {
Expand Down Expand Up @@ -1231,19 +1237,34 @@ export const layer: Layer.Layer<

const retryStillAllowed = Effect.fn("SessionProcessor.retryStillAllowed")(function* (stage: string) {
const lifecycleAction = currentLifecycleCloseAction(ctx.directory)
if (!lifecycleAction) return { allowed: true as const }
ctx.runTrace.recordScopeClosed({
at: Date.now(),
monotonicMs: performance.now(),
source: `session.processor.safe_recovery.${stage}`,
reason: "lifecycle_close_before_auto_retry",
propagationPoint: "session.processor.safe_recovery",
...lifecycleCloseActionMeta(lifecycleAction),
})
return {
allowed: false as const,
interruptionMessage: LOCAL_LIFECYCLE_CLOSE_INTERRUPTION_MESSAGE,
if (lifecycleAction) {
ctx.runTrace.recordScopeClosed({
at: Date.now(),
monotonicMs: performance.now(),
source: `session.processor.safe_recovery.${stage}`,
reason: "lifecycle_close_before_auto_retry",
propagationPoint: "session.processor.safe_recovery",
...lifecycleCloseActionMeta(lifecycleAction),
})
return {
allowed: false as const,
interruptionMessage: LOCAL_LIFECYCLE_CLOSE_INTERRUPTION_MESSAGE,
}
}
if (isLifecycleClosing(ctx.directory)) {
ctx.runTrace.recordScopeClosed({
at: Date.now(),
monotonicMs: performance.now(),
source: `session.processor.safe_recovery.${stage}`,
reason: "maintenance_lifecycle_close_pending",
propagationPoint: "session.processor.safe_recovery",
})
return {
allowed: false as const,
interruptionMessage: LOCAL_LIFECYCLE_CLOSE_INTERRUPTION_MESSAGE,
}
}
return { allowed: true as const }
})

const retrySignalFor = (error: unknown) => {
Expand Down Expand Up @@ -1311,6 +1332,7 @@ export const layer: Layer.Layer<
presentation: info.presentation,
reason: info.reason,
}),
delay: input.safeRecoveryDelay,
}),
)

Expand Down Expand Up @@ -1428,15 +1450,34 @@ export const layer: Layer.Layer<
if (beforeRetry.allowed) {
automaticStreamRetriesUsed += 1
yield* removeReasoningForAttempt(attemptID)
const safeRecoveryScheduled = yield* safeRecoveryStep(undefined).pipe(
Effect.as(true),
Effect.catchCause(() => Effect.succeed(false)),
const lifecycleCloseWatch = Effect.callback<"lifecycle_close">((resume) => {
const handle = whenLifecycleCloseBegins(ctx.directory)
handle.promise.then(() => resume(Effect.succeed("lifecycle_close" as const)))
return Effect.sync(() => handle.cancel())
})
const backoffResult = yield* Effect.race(
safeRecoveryStep(undefined).pipe(Effect.as("scheduled" as const)),
lifecycleCloseWatch,
).pipe(
Effect.onInterrupt(() => recordProcessInterrupt(attemptID)),
Comment thread
Astro-Han marked this conversation as resolved.
Effect.catchCause((cause) =>
Cause.hasInterruptsOnly(cause)
? Effect.interrupt
: Effect.succeed("exhausted" as const),
),
)
if (!safeRecoveryScheduled) {
if (backoffResult === "exhausted") {
yield* writeSafeRetryFailedNotice(attemptID)
break
}
if (backoffResult === "lifecycle_close") {
const closeCheck = yield* retryStillAllowed("during_backoff")
yield* halt(result.error, attemptID, {
recordFailure: false,
interruptionMessage: closeCheck.allowed ? LOCAL_LIFECYCLE_CLOSE_INTERRUPTION_MESSAGE : closeCheck.interruptionMessage,
})
break
}
const afterRetry = yield* retryStillAllowed("after_backoff")
if (afterRetry.allowed) {
ctx.runTrace.recordAutoRetryAttempted({
Expand Down
10 changes: 6 additions & 4 deletions packages/opencode/src/session/retry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,8 +15,7 @@ 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
export const SAFE_RECOVERY_MAX_ATTEMPTS = 3

function cap(ms: number) {
return Math.min(ms, RETRY_MAX_DELAY)
Expand Down Expand Up @@ -185,20 +184,23 @@ export function safeRecoveryPolicy(opts: {
presentation: "recovery"
reason: "network_connection_dropped"
}) => Effect.Effect<void>
delay?: (attempt: number) => number
}) {
const delayFn = opts.delay ?? delay
return Schedule.fromStepWithMetadata(
Effect.succeed((meta: Schedule.InputMetadata<unknown>) => {
if (meta.attempt > SAFE_RECOVERY_MAX_ATTEMPTS) return Cause.done(meta.attempt)
return Effect.gen(function* () {
const wait = delayFn(meta.attempt)
const now = yield* Clock.currentTimeMillis
yield* opts.set({
attempt: meta.attempt,
message: "",
next: now + SAFE_RECOVERY_REPLAY_DELAY,
next: now + wait,
presentation: "recovery",
reason: "network_connection_dropped",
})
return [meta.attempt, Duration.millis(SAFE_RECOVERY_REPLAY_DELAY)] as [number, Duration.Duration]
return [meta.attempt, Duration.millis(wait)] as [number, Duration.Duration]
})
}),
)
Expand Down
18 changes: 12 additions & 6 deletions packages/opencode/src/session/run-incident/policy.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,12 @@
import { allowsBeforeProgressRetry as boundaryAllowsBeforeProgressRetry } from "../run-observability/boundary"
import { RETRY_INITIAL_DELAY, SAFE_RECOVERY_MAX_ATTEMPTS } from "../retry"
import type { IncidentFacts, RecoveryDecision, TerminalCause } from "./types"

const SAFE_RECOVERY_AUTO_RETRY = {
max_attempts: SAFE_RECOVERY_MAX_ATTEMPTS,
backoff_ms: RETRY_INITIAL_DELAY,
} as const

export function recoveryFor(input: {
cause: TerminalCause
facts: IncidentFacts
Expand Down Expand Up @@ -35,10 +41,10 @@ export function recoveryFor(input: {
) {
return {
...base,
recommendation: "auto_retry_once",
recommendation: "auto_retry",
confidence: "high",
reason: "no_visible_output_or_tool_execution",
auto_retry: { max_attempts: 1, backoff_ms: 1_000 },
auto_retry: SAFE_RECOVERY_AUTO_RETRY,
}
}
if (isBeforeFirstProviderProgressCause(input.cause) && beforeProgressBoundaryEvidenceBlocksRetry(terminalFacts)) {
Expand All @@ -60,10 +66,10 @@ export function recoveryFor(input: {
) {
return {
...base,
recommendation: "auto_retry_once",
recommendation: "auto_retry",
confidence: "high",
reason: "reasoning_only_without_final_text_or_tool_activity",
auto_retry: { max_attempts: 1, backoff_ms: 1_000 },
auto_retry: SAFE_RECOVERY_AUTO_RETRY,
}
}
if (!terminalFacts.side_effect_facts_complete) {
Expand Down Expand Up @@ -147,10 +153,10 @@ export function recoveryFor(input: {
if (retryableTransport) {
return {
...base,
recommendation: "auto_retry_once",
recommendation: "auto_retry",
confidence: "medium",
reason: "no_visible_output_or_tool_execution",
auto_retry: { max_attempts: 1, backoff_ms: 1_000 },
auto_retry: SAFE_RECOVERY_AUTO_RETRY,
}
}
return { ...base, recommendation: "unknown", confidence: "low", reason: "unknown" }
Expand Down
2 changes: 1 addition & 1 deletion packages/opencode/src/session/run-incident/presentation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ export function plainSummary(input: { cause: TerminalCause; facts: IncidentFacts
}

function actionKey(recovery: RecoveryDecision) {
if (recovery.recommendation === "auto_retry_once") return "run_incident.action.retry"
if (recovery.recommendation === "auto_retry") return "run_incident.action.retry"
if (recovery.recommendation === "offer_continue") return "run_incident.action.continue"
if (recovery.recommendation === "offer_resume_with_confirmation") return "run_incident.action.confirm_continue"
if (recovery.recommendation === "ask_user_before_retry") return "run_incident.action.confirm_retry"
Expand Down
2 changes: 1 addition & 1 deletion packages/opencode/src/session/run-incident/safety-gate.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ export function evaluateReplaySafety(input: {
}): ReplaySafetyDecision {
const safety = input.recovery

if (safety.recommendation === "auto_retry_once") {
if (safety.recommendation === "auto_retry") {
const maxAttempts = safety.auto_retry?.max_attempts ?? 1
if (input.safeRecoveryAttempt < maxAttempts) {
return {
Expand Down
4 changes: 2 additions & 2 deletions packages/opencode/src/session/run-incident/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -168,7 +168,7 @@ export type MaterializedToolBoundary = {

export type RecoveryDecision = {
recommendation:
| "auto_retry_once"
| "auto_retry"
| "offer_continue"
| "offer_resume_with_confirmation"
| "ask_user_before_retry"
Expand All @@ -188,7 +188,7 @@ export type RecoveryDecision = {
| "local_lifecycle_close"
| "user_cancel"
| "unknown"
auto_retry?: { max_attempts: 1; backoff_ms: number; attempted_at?: number }
auto_retry?: { max_attempts: number; backoff_ms: number; attempted_at?: number }
user_action?: { kind: "continue" | "resume" | "retry" | "confirm_continue" | "dismiss"; idempotency_key: string }
safety_scope: "visible_output_and_tool_side_effects"
}
Expand Down
6 changes: 5 additions & 1 deletion packages/opencode/src/session/run-observability/recorder.ts
Original file line number Diff line number Diff line change
Expand Up @@ -534,7 +534,11 @@ export function createRecorder(input: RecorderInput): Recorder {
return deriveCurrentIncident(next.at, { includeRecoveredTerminal: true })?.recovery ?? unknownRecovery()
},
recordRecoveryDecision(next) {
if (recoveryDecision?.retry_attempted && next.safe_recovery_attempt > recoveryDecision.safe_recovery_attempt) {
if (
recoveryDecision?.retry_attempted &&
next.safe_recovery_attempt > recoveryDecision.safe_recovery_attempt &&
next.recovery_mode !== "auto_replay_blocked"
) {
rememberEvent(next.monotonicMs)
return
}
Expand Down
4 changes: 2 additions & 2 deletions packages/opencode/test/session/export.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2264,7 +2264,7 @@ describe("redactPart", () => {
monotonicMs: 125,
technical_retryable: true,
safety_gate_decision: {
recommendation: "auto_retry_once",
recommendation: "auto_retry",
confidence: "high",
reason: "no_visible_output_or_tool_execution",
safety_scope: "visible_output_and_tool_side_effects",
Expand Down Expand Up @@ -2321,7 +2321,7 @@ describe("redactPart", () => {
})
expect(sanitized.diagnostics.run_observability?.[0]?.recovery_decision).toMatchObject({
technical_retryable: true,
safety_gate_recommendation: "auto_retry_once",
safety_gate_recommendation: "auto_retry",
recovery_mode: "replay",
attempt_kind: "safe_recovery_replay",
timeout_policy: "reasoning_first_attempt",
Expand Down
Loading
Loading