mirror of
https://github.com/anomalyco/opencode.git
synced 2026-08-05 01:43:27 -04:00
Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 5562652f63 | |||
| 031510e1c3 | |||
| f0afb6750e | |||
| 703d09f306 | |||
| aefaf140c1 | |||
| 44614c79c4 |
@@ -40,7 +40,10 @@ export class Subscription {
|
||||
private readonly abort = new AbortController()
|
||||
private readonly shellSnapshots = new Map<string, string>()
|
||||
private readonly toolStarts = new Set<string>()
|
||||
private readonly connectionWaiters = new Set<() => void>()
|
||||
private readonly idleWaiters = new Map<string, Set<ReturnType<typeof signal>>>()
|
||||
private readonly permission: ACPPermission.Handler
|
||||
private connected = false
|
||||
private started = false
|
||||
|
||||
constructor(
|
||||
@@ -63,10 +66,35 @@ export class Subscription {
|
||||
|
||||
stop() {
|
||||
this.abort.abort()
|
||||
this.disconnected()
|
||||
for (const resolve of this.connectionWaiters) resolve()
|
||||
this.connectionWaiters.clear()
|
||||
}
|
||||
|
||||
async runUntilIdle<A>(sessionId: string, request: () => Promise<A>) {
|
||||
await this.waitUntilConnected()
|
||||
const waiter = signal()
|
||||
const waiters = this.idleWaiters.get(sessionId) ?? new Set()
|
||||
waiters.add(waiter)
|
||||
this.idleWaiters.set(sessionId, waiters)
|
||||
|
||||
try {
|
||||
// Idle is queued after the turn's events, and this subscription awaits each update in order.
|
||||
void waiter.promise.catch(() => {})
|
||||
const response = await request()
|
||||
await waiter.promise
|
||||
return response
|
||||
} finally {
|
||||
waiters.delete(waiter)
|
||||
if (waiters.size === 0) this.idleWaiters.delete(sessionId)
|
||||
}
|
||||
}
|
||||
|
||||
async handle(event: Event) {
|
||||
switch (event.type) {
|
||||
case "session.status":
|
||||
if (event.properties.status.type === "idle") this.idle(event.properties.sessionID)
|
||||
return
|
||||
case "permission.asked":
|
||||
this.permission.handle(event)
|
||||
return
|
||||
@@ -115,19 +143,51 @@ export class Subscription {
|
||||
|
||||
private async run() {
|
||||
while (!this.abort.signal.aborted) {
|
||||
const events = (await this.input.sdk.global.event({
|
||||
signal: this.abort.signal,
|
||||
})) as GlobalEventStream
|
||||
|
||||
for await (const event of events.stream) {
|
||||
if (this.abort.signal.aborted) return
|
||||
if (!event.payload) continue
|
||||
await this.handle(event.payload).catch(() => {})
|
||||
}
|
||||
await this.consume().catch(() => {})
|
||||
this.disconnected()
|
||||
if (!this.abort.signal.aborted) await new Promise((resolve) => setTimeout(resolve, 1000))
|
||||
}
|
||||
}
|
||||
|
||||
private async consume() {
|
||||
const events = (await this.input.sdk.global.event({
|
||||
signal: this.abort.signal,
|
||||
})) as GlobalEventStream
|
||||
this.connected = true
|
||||
for (const resolve of this.connectionWaiters) resolve()
|
||||
this.connectionWaiters.clear()
|
||||
|
||||
for await (const event of events.stream) {
|
||||
if (this.abort.signal.aborted) return
|
||||
if (!event.payload) continue
|
||||
await this.handle(event.payload).catch(() => {})
|
||||
}
|
||||
}
|
||||
|
||||
private async waitUntilConnected() {
|
||||
while (!this.connected) {
|
||||
if (this.abort.signal.aborted) throw new Error("ACP event subscription stopped")
|
||||
await new Promise<void>((resolve) => this.connectionWaiters.add(resolve))
|
||||
}
|
||||
}
|
||||
|
||||
private disconnected() {
|
||||
if (!this.connected) return
|
||||
this.connected = false
|
||||
const error = new Error("ACP event stream disconnected")
|
||||
for (const waiters of this.idleWaiters.values()) {
|
||||
for (const waiter of waiters) waiter.reject(error)
|
||||
}
|
||||
this.idleWaiters.clear()
|
||||
}
|
||||
|
||||
private idle(sessionId: string) {
|
||||
const waiters = this.idleWaiters.get(sessionId)
|
||||
if (!waiters) return
|
||||
this.idleWaiters.delete(sessionId)
|
||||
for (const waiter of waiters) waiter.resolve()
|
||||
}
|
||||
|
||||
private async handlePartUpdated(event: EventMessagePartUpdated) {
|
||||
const part = event.properties.part
|
||||
const sessionId = part.sessionID || event.properties.sessionID
|
||||
@@ -339,4 +399,23 @@ export class Subscription {
|
||||
}
|
||||
}
|
||||
|
||||
function signal() {
|
||||
const state: {
|
||||
resolve: () => void
|
||||
reject: (reason?: unknown) => void
|
||||
} = {
|
||||
resolve: () => {},
|
||||
reject: () => {},
|
||||
}
|
||||
const promise = new Promise<void>((resolve, reject) => {
|
||||
state.resolve = resolve
|
||||
state.reject = reject
|
||||
})
|
||||
return {
|
||||
promise,
|
||||
resolve: () => state.resolve(),
|
||||
reject: (reason?: unknown) => state.reject(reason),
|
||||
}
|
||||
}
|
||||
|
||||
export * as ACPEvent from "./event"
|
||||
|
||||
@@ -88,6 +88,8 @@ export function make(input: {
|
||||
? ACPEvent.start({ sdk: input.sdk, connection: input.connection, session })
|
||||
: undefined
|
||||
if (events) input.eventSubscription?.(events)
|
||||
const runUntilIdle = <A>(sessionId: string, fn: () => Promise<A>) =>
|
||||
events ? events.runUntilIdle(sessionId, fn) : fn()
|
||||
|
||||
const initialize = Effect.fn("ACP.initialize")(function* (params: InitializeRequest) {
|
||||
const started = performance.now()
|
||||
@@ -504,19 +506,21 @@ export function make(input: {
|
||||
if (!command) {
|
||||
const response = yield* request(
|
||||
() =>
|
||||
input.sdk.session.prompt(
|
||||
{
|
||||
sessionID: current.id,
|
||||
model: {
|
||||
providerID: selected.providerID,
|
||||
modelID: selected.modelID,
|
||||
runUntilIdle(current.id, () =>
|
||||
input.sdk.session.prompt(
|
||||
{
|
||||
sessionID: current.id,
|
||||
model: {
|
||||
providerID: selected.providerID,
|
||||
modelID: selected.modelID,
|
||||
},
|
||||
...(variant ? { variant } : {}),
|
||||
parts,
|
||||
...(modeId ? { agent: modeId } : {}),
|
||||
directory: current.cwd,
|
||||
},
|
||||
...(variant ? { variant } : {}),
|
||||
parts,
|
||||
...(modeId ? { agent: modeId } : {}),
|
||||
directory: current.cwd,
|
||||
},
|
||||
{ throwOnError: true },
|
||||
{ throwOnError: true },
|
||||
),
|
||||
),
|
||||
"session",
|
||||
)
|
||||
@@ -528,17 +532,19 @@ export function make(input: {
|
||||
if (known) {
|
||||
const response = yield* request(
|
||||
() =>
|
||||
input.sdk.session.command(
|
||||
{
|
||||
sessionID: current.id,
|
||||
command: known.name,
|
||||
arguments: command.args,
|
||||
model: `${selected.providerID}/${selected.modelID}`,
|
||||
...(variant ? { variant } : {}),
|
||||
...(modeId ? { agent: modeId } : {}),
|
||||
directory: current.cwd,
|
||||
},
|
||||
{ throwOnError: true },
|
||||
runUntilIdle(current.id, () =>
|
||||
input.sdk.session.command(
|
||||
{
|
||||
sessionID: current.id,
|
||||
command: known.name,
|
||||
arguments: command.args,
|
||||
model: `${selected.providerID}/${selected.modelID}`,
|
||||
...(variant ? { variant } : {}),
|
||||
...(modeId ? { agent: modeId } : {}),
|
||||
directory: current.cwd,
|
||||
},
|
||||
{ throwOnError: true },
|
||||
),
|
||||
),
|
||||
"session",
|
||||
)
|
||||
@@ -549,14 +555,16 @@ export function make(input: {
|
||||
if (command.name === "compact") {
|
||||
yield* request(
|
||||
() =>
|
||||
input.sdk.session.summarize(
|
||||
{
|
||||
sessionID: current.id,
|
||||
directory: current.cwd,
|
||||
providerID: selected.providerID,
|
||||
modelID: selected.modelID,
|
||||
},
|
||||
{ throwOnError: true },
|
||||
runUntilIdle(current.id, () =>
|
||||
input.sdk.session.summarize(
|
||||
{
|
||||
sessionID: current.id,
|
||||
directory: current.cwd,
|
||||
providerID: selected.providerID,
|
||||
modelID: selected.modelID,
|
||||
},
|
||||
{ throwOnError: true },
|
||||
),
|
||||
),
|
||||
"session",
|
||||
)
|
||||
|
||||
@@ -97,6 +97,29 @@ export function http(
|
||||
headers.delete("content-encoding")
|
||||
headers.delete("content-length")
|
||||
|
||||
// An upstream 5xx from a remote workspace sandbox arrives here as an opaque
|
||||
// status — its real cause (and log line) live only inside the sandbox. Buffer
|
||||
// the small error body, log it locally so it shows up in the host's log, and
|
||||
// forward it unchanged (preserving content-type so the client can still parse
|
||||
// the structured error, e.g. its `ref`).
|
||||
if (response.status >= 500) {
|
||||
const body = yield* response.text.pipe(Effect.catch(() => Effect.succeed("")))
|
||||
const contentType = response.headers["content-type"] ?? "application/json"
|
||||
headers.delete("content-type")
|
||||
yield* Effect.logError("workspace proxy upstream error", {
|
||||
url: url.toString(),
|
||||
method: request.method,
|
||||
status: response.status,
|
||||
body: body.slice(0, 2000),
|
||||
})
|
||||
return HttpServerResponse.text(body, {
|
||||
status: response.status,
|
||||
statusText: statusText(response),
|
||||
headers,
|
||||
contentType,
|
||||
})
|
||||
}
|
||||
|
||||
return HttpServerResponse.stream(response.stream.pipe(Stream.catchCause(() => Stream.empty)), {
|
||||
status: response.status,
|
||||
statusText: statusText(response),
|
||||
|
||||
@@ -34,5 +34,12 @@ export function workspaceProxyURL(target: string | URL, requestURL: URL) {
|
||||
proxyURL.search = requestURL.search
|
||||
proxyURL.hash = requestURL.hash
|
||||
proxyURL.searchParams.delete("workspace")
|
||||
// The `directory` param is the *host's* working directory (e.g. a Windows
|
||||
// path like `F:\proj`). It is meaningless — and dangerous — on the remote:
|
||||
// the sandbox would `path.resolve` it against its own cwd, producing a bogus
|
||||
// path like `/home/daytona/workspace/repo/F:\proj` that does not exist and
|
||||
// crashes prompt handling. Drop it so the remote falls back to its own
|
||||
// project root. This mirrors ProxyUtil.headers stripping `x-opencode-directory`.
|
||||
proxyURL.searchParams.delete("directory")
|
||||
return proxyURL
|
||||
}
|
||||
|
||||
@@ -636,14 +636,30 @@ 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) => handleEvent(event)),
|
||||
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.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 } from "effect"
|
||||
import { Cause, Clock, Duration, Effect, Schedule, Schema } from "effect"
|
||||
import { MessageV2 } from "./message-v2"
|
||||
import { iife } from "@/util/iife"
|
||||
import { isRecord } from "@/util/record"
|
||||
@@ -23,6 +23,10 @@ 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
|
||||
@@ -181,7 +185,8 @@ export function policy(opts: {
|
||||
return Schedule.fromStepWithMetadata(
|
||||
Effect.succeed((meta: Schedule.InputMetadata<unknown>) => {
|
||||
const error = opts.parse(meta.input)
|
||||
const retry = retryable(error, opts.provider)
|
||||
const retry =
|
||||
meta.input instanceof EmptyResponseError ? { message: meta.input.message } : 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)
|
||||
|
||||
@@ -10,7 +10,7 @@ import type {
|
||||
SessionConfigSelectOption,
|
||||
SetSessionConfigOptionResponse,
|
||||
} from "@agentclientprotocol/sdk"
|
||||
import type { AssistantMessage, OpencodeClient } from "@opencode-ai/sdk/v2"
|
||||
import type { AssistantMessage, Event, OpencodeClient } from "@opencode-ai/sdk/v2"
|
||||
import { ProviderV2 } from "@opencode-ai/core/provider"
|
||||
import { ModelV2 } from "@opencode-ai/core/model"
|
||||
import { Effect } from "effect"
|
||||
@@ -24,6 +24,54 @@ const modelID = ModelV2.ID.make("test-model")
|
||||
const configuredModelID = ModelV2.ID.make("configured-model")
|
||||
const secondModelID = ModelV2.ID.make("second-model")
|
||||
|
||||
function createEventStream() {
|
||||
const queue: Event[] = []
|
||||
const waiters: Array<(event: Event | undefined) => void> = []
|
||||
const push = (event: Event) => {
|
||||
const waiter = waiters.shift()
|
||||
if (waiter) return waiter(event)
|
||||
queue.push(event)
|
||||
}
|
||||
const stream = async function* (signal?: AbortSignal) {
|
||||
while (!signal?.aborted) {
|
||||
const event = queue.shift()
|
||||
if (event) {
|
||||
yield { payload: event }
|
||||
continue
|
||||
}
|
||||
const next = await new Promise<Event | undefined>((resolve) => {
|
||||
waiters.push(resolve)
|
||||
signal?.addEventListener("abort", () => resolve(undefined), { once: true })
|
||||
})
|
||||
if (!next) return
|
||||
yield { payload: next }
|
||||
}
|
||||
}
|
||||
return { push, stream }
|
||||
}
|
||||
|
||||
function idleEvent(sessionID: string): Event {
|
||||
return {
|
||||
id: `evt_idle_${sessionID}`,
|
||||
type: "session.status",
|
||||
properties: {
|
||||
sessionID,
|
||||
status: { type: "idle" },
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
function deferred<A>() {
|
||||
const state: { resolve?: (value: A) => void } = {}
|
||||
const promise = new Promise<A>((resolve) => {
|
||||
state.resolve = resolve
|
||||
})
|
||||
return {
|
||||
promise,
|
||||
resolve: (value: A) => state.resolve?.(value),
|
||||
}
|
||||
}
|
||||
|
||||
const provider: Provider.Info = {
|
||||
id: providerID,
|
||||
name: "Test",
|
||||
@@ -147,6 +195,7 @@ describe("ACP service sessions", () => {
|
||||
options?: {
|
||||
abort?: (input: { sessionID: string }) => Promise<{ data: boolean }>
|
||||
prompt?: (input: unknown) => Promise<{ data: { info: ReturnType<typeof assistantInfo> } }>
|
||||
sessionUpdate?: (update: SessionNotification) => Promise<void>
|
||||
},
|
||||
) => {
|
||||
const updates: SessionNotification[] = []
|
||||
@@ -157,6 +206,7 @@ describe("ACP service sessions", () => {
|
||||
const commands: unknown[] = []
|
||||
const summarizes: unknown[] = []
|
||||
const usageUpdates: string[] = []
|
||||
const events = createEventStream()
|
||||
const sessions = Array.from({ length: 102 }, (_, index) => ({
|
||||
id: `ses_${index + 1}`,
|
||||
directory: index % 2 === 0 ? "/workspace" : "/other",
|
||||
@@ -164,6 +214,9 @@ describe("ACP service sessions", () => {
|
||||
time: { created: index + 1, updated: index + 1 },
|
||||
}))
|
||||
const sdk = {
|
||||
global: {
|
||||
event: (input?: { signal?: AbortSignal }) => Promise.resolve({ stream: events.stream(input?.signal) }),
|
||||
},
|
||||
config: {
|
||||
providers: () => Promise.resolve({ data: { providers: [provider], default: { test: modelID } } }),
|
||||
get: () => Promise.resolve({ data: {} }),
|
||||
@@ -196,11 +249,9 @@ describe("ACP service sessions", () => {
|
||||
data: input.directory ? sessions.filter((session) => session.directory === input.directory) : sessions,
|
||||
}),
|
||||
messages: () => Promise.resolve({ data: messages }),
|
||||
prompt:
|
||||
options?.prompt ??
|
||||
((input: unknown) => {
|
||||
prompts.push(input)
|
||||
return Promise.resolve({
|
||||
prompt: async (input: { sessionID: string }) => {
|
||||
const response = await (options?.prompt?.(input) ??
|
||||
Promise.resolve({
|
||||
data: {
|
||||
info: assistantInfo({
|
||||
input: 100,
|
||||
@@ -209,10 +260,14 @@ describe("ACP service sessions", () => {
|
||||
cache: { read: 11, write: 13 },
|
||||
}),
|
||||
},
|
||||
})
|
||||
}),
|
||||
command: (input: unknown) => {
|
||||
}))
|
||||
prompts.push(input)
|
||||
events.push(idleEvent(input.sessionID))
|
||||
return response
|
||||
},
|
||||
command: (input: { sessionID: string }) => {
|
||||
commands.push(input)
|
||||
events.push(idleEvent(input.sessionID))
|
||||
return Promise.resolve({
|
||||
data: {
|
||||
info: assistantInfo({
|
||||
@@ -224,8 +279,9 @@ describe("ACP service sessions", () => {
|
||||
},
|
||||
})
|
||||
},
|
||||
summarize: (input: unknown) => {
|
||||
summarize: (input: { sessionID: string }) => {
|
||||
summarizes.push(input)
|
||||
events.push(idleEvent(input.sessionID))
|
||||
return Promise.resolve({ data: true })
|
||||
},
|
||||
abort:
|
||||
@@ -249,7 +305,7 @@ describe("ACP service sessions", () => {
|
||||
const connection = {
|
||||
sessionUpdate: (update: SessionNotification) => {
|
||||
updates.push(update)
|
||||
return Promise.resolve()
|
||||
return options?.sessionUpdate?.(update) ?? Promise.resolve()
|
||||
},
|
||||
} as Pick<AgentSideConnection, "sessionUpdate">
|
||||
const usage = UsageService.Service.of({
|
||||
@@ -273,6 +329,7 @@ describe("ACP service sessions", () => {
|
||||
commands,
|
||||
summarizes,
|
||||
usageUpdates,
|
||||
events,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1018,6 +1075,75 @@ describe("ACP service sessions", () => {
|
||||
expect(usageUpdates).toEqual([session.sessionId])
|
||||
})
|
||||
|
||||
it("waits for queued session updates before returning end_turn", async () => {
|
||||
const called = deferred<void>()
|
||||
const response = deferred<{ data: { info: ReturnType<typeof assistantInfo> } }>()
|
||||
const update = deferred<void>()
|
||||
const release = deferred<void>()
|
||||
const order: string[] = []
|
||||
const fixture = makeService([], {
|
||||
prompt: () => {
|
||||
called.resolve(undefined)
|
||||
return response.promise
|
||||
},
|
||||
sessionUpdate: (notification) => {
|
||||
if (notification.update.sessionUpdate !== "agent_thought_chunk") return Promise.resolve()
|
||||
update.resolve(undefined)
|
||||
return release.promise.then(() => {
|
||||
order.push("update")
|
||||
})
|
||||
},
|
||||
})
|
||||
const session = await Effect.runPromise(fixture.service.newSession({ cwd: "/workspace", mcpServers: [] }))
|
||||
const result = Effect.runPromise(
|
||||
fixture.service.prompt({ sessionId: session.sessionId, prompt: [{ type: "text", text: "hello" }] }),
|
||||
).then((value) => {
|
||||
order.push("response")
|
||||
return value
|
||||
})
|
||||
|
||||
await called.promise
|
||||
fixture.events.push({
|
||||
id: "evt_part",
|
||||
type: "message.part.updated",
|
||||
properties: {
|
||||
sessionID: session.sessionId,
|
||||
time: Date.now(),
|
||||
part: {
|
||||
id: "part_reasoning",
|
||||
sessionID: session.sessionId,
|
||||
messageID: "msg_assistant",
|
||||
type: "reasoning",
|
||||
text: "",
|
||||
time: { start: Date.now() },
|
||||
},
|
||||
},
|
||||
})
|
||||
fixture.events.push({
|
||||
id: "evt_delta",
|
||||
type: "message.part.delta",
|
||||
properties: {
|
||||
sessionID: session.sessionId,
|
||||
messageID: "msg_assistant",
|
||||
partID: "part_reasoning",
|
||||
field: "text",
|
||||
delta: "thinking",
|
||||
},
|
||||
})
|
||||
response.resolve({
|
||||
data: {
|
||||
info: assistantInfo({ input: 1, output: 1, reasoning: 1, cache: { read: 0, write: 0 } }),
|
||||
},
|
||||
})
|
||||
|
||||
await update.promise
|
||||
expect(order).toEqual([])
|
||||
|
||||
release.resolve(undefined)
|
||||
expect((await result).stopReason).toBe("end_turn")
|
||||
expect(order).toEqual(["update", "response"])
|
||||
})
|
||||
|
||||
it("maps assistant prompt errors to request errors instead of end turn", async () => {
|
||||
const { service } = makeService([], {
|
||||
prompt: () =>
|
||||
|
||||
@@ -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 unknown
|
||||
// finish, not a fatal provider/session error. Lock that distinction in so it
|
||||
// is not accidentally used as the failure compatibility oracle.
|
||||
// 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.
|
||||
cliIt.concurrent(
|
||||
"unknown stream finish preserves partial output and exits 0",
|
||||
"empty unknown stream finish retries and preserves partial output",
|
||||
({ llm, opencode }) =>
|
||||
Effect.gen(function* () {
|
||||
yield* llm.push(
|
||||
@@ -95,9 +95,10 @@ 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\n")
|
||||
expect(result.stdout).toBe("partial response\nrecovered response\n")
|
||||
expect(result.stderr).not.toContain("upstream provider exploded mid-stream")
|
||||
}),
|
||||
60_000,
|
||||
@@ -213,7 +214,7 @@ describe("opencode run (non-interactive subprocess)", () => {
|
||||
)
|
||||
|
||||
cliIt.concurrent(
|
||||
"--format json records partial output for an unknown stream finish",
|
||||
"--format json records an empty unknown stream retry",
|
||||
({ llm, opencode }) =>
|
||||
Effect.gen(function* () {
|
||||
yield* llm.push(
|
||||
@@ -223,6 +224,7 @@ 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)
|
||||
@@ -234,9 +236,13 @@ 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(-1)?.part).toEqual(expect.objectContaining({ type: "step-finish", reason: "unknown" }))
|
||||
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" }))
|
||||
}),
|
||||
60_000,
|
||||
)
|
||||
|
||||
@@ -80,6 +80,13 @@ describe("workspaceProxyURL", () => {
|
||||
expect(result.searchParams.get("keep")).toBe("yes")
|
||||
})
|
||||
|
||||
test("strips the host directory param so the remote resolves its own root", () => {
|
||||
const url = new URL("http://localhost/session/abc?directory=F%3A%5Cproj&keep=yes")
|
||||
const result = workspaceProxyURL("http://remote:8080/base", url)
|
||||
expect(result.searchParams.get("directory")).toBeNull()
|
||||
expect(result.searchParams.get("keep")).toBe("yes")
|
||||
})
|
||||
|
||||
test("preserves hash from request", () => {
|
||||
const url = new URL("http://localhost/page#section")
|
||||
const result = workspaceProxyURL("http://remote:8080", url)
|
||||
|
||||
@@ -604,6 +604,68 @@ 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