diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index 478fee4e9720..4bac5fb5965e 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -646,6 +646,16 @@ const layer = Layer.effect( ctx.toolcalls = {} ctx.assistantMessage.time.completed = Date.now() yield* session.updateMessage(ctx.assistantMessage) + yield* status.setMessage(ctx.sessionID, ctx.assistantMessage.id) + if (ctx.assistantMessage.error) { + // Terminal listeners must see the failed message and its final parts before settling delivery. + yield* events.publish(Session.Event.Error, { + sessionID: ctx.sessionID, + messageID: ctx.assistantMessage.id, + error: ctx.assistantMessage.error, + }) + yield* status.set(ctx.sessionID, { type: "idle" }) + } }) const halt = Effect.fn("SessionProcessor.halt")(function* (e: unknown) { @@ -660,8 +670,6 @@ const layer = Layer.effect( if ((yield* config.get()).compaction?.auto === false && !ctx.assistantMessage.summary) { ctx.assistantMessage.error = error ctx.assistantMessage.finish = "error" - yield* events.publish(Session.Event.Error, { sessionID: ctx.sessionID, error }) - yield* status.set(ctx.sessionID, { type: "idle" }) return } ctx.needsCompaction = true @@ -669,11 +677,6 @@ const layer = Layer.effect( return } ctx.assistantMessage.error = error - yield* events.publish(Session.Event.Error, { - sessionID: ctx.assistantMessage.sessionID, - error: ctx.assistantMessage.error, - }) - yield* status.set(ctx.sessionID, { type: "idle" }) }) const process = Effect.fn("SessionProcessor.process")(function* (streamInput: ProcessInput) { diff --git a/packages/opencode/src/session/retry.ts b/packages/opencode/src/session/retry.ts index 7ac151791212..6068ca085766 100644 --- a/packages/opencode/src/session/retry.ts +++ b/packages/opencode/src/session/retry.ts @@ -14,6 +14,7 @@ export type RetryReason = "free_tier_limit" | "account_rate_limit" | (string & { export type Retryable = { message: string + retries?: number action?: { reason: RetryReason provider: string @@ -66,9 +67,13 @@ export function delay(attempt: number, error?: SessionV1.APIError) { return cap(Math.min(RETRY_INITIAL_DELAY * Math.pow(RETRY_BACKOFF_FACTOR, attempt - 1), RETRY_MAX_DELAY_NO_HEADERS)) } -export function retryable(error: Err, provider: string) { +export function retryable(error: Err, provider: string): Retryable | undefined { // context overflow errors should not be retried if (SessionV1.ContextOverflowError.isInstance(error)) return undefined + const message = isRecord(error.data) ? error.data.message : undefined + if (typeof message === "string" && /\bgetaddrinfo (?:ETIMEOUT|ETIMEDOUT|EAI_AGAIN)\b/.test(message)) { + return { message: "Connection lookup failed. Retrying", retries: 2 } + } if (SessionV1.APIError.isInstance(error)) { const status = error.data.statusCode // 5xx errors are transient server failures and should always be retried, @@ -198,7 +203,8 @@ export function policy(opts: { const retry = retryable(error, opts.provider()) // Backoff and the retry budget restart on every swap; the status keeps counting across swaps. const sinceSwap = meta.attempt - swapAt - const exhausted = opts.retries !== undefined && sinceSwap > opts.retries + const retries = Math.min(opts.retries ?? Infinity, retry?.retries ?? Infinity) + const exhausted = sinceSwap > retries if (retry && !exhausted) { return Effect.gen(function* () { const wait = delay(sinceSwap, SessionV1.APIError.isInstance(error) ? error : undefined) diff --git a/packages/opencode/src/session/run-state.ts b/packages/opencode/src/session/run-state.ts index 5cefdd04a3f3..af12491fcfe6 100644 --- a/packages/opencode/src/session/run-state.ts +++ b/packages/opencode/src/session/run-state.ts @@ -90,7 +90,13 @@ const layer = Layer.effect( onInterrupt: Effect.Effect, work: Effect.Effect, ) { - return yield* (yield* runner(sessionID, onInterrupt)).ensureRunning(work) + return yield* (yield* runner(sessionID, onInterrupt)).ensureRunning( + work.pipe( + Effect.tap((result) => + result.info.role === "assistant" ? status.setMessage(sessionID, result.info.id) : Effect.void, + ), + ), + ) }) const startShell = Effect.fn("SessionRunState.startShell")(function* ( @@ -100,7 +106,14 @@ const layer = Layer.effect( ready?: Latch.Latch, ) { return yield* (yield* runner(sessionID, onInterrupt)) - .startShell(work, ready) + .startShell( + work.pipe( + Effect.tap((result) => + result.info.role === "assistant" ? status.setMessage(sessionID, result.info.id) : Effect.void, + ), + ), + ready, + ) .pipe(Effect.catchTag("RunnerBusy", () => Effect.fail(busyError(sessionID)))) }) diff --git a/packages/opencode/src/session/status.ts b/packages/opencode/src/session/status.ts index 11140acfeef5..ed335a17c80a 100644 --- a/packages/opencode/src/session/status.ts +++ b/packages/opencode/src/session/status.ts @@ -1,6 +1,6 @@ import { LayerNode } from "@opencode-ai/core/effect/layer-node" import { InstanceState } from "@/effect/instance-state" -import { SessionID } from "./schema" +import { MessageID, SessionID } from "./schema" import { Effect, Layer, Context } from "effect" import { EventV2Bridge } from "@/event-v2-bridge" import { SessionStatusEvent } from "@opencode-ai/schema/session-status-event" @@ -14,6 +14,7 @@ export interface Interface { readonly get: (sessionID: SessionID) => Effect.Effect readonly list: () => Effect.Effect> readonly set: (sessionID: SessionID, status: Info) => Effect.Effect + readonly setMessage: (sessionID: SessionID, messageID: MessageID) => Effect.Effect } export class Service extends Context.Service()("@opencode/SessionStatus") {} @@ -24,30 +25,38 @@ const layer = Layer.effect( const events = yield* EventV2Bridge.Service const state = yield* InstanceState.make( - Effect.fn("SessionStatus.state")(() => Effect.succeed(new Map())), + Effect.fn("SessionStatus.state")(() => + Effect.succeed({ statuses: new Map(), messages: new Map() }), + ), ) const get = Effect.fn("SessionStatus.get")(function* (sessionID: SessionID) { const data = yield* InstanceState.get(state) - return data.get(sessionID) ?? { type: "idle" as const } + return data.statuses.get(sessionID) ?? { type: "idle" as const } }) const list = Effect.fn("SessionStatus.list")(function* () { - return new Map(yield* InstanceState.get(state)) + return new Map((yield* InstanceState.get(state)).statuses) }) const set = Effect.fn("SessionStatus.set")(function* (sessionID: SessionID, status: Info) { const data = yield* InstanceState.get(state) yield* events.publish(Event.Status, { sessionID, status }) if (status.type === "idle") { - yield* events.publish(Event.Idle, { sessionID }) - data.delete(sessionID) + yield* events.publish(Event.Idle, { sessionID, messageID: data.messages.get(sessionID) }) + data.statuses.delete(sessionID) return } - data.set(sessionID, status) + if (!data.statuses.has(sessionID)) data.messages.delete(sessionID) + data.statuses.set(sessionID, status) }) - return Service.of({ get, list, set }) + const setMessage = Effect.fn("SessionStatus.setMessage")(function* (sessionID: SessionID, messageID: MessageID) { + const data = yield* InstanceState.get(state) + data.messages.set(sessionID, messageID) + }) + + return Service.of({ get, list, set, setMessage }) }), ) diff --git a/packages/opencode/test/session/processor-effect.test.ts b/packages/opencode/test/session/processor-effect.test.ts index 40f6a28b7369..45f3646c4b63 100644 --- a/packages/opencode/test/session/processor-effect.test.ts +++ b/packages/opencode/test/session/processor-effect.test.ts @@ -4,7 +4,7 @@ import { LayerNode } from "@opencode-ai/core/effect/layer-node" import { EventV2Bridge } from "@/event-v2-bridge" import { afterEach, expect } from "bun:test" import { tool } from "ai" -import { Cause, Effect, Exit, Fiber, Layer, Stream } from "effect" +import { Cause, Effect, Exit, Fiber, Layer, Schema, Stream } from "effect" import path from "path" import z from "zod" import type { Agent } from "../../src/agent/agent" @@ -1221,12 +1221,26 @@ itFragmentFailure.live("session.processor effect tests retain partial legacy par const chat = yield* session.create({}) const parent = yield* user(chat.id, "provider failure") + const database = yield* Database.Service const msg = yield* assistant(chat.id, parent.id, path.resolve(dir)) const mdl = yield* provider.getModel(ref.providerID, ref.modelID) const seen: string[] = [] + const terminal: { type: string; completed?: number; error?: string; messageID?: string }[] = [] const off = yield* events.listen((event) => { seen.push(event.type) - return Effect.void + if (event.type !== Session.Event.Error.type && event.type !== SessionStatus.Event.Idle.type) + return Effect.void + return Effect.gen(function* () { + const stored = yield* MessageV2.get({ sessionID: chat.id, messageID: msg.id }) + if (stored.info.role !== "assistant") throw new Error("Expected assistant message") + const data = Schema.decodeUnknownSync(Session.Event.Error.data)(event.data) + terminal.push({ + type: event.type, + completed: stored.info.time.completed, + error: stored.info.error?.name, + messageID: data.messageID, + }) + }).pipe(Effect.provideService(Database.Service, database), Effect.orDie) }) const handle = yield* processors.create({ assistantMessage: msg, sessionID: chat.id, model: mdl }) @@ -1259,6 +1273,12 @@ itFragmentFailure.live("session.processor effect tests retain partial legacy par ) expect(seen).toContain(MessageV2.Event.PartUpdated.type) expect(seen).toContain(Session.Event.Error.type) + expect(terminal.map((event) => event.type)).toEqual([Session.Event.Error.type, SessionStatus.Event.Idle.type]) + for (const event of terminal) { + expect(event.completed).toBeNumber() + expect(event.error).toBeDefined() + expect(event.messageID).toBe(msg.id) + } expect(seen.filter((type) => type.startsWith("session.next."))).toEqual([]) }), { config: cfg }, diff --git a/packages/opencode/test/session/retry.test.ts b/packages/opencode/test/session/retry.test.ts index dabac95c9379..e94b5dc63f4f 100644 --- a/packages/opencode/test/session/retry.test.ts +++ b/packages/opencode/test/session/retry.test.ts @@ -160,6 +160,36 @@ describe("session.retry.delay", () => { }) describe("session.retry.retryable", () => { + test.each(["ETIMEOUT", "ETIMEDOUT", "EAI_AGAIN"])("retries transient DNS %s failures", (code) => { + const error = MessageV2.fromError(new TypeError(`getaddrinfo ${code} proxy.example.com`), { providerID }) + expect(SessionRetry.retryable(error, retryProvider)).toMatchObject({ retries: 2 }) + }) + + test.each(["getaddrinfo ENOTFOUND proxy.example.com", "connect ETIMEDOUT", "fetch failed"])( + "does not classify %s as a transient DNS failure", + (message) => { + expect(SessionRetry.retryable(wrap(message), retryProvider)).toBeUndefined() + }, + ) + + test("stops transient DNS retries after two attempts even without a configured budget", async () => { + await Effect.runPromise( + Effect.gen(function* () { + const error = wrap("getaddrinfo ETIMEOUT proxy.example.com") + const step = yield* Schedule.toStepWithMetadata( + SessionRetry.policy({ + provider: () => retryProvider, + parse: () => error, + set: () => Effect.void, + }), + ) + expect(Duration.toMillis((yield* step(error)).duration)).toBe(2000) + expect(Duration.toMillis((yield* step(error)).duration)).toBe(4000) + expect(Exit.isFailure(yield* Effect.exit(step(error)))).toBe(true) + }), + ) + }, 10_000) + test("maps too_many_requests json messages", () => { const error = wrap(JSON.stringify({ type: "error", error: { type: "too_many_requests" } })) expect(SessionRetry.retryable(error, retryProvider)).toEqual({ message: "Too Many Requests" }) diff --git a/packages/schema/src/session-status-event.ts b/packages/schema/src/session-status-event.ts index f6a3022bcb1d..aa2307cf3e81 100644 --- a/packages/schema/src/session-status-event.ts +++ b/packages/schema/src/session-status-event.ts @@ -5,6 +5,7 @@ import { optional } from "./schema" import { Event } from "./event" import { NonNegativeInt } from "./schema" import { SessionID } from "./session-id" +import { MessageID } from "./v1/session" export const Info = Schema.Union([ Schema.Struct({ @@ -45,6 +46,7 @@ export const Idle = Event.define({ type: "session.idle", schema: { sessionID: SessionID, + messageID: optional(MessageID), }, }) diff --git a/packages/schema/src/v1/session.ts b/packages/schema/src/v1/session.ts index 75e9282f117c..c0c18f2e281b 100644 --- a/packages/schema/src/v1/session.ts +++ b/packages/schema/src/v1/session.ts @@ -652,6 +652,7 @@ export const Error = define({ type: "session.error", schema: { sessionID: Schema.optional(SessionID), + messageID: Schema.optional(MessageID), error: Assistant.fields.error, }, }) diff --git a/packages/sdk/js/src/v2/gen/types.gen.ts b/packages/sdk/js/src/v2/gen/types.gen.ts index 5e067f3afb23..f3f1c181a973 100644 --- a/packages/sdk/js/src/v2/gen/types.gen.ts +++ b/packages/sdk/js/src/v2/gen/types.gen.ts @@ -1214,6 +1214,7 @@ export type GlobalEvent = { type: "session.error" properties: { sessionID?: string + messageID?: string error?: | ProviderAuthError | UnknownError @@ -1503,6 +1504,7 @@ export type GlobalEvent = { type: "session.idle" properties: { sessionID: string + messageID?: string } } | { @@ -1746,6 +1748,7 @@ export type ProviderConfig = { baseURL?: string enterpriseUrl?: string setCacheKey?: boolean + promptCacheKey?: string /** * Timeout in milliseconds for full requests to this provider. Set to false to disable timeout. */ @@ -5348,6 +5351,7 @@ export type SessionError = { location?: LocationRef data: { sessionID?: string + messageID?: string error?: | ProviderAuthError | UnknownError @@ -5919,6 +5923,7 @@ export type SessionIdle = { location?: LocationRef data: { sessionID: string + messageID?: string } } @@ -6673,6 +6678,7 @@ export type EventSessionError = { type: "session.error" properties: { sessionID?: string + messageID?: string error?: | ProviderAuthError | UnknownError @@ -6937,6 +6943,7 @@ export type EventSessionIdle = { type: "session.idle" properties: { sessionID: string + messageID?: string } }