refactor(httpapi): preserve typed errors in session prompt handlers (#25181)
This commit is contained in:
@@ -1,4 +1,4 @@
|
|||||||
import { Effect, Fiber } from "effect"
|
import { Effect, Exit, Fiber } from "effect"
|
||||||
import { WorkspaceContext } from "@/control-plane/workspace-context"
|
import { WorkspaceContext } from "@/control-plane/workspace-context"
|
||||||
import { Instance, type InstanceContext } from "@/project/instance"
|
import { Instance, type InstanceContext } from "@/project/instance"
|
||||||
import type { WorkspaceID } from "@/control-plane/schema"
|
import type { WorkspaceID } from "@/control-plane/schema"
|
||||||
@@ -9,6 +9,7 @@ import { attachWith } from "./run-service"
|
|||||||
export interface Shape {
|
export interface Shape {
|
||||||
readonly promise: <A, E, R>(effect: Effect.Effect<A, E, R>) => Promise<A>
|
readonly promise: <A, E, R>(effect: Effect.Effect<A, E, R>) => Promise<A>
|
||||||
readonly fork: <A, E, R>(effect: Effect.Effect<A, E, R>) => Fiber.Fiber<A, E>
|
readonly fork: <A, E, R>(effect: Effect.Effect<A, E, R>) => Fiber.Fiber<A, E>
|
||||||
|
readonly run: <A, E, R>(effect: Effect.Effect<A, E, R>) => Effect.Effect<A, E>
|
||||||
}
|
}
|
||||||
|
|
||||||
function restore<R>(instance: InstanceContext | undefined, workspace: WorkspaceID | undefined, fn: () => R): R {
|
function restore<R>(instance: InstanceContext | undefined, workspace: WorkspaceID | undefined, fn: () => R): R {
|
||||||
@@ -43,6 +44,14 @@ export function make(): Effect.Effect<Shape> {
|
|||||||
restore(instance, workspace, () => Effect.runPromise(wrap(effect))),
|
restore(instance, workspace, () => Effect.runPromise(wrap(effect))),
|
||||||
fork: <A, E, R>(effect: Effect.Effect<A, E, R>) =>
|
fork: <A, E, R>(effect: Effect.Effect<A, E, R>) =>
|
||||||
restore(instance, workspace, () => Effect.runFork(wrap(effect))),
|
restore(instance, workspace, () => Effect.runFork(wrap(effect))),
|
||||||
|
run: <A, E, R>(effect: Effect.Effect<A, E, R>) =>
|
||||||
|
Effect.callback<A, E>((resume) => {
|
||||||
|
restore(instance, workspace, () =>
|
||||||
|
Effect.runPromiseExit(wrap(effect)).then((exit) =>
|
||||||
|
resume(Exit.isSuccess(exit) ? Effect.succeed(exit.value) : Effect.failCause(exit.cause)),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
}),
|
||||||
} satisfies Shape
|
} satisfies Shape
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -18,9 +18,8 @@ import { SessionSummary } from "@/session/summary"
|
|||||||
import { Todo } from "@/session/todo"
|
import { Todo } from "@/session/todo"
|
||||||
import { MessageID, PartID, SessionID } from "@/session/schema"
|
import { MessageID, PartID, SessionID } from "@/session/schema"
|
||||||
import { NotFoundError } from "@/storage/storage"
|
import { NotFoundError } from "@/storage/storage"
|
||||||
import * as Log from "@opencode-ai/core/util/log"
|
|
||||||
import { NamedError } from "@opencode-ai/core/util/error"
|
import { NamedError } from "@opencode-ai/core/util/error"
|
||||||
import { Effect, Schema } from "effect"
|
import { Cause, Effect, Schema } from "effect"
|
||||||
import * as Stream from "effect/Stream"
|
import * as Stream from "effect/Stream"
|
||||||
import { HttpServerRequest, HttpServerResponse } from "effect/unstable/http"
|
import { HttpServerRequest, HttpServerResponse } from "effect/unstable/http"
|
||||||
import { HttpApiBuilder, HttpApiError, HttpApiSchema } from "effect/unstable/httpapi"
|
import { HttpApiBuilder, HttpApiError, HttpApiSchema } from "effect/unstable/httpapi"
|
||||||
@@ -40,8 +39,6 @@ import {
|
|||||||
UpdatePayload,
|
UpdatePayload,
|
||||||
} from "../groups/session"
|
} from "../groups/session"
|
||||||
|
|
||||||
const log = Log.create({ service: "server" })
|
|
||||||
|
|
||||||
const mapNotFound = <A, E, R>(self: Effect.Effect<A, E, R>) =>
|
const mapNotFound = <A, E, R>(self: Effect.Effect<A, E, R>) =>
|
||||||
self.pipe(
|
self.pipe(
|
||||||
Effect.catchIf(NotFoundError.isInstance, () => Effect.fail(new HttpApiError.NotFound({}))),
|
Effect.catchIf(NotFoundError.isInstance, () => Effect.fail(new HttpApiError.NotFound({}))),
|
||||||
@@ -63,6 +60,7 @@ export const sessionHandlers = HttpApiBuilder.group(InstanceHttpApi, "session",
|
|||||||
const statusSvc = yield* SessionStatus.Service
|
const statusSvc = yield* SessionStatus.Service
|
||||||
const todoSvc = yield* Todo.Service
|
const todoSvc = yield* Todo.Service
|
||||||
const summary = yield* SessionSummary.Service
|
const summary = yield* SessionSummary.Service
|
||||||
|
const bus = yield* Bus.Service
|
||||||
|
|
||||||
const list = Effect.fn("SessionHttpApi.list")(function* (ctx: { query: typeof ListQuery.Type }) {
|
const list = Effect.fn("SessionHttpApi.list")(function* (ctx: { query: typeof ListQuery.Type }) {
|
||||||
const instance = yield* InstanceState.context
|
const instance = yield* InstanceState.context
|
||||||
@@ -264,13 +262,11 @@ export const sessionHandlers = HttpApiBuilder.group(InstanceHttpApi, "session",
|
|||||||
const bridge = yield* EffectBridge.make()
|
const bridge = yield* EffectBridge.make()
|
||||||
return HttpServerResponse.stream(
|
return HttpServerResponse.stream(
|
||||||
Stream.fromEffect(
|
Stream.fromEffect(
|
||||||
Effect.promise(() =>
|
bridge.run(
|
||||||
bridge.promise(
|
promptSvc.prompt({
|
||||||
promptSvc.prompt({
|
...ctx.payload,
|
||||||
...ctx.payload,
|
sessionID: ctx.params.sessionID,
|
||||||
sessionID: ctx.params.sessionID,
|
}),
|
||||||
}),
|
|
||||||
),
|
|
||||||
),
|
),
|
||||||
).pipe(
|
).pipe(
|
||||||
Stream.map((message) => JSON.stringify(message)),
|
Stream.map((message) => JSON.stringify(message)),
|
||||||
@@ -288,12 +284,12 @@ export const sessionHandlers = HttpApiBuilder.group(InstanceHttpApi, "session",
|
|||||||
yield* Effect.sync(() => {
|
yield* Effect.sync(() => {
|
||||||
bridge.fork(
|
bridge.fork(
|
||||||
promptSvc.prompt({ ...ctx.payload, sessionID: ctx.params.sessionID }).pipe(
|
promptSvc.prompt({ ...ctx.payload, sessionID: ctx.params.sessionID }).pipe(
|
||||||
Effect.catchCause((error) =>
|
Effect.catchCause((cause) =>
|
||||||
Effect.sync(() => {
|
Effect.gen(function* () {
|
||||||
log.error("prompt_async failed", { sessionID: ctx.params.sessionID, error })
|
yield* Effect.logError("prompt_async failed", { sessionID: ctx.params.sessionID, cause })
|
||||||
void Bus.publish(Session.Event.Error, {
|
yield* bus.publish(Session.Event.Error, {
|
||||||
sessionID: ctx.params.sessionID,
|
sessionID: ctx.params.sessionID,
|
||||||
error: new NamedError.Unknown({ message: String(error) }).toObject(),
|
error: new NamedError.Unknown({ message: Cause.pretty(cause) }).toObject(),
|
||||||
})
|
})
|
||||||
}),
|
}),
|
||||||
),
|
),
|
||||||
|
|||||||
Reference in New Issue
Block a user