mirror of
https://github.com/anomalyco/opencode.git
synced 2026-08-05 01:43:27 -04:00
Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| b50cd706ce |
+21
-15
@@ -165,6 +165,7 @@ export const allBounded = (events: Interface, capacity: number) =>
|
||||
|
||||
export interface LayerOptions {
|
||||
readonly beforeAggregateRead?: (aggregateID: string) => Effect.Effect<void>
|
||||
readonly persistDurableEvents?: boolean
|
||||
}
|
||||
|
||||
export const layerWith = (options?: LayerOptions) =>
|
||||
@@ -180,6 +181,7 @@ export const layerWith = (options?: LayerOptions) =>
|
||||
// TODO: Bind durable projectors to exact type+version before supporting incompatible historical payloads.
|
||||
const listeners = new Array<Subscriber>()
|
||||
const { db } = yield* Database.Service
|
||||
const persistDurableEvents = options?.persistDurableEvents ?? true
|
||||
|
||||
const getOrCreate = (definition: Definition) =>
|
||||
Effect.gen(function* () {
|
||||
@@ -333,19 +335,20 @@ export const layerWith = (options?: LayerOptions) =>
|
||||
})
|
||||
.run()
|
||||
.pipe(Effect.orDie)
|
||||
yield* db
|
||||
.insert(EventTable)
|
||||
.values([
|
||||
{
|
||||
id: event.id,
|
||||
aggregate_id: aggregateID,
|
||||
seq,
|
||||
type: versionedType(definition.type, durable.version),
|
||||
data: encoded,
|
||||
},
|
||||
])
|
||||
.run()
|
||||
.pipe(Effect.orDie)
|
||||
if (persistDurableEvents)
|
||||
yield* db
|
||||
.insert(EventTable)
|
||||
.values([
|
||||
{
|
||||
id: event.id,
|
||||
aggregate_id: aggregateID,
|
||||
seq,
|
||||
type: versionedType(definition.type, durable.version),
|
||||
data: encoded,
|
||||
},
|
||||
])
|
||||
.run()
|
||||
.pipe(Effect.orDie)
|
||||
return { aggregateID, seq }
|
||||
}),
|
||||
{ behavior: "immediate" },
|
||||
@@ -634,5 +637,8 @@ export const layerWith = (options?: LayerOptions) =>
|
||||
}),
|
||||
)
|
||||
|
||||
const layer = layerWith()
|
||||
export const node = makeGlobalNode({ service: Service, layer: layer, deps: [Database.node] })
|
||||
export const nodeWith = (options?: LayerOptions) =>
|
||||
makeGlobalNode({ service: Service, layer: layerWith(options), deps: [Database.node] })
|
||||
|
||||
export const node = nodeWith()
|
||||
export const sequenceOnlyNode = nodeWith({ persistDurableEvents: false })
|
||||
|
||||
@@ -51,6 +51,7 @@ import { memoMap } from "@opencode-ai/core/effect/memo-map"
|
||||
import { BackgroundJob } from "@/background/job"
|
||||
import { RuntimeFlags } from "@/effect/runtime-flags"
|
||||
import { EventV2Bridge } from "@/event-v2-bridge"
|
||||
import { EventV2 } from "@opencode-ai/core/event"
|
||||
import { LayerNode } from "@opencode-ai/core/effect/layer-node"
|
||||
import { AppNodeBuilderV1 } from "./app-node-builder-v1"
|
||||
import { SessionProjector } from "@opencode-ai/core/session/projector"
|
||||
@@ -106,6 +107,7 @@ export const AppLayer = AppNodeBuilderV1.build(
|
||||
ShareNext.node,
|
||||
SessionShare.node,
|
||||
]),
|
||||
[[EventV2.node, EventV2.sequenceOnlyNode]],
|
||||
).pipe(Layer.provideMerge(AppNodeBuilderV1.build(Ripgrep.node)), Layer.provideMerge(Observability.layer))
|
||||
|
||||
const rt = ManagedRuntime.make(AppLayer, { memoMap })
|
||||
|
||||
@@ -270,8 +270,10 @@ const app = LayerNode.group([
|
||||
|
||||
export function createRoutes(
|
||||
corsOptions?: CorsOptions,
|
||||
persistDurableEvents = false,
|
||||
): Layer.Layer<never, EffectConfig.ConfigError, RouteRequirements> {
|
||||
const locationServiceMapV2 = buildLocationServiceMap()
|
||||
const eventNode = persistDurableEvents ? EventV2.node : EventV2.sequenceOnlyNode
|
||||
|
||||
return Layer.mergeAll(
|
||||
rootApiRoutes,
|
||||
@@ -288,7 +290,10 @@ export function createRoutes(
|
||||
corsVaryFix,
|
||||
fenceLayer,
|
||||
cors(corsOptions),
|
||||
AppNodeBuilderV1.build(MoveSession.node, [[LocationServiceMap.node, locationServiceMapV2]]),
|
||||
AppNodeBuilderV1.build(MoveSession.node, [
|
||||
[LocationServiceMap.node, locationServiceMapV2],
|
||||
[EventV2.node, eventNode],
|
||||
]),
|
||||
HttpServer.layerServices,
|
||||
]),
|
||||
Layer.provide(Layer.succeed(CorsConfig)(corsOptions)),
|
||||
@@ -299,11 +304,12 @@ export function createRoutes(
|
||||
AppNodeBuilderV1.build(SessionV2.node, [
|
||||
[LocationServiceMap.node, locationServiceMapV2],
|
||||
[SessionExecution.node, SessionExecutionLocal.node],
|
||||
[EventV2.node, eventNode],
|
||||
]),
|
||||
),
|
||||
Layer.provide(locationServiceMapV2),
|
||||
|
||||
Layer.provide(AppNodeBuilderV1.build(app)),
|
||||
Layer.provide(AppNodeBuilderV1.build(app, [[EventV2.node, eventNode]])),
|
||||
// Must stay last: layers provided later in this pipe build beneath earlier ones,
|
||||
// so Observability must come after every service graph. Otherwise eagerly forked
|
||||
// fibers (e.g. the ModelsDev background refresh) capture Effect's default stdout
|
||||
|
||||
@@ -636,30 +636,14 @@ const layer = Layer.effect(
|
||||
yield* Effect.gen(function* () {
|
||||
ctx.currentText = undefined
|
||||
ctx.reasoningMap = {}
|
||||
let generated = false
|
||||
yield* status.set(ctx.sessionID, { type: "busy" })
|
||||
const stream = llm.stream(streamInput)
|
||||
|
||||
yield* stream.pipe(
|
||||
Stream.tap((event) => {
|
||||
if (
|
||||
(event.type === "text-delta" && event.text.length > 0) ||
|
||||
(event.type === "reasoning-delta" && event.text.length > 0) ||
|
||||
event.type === "tool-input-start" ||
|
||||
event.type === "tool-call"
|
||||
) {
|
||||
generated = true
|
||||
}
|
||||
return handleEvent(event)
|
||||
}),
|
||||
Stream.tap((event) => handleEvent(event)),
|
||||
Stream.takeUntil(() => ctx.needsCompaction),
|
||||
Stream.runDrain,
|
||||
)
|
||||
if (ctx.assistantMessage.finish === "unknown" && !generated) {
|
||||
yield* new SessionRetry.EmptyResponseError({
|
||||
message: "The model returned an empty response with an unknown finish reason",
|
||||
})
|
||||
}
|
||||
}).pipe(
|
||||
Effect.onInterrupt(() =>
|
||||
Effect.gen(function* () {
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import type { NamedError } from "@opencode-ai/core/util/error"
|
||||
import { SessionV1 } from "@opencode-ai/core/v1/session"
|
||||
import { Cause, Clock, Duration, Effect, Schedule, Schema } from "effect"
|
||||
import { Cause, Clock, Duration, Effect, Schedule } from "effect"
|
||||
import { MessageV2 } from "./message-v2"
|
||||
import { iife } from "@/util/iife"
|
||||
import { isRecord } from "@/util/record"
|
||||
@@ -23,10 +23,6 @@ export type Retryable = {
|
||||
}
|
||||
}
|
||||
|
||||
export class EmptyResponseError extends Schema.TaggedErrorClass<EmptyResponseError>()("SessionEmptyResponseError", {
|
||||
message: Schema.String,
|
||||
}) {}
|
||||
|
||||
export const RETRY_INITIAL_DELAY = 2000
|
||||
export const RETRY_BACKOFF_FACTOR = 2
|
||||
export const RETRY_MAX_DELAY_NO_HEADERS = 30_000 // 30 seconds
|
||||
@@ -185,8 +181,7 @@ export function policy(opts: {
|
||||
return Schedule.fromStepWithMetadata(
|
||||
Effect.succeed((meta: Schedule.InputMetadata<unknown>) => {
|
||||
const error = opts.parse(meta.input)
|
||||
const retry =
|
||||
meta.input instanceof EmptyResponseError ? { message: meta.input.message } : retryable(error, opts.provider)
|
||||
const retry = retryable(error, opts.provider)
|
||||
if (!retry) return Cause.done(meta.attempt)
|
||||
return Effect.gen(function* () {
|
||||
const wait = delay(meta.attempt, SessionV1.APIError.isInstance(error) ? error : undefined)
|
||||
|
||||
@@ -81,11 +81,11 @@ describe("opencode run (non-interactive subprocess)", () => {
|
||||
30_000,
|
||||
)
|
||||
|
||||
// The test provider's SSE error item is interpreted by the SDK as an empty
|
||||
// response with an unknown finish. That attempt should retry while preserving
|
||||
// output from the preceding tool-call step.
|
||||
// The test provider's SSE error item is interpreted by the SDK as an unknown
|
||||
// finish, not a fatal provider/session error. Lock that distinction in so it
|
||||
// is not accidentally used as the failure compatibility oracle.
|
||||
cliIt.concurrent(
|
||||
"empty unknown stream finish retries and preserves partial output",
|
||||
"unknown stream finish preserves partial output and exits 0",
|
||||
({ llm, opencode }) =>
|
||||
Effect.gen(function* () {
|
||||
yield* llm.push(
|
||||
@@ -95,10 +95,9 @@ describe("opencode run (non-interactive subprocess)", () => {
|
||||
}),
|
||||
)
|
||||
yield* llm.fail("upstream provider exploded mid-stream")
|
||||
yield* llm.text("recovered response")
|
||||
const result = yield* opencode.run("trigger midstream error", { timeoutMs: 30_000 })
|
||||
expect(result.exitCode).toBe(0)
|
||||
expect(result.stdout).toBe("partial response\nrecovered response\n")
|
||||
expect(result.stdout).toBe("partial response\n")
|
||||
expect(result.stderr).not.toContain("upstream provider exploded mid-stream")
|
||||
}),
|
||||
60_000,
|
||||
@@ -214,7 +213,7 @@ describe("opencode run (non-interactive subprocess)", () => {
|
||||
)
|
||||
|
||||
cliIt.concurrent(
|
||||
"--format json records an empty unknown stream retry",
|
||||
"--format json records partial output for an unknown stream finish",
|
||||
({ llm, opencode }) =>
|
||||
Effect.gen(function* () {
|
||||
yield* llm.push(
|
||||
@@ -224,7 +223,6 @@ describe("opencode run (non-interactive subprocess)", () => {
|
||||
}),
|
||||
)
|
||||
yield* llm.fail("provider failed")
|
||||
yield* llm.text("recovered json")
|
||||
const result = yield* opencode.run("fail after output", { format: "json" })
|
||||
|
||||
const events = opencode.parseJsonEvents(result.stdout)
|
||||
@@ -236,13 +234,9 @@ describe("opencode run (non-interactive subprocess)", () => {
|
||||
"step_finish",
|
||||
"step_start",
|
||||
"step_finish",
|
||||
"step_start",
|
||||
"text",
|
||||
"step_finish",
|
||||
])
|
||||
expect(events[1]?.part).toEqual(expect.objectContaining({ type: "text", text: "partial json" }))
|
||||
expect(events.at(-2)?.part).toEqual(expect.objectContaining({ type: "text", text: "recovered json" }))
|
||||
expect(events.at(-1)?.part).toEqual(expect.objectContaining({ type: "step-finish", reason: "stop" }))
|
||||
expect(events.at(-1)?.part).toEqual(expect.objectContaining({ type: "step-finish", reason: "unknown" }))
|
||||
}),
|
||||
60_000,
|
||||
)
|
||||
|
||||
@@ -48,7 +48,7 @@ const appLayer = AppNodeBuilder.build(
|
||||
[[InstanceStore.bootstrapNode, noopBootstrapLayer]],
|
||||
)
|
||||
const servedRoutes: Layer.Layer<never, Config.ConfigError, HttpServer.HttpServer> = HttpRouter.serve(
|
||||
HttpApiApp.routes,
|
||||
HttpApiApp.createRoutes(undefined, true),
|
||||
{
|
||||
disableListenLog: true,
|
||||
disableLogger: true,
|
||||
|
||||
@@ -604,68 +604,6 @@ it.live("session.processor effect tests retry recognized structured json errors"
|
||||
),
|
||||
)
|
||||
|
||||
it.live("session.processor effect tests retry empty responses with unknown finish reasons", () =>
|
||||
provideTmpdirServer(
|
||||
({ dir, llm }) =>
|
||||
Effect.gen(function* () {
|
||||
const { processors, session, provider } = yield* boot()
|
||||
|
||||
yield* llm.push(
|
||||
raw({
|
||||
chunks: [
|
||||
{
|
||||
id: "chatcmpl-test",
|
||||
object: "chat.completion.chunk",
|
||||
choices: [{ delta: { role: "assistant" }, finish_reason: null }],
|
||||
},
|
||||
{
|
||||
id: "chatcmpl-test",
|
||||
object: "chat.completion.chunk",
|
||||
choices: [{ delta: {}, finish_reason: "unknown_reason" }],
|
||||
},
|
||||
],
|
||||
}),
|
||||
reply().text("after").stop(),
|
||||
)
|
||||
|
||||
const chat = yield* session.create({})
|
||||
const parent = yield* user(chat.id, "retry empty")
|
||||
const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
|
||||
const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
|
||||
const handle = yield* processors.create({
|
||||
assistantMessage: msg,
|
||||
sessionID: chat.id,
|
||||
model: mdl,
|
||||
})
|
||||
|
||||
const value = yield* handle.process({
|
||||
user: {
|
||||
id: parent.id,
|
||||
sessionID: chat.id,
|
||||
role: "user",
|
||||
time: parent.time,
|
||||
agent: parent.agent,
|
||||
model: { providerID: ref.providerID, modelID: ref.modelID },
|
||||
} satisfies SessionV1.User,
|
||||
sessionID: chat.id,
|
||||
model: mdl,
|
||||
agent: agent(),
|
||||
system: [],
|
||||
messages: [{ role: "user", content: "retry empty" }],
|
||||
tools: {},
|
||||
})
|
||||
|
||||
const parts = yield* MessageV2.parts(msg.id)
|
||||
|
||||
expect(value).toBe("continue")
|
||||
expect(yield* llm.calls).toBe(2)
|
||||
expect(parts.some((part) => part.type === "text" && part.text === "after")).toBe(true)
|
||||
expect(handle.message.error).toBeUndefined()
|
||||
}),
|
||||
{ config: (url) => providerCfg(url) },
|
||||
),
|
||||
)
|
||||
|
||||
it.live("session.processor effect tests publish retry status updates", () =>
|
||||
provideTmpdirServer(
|
||||
({ dir, llm }) =>
|
||||
|
||||
Reference in New Issue
Block a user