Compare commits

..

4 Commits

Author SHA1 Message Date
Kit Langton bba1c5bbf4 chore(client): generate provider error types 2026-06-26 00:21:46 -04:00
Kit Langton 007cdd6d70 fix(core): preserve provider session failures 2026-06-25 23:58:00 -04:00
opencode-agent[bot] 41f2705e72 chore: generate 2026-06-26 03:35:39 +00:00
Kit Langton ef5c9f4931 feat(sdk): expose active sessions (#33991) 2026-06-26 03:33:35 +00:00
31 changed files with 834 additions and 303 deletions
+1
View File
@@ -174,6 +174,7 @@ _Avoid_: Response envelope
- `sessions.messages(...)` returns a **Page** and uses the same cursor discipline as `sessions.list(...)`: the initial request supplies `sessionID`, ordering, and page size; continuation supplies `sessionID` plus only an opaque branded message cursor carrying ordering, page size, direction, and message anchor. Using a cursor with another Session is invalid.
- `sessions.message({ sessionID, messageID })` is a required resource lookup. An unknown Session fails with `SessionNotFoundError`; a known Session with an absent or differently owned message fails with `MessageNotFoundError` without disclosing cross-Session ownership. Absence is not represented as `undefined` across the public HTTP boundary.
- `sessions.interrupt({ sessionID })` first verifies that the durable Session exists, failing with `SessionNotFoundError` otherwise. For a known Session, interruption is idempotent: idle, already-settled, or locally unowned execution is a no-op.
- `sessions.active()` snapshots the current process's foreground Session drain registry as a record of Session IDs to `{ type: "running" }`. Missing IDs are inactive; background subagents and tasks do not make their parent Session active, and process restart clears the registry.
- `sessions.context({ sessionID })` preserves the existing message-only operation. It returns projected conversational messages selected as Session context; it does not include or represent the complete provider request context, whose baseline system context and other contributions remain separate.
- **Open question**: Should a future, separately named operation expose the complete provider request context, including baseline system context, selected source contributions, and context-epoch metadata?
- `sessions.prompt(...)` exposes `resume?: boolean`. Omitting it preserves durable admission followed by an advisory execution wake; `resume: false` requests durable admit-only behavior.
+74 -67
View File
@@ -55,43 +55,49 @@ const Endpoint0_1 = (raw: RawClient["server.session"]) => (input?: Endpoint0_1In
Effect.map((value) => value.data),
)
type Endpoint0_2Request = Parameters<RawClient["server.session"]["session.get"]>[0]
type Endpoint0_2Input = { readonly sessionID: Endpoint0_2Request["params"]["sessionID"] }
const Endpoint0_2 = (raw: RawClient["server.session"]) => (input: Endpoint0_2Input) =>
const Endpoint0_2 = (raw: RawClient["server.session"]) => () =>
raw["session.active"]({}).pipe(
Effect.mapError(mapClientError),
Effect.map((value) => value.data),
)
type Endpoint0_3Request = Parameters<RawClient["server.session"]["session.get"]>[0]
type Endpoint0_3Input = { readonly sessionID: Endpoint0_3Request["params"]["sessionID"] }
const Endpoint0_3 = (raw: RawClient["server.session"]) => (input: Endpoint0_3Input) =>
raw["session.get"]({ params: { sessionID: input.sessionID } }).pipe(
Effect.mapError(mapClientError),
Effect.map((value) => value.data),
)
type Endpoint0_3Request = Parameters<RawClient["server.session"]["session.switchAgent"]>[0]
type Endpoint0_3Input = {
readonly sessionID: Endpoint0_3Request["params"]["sessionID"]
readonly agent: Endpoint0_3Request["payload"]["agent"]
type Endpoint0_4Request = Parameters<RawClient["server.session"]["session.switchAgent"]>[0]
type Endpoint0_4Input = {
readonly sessionID: Endpoint0_4Request["params"]["sessionID"]
readonly agent: Endpoint0_4Request["payload"]["agent"]
}
const Endpoint0_3 = (raw: RawClient["server.session"]) => (input: Endpoint0_3Input) =>
const Endpoint0_4 = (raw: RawClient["server.session"]) => (input: Endpoint0_4Input) =>
raw["session.switchAgent"]({ params: { sessionID: input.sessionID }, payload: { agent: input.agent } }).pipe(
Effect.mapError(mapClientError),
)
type Endpoint0_4Request = Parameters<RawClient["server.session"]["session.switchModel"]>[0]
type Endpoint0_4Input = {
readonly sessionID: Endpoint0_4Request["params"]["sessionID"]
readonly model: Endpoint0_4Request["payload"]["model"]
type Endpoint0_5Request = Parameters<RawClient["server.session"]["session.switchModel"]>[0]
type Endpoint0_5Input = {
readonly sessionID: Endpoint0_5Request["params"]["sessionID"]
readonly model: Endpoint0_5Request["payload"]["model"]
}
const Endpoint0_4 = (raw: RawClient["server.session"]) => (input: Endpoint0_4Input) =>
const Endpoint0_5 = (raw: RawClient["server.session"]) => (input: Endpoint0_5Input) =>
raw["session.switchModel"]({ params: { sessionID: input.sessionID }, payload: { model: input.model } }).pipe(
Effect.mapError(mapClientError),
)
type Endpoint0_5Request = Parameters<RawClient["server.session"]["session.prompt"]>[0]
type Endpoint0_5Input = {
readonly sessionID: Endpoint0_5Request["params"]["sessionID"]
readonly id?: Endpoint0_5Request["payload"]["id"]
readonly prompt: Endpoint0_5Request["payload"]["prompt"]
readonly delivery?: Endpoint0_5Request["payload"]["delivery"]
readonly resume?: Endpoint0_5Request["payload"]["resume"]
type Endpoint0_6Request = Parameters<RawClient["server.session"]["session.prompt"]>[0]
type Endpoint0_6Input = {
readonly sessionID: Endpoint0_6Request["params"]["sessionID"]
readonly id?: Endpoint0_6Request["payload"]["id"]
readonly prompt: Endpoint0_6Request["payload"]["prompt"]
readonly delivery?: Endpoint0_6Request["payload"]["delivery"]
readonly resume?: Endpoint0_6Request["payload"]["resume"]
}
const Endpoint0_5 = (raw: RawClient["server.session"]) => (input: Endpoint0_5Input) =>
const Endpoint0_6 = (raw: RawClient["server.session"]) => (input: Endpoint0_6Input) =>
raw["session.prompt"]({
params: { sessionID: input.sessionID },
payload: { id: input.id, prompt: input.prompt, delivery: input.delivery, resume: input.resume },
@@ -100,23 +106,23 @@ const Endpoint0_5 = (raw: RawClient["server.session"]) => (input: Endpoint0_5Inp
Effect.map((value) => value.data),
)
type Endpoint0_6Request = Parameters<RawClient["server.session"]["session.compact"]>[0]
type Endpoint0_6Input = { readonly sessionID: Endpoint0_6Request["params"]["sessionID"] }
const Endpoint0_6 = (raw: RawClient["server.session"]) => (input: Endpoint0_6Input) =>
raw["session.compact"]({ params: { sessionID: input.sessionID } }).pipe(Effect.mapError(mapClientError))
type Endpoint0_7Request = Parameters<RawClient["server.session"]["session.wait"]>[0]
type Endpoint0_7Request = Parameters<RawClient["server.session"]["session.compact"]>[0]
type Endpoint0_7Input = { readonly sessionID: Endpoint0_7Request["params"]["sessionID"] }
const Endpoint0_7 = (raw: RawClient["server.session"]) => (input: Endpoint0_7Input) =>
raw["session.compact"]({ params: { sessionID: input.sessionID } }).pipe(Effect.mapError(mapClientError))
type Endpoint0_8Request = Parameters<RawClient["server.session"]["session.wait"]>[0]
type Endpoint0_8Input = { readonly sessionID: Endpoint0_8Request["params"]["sessionID"] }
const Endpoint0_8 = (raw: RawClient["server.session"]) => (input: Endpoint0_8Input) =>
raw["session.wait"]({ params: { sessionID: input.sessionID } }).pipe(Effect.mapError(mapClientError))
type Endpoint0_8Request = Parameters<RawClient["server.session"]["session.revert.stage"]>[0]
type Endpoint0_8Input = {
readonly sessionID: Endpoint0_8Request["params"]["sessionID"]
readonly messageID: Endpoint0_8Request["payload"]["messageID"]
readonly files?: Endpoint0_8Request["payload"]["files"]
type Endpoint0_9Request = Parameters<RawClient["server.session"]["session.revert.stage"]>[0]
type Endpoint0_9Input = {
readonly sessionID: Endpoint0_9Request["params"]["sessionID"]
readonly messageID: Endpoint0_9Request["payload"]["messageID"]
readonly files?: Endpoint0_9Request["payload"]["files"]
}
const Endpoint0_8 = (raw: RawClient["server.session"]) => (input: Endpoint0_8Input) =>
const Endpoint0_9 = (raw: RawClient["server.session"]) => (input: Endpoint0_9Input) =>
raw["session.revert.stage"]({
params: { sessionID: input.sessionID },
payload: { messageID: input.messageID, files: input.files },
@@ -125,30 +131,30 @@ const Endpoint0_8 = (raw: RawClient["server.session"]) => (input: Endpoint0_8Inp
Effect.map((value) => value.data),
)
type Endpoint0_9Request = Parameters<RawClient["server.session"]["session.revert.clear"]>[0]
type Endpoint0_9Input = { readonly sessionID: Endpoint0_9Request["params"]["sessionID"] }
const Endpoint0_9 = (raw: RawClient["server.session"]) => (input: Endpoint0_9Input) =>
raw["session.revert.clear"]({ params: { sessionID: input.sessionID } }).pipe(Effect.mapError(mapClientError))
type Endpoint0_10Request = Parameters<RawClient["server.session"]["session.revert.commit"]>[0]
type Endpoint0_10Request = Parameters<RawClient["server.session"]["session.revert.clear"]>[0]
type Endpoint0_10Input = { readonly sessionID: Endpoint0_10Request["params"]["sessionID"] }
const Endpoint0_10 = (raw: RawClient["server.session"]) => (input: Endpoint0_10Input) =>
raw["session.revert.commit"]({ params: { sessionID: input.sessionID } }).pipe(Effect.mapError(mapClientError))
raw["session.revert.clear"]({ params: { sessionID: input.sessionID } }).pipe(Effect.mapError(mapClientError))
type Endpoint0_11Request = Parameters<RawClient["server.session"]["session.context"]>[0]
type Endpoint0_11Request = Parameters<RawClient["server.session"]["session.revert.commit"]>[0]
type Endpoint0_11Input = { readonly sessionID: Endpoint0_11Request["params"]["sessionID"] }
const Endpoint0_11 = (raw: RawClient["server.session"]) => (input: Endpoint0_11Input) =>
raw["session.revert.commit"]({ params: { sessionID: input.sessionID } }).pipe(Effect.mapError(mapClientError))
type Endpoint0_12Request = Parameters<RawClient["server.session"]["session.context"]>[0]
type Endpoint0_12Input = { readonly sessionID: Endpoint0_12Request["params"]["sessionID"] }
const Endpoint0_12 = (raw: RawClient["server.session"]) => (input: Endpoint0_12Input) =>
raw["session.context"]({ params: { sessionID: input.sessionID } }).pipe(
Effect.mapError(mapClientError),
Effect.map((value) => value.data),
)
type Endpoint0_12Request = Parameters<RawClient["server.session"]["session.events"]>[0]
type Endpoint0_12Input = {
readonly sessionID: Endpoint0_12Request["params"]["sessionID"]
readonly after?: Endpoint0_12Request["query"]["after"]
type Endpoint0_13Request = Parameters<RawClient["server.session"]["session.events"]>[0]
type Endpoint0_13Input = {
readonly sessionID: Endpoint0_13Request["params"]["sessionID"]
readonly after?: Endpoint0_13Request["query"]["after"]
}
const Endpoint0_12 = (raw: RawClient["server.session"]) => (input: Endpoint0_12Input) =>
const Endpoint0_13 = (raw: RawClient["server.session"]) => (input: Endpoint0_13Input) =>
Stream.unwrap(
raw["session.events"]({ params: { sessionID: input.sessionID }, query: { after: input.after } }).pipe(
Effect.mapError(mapClientError),
@@ -156,17 +162,17 @@ const Endpoint0_12 = (raw: RawClient["server.session"]) => (input: Endpoint0_12I
),
)
type Endpoint0_13Request = Parameters<RawClient["server.session"]["session.interrupt"]>[0]
type Endpoint0_13Input = { readonly sessionID: Endpoint0_13Request["params"]["sessionID"] }
const Endpoint0_13 = (raw: RawClient["server.session"]) => (input: Endpoint0_13Input) =>
type Endpoint0_14Request = Parameters<RawClient["server.session"]["session.interrupt"]>[0]
type Endpoint0_14Input = { readonly sessionID: Endpoint0_14Request["params"]["sessionID"] }
const Endpoint0_14 = (raw: RawClient["server.session"]) => (input: Endpoint0_14Input) =>
raw["session.interrupt"]({ params: { sessionID: input.sessionID } }).pipe(Effect.mapError(mapClientError))
type Endpoint0_14Request = Parameters<RawClient["server.session"]["session.message"]>[0]
type Endpoint0_14Input = {
readonly sessionID: Endpoint0_14Request["params"]["sessionID"]
readonly messageID: Endpoint0_14Request["params"]["messageID"]
type Endpoint0_15Request = Parameters<RawClient["server.session"]["session.message"]>[0]
type Endpoint0_15Input = {
readonly sessionID: Endpoint0_15Request["params"]["sessionID"]
readonly messageID: Endpoint0_15Request["params"]["messageID"]
}
const Endpoint0_14 = (raw: RawClient["server.session"]) => (input: Endpoint0_14Input) =>
const Endpoint0_15 = (raw: RawClient["server.session"]) => (input: Endpoint0_15Input) =>
raw["session.message"]({ params: { sessionID: input.sessionID, messageID: input.messageID } }).pipe(
Effect.mapError(mapClientError),
Effect.map((value) => value.data),
@@ -175,19 +181,20 @@ const Endpoint0_14 = (raw: RawClient["server.session"]) => (input: Endpoint0_14I
const adaptGroup0 = (raw: RawClient["server.session"]) => ({
list: Endpoint0_0(raw),
create: Endpoint0_1(raw),
get: Endpoint0_2(raw),
switchAgent: Endpoint0_3(raw),
switchModel: Endpoint0_4(raw),
prompt: Endpoint0_5(raw),
compact: Endpoint0_6(raw),
wait: Endpoint0_7(raw),
stage: Endpoint0_8(raw),
clear: Endpoint0_9(raw),
commit: Endpoint0_10(raw),
context: Endpoint0_11(raw),
events: Endpoint0_12(raw),
interrupt: Endpoint0_13(raw),
message: Endpoint0_14(raw),
active: Endpoint0_2(raw),
get: Endpoint0_3(raw),
switchAgent: Endpoint0_4(raw),
switchModel: Endpoint0_5(raw),
prompt: Endpoint0_6(raw),
compact: Endpoint0_7(raw),
wait: Endpoint0_8(raw),
stage: Endpoint0_9(raw),
clear: Endpoint0_10(raw),
commit: Endpoint0_11(raw),
context: Endpoint0_12(raw),
events: Endpoint0_13(raw),
interrupt: Endpoint0_14(raw),
message: Endpoint0_15(raw),
})
const adaptClient = (raw: RawClient) => ({ sessions: adaptGroup0(raw["server.session"]) })
+12
View File
@@ -3,6 +3,7 @@ import type {
SessionsListOutput,
SessionsCreateInput,
SessionsCreateOutput,
SessionsActiveOutput,
SessionsGetInput,
SessionsGetOutput,
SessionsSwitchAgentInput,
@@ -198,6 +199,17 @@ export function make(options: ClientOptions) {
},
requestOptions,
).then((value) => value.data),
active: (requestOptions?: RequestOptions) =>
request<{ readonly data: SessionsActiveOutput }>(
{
method: "GET",
path: `/api/session/active`,
successStatus: 200,
declaredStatuses: [401, 400],
empty: false,
},
requestOptions,
).then((value) => value.data),
get: (input: SessionsGetInput, requestOptions?: RequestOptions) =>
request<{ readonly data: SessionsGetOutput }>(
{
+62 -3
View File
@@ -243,6 +243,8 @@ export type SessionsCreateOutput = {
}
}["data"]
export type SessionsActiveOutput = { readonly data: { readonly [x: string]: { readonly type: "running" } } }["data"]
export type SessionsGetInput = { readonly sessionID: { readonly sessionID: string }["sessionID"] }
export type SessionsGetOutput = {
@@ -579,7 +581,26 @@ export type SessionsContextOutput = {
readonly reasoning: number
readonly cache: { readonly read: number; readonly write: number }
}
readonly error?: { readonly type: "unknown"; readonly message: string }
readonly error?:
| { readonly type: "unknown"; readonly message: string }
| {
readonly type: "provider"
readonly category:
| "invalid-request"
| "no-route"
| "authentication"
| "rate-limit"
| "quota-exceeded"
| "content-policy"
| "provider-internal"
| "transport"
| "invalid-provider-output"
| "unknown"
readonly message: string
readonly status?: number | null
readonly retryable: boolean
readonly retryAfterMs?: number | null
}
}
| {
readonly type: "compaction"
@@ -792,7 +813,26 @@ export type SessionsEventsOutput =
readonly timestamp: number
readonly sessionID: string
readonly assistantMessageID: string
readonly error: { readonly type: "unknown"; readonly message: string }
readonly error:
| { readonly type: "unknown"; readonly message: string }
| {
readonly type: "provider"
readonly category:
| "invalid-request"
| "no-route"
| "authentication"
| "rate-limit"
| "quota-exceeded"
| "content-policy"
| "provider-internal"
| "transport"
| "invalid-provider-output"
| "unknown"
readonly message: string
readonly status?: number | undefined
readonly retryable: boolean
readonly retryAfterMs?: number | undefined
}
}
}
| {
@@ -1196,7 +1236,26 @@ export type SessionsMessageOutput = {
readonly reasoning: number
readonly cache: { readonly read: number; readonly write: number }
}
readonly error?: { readonly type: "unknown"; readonly message: string }
readonly error?:
| { readonly type: "unknown"; readonly message: string }
| {
readonly type: "provider"
readonly category:
| "invalid-request"
| "no-route"
| "authentication"
| "rate-limit"
| "quota-exceeded"
| "content-policy"
| "provider-internal"
| "transport"
| "invalid-provider-output"
| "unknown"
readonly message: string
readonly status?: number | null
readonly retryable: boolean
readonly retryAfterMs?: number | null
}
}
| {
readonly type: "compaction"
+8 -1
View File
@@ -37,6 +37,11 @@ test("session methods retain decoded Effect inputs and outputs", async () => {
if (url.includes("/message/")) {
return Effect.succeed(HttpClientResponse.fromWeb(request, Response.json({ data: modelSwitchedMessage })))
}
if (url.endsWith("/api/session/active")) {
return Effect.succeed(
HttpClientResponse.fromWeb(request, Response.json({ data: { ses_test: { type: "running" } } })),
)
}
if (request.method === "POST" && url.endsWith("/api/session")) {
return Effect.succeed(HttpClientResponse.fromWeb(request, Response.json(session)))
}
@@ -50,6 +55,7 @@ test("session methods retain decoded Effect inputs and outputs", async () => {
const result = await Effect.gen(function* () {
const client = yield* OpenCode.make({ baseUrl: "http://localhost:3000" })
const page = yield* client.sessions.list({ limit: 10 })
const active = yield* client.sessions.active()
const created = yield* client.sessions.create({
location: Location.Ref.make({ directory: AbsolutePath.make("/tmp/project") }),
})
@@ -74,10 +80,11 @@ test("session methods retain decoded Effect inputs and outputs", async () => {
sessionID: Session.ID.make("ses_test"),
messageID: SessionMessage.ID.make("msg_model"),
})
return { page, created, admitted, context, events, message }
return { page, active, created, admitted, context, events, message }
}).pipe(Effect.provideService(HttpClient.HttpClient, httpClient), Effect.runPromise)
expect(DateTime.toEpochMillis(result.page.data[0].time.created)).toBe(1_717_171_717_000)
expect(result.active).toEqual({ ses_test: { type: "running" } })
expect(Object.getPrototypeOf(result.page.data[0])).toBe(Object.prototype)
expect(Object.getPrototypeOf(result.created)).toBe(Object.prototype)
expect(result.created.id).toBe("ses_test")
+5 -1
View File
@@ -32,6 +32,7 @@ test("session methods use the public HTTP contract", async () => {
if (url.includes("/prompt")) return Response.json(admission)
if (url.includes("/context")) return Response.json({ data: [] })
if (url.includes("/message/")) return Response.json({ data: modelSwitchedMessage })
if (url.endsWith("/api/session/active")) return Response.json({ data: { ses_test: { type: "running" } } })
if (init?.method === "POST" && url.endsWith("/api/session")) return Response.json(session)
if (init?.method === "POST") return new Response(null, { status: 204 })
return Response.json({ data: [session.data], cursor: { next: "next" } })
@@ -39,6 +40,7 @@ test("session methods use the public HTTP contract", async () => {
})
const page = await client.sessions.list({ limit: "10", order: "desc" })
const active = await client.sessions.active()
const created = await client.sessions.create({ location: { directory: "/tmp/project" } })
await client.sessions.switchAgent({ sessionID: "ses_test", agent: "build" })
await client.sessions.switchModel({
@@ -59,6 +61,7 @@ test("session methods use the public HTTP contract", async () => {
const message = await client.sessions.message({ sessionID: "ses_test", messageID: "msg_model" })
expect(page.cursor.next).toBe("next")
expect(active).toEqual({ ses_test: { type: "running" } })
expect(created.id).toBe("ses_test")
expect(admitted.id).toBe("msg_test")
expect(context).toEqual([])
@@ -66,6 +69,7 @@ test("session methods use the public HTTP contract", async () => {
expect(message).toEqual(modelSwitchedMessage)
expect(requests.map((request) => [request.init?.method, request.url])).toEqual([
["GET", "http://localhost:3000/api/session?limit=10&order=desc"],
["GET", "http://localhost:3000/api/session/active"],
["POST", "http://localhost:3000/api/session"],
["POST", "http://localhost:3000/api/session/ses_test/agent"],
["POST", "http://localhost:3000/api/session/ses_test/model"],
@@ -77,7 +81,7 @@ test("session methods use the public HTTP contract", async () => {
["POST", "http://localhost:3000/api/session/ses_test/interrupt"],
["GET", "http://localhost:3000/api/session/ses_test/message/msg_model"],
])
const body = requests[4]?.init?.body
const body = requests.find((request) => request.url.endsWith("/api/session/ses_test/prompt"))?.init?.body
if (typeof body !== "string") throw new Error("Expected JSON request body")
expect(JSON.parse(body)).toEqual({
prompt: { text: "Hello" },
+2
View File
@@ -155,6 +155,7 @@ export interface Interface {
}) => Effect.Effect<void, OperationUnavailableError>
readonly compact: (input: CompactInput) => Effect.Effect<void, NotFoundError | OperationUnavailableError>
readonly wait: (id: SessionSchema.ID) => Effect.Effect<void, NotFoundError | OperationUnavailableError>
readonly active: Effect.Effect<ReadonlySet<SessionSchema.ID>>
readonly resume: (sessionID: SessionSchema.ID) => Effect.Effect<void, NotFoundError | SessionRunner.RunError>
readonly interrupt: (sessionID: SessionSchema.ID) => Effect.Effect<void>
readonly revert: {
@@ -402,6 +403,7 @@ export const layer = Layer.unwrap(
yield* result.get(sessionID)
return yield* new OperationUnavailableError({ operation: "wait" })
}),
active: execution.active,
resume: Effect.fn("V2Session.resume")(function* (sessionID) {
yield* result.get(sessionID)
yield* execution.resume(sessionID)
+8 -1
View File
@@ -5,6 +5,8 @@ import { SessionRunner } from "./runner/index"
import { SessionSchema } from "./schema"
export interface Interface {
/** Snapshots active execution owned by this process. */
readonly active: Effect.Effect<ReadonlySet<SessionSchema.ID>>
/** Starts execution while idle or joins the active execution. */
readonly resume: (sessionID: SessionSchema.ID) => Effect.Effect<void, SessionRunner.RunError>
/** Registers newly recorded work. Repeated wakeups may coalesce. */
@@ -19,5 +21,10 @@ export class Service extends Context.Service<Service, Interface>()("@opencode/v2
/** Low-level compatibility layer for callers that only need durable Session recording. */
export const noopLayer = Layer.succeed(
Service,
Service.of({ resume: () => Effect.void, wake: () => Effect.void, interrupt: () => Effect.void }),
Service.of({
active: Effect.succeed(new Set()),
resume: () => Effect.void,
wake: () => Effect.void,
interrupt: () => Effect.void,
}),
)
@@ -28,6 +28,7 @@ export const layer = Layer.effect(
})
return SessionExecution.Service.of({
active: coordinator.active,
interrupt: coordinator.interrupt,
resume: coordinator.run,
wake: coordinator.wake,
+3 -1
View File
@@ -4,6 +4,8 @@ import { Deferred, Effect, Exit, Fiber, FiberSet, Scope } from "effect"
/** Serializes execution for each key while allowing different keys to run concurrently. */
export interface Coordinator<Key, E> {
/** Snapshots keys with an execution owned by this coordinator. */
readonly active: Effect.Effect<ReadonlySet<Key>>
/** Starts execution while idle or joins the active execution. */
readonly run: (key: Key) => Effect.Effect<void, E>
/** Registers one coalesced follow-up after newly recorded work. */
@@ -98,5 +100,5 @@ export const make = <Key, E>(options: {
return Fiber.interrupt(entry.owner)
})
return { run, wake, interrupt }
return { active: Effect.sync(() => new Set(active.keys())), run, wake, interrupt }
})
+69 -66
View File
@@ -6,6 +6,7 @@ import {
Message,
SystemPart,
isContextOverflowFailure,
type LLMErrorReason,
type ProviderErrorEvent,
} from "@opencode-ai/llm"
import { Cause, DateTime, Effect, FiberSet, Layer, Option, Semaphore, Stream } from "effect"
@@ -32,14 +33,11 @@ import { SessionSchema } from "../schema"
import { SessionStore } from "../store"
import { type RunError, Service } from "./index"
import { SessionRunnerModel } from "./model"
import { createLLMEventPublisher } from "./publish-llm-event"
import { createLLMEventPublisher, providerError } from "./publish-llm-event"
import { toLLMMessages } from "./to-llm-message"
import { MAX_STEPS_PROMPT } from "./max-steps"
import { Snapshot } from "../../snapshot"
const MAX_PROVIDER_RETRIES = 2
const PROVIDER_RETRY_DELAY_MS = 500
/**
* Runs one durable coding-agent Session until it settles.
*
@@ -90,6 +88,29 @@ const PROVIDER_RETRY_DELAY_MS = 500
* explicit loop starts the next provider turn after local settlement. Configured agent step limits bound the loop.
*/
const providerCategories = {
InvalidRequest: "invalid-request",
NoRoute: "no-route",
Authentication: "authentication",
RateLimit: "rate-limit",
QuotaExceeded: "quota-exceeded",
ContentPolicy: "content-policy",
ProviderInternal: "provider-internal",
Transport: "transport",
InvalidProviderOutput: "invalid-provider-output",
UnknownProvider: "unknown",
} as const satisfies Record<LLMErrorReason["_tag"], Parameters<typeof providerError>[0]["category"]>
const toSessionProviderError = (reason: LLMErrorReason) =>
providerError({
category: providerCategories[reason._tag],
status:
("http" in reason ? reason.http?.response?.status : undefined) ??
("status" in reason ? reason.status : undefined),
retryable: reason.retryable,
retryAfterMs: "retryAfterMs" in reason ? reason.retryAfterMs : undefined,
})
export const layer = Layer.effect(
Service,
Effect.gen(function* () {
@@ -225,72 +246,54 @@ export const layer = Layer.effect(
const publish = (event: LLMEvent, outputPaths: ReadonlyArray<string> = []) =>
withPublication(publisher.publish(event, outputPaths))
let overflowFailure: ProviderErrorEvent | undefined
const runProvider: (remaining: number, attempt: number) => Effect.Effect<void, LLMError> = Effect.fnUntraced(
function* (remaining: number, attempt: number) {
let retryFailure: ProviderErrorEvent | undefined
yield* llm.stream(request).pipe(
Stream.runForEach((event) =>
Effect.gen(function* () {
if (overflowFailure || retryFailure || publisher.hasProviderError()) return
if (LLMEvent.is.providerError(event) && !publisher.hasAssistantStarted()) {
if (isContextOverflowFailure(event)) {
overflowFailure = event
return
}
if (event.retryable) {
retryFailure = event
return
}
}
yield* publish(event)
if (event.type !== "tool-call" || event.providerExecuted) return
if (!toolMaterialization) {
yield* withPublication(
publisher.failUnsettledTools("Tools are disabled after the maximum agent steps"),
)
return
}
needsContinuation = true
const assistantMessageID = yield* publisher.assistantMessageID(event.id)
yield* Effect.uninterruptibleMask((restore) =>
restore(
toolMaterialization.settle({
sessionID: session.id,
agent: agent.id,
assistantMessageID,
call: event,
const providerStream = llm.stream(request).pipe(
Stream.runForEach((event) =>
Effect.gen(function* () {
if (overflowFailure || publisher.hasProviderError()) return
if (LLMEvent.is.providerError(event)) {
if (isContextOverflowFailure(event) && !publisher.hasAssistantStarted()) {
overflowFailure = event
return
}
}
yield* publish(event)
if (event.type !== "tool-call" || event.providerExecuted) return
if (!toolMaterialization) {
yield* withPublication(publisher.failUnsettledTools("Tools are disabled after the maximum agent steps"))
return
}
needsContinuation = true
const assistantMessageID = yield* publisher.assistantMessageID(event.id)
yield* Effect.uninterruptibleMask((restore) =>
restore(
toolMaterialization.settle({
sessionID: session.id,
agent: agent.id,
assistantMessageID,
call: event,
}),
).pipe(
Effect.flatMap((settlement) =>
publish(
LLMEvent.toolResult({
id: event.id,
name: event.name,
result: settlement.result,
output: settlement.output,
}),
).pipe(
Effect.flatMap((settlement) =>
publish(
LLMEvent.toolResult({
id: event.id,
name: event.name,
result: settlement.result,
output: settlement.output,
}),
settlement.outputPaths ?? [],
),
),
settlement.outputPaths ?? [],
),
).pipe(FiberSet.run(toolFibers))
}),
),
Effect.ensuring(withPublication(publisher.flush())),
)
if (!retryFailure) return
if (remaining === 0) {
yield* publish(retryFailure)
return
}
yield* Effect.sleep(`${PROVIDER_RETRY_DELAY_MS * 2 ** attempt} millis`)
yield* runProvider(remaining - 1, attempt + 1)
},
),
),
).pipe(FiberSet.run(toolFibers))
}),
),
Effect.ensuring(withPublication(publisher.flush())),
)
return yield* Effect.uninterruptibleMask((restore) =>
Effect.gen(function* () {
const stream = yield* restore(runProvider(MAX_PROVIDER_RETRIES, 0)).pipe(Effect.exit)
const stream = yield* restore(providerStream).pipe(Effect.exit)
const failure =
stream._tag === "Failure" ? Option.getOrUndefined(Cause.findErrorOption(stream.cause)) : undefined
if (
@@ -304,7 +307,7 @@ export const layer = Layer.effect(
const llmFailure = failure instanceof LLMError ? failure : undefined
if (llmFailure && !publisher.hasProviderError()) {
yield* withPublication(publisher.failUnsettledTools("Provider did not return a tool result", true))
yield* withPublication(publisher.failAssistant(llmFailure.reason.message))
yield* withPublication(publisher.failAssistant(toSessionProviderError(llmFailure.reason)))
}
if (stream._tag === "Failure" && Cause.hasInterrupts(stream.cause)) yield* FiberSet.clear(toolFibers)
const settled = yield* restore(awaitToolFibers(toolFibers)).pipe(Effect.exit)
@@ -320,7 +323,7 @@ export const layer = Layer.effect(
yield* FiberSet.clear(toolFibers)
yield* withPublication(publisher.failUnsettledTools("Tool execution interrupted"))
if (publisher.hasActiveAssistant())
yield* withPublication(publisher.failAssistant("Provider turn interrupted"))
yield* withPublication(publisher.failAssistant({ type: "unknown", message: "Provider turn interrupted" }))
}
if (settled._tag === "Failure" && !Cause.hasInterrupts(settled.cause)) {
const failure = Cause.squash(settled.cause)
@@ -50,6 +50,33 @@ const settledOutput = (value: ToolOutput | undefined, result: ToolResultValue):
return { structured: record(settled.structured), content: settled.content }
}
const providerMessages: Record<SessionMessage.ProviderErrorCategory, string> = {
"invalid-request": "Provider rejected the request",
"no-route": "Provider unavailable: no provider route",
authentication: "Provider authentication failed",
"rate-limit": "Provider rate limit exceeded",
"quota-exceeded": "Provider quota exceeded",
"content-policy": "Provider rejected the invalid request due to content policy",
"provider-internal": "Provider service unavailable",
transport: "Provider connection failed",
"invalid-provider-output": "Provider returned an invalid response",
unknown: "Provider request failed",
}
export const providerError = (input: {
readonly category: SessionMessage.ProviderErrorCategory
readonly status?: number
readonly retryable?: boolean
readonly retryAfterMs?: number
}): SessionMessage.ProviderError => ({
type: "provider",
category: input.category,
message: providerMessages[input.category],
status: input.status,
retryable: input.retryable ?? (input.category === "rate-limit" || input.category === "provider-internal"),
retryAfterMs: input.retryAfterMs,
})
/** Persist one provider turn without executing tools or starting a continuation turn. */
export const createLLMEventPublisher = (events: EventV2.Interface, input: Input) => {
const tools = new Map<
@@ -67,7 +94,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
const timestamp = DateTime.now
let assistantMessageID: SessionMessage.ID | undefined
let assistantActive = false
let assistantFailed = false
let providerFailed = false
let stepSettlement: { readonly finish: string; readonly tokens: ReturnType<typeof tokens> } | undefined
@@ -196,17 +222,17 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
yield* flushFragments()
})
const failAssistant = Effect.fnUntraced(function* (message: string) {
if (assistantFailed) return
const failAssistant = Effect.fnUntraced(function* (error: SessionMessage.Error) {
if (providerFailed) return
providerFailed = true
yield* flush()
const assistantMessageID = yield* startAssistant()
assistantActive = false
assistantFailed = true
yield* events.publish(SessionEvent.Step.Failed, {
sessionID: input.sessionID,
timestamp: yield* timestamp,
assistantMessageID,
error: { type: "unknown", message },
error,
})
})
@@ -402,8 +428,13 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
case "finish":
return
case "provider-error":
providerFailed = true
yield* failAssistant(event.message)
yield* failAssistant(
providerError({
category: event.category ?? (event.classification === "context-overflow" ? "invalid-request" : "unknown"),
status: event.status,
retryable: event.retryable,
}),
)
return
}
})
@@ -22,9 +22,11 @@ import { testEffect } from "./lib/effect"
const executionCalls: SessionV2.ID[] = []
const interruptCalls: SessionV2.ID[] = []
const wakeCalls: SessionV2.ID[] = []
const activeSessions = new Set<SessionV2.ID>()
const execution = Layer.succeed(
SessionExecution.Service,
SessionExecution.Service.of({
active: Effect.sync(() => new Set(activeSessions)),
resume: (sessionID) =>
Effect.sync(() => {
executionCalls.push(sessionID)
@@ -108,6 +110,13 @@ const eventCount = (type: string) =>
)
describe("SessionV2.prompt", () => {
it.effect("exposes the execution registry", () =>
Effect.gen(function* () {
activeSessions.add(sessionID)
expect(Array.from(yield* (yield* SessionV2.Service).active)).toEqual([sessionID])
}).pipe(Effect.ensuring(Effect.sync(() => activeSessions.clear()))),
)
it.effect("delegates execution continuation through SessionExecution", () =>
Effect.gen(function* () {
yield* setup
@@ -66,6 +66,78 @@ describe("SessionRunCoordinator", () => {
),
)
it.effect("snapshots only active executions", () =>
Effect.scoped(
Effect.gen(function* () {
const firstStarted = yield* Deferred.make<void>()
const secondStarted = yield* Deferred.make<void>()
const firstGate = yield* Deferred.make<void>()
const secondGate = yield* Deferred.make<void>()
const coordinator = yield* SessionRunCoordinator.make({
drain: (key: string) =>
Deferred.succeed(key === "first" ? firstStarted : secondStarted, undefined).pipe(
Effect.andThen(Deferred.await(key === "first" ? firstGate : secondGate)),
),
})
expect(Array.from(yield* coordinator.active)).toEqual([])
const first = yield* coordinator.run("first").pipe(Effect.forkChild)
yield* Deferred.await(firstStarted)
expect(Array.from(yield* coordinator.active)).toEqual(["first"])
const second = yield* coordinator.run("second").pipe(Effect.forkChild)
yield* Deferred.await(secondStarted)
expect(Array.from(yield* coordinator.active)).toEqual(["first", "second"])
yield* Deferred.succeed(firstGate, undefined)
yield* Fiber.join(first)
expect(Array.from(yield* coordinator.active)).toEqual(["second"])
yield* Deferred.succeed(secondGate, undefined)
yield* Fiber.join(second)
expect(Array.from(yield* coordinator.active)).toEqual([])
}),
),
)
it.effect("cleans active executions after failure and defect", () =>
Effect.scoped(
Effect.gen(function* () {
const failure = new Error("failed")
const defect = new Error("defect")
const coordinator = yield* SessionRunCoordinator.make({
drain: (key: string) => (key === "failure" ? Effect.fail(failure) : Effect.die(defect)),
})
const failed = yield* coordinator.run("failure").pipe(Effect.exit)
expect(Exit.isFailure(failed) && Cause.hasFails(failed.cause)).toBeTrue()
expect(Array.from(yield* coordinator.active)).toEqual([])
const died = yield* coordinator.run("defect").pipe(Effect.exit)
expect(Exit.isFailure(died) && Cause.hasDies(died.cause)).toBeTrue()
expect(Array.from(yield* coordinator.active)).toEqual([])
}),
),
)
it.effect("cleans active executions when its scope closes", () =>
Effect.gen(function* () {
const started = yield* Deferred.make<void>()
const coordinator = yield* Effect.scoped(
Effect.gen(function* () {
const coordinator = yield* SessionRunCoordinator.make({
drain: () => Deferred.succeed(started, undefined).pipe(Effect.andThen(Effect.never)),
})
yield* coordinator.wake("session")
yield* Deferred.await(started)
expect(Array.from(yield* coordinator.active)).toEqual(["session"])
return coordinator
}),
)
expect(Array.from(yield* coordinator.active)).toEqual([])
}),
)
it.effect("coalesces wakes received during active execution", () =>
Effect.scoped(
Effect.gen(function* () {
@@ -166,6 +238,7 @@ describe("SessionRunCoordinator", () => {
const exit = yield* Fiber.await(resumed)
expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBeTrue()
expect(Array.from(yield* coordinator.active)).toEqual([])
expect(runs).toBe(1)
}),
),
@@ -95,6 +95,7 @@ const execution = Layer.effect(
drain: (sessionID, force) => sessionRunner.run({ sessionID, force }),
})
return SessionExecution.Service.of({
active: coordinator.active,
resume: coordinator.run,
wake: coordinator.wake,
interrupt: coordinator.interrupt,
+209 -67
View File
@@ -4,6 +4,11 @@ import {
LLMError,
LLMEvent,
Model,
HttpContext,
HttpRateLimitDetails,
HttpRequestDetails,
HttpResponseDetails,
RateLimitReason,
TransportReason,
InvalidRequestReason,
type LLMClientShape,
@@ -54,7 +59,6 @@ import { ModelV2 } from "@opencode-ai/core/model"
import { Location } from "@opencode-ai/core/location"
import { ProviderV2 } from "@opencode-ai/core/provider"
import { Cause, DateTime, Deferred, Effect, Exit, Fiber, Layer, Schema, Stream } from "effect"
import * as TestClock from "effect/testing/TestClock"
import { asc, eq } from "drizzle-orm"
import { testEffect } from "./lib/effect"
@@ -255,6 +259,7 @@ const execution = Layer.effect(
drain: (sessionID, force) => sessionRunner.run({ sessionID, force }),
})
return SessionExecution.Service.of({
active: coordinator.active,
resume: coordinator.run,
wake: coordinator.wake,
interrupt: coordinator.interrupt,
@@ -515,13 +520,21 @@ const verifyPartialFlushOnFailure = (kind: FragmentKind) =>
responseStream = Stream.concat(Stream.fromIterable(fixture.partialEvents), Stream.fail(failure))
expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure)
const expectedContent =
kind === "tool input"
? {
type: "tool",
id: fragmentID(kind, "partial"),
state: { status: "error", error: { message: "Tool execution interrupted" } },
}
: fixture.expectedContent
expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user", text: prompt },
{
type: "assistant",
finish: "error",
error: { type: "unknown", message: "Provider unavailable" },
content: [fixture.expectedContent],
error: { type: "provider", category: "transport", message: "Provider connection failed", retryable: false },
content: [expectedContent],
},
])
})
@@ -1196,7 +1209,16 @@ describe("SessionRunnerLLM", () => {
expect(requests).toHaveLength(3)
expect(yield* session.context(sessionID)).toMatchObject([
{ type: "compaction" },
{ type: "assistant", finish: "error", error: { message: "prompt too long" } },
{
type: "assistant",
finish: "error",
error: {
type: "provider",
category: "invalid-request",
message: "Provider rejected the request",
retryable: false,
},
},
])
}),
)
@@ -1244,7 +1266,16 @@ describe("SessionRunnerLLM", () => {
expect(context.some((message) => message.type === "compaction")).toBe(false)
expect(context.slice(-2)).toMatchObject([
{ type: "user", text: "Continue" },
{ type: "assistant", finish: "error", error: { message: "prompt too long" } },
{
type: "assistant",
finish: "error",
error: {
type: "provider",
category: "invalid-request",
message: "Provider rejected the request",
retryable: false,
},
},
])
}),
)
@@ -2920,7 +2951,16 @@ describe("SessionRunnerLLM", () => {
expect(requests).toHaveLength(1)
expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user", text: "Fail durably" },
{ type: "assistant", finish: "error", error: { type: "unknown", message: "Provider unavailable" } },
{
type: "assistant",
finish: "error",
error: {
type: "provider",
category: "unknown",
message: "Provider request failed",
retryable: false,
},
},
])
}),
)
@@ -2939,63 +2979,16 @@ describe("SessionRunnerLLM", () => {
expect(requests).toHaveLength(1)
expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user", text: "Fail before step" },
{ type: "assistant", finish: "error", error: { type: "unknown", message: "Provider unavailable" } },
])
}),
)
it.effect("retries transient provider errors before durable assistant output", () =>
Effect.gen(function* () {
yield* setup
const session = yield* SessionV2.Service
yield* session.prompt({ sessionID, prompt: Prompt.make({ text: "Recover transient stream" }), resume: false })
requests.length = 0
responses = [
[LLMEvent.providerError({ message: "stream_read_error", retryable: true })],
[
LLMEvent.stepStart({ index: 0 }),
LLMEvent.textStart({ id: "text-after-retry" }),
LLMEvent.textDelta({ id: "text-after-retry", text: "Recovered" }),
LLMEvent.textEnd({ id: "text-after-retry" }),
LLMEvent.stepFinish({ index: 0, reason: "stop" }),
LLMEvent.finish({ reason: "stop" }),
],
]
const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
yield* Effect.yieldNow
yield* TestClock.adjust("500 millis")
yield* Fiber.join(run)
expect(requests).toHaveLength(2)
expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user", text: "Recover transient stream" },
{ type: "assistant", finish: "stop", content: [{ type: "text", text: "Recovered" }] },
])
}),
)
it.effect("bounds transient provider retries", () =>
Effect.gen(function* () {
yield* setup
const session = yield* SessionV2.Service
yield* session.prompt({ sessionID, prompt: Prompt.make({ text: "Bound transient retries" }), resume: false })
requests.length = 0
responses = Array.from({ length: 3 }, () => [
LLMEvent.providerError({ message: "server_error", retryable: true }),
])
const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
yield* Effect.yieldNow
yield* TestClock.adjust("1500 millis")
yield* Fiber.join(run)
expect(requests).toHaveLength(3)
expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user", text: "Bound transient retries" },
{ type: "assistant", finish: "error", error: { type: "unknown", message: "server_error" } },
{
type: "assistant",
finish: "error",
error: {
type: "provider",
category: "unknown",
message: "Provider request failed",
retryable: false,
},
},
])
}),
)
@@ -3022,7 +3015,12 @@ describe("SessionRunnerLLM", () => {
{
type: "assistant",
finish: "error",
error: { message: "prompt too long" },
error: {
type: "provider",
category: "invalid-request",
message: "Provider rejected the request",
retryable: false,
},
content: [{ type: "text", text: "Partial" }],
},
])
@@ -3041,11 +3039,138 @@ describe("SessionRunnerLLM", () => {
yield* replaySessionProjection(sessionID)
expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user", text: "Fail raw stream durably" },
{ type: "assistant", finish: "error", error: { type: "unknown", message: "Provider unavailable" } },
{
type: "assistant",
finish: "error",
error: { type: "provider", category: "transport", message: "Provider connection failed", retryable: false },
},
])
}),
)
it.effect("does not publish step ended when the provider throws after step finish", () =>
Effect.gen(function* () {
yield* setup
const session = yield* SessionV2.Service
const db = yield* Database.Service
yield* session.prompt({ sessionID, prompt: Prompt.make({ text: "Fail after finish" }), resume: false })
const failure = providerUnavailable()
responseStream = Stream.concat(
Stream.fromIterable([LLMEvent.stepFinish({ index: 0, reason: "stop" })]),
Stream.fail(failure),
)
expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure)
const types = yield* db.db.select({ type: EventTable.type }).from(EventTable).all().pipe(Effect.orDie)
expect(types).toContainEqual({
type: EventV2.versionedType(SessionEvent.Step.Failed.type, 2),
})
expect(types).not.toContainEqual({
type: EventV2.versionedType(SessionEvent.Step.Ended.type, 2),
})
expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user", text: "Fail after finish" },
{
type: "assistant",
finish: "error",
error: { type: "provider", category: "transport", message: "Provider connection failed" },
},
])
}),
)
it.effect("preserves sanitized HTTP rate-limit details in the durable event and projection", () =>
Effect.gen(function* () {
yield* setup
const session = yield* SessionV2.Service
const db = yield* Database.Service
yield* session.prompt({ sessionID, prompt: Prompt.make({ text: "Fail with rate limit" }), resume: false })
const failure = new LLMError({
module: "RequestExecutor",
method: "execute",
reason: new RateLimitReason({
message: "secret provider response message",
retryAfterMs: 12_000,
rateLimit: new HttpRateLimitDetails({ retryAfterMs: 12_000, limit: { requests: "secret-limit" } }),
http: new HttpContext({
request: new HttpRequestDetails({
method: "POST",
url: "https://secret.example/v1/responses?api_key=credential",
headers: { authorization: "Bearer credential" },
}),
response: new HttpResponseDetails({ status: 429, headers: { "x-secret": "secret-header" } }),
body: '{"secret":"provider body"}',
requestId: "secret-request-id",
}),
}),
})
responseStream = Stream.fail(failure)
expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure)
const event = yield* db.db
.select({ data: EventTable.data })
.from(EventTable)
.where(eq(EventTable.type, EventV2.versionedType(SessionEvent.Step.Failed.type, 2)))
.get()
.pipe(Effect.orDie)
const expected = {
type: "provider",
category: "rate-limit",
message: "Provider rate limit exceeded",
status: 429,
retryable: true,
retryAfterMs: 12_000,
}
expect(event?.data.error).toEqual(expected)
expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user", text: "Fail with rate limit" },
{ type: "assistant", finish: "error", error: expected },
])
expect(JSON.stringify(event)).not.toMatch(
/secret provider|secret\.example|credential|secret-header|provider body|secret-request-id|secret-limit/,
)
expect(JSON.stringify(yield* session.context(sessionID))).not.toMatch(
/secret provider|secret\.example|credential|secret-header|provider body|secret-request-id|secret-limit/,
)
}),
)
it.effect("projects categorized in-band failures without fabricating an HTTP status", () =>
Effect.gen(function* () {
yield* setup
const session = yield* SessionV2.Service
yield* session.prompt({ sessionID, prompt: Prompt.make({ text: "Fail in band" }), resume: false })
response = [
LLMEvent.providerError({
message: "rate_limit_exceeded: secret provider message",
category: "rate-limit",
retryable: false,
}),
]
yield* session.resume(sessionID)
expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user", text: "Fail in band" },
{
type: "assistant",
finish: "error",
error: {
type: "provider",
category: "rate-limit",
message: "Provider rate limit exceeded",
retryable: false,
},
},
])
const assistant = (yield* session.context(sessionID))[1]
expect(assistant?.type === "assistant" ? assistant.error : undefined).not.toHaveProperty("status")
expect(JSON.stringify(assistant)).not.toContain("secret provider message")
}),
)
it.effect("does not continue automatically after a provider error follows a local tool call", () =>
Effect.gen(function* () {
yield* setup
@@ -3061,13 +3186,25 @@ describe("SessionRunnerLLM", () => {
response = [
LLMEvent.stepStart({ index: 0 }),
LLMEvent.toolCall({ id: "call-before-provider-error", name: "echo", input: { text: "settled" } }),
LLMEvent.providerError({ message: "Provider unavailable" }),
LLMEvent.providerError({ message: "secret provider failure" }),
]
yield* session.resume(sessionID)
expect(requests).toHaveLength(1)
expect(executions.slice(executionCount)).toEqual(["settled"])
const assistant = (yield* session.context(sessionID))[1]
expect(assistant).toMatchObject({
type: "assistant",
finish: "error",
error: {
type: "provider",
category: "unknown",
message: "Provider request failed",
retryable: false,
},
})
expect(JSON.stringify(assistant)).not.toContain("secret provider failure")
}),
)
@@ -3157,7 +3294,12 @@ describe("SessionRunnerLLM", () => {
{
type: "assistant",
finish: "error",
error: { type: "unknown", message: "Provider unavailable" },
error: {
type: "provider",
category: "transport",
message: "Provider connection failed",
retryable: false,
},
content: [{ type: "tool", id: "call-hosted-raw-failure", state: { status: "error" } }],
},
])
+16
View File
@@ -144,6 +144,8 @@ test("Core reuses the canonical shared schemas", async () => {
[coreSessionInput.Admitted, SessionInput.Admitted],
[coreSessionMessage.ID, SessionMessage.ID],
[coreSessionMessage.UnknownError, SessionMessage.UnknownError],
[coreSessionMessage.ProviderError, SessionMessage.ProviderError],
[coreSessionMessage.Error, SessionMessage.Error],
[coreSessionMessage.AgentSwitched, SessionMessage.AgentSwitched],
[coreSessionMessage.ModelSwitched, SessionMessage.ModelSwitched],
[coreSessionMessage.User, SessionMessage.User],
@@ -204,3 +206,17 @@ test("shared record schemas construct and decode plain objects", () => {
expect(Prompt.fromUserMessage({ text: "hello" })).toEqual(made)
expect(Workspace.ID.ascending("")).toStartWith("wrk_")
})
test("assistant errors retain legacy unknown decode compatibility", () => {
const assistant = Schema.decodeUnknownSync(SessionMessage.Assistant)({
id: "msg_legacy_error",
type: "assistant",
agent: "build",
model: { id: "model", providerID: "provider" },
content: [],
error: { type: "unknown", message: "Legacy failure" },
time: { created: 0 },
})
expect(assistant.error).toEqual({ type: "unknown", message: "Legacy failure" })
})
+6 -23
View File
@@ -201,7 +201,6 @@ type OpenAIResponsesStreamItem = Schema.Schema.Type<typeof OpenAIResponsesStream
// `response.failed` carries them under `response.error`. We capture both so
// the parser can surface a useful provider-error message in either path.
const OpenAIResponsesErrorPayload = Schema.Struct({
type: optionalNull(Schema.String),
code: optionalNull(Schema.String),
message: optionalNull(Schema.String),
param: optionalNull(Schema.String),
@@ -225,7 +224,6 @@ const OpenAIResponsesEvent = Schema.Struct({
[Schema.Record(Schema.String, Schema.Unknown)],
),
),
error: optionalNull(OpenAIResponsesErrorPayload),
code: Schema.optional(Schema.String),
message: Schema.optional(Schema.String),
param: Schema.optional(Schema.String),
@@ -595,10 +593,10 @@ type StepResult = readonly [ParserState, ReadonlyArray<LLMEvent>]
const NO_EVENTS: StepResult["1"] = []
// `response.completed` / `response.incomplete` are clean finishes that emit a
// `finish` event; `response.failed` and `error` are hard failures that emit a
// `provider-error`. All four end the stream — kept in one set so `step` and
// `finish` event; `response.failed` is a hard failure that emits a
// `provider-error`. All three end the stream — kept in one set so `step` and
// the protocol's `terminal` predicate stay in sync.
const TERMINAL_TYPES = new Set(["response.completed", "response.incomplete", "response.failed", "error"])
const TERMINAL_TYPES = new Set(["response.completed", "response.incomplete", "response.failed"])
const onOutputTextDelta = (state: ParserState, event: OpenAIResponsesEvent): StepResult => {
if (!event.delta) return [state, NO_EVENTS]
@@ -882,7 +880,7 @@ const onResponseFinish = (state: ParserState, event: OpenAIResponsesEvent): Step
// the bare message — production rate limits and context-length failures used
// to be indistinguishable from generic stream drops.
const providerErrorMessage = (event: OpenAIResponsesEvent, fallback: string): string => {
const nested = event.error ?? event.response?.error ?? undefined
const nested = event.response?.error ?? undefined
const message = event.message || nested?.message || undefined
const code = event.code || nested?.code || undefined
if (message && code) return `${code}: ${message}`
@@ -890,26 +888,11 @@ const providerErrorMessage = (event: OpenAIResponsesEvent, fallback: string): st
}
const providerError = (event: OpenAIResponsesEvent, fallback: string) => {
const nested = event.error ?? event.response?.error ?? undefined
const code = event.code || nested?.code || undefined
const type = nested?.type || undefined
const code = event.code || event.response?.error?.code || undefined
const message = providerErrorMessage(event, fallback)
const retryable = [
"internal_error",
"rate_limit_error",
"rate_limit_exceeded",
"server_error",
"server_is_overloaded",
"service_unavailable_error",
"stream_read_error",
"upstream_error",
].some((value) => value === code || value === type)
return LLMEvent.providerError({
message,
...(code === "context_length_exceeded" || isContextOverflow(message)
? { classification: "context-overflow" as const }
: {}),
...(retryable ? { retryable: true } : {}),
classification: code === "context_length_exceeded" || isContextOverflow(message) ? "context-overflow" : undefined,
})
}
+14
View File
@@ -4,6 +4,20 @@ import { ModelID, ProviderID, ProviderMetadata, RouteID } from "./ids"
export const ProviderFailureClassification = Schema.Literal("context-overflow")
export type ProviderFailureClassification = typeof ProviderFailureClassification.Type
export const ProviderFailureCategory = Schema.Literals([
"invalid-request",
"no-route",
"authentication",
"rate-limit",
"quota-exceeded",
"content-policy",
"provider-internal",
"transport",
"invalid-provider-output",
"unknown",
])
export type ProviderFailureCategory = typeof ProviderFailureCategory.Type
export class HttpRequestDetails extends Schema.Class<HttpRequestDetails>("LLM.HttpRequestDetails")({
method: Schema.String,
url: Schema.String,
+3 -1
View File
@@ -2,7 +2,7 @@ import { Schema } from "effect"
import { ContentBlockID, FinishReason, ProtocolID, ProviderMetadata, RouteID, ToolCallID } from "./ids"
import { ModelSchema } from "./options"
import { ToolOutput, ToolResultValue } from "./messages"
import { ProviderFailureClassification } from "./errors"
import { ProviderFailureCategory, ProviderFailureClassification } from "./errors"
/**
* Token usage reported by an LLM provider.
@@ -201,6 +201,8 @@ export const ProviderErrorEvent = Schema.Struct({
type: Schema.tag("provider-error"),
message: Schema.String,
classification: Schema.optional(ProviderFailureClassification),
category: Schema.optional(ProviderFailureCategory),
status: Schema.optional(Schema.Number),
retryable: Schema.optional(Schema.Boolean),
providerMetadata: Schema.optional(ProviderMetadata),
}).annotate({ identifier: "LLM.Event.ProviderError" })
@@ -1327,9 +1327,7 @@ describe("OpenAI Responses route", () => {
// sometimes-generic provider message. The bare message alone meant
// production errors like rate limits were indistinguishable from
// unrelated stream failures.
expect(response.events).toEqual([
{ type: "provider-error", message: "rate_limit_exceeded: Slow down", retryable: true },
])
expect(response.events).toEqual([{ type: "provider-error", message: "rate_limit_exceeded: Slow down" }])
}),
)
@@ -1339,7 +1337,7 @@ describe("OpenAI Responses route", () => {
Effect.provide(fixedResponse(sseEvents({ type: "error", code: "internal_error" }))),
)
expect(response.events).toEqual([{ type: "provider-error", message: "internal_error", retryable: true }])
expect(response.events).toEqual([{ type: "provider-error", message: "internal_error" }])
}),
)
@@ -1349,7 +1347,7 @@ describe("OpenAI Responses route", () => {
Effect.provide(fixedResponse(sseEvents({ type: "error", code: "internal_error", message: "" }))),
)
expect(response.events).toEqual([{ type: "provider-error", message: "internal_error", retryable: true }])
expect(response.events).toEqual([{ type: "provider-error", message: "internal_error" }])
}),
)
@@ -1373,9 +1371,7 @@ describe("OpenAI Responses route", () => {
),
)
expect(response.events).toEqual([
{ type: "provider-error", message: "server_error: Upstream model unavailable", retryable: true },
])
expect(response.events).toEqual([{ type: "provider-error", message: "server_error: Upstream model unavailable" }])
}),
)
@@ -1424,55 +1420,6 @@ describe("OpenAI Responses route", () => {
}),
)
it.effect("surfaces and marks transient nested error envelopes retryable", () =>
Effect.gen(function* () {
const response = yield* LLMClient.generate(request).pipe(
Effect.provide(
fixedResponse(
sseEvents({
type: "error",
error: {
type: "upstream_error",
code: "stream_read_error",
message: "The upstream stream ended unexpectedly",
},
}),
),
),
)
expect(response.events).toEqual([
{
type: "provider-error",
message: "stream_read_error: The upstream stream ended unexpectedly",
retryable: true,
},
])
}),
)
it.effect("stops parsing after a terminal error event", () =>
Effect.gen(function* () {
const response = yield* LLMClient.generate(request).pipe(
Effect.provide(
fixedResponse(
sseEvents(
{
type: "error",
error: { type: "server_error", code: "server_error", message: "Transient failure" },
},
{ type: "response.output_text.delta", item_id: "ignored", delta: "must not publish" },
),
),
),
)
expect(response.events).toEqual([
{ type: "provider-error", message: "server_error: Transient failure", retryable: true },
])
}),
)
it.effect("falls back to a stable default when both error and response are absent", () =>
Effect.gen(function* () {
const response = yield* LLMClient.generate(request).pipe(
@@ -964,6 +964,7 @@ const scenarios: Scenario[] = [
headers: ctx.headers(),
}))
.status(400, undefined, "none"),
http.protected.get("/api/session/active", "v2.session.active").json(200, data(object), "none"),
http.protected
.post("/api/session", "v2.session.create")
.at((ctx) => ({
+16
View File
@@ -73,6 +73,10 @@ export const SessionsCursor = Schema.String.pipe(
)
export type SessionsCursor = typeof SessionsCursor.Type
const SessionActive = Schema.Struct({
type: Schema.Literal("running"),
}).annotate({ identifier: "SessionActive" })
const SessionsQueryCursor = SessionsCursor.annotate({
description: "Opaque pagination cursor returned as cursor.previous or cursor.next in the previous response.",
})
@@ -124,6 +128,18 @@ export const makeSessionGroup = <I extends HttpApiMiddleware.AnyId, S>(sessionLo
}),
),
)
.add(
HttpApiEndpoint.get("session.active", "/api/session/active", {
success: Schema.Struct({ data: Schema.Record(Session.ID, SessionActive) }),
}).annotateMerge(
OpenApi.annotations({
identifier: "v2.session.active",
summary: "List active sessions",
description:
"Retrieve foreground Session drains currently owned by this OpenCode process. Sessions absent from the result are inactive.",
}),
),
)
.add(
HttpApiEndpoint.get("session.get", "/api/session/:sessionID", {
params: { sessionID: Session.ID },
+5 -1
View File
@@ -50,6 +50,10 @@ const stepSettlementOptions = {
export const UnknownError = SessionMessage.UnknownError
export type UnknownError = SessionMessage.UnknownError
export const ProviderError = SessionMessage.ProviderError
export type ProviderError = SessionMessage.ProviderError
export const Error = SessionMessage.Error
export type Error = SessionMessage.Error
export const AgentSwitched = Event.define({
type: "session.next.agent.switched",
@@ -188,7 +192,7 @@ export namespace Step {
schema: {
...Base,
assistantMessageID: SessionMessage.ID,
error: UnknownError,
error: Error,
},
})
export type Failed = typeof Failed.Type
+28 -1
View File
@@ -21,6 +21,33 @@ export const UnknownError = Schema.Struct({
message: Schema.String,
}).annotate({ identifier: "Session.Error.Unknown" })
export const ProviderErrorCategory = Schema.Literals([
"invalid-request",
"no-route",
"authentication",
"rate-limit",
"quota-exceeded",
"content-policy",
"provider-internal",
"transport",
"invalid-provider-output",
"unknown",
])
export type ProviderErrorCategory = typeof ProviderErrorCategory.Type
export interface ProviderError extends Schema.Schema.Type<typeof ProviderError> {}
export const ProviderError = Schema.Struct({
type: Schema.Literal("provider"),
category: ProviderErrorCategory,
message: Schema.String,
status: Schema.Finite.pipe(Schema.optional),
retryable: Schema.Boolean,
retryAfterMs: Schema.Finite.pipe(Schema.optional),
}).annotate({ identifier: "Session.Error.Provider" })
export const Error = Schema.Union([UnknownError, ProviderError]).pipe(Schema.toTaggedUnion("type"))
export type Error = UnknownError | ProviderError
const Base = {
id: ID,
metadata: Schema.Record(Schema.String, Schema.Unknown).pipe(optional),
@@ -177,7 +204,7 @@ export const Assistant = Schema.Struct({
reasoning: Schema.Finite,
cache: Schema.Struct({ read: Schema.Finite, write: Schema.Finite }),
}).pipe(optional),
error: UnknownError.pipe(optional),
error: Error.pipe(optional),
time: Schema.Struct({
created: DateTimeUtcFromMillis,
completed: DateTimeUtcFromMillis.pipe(optional),
+2
View File
@@ -33,6 +33,7 @@ test("embedded client uses the real router and handlers", async () => {
yield* opencode.sessions.switchModel({ sessionID, model })
const selected = yield* opencode.sessions.get({ sessionID })
const page = yield* opencode.sessions.list({ directory: AbsolutePath.make(directory) })
const active = yield* opencode.sessions.active()
const admitted = yield* opencode.sessions.prompt({
sessionID,
prompt: Prompt.make({ text: "Do not run" }),
@@ -81,6 +82,7 @@ test("embedded client uses the real router and handlers", async () => {
expect(selected.model?.id).toBe(model.id)
expect(selected.model?.providerID).toBe(model.providerID)
expect(page.data.some((session) => session.id === sessionID)).toBe(true)
expect(active).toEqual({})
expect(admitted.sessionID).toBe(sessionID)
expect(prompted.type).toBe("session.next.prompted")
expect(wakeContext).toContainEqual(expect.objectContaining({ id: wake.id, type: "user" }))
+14
View File
@@ -333,6 +333,8 @@ import type {
V2QuestionRequestListResponses,
V2ReferenceListErrors,
V2ReferenceListResponses,
V2SessionActiveErrors,
V2SessionActiveResponses,
V2SessionCompactErrors,
V2SessionCompactResponses,
V2SessionContextErrors,
@@ -5501,6 +5503,18 @@ export class Session3 extends HeyApiClient {
})
}
/**
* List active sessions
*
* Retrieve foreground Session drains currently owned by this OpenCode process. Sessions absent from the result are inactive.
*/
public active<ThrowOnError extends boolean = false>(options?: Options<never, ThrowOnError>) {
return (options?.client ?? this.client).get<V2SessionActiveResponses, V2SessionActiveErrors, ThrowOnError>({
url: "/api/session/active",
...options,
})
}
/**
* Get session
*
+62 -6
View File
@@ -952,7 +952,7 @@ export type GlobalEvent = {
timestamp: number
sessionID: string
assistantMessageID: string
error: SessionErrorUnknown
error: SessionErrorUnknown | SessionErrorProvider
}
}
| {
@@ -2690,6 +2690,10 @@ export type InvalidCursorError = {
message: string
}
export type SessionActive = {
type: "running"
}
export type SessionNotFoundError = {
_tag: "SessionNotFoundError"
sessionID: string
@@ -2951,6 +2955,25 @@ export type SessionErrorUnknown = {
message: string
}
export type SessionErrorProvider = {
type: "provider"
category:
| "invalid-request"
| "no-route"
| "authentication"
| "rate-limit"
| "quota-exceeded"
| "content-policy"
| "provider-internal"
| "transport"
| "invalid-provider-output"
| "unknown"
message: string
status?: number
retryable: boolean
retryAfterMs?: number
}
export type LlmProviderMetadata = {
[key: string]: {
[key: string]: unknown
@@ -3395,7 +3418,7 @@ export type SyncEventSessionNextStepFailed = {
timestamp: number
sessionID: string
assistantMessageID: string
error: SessionErrorUnknown
error: SessionErrorUnknown | SessionErrorProvider
}
}
}
@@ -4005,7 +4028,7 @@ export type SessionMessageAssistant = {
write: number
}
}
error?: SessionErrorUnknown
error?: SessionErrorUnknown | SessionErrorProvider
}
export type SessionMessageCompaction = {
@@ -4285,7 +4308,7 @@ export type SessionNextStepFailed = {
timestamp: number
sessionID: string
assistantMessageID: string
error: SessionErrorUnknown
error: SessionErrorUnknown | SessionErrorProvider
}
}
@@ -5340,7 +5363,7 @@ export type V2EventSessionNextStepFailed = {
timestamp: number
sessionID: string
assistantMessageID: string
error: SessionErrorUnknown
error: SessionErrorUnknown | SessionErrorProvider
}
}
@@ -6927,7 +6950,7 @@ export type EventSessionNextStepFailed = {
timestamp: number
sessionID: string
assistantMessageID: string
error: SessionErrorUnknown
error: SessionErrorUnknown | SessionErrorProvider
}
}
@@ -11940,6 +11963,39 @@ export type V2SessionCreateResponses = {
export type V2SessionCreateResponse = V2SessionCreateResponses[keyof V2SessionCreateResponses]
export type V2SessionActiveData = {
body?: never
path?: never
query?: never
url: "/api/session/active"
}
export type V2SessionActiveErrors = {
/**
* InvalidRequestError
*/
400: InvalidRequestError
/**
* UnauthorizedError
*/
401: UnauthorizedError
}
export type V2SessionActiveError = V2SessionActiveErrors[keyof V2SessionActiveErrors]
export type V2SessionActiveResponses = {
/**
* Success
*/
200: {
data: {
[key: string]: unknown | SessionActive
}
}
}
export type V2SessionActiveResponse = V2SessionActiveResponses[keyof V2SessionActiveResponses]
export type V2SessionGetData = {
body?: never
path: {
+71
View File
@@ -10187,6 +10187,66 @@
]
}
},
"/api/session/active": {
"get": {
"tags": ["sessions"],
"operationId": "v2.session.active",
"parameters": [],
"security": [],
"responses": {
"200": {
"description": "Success",
"content": {
"application/json": {
"schema": {
"type": "object",
"properties": {
"data": {
"type": "object",
"patternProperties": {
"^ses": {
"$ref": "#/components/schemas/SessionActive"
}
}
}
},
"required": ["data"],
"additionalProperties": false
}
}
}
},
"400": {
"description": "InvalidRequestError",
"content": {
"application/json": {
"schema": {
"$ref": "#/components/schemas/InvalidRequestError"
}
}
}
},
"401": {
"description": "UnauthorizedError",
"content": {
"application/json": {
"schema": {
"$ref": "#/components/schemas/UnauthorizedError"
}
}
}
}
},
"description": "Retrieve foreground Session drains currently owned by this OpenCode process. Sessions absent from the result are inactive.",
"summary": "List active sessions",
"x-codeSamples": [
{
"lang": "js",
"source": "import { createOpencodeClient } from \"@opencode-ai/sdk\n\nconst client = createOpencodeClient()\nawait client.v2.session.active({\n ...\n})"
}
]
}
},
"/api/session/{sessionID}": {
"get": {
"tags": ["sessions"],
@@ -23445,6 +23505,17 @@
"required": ["_tag", "message"],
"additionalProperties": false
},
"SessionActive": {
"type": "object",
"properties": {
"type": {
"type": "string",
"enum": ["running"]
}
},
"required": ["type"],
"additionalProperties": false
},
"SessionNotFoundError": {
"type": "object",
"properties": {
+10
View File
@@ -76,6 +76,16 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl
}
}),
)
.handle(
"session.active",
Effect.fn(function* () {
return {
data: Object.fromEntries(
Array.from(yield* session.active, (sessionID) => [sessionID, { type: "running" as const }]),
),
}
}),
)
.handle(
"session.get",
Effect.fn(function* (ctx) {
+7
View File
@@ -25,6 +25,11 @@ sessions.interrupt(sessionID)
-> clears a coalesced follow-up wake already registered with this coordinator
-> preserves durable inbox rows for a later wake or resume
-> idle or missing Session is a no-op
sessions.active()
-> snapshots foreground Session drains owned by this process
-> returns only active Session IDs with { type: "running" }
-> absence means inactive; activity is not durable across process restarts
```
`session_input` is the durable admission inbox. `PromptAdmitted` records and projects accepted input so pending queue state can be replayed, replicated, and observed by clients. Admitted inputs remain outside model-visible Session history until the serialized runner publishes `Prompted`. Its projector atomically writes the visible user message and marks the inbox row promoted in the same event transaction. The V1-to-V2 shadow bridge publishes the same `Prompted` event for already-visible V1 prompts.
@@ -161,6 +166,8 @@ Post-crash continuation recovery is intentionally deferred. A wake does not infe
A process-global `SessionRunCoordinator` serializes execution for each local Session while allowing different Sessions to run concurrently. Resumes join active execution, overlapping wakes coalesce into one follow-up, and interruption stops current process-local execution without deleting durable inbox work. The runner enters the Session's current Location when execution starts and fences each new provider turn against that Location.
The coordinator's active registry is also the source for `sessions.active()`. It represents only foreground Session drains owned by the current process; background subagents and tasks do not add parent Sessions to this registry. The snapshot is runtime state and is empty after a process restart.
Inbox promotion coalesces pending steers in durable admission order. Once continuation would otherwise end, it promotes one queued input at a time in FIFO order. Add explicit inbox backlog and steering-batch limits before exposing broad multi-caller admission or untrusted queue growth.
Eager local-tool execution is intentionally unbounded in the current local slice. This minimizes tool latency but does not increase SQLite settlement throughput: Session-event publication remains serialized per provider turn. Before broadening exposure, revisit per-turn call limits, output truncation, and operational backpressure using observed workloads. The `session.next.*` event schemas remain experimental and unshipped; databases created by earlier experimental builds are disposable rather than compatibility targets.