Skip to content
Draft
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
17 changes: 10 additions & 7 deletions packages/opencode/src/session/processor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -660,20 +670,13 @@ 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
yield* events.publish(Session.Event.Error, { sessionID: ctx.sessionID, error })
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) {
Expand Down
10 changes: 8 additions & 2 deletions packages/opencode/src/session/retry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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)
Expand Down
17 changes: 15 additions & 2 deletions packages/opencode/src/session/run-state.ts
Original file line number Diff line number Diff line change
Expand Up @@ -90,7 +90,13 @@ const layer = Layer.effect(
onInterrupt: Effect.Effect<SessionV1.WithParts>,
work: Effect.Effect<SessionV1.WithParts>,
) {
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* (
Expand All @@ -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))))
})

Expand Down
25 changes: 17 additions & 8 deletions packages/opencode/src/session/status.ts
Original file line number Diff line number Diff line change
@@ -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"
Expand All @@ -14,6 +14,7 @@ export interface Interface {
readonly get: (sessionID: SessionID) => Effect.Effect<Info>
readonly list: () => Effect.Effect<Map<SessionID, Info>>
readonly set: (sessionID: SessionID, status: Info) => Effect.Effect<void>
readonly setMessage: (sessionID: SessionID, messageID: MessageID) => Effect.Effect<void>
}

export class Service extends Context.Service<Service, Interface>()("@opencode/SessionStatus") {}
Expand All @@ -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<SessionID, Info>())),
Effect.fn("SessionStatus.state")(() =>
Effect.succeed({ statuses: new Map<SessionID, Info>(), messages: new Map<SessionID, MessageID>() }),
),
)

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 })
}),
)

Expand Down
24 changes: 22 additions & 2 deletions packages/opencode/test/session/processor-effect.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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 })

Expand Down Expand Up @@ -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 },
Expand Down
30 changes: 30 additions & 0 deletions packages/opencode/test/session/retry.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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" })
Expand Down
2 changes: 2 additions & 0 deletions packages/schema/src/session-status-event.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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({
Expand Down Expand Up @@ -45,6 +46,7 @@ export const Idle = Event.define({
type: "session.idle",
schema: {
sessionID: SessionID,
messageID: optional(MessageID),
},
})

Expand Down
1 change: 1 addition & 0 deletions packages/schema/src/v1/session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -652,6 +652,7 @@ export const Error = define({
type: "session.error",
schema: {
sessionID: Schema.optional(SessionID),
messageID: Schema.optional(MessageID),
error: Assistant.fields.error,
},
})
Expand Down
7 changes: 7 additions & 0 deletions packages/sdk/js/src/v2/gen/types.gen.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1214,6 +1214,7 @@ export type GlobalEvent = {
type: "session.error"
properties: {
sessionID?: string
messageID?: string
error?:
| ProviderAuthError
| UnknownError
Expand Down Expand Up @@ -1503,6 +1504,7 @@ export type GlobalEvent = {
type: "session.idle"
properties: {
sessionID: string
messageID?: string
}
}
| {
Expand Down Expand Up @@ -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.
*/
Expand Down Expand Up @@ -5348,6 +5351,7 @@ export type SessionError = {
location?: LocationRef
data: {
sessionID?: string
messageID?: string
error?:
| ProviderAuthError
| UnknownError
Expand Down Expand Up @@ -5919,6 +5923,7 @@ export type SessionIdle = {
location?: LocationRef
data: {
sessionID: string
messageID?: string
}
}

Expand Down Expand Up @@ -6673,6 +6678,7 @@ export type EventSessionError = {
type: "session.error"
properties: {
sessionID?: string
messageID?: string
error?:
| ProviderAuthError
| UnknownError
Expand Down Expand Up @@ -6937,6 +6943,7 @@ export type EventSessionIdle = {
type: "session.idle"
properties: {
sessionID: string
messageID?: string
}
}

Expand Down
Loading