Compare commits

..

1 Commits

Author SHA1 Message Date
Kit Langton 8272d68335 fix(console): harden zen stream relay diagnostics
Add stream lifecycle metrics and extract the Zen relay loop so missing bodies, decoder stalls, EOF framing, and postprocess failures are easier to diagnose. Cover the new relay behavior with focused repro tests for the streaming edge cases we suspect in production.
2026-03-26 11:24:56 -04:00
9 changed files with 564 additions and 77 deletions
@@ -39,6 +39,7 @@ import { createRateLimiter } from "./rateLimiter"
import { createDataDumper } from "./dataDumper"
import { createTrialLimiter } from "./trialLimiter"
import { createStickyTracker } from "./stickyProviderTracker"
import { relay } from "./stream"
import { LiteData } from "@opencode-ai/console-core/lite.js"
import { Resource } from "@opencode-ai/console-resource"
import { i18n, type Key } from "~/i18n"
@@ -254,80 +255,41 @@ export async function handler(
const streamConverter = createStreamPartConverter(providerInfo.format, opts.format)
const usageParser = providerInfo.createUsageParser()
const binaryDecoder = providerInfo.createBinaryStreamDecoder()
const stream = new ReadableStream({
start(c) {
const reader = res.body?.getReader()
const decoder = new TextDecoder()
const encoder = new TextEncoder()
let buffer = ""
let responseLength = 0
function pump(): Promise<void> {
return (
reader?.read().then(async ({ done, value: rawValue }) => {
if (done) {
logger.metric({
response_length: responseLength,
"timestamp.last_byte": Date.now(),
})
dataDumper?.flush()
await rateLimiter?.track()
const usage = usageParser.retrieve()
if (usage) {
const usageInfo = providerInfo.normalizeUsage(usage)
const costInfo = calculateCost(modelInfo, usageInfo)
await trialLimiter?.track(usageInfo)
await trackUsage(sessionId, billingSource, authInfo, modelInfo, providerInfo, usageInfo, costInfo)
await reload(billingSource, authInfo, costInfo)
const cost = calculateOccurredCost(billingSource, costInfo)
c.enqueue(encoder.encode(buildCostChunk(opts.format, cost)))
}
c.close()
return
}
if (responseLength === 0) {
const now = Date.now()
logger.metric({
time_to_first_byte: now - startTimestamp,
"timestamp.first_byte": now,
})
}
const value = binaryDecoder ? binaryDecoder(rawValue) : rawValue
if (!value) return
responseLength += value.length
buffer += decoder.decode(value, { stream: true })
dataDumper?.provideStream(buffer)
const parts = buffer.split(providerInfo.streamSeparator)
buffer = parts.pop() ?? ""
for (let part of parts) {
logger.debug("PART: " + part)
part = part.trim()
usageParser.parse(part)
if (providerInfo.format !== opts.format) {
part = streamConverter(part)
c.enqueue(encoder.encode(part + "\n\n"))
}
}
if (providerInfo.format === opts.format) {
c.enqueue(value)
}
return pump()
}) || Promise.resolve()
)
}
return pump()
const metric = (values: Record<string, unknown>) =>
logger.metric({
request: requestId,
session: sessionId,
client: ocClient,
provider: providerInfo.id,
model: modelInfo.id,
...values,
})
const stream = relay({
body: res.body,
separator: providerInfo.streamSeparator,
signal: input.request.signal,
start: startTimestamp,
same: providerInfo.format === opts.format,
binary: binaryDecoder,
parse: (part) => {
logger.debug("PART: " + part)
usageParser.parse(part)
},
convert: streamConverter,
tail: async () => {
await rateLimiter?.track()
const usage = usageParser.retrieve()
if (!usage) return
const usageInfo = providerInfo.normalizeUsage(usage)
const costInfo = calculateCost(modelInfo, usageInfo)
await trialLimiter?.track(usageInfo)
await trackUsage(sessionId, billingSource, authInfo, modelInfo, providerInfo, usageInfo, costInfo)
await reload(billingSource, authInfo, costInfo)
const cost = calculateOccurredCost(billingSource, costInfo)
return buildCostChunk(opts.format, cost)
},
metric,
dump: dataDumper,
})
return new Response(stream, {
status: resStatus,
@@ -0,0 +1,196 @@
const done = new Set(["[DONE]", "message_stop", "response.completed"])
type Dump = {
provideStream: (chunk: string) => void
flush: () => void
}
type Opts = {
body: ReadableStream<Uint8Array> | null | undefined
separator: string
signal: AbortSignal
start: number
same: boolean
binary?: (chunk: Uint8Array) => Uint8Array | undefined
parse: (part: string) => void
convert: (part: string) => string
tail: () => Promise<string | undefined>
metric: (values: Record<string, unknown>) => void
dump?: Dump
}
export const eventName = (part: string) => {
const line = part.split("\n", 1)[0]?.trim() ?? ""
if (line.startsWith("event:")) return line.slice(6).trim() || "message"
if (part.includes("[DONE]")) return "[DONE]"
if (line.startsWith("data:")) return "message"
return "unknown"
}
const errInfo = (err: unknown) => {
if (err instanceof Error) {
return {
"stream.error_type": err.constructor.name,
"stream.error_message": err.message,
}
}
return {
"stream.error_type": typeof err,
"stream.error_message": String(err),
}
}
const stats = (len: number, cnt: number, seen: number, buf: string, end: string | undefined, gap: number) => ({
"stream.response_length": len,
"stream.chunk_count": cnt,
"stream.event_count": seen,
"stream.pending_length": buf.length,
"stream.last_event": end,
"stream.max_gap_ms": gap || undefined,
})
export const relay = (opts: Opts) =>
new ReadableStream({
async start(c) {
let phase = "start"
let len = 0
let cnt = 0
let seen = 0
let end: string | undefined
let gap = 0
let prev: number | undefined
let completed = false
let aborted = false
let buf = ""
if (!opts.body) {
opts.metric({
"stream.event": "missing_body",
"stream.phase": phase,
})
opts.dump?.flush()
c.close()
return
}
const reader = opts.body.getReader()
const dec = new TextDecoder()
const enc = new TextEncoder()
const names: Record<string, number> = {}
const abort = () => {
aborted = true
reader.cancel().catch(() => undefined)
}
const note = (part: string) => {
const name = eventName(part)
end = name
seen += 1
names[name] = (names[name] ?? 0) + 1
if (done.has(name)) completed = true
}
opts.signal.addEventListener("abort", abort)
opts.metric({
"stream.event": "started",
})
try {
while (true) {
phase = "read"
const raw = await reader.read()
if (raw.done) break
if (len === 0) {
const now = Date.now()
opts.metric({
time_to_first_byte: now - opts.start,
"timestamp.first_byte": now,
})
}
const value = opts.binary ? opts.binary(raw.value) : raw.value
if (!value) continue
cnt += 1
len += value.length
const now = Date.now()
if (prev !== undefined) gap = Math.max(gap, now - prev)
prev = now
const text = dec.decode(value, { stream: true })
buf += text
opts.dump?.provideStream(text)
const parts = buf.split(opts.separator)
buf = parts.pop() ?? ""
for (let part of parts) {
part = part.trim()
if (!part) continue
note(part)
phase = "parse"
opts.parse(part)
if (opts.same) continue
phase = "convert"
c.enqueue(enc.encode(opts.convert(part) + "\n\n"))
}
if (opts.same) c.enqueue(value)
}
const tail = dec.decode()
if (tail) {
buf += tail
opts.dump?.provideStream(tail)
}
if (buf.trim()) {
const part = buf.trim()
note(part)
phase = "parse"
opts.parse(part)
if (!opts.same) {
phase = "convert"
c.enqueue(enc.encode(opts.convert(part) + "\n\n"))
}
buf = ""
}
opts.metric({
response_length: len,
"timestamp.last_byte": Date.now(),
})
opts.dump?.flush()
phase = "tail"
const chunk = await opts.tail()
if (chunk) c.enqueue(enc.encode(chunk))
c.close()
opts.metric({
"stream.event": "finished",
"stream.phase": "done",
"stream.duration_ms": Date.now() - opts.start,
"stream.saw_completed": completed,
"stream.events": JSON.stringify(names),
...stats(len, cnt, seen, buf, end, gap),
})
} catch (err) {
opts.metric({
"stream.event": aborted ? "aborted" : "error",
"stream.phase": phase,
"stream.duration_ms": Date.now() - opts.start,
"stream.saw_completed": completed,
"stream.events": JSON.stringify(names),
...stats(len, cnt, seen, buf, end, gap),
...errInfo(err),
})
c.error(err)
} finally {
opts.signal.removeEventListener("abort", abort)
}
},
})
+139
View File
@@ -0,0 +1,139 @@
import { describe, expect, test } from "bun:test"
import { eventName, relay } from "../src/routes/zen/util/stream"
const enc = new TextEncoder()
const read = (stream: ReadableStream<Uint8Array>) => new Response(stream).text()
const body = (parts: string[]) =>
new ReadableStream<Uint8Array>({
async start(c) {
for (const part of parts) c.enqueue(enc.encode(part))
c.close()
},
})
describe("zen stream", () => {
test("parses known event names", () => {
expect(eventName("event: response.created\ndata: {}")).toBe("response.created")
expect(eventName('data: {"ok":true}')).toBe("message")
expect(eventName("data: [DONE]")).toBe("[DONE]")
})
test("relays split OpenAI responses and logs completion", async () => {
const seen: string[] = []
const logs: Array<Record<string, unknown>> = []
const stream = relay({
body: body([
"event: response.created\n",
'data: {"type":"response.created"}\n\n',
"event: response.completed\n",
'data: {"response":{"usage":{"input_tokens":1,"output_tokens":2}}}\n\n',
]),
separator: "\n\n",
signal: new AbortController().signal,
start: Date.now(),
same: true,
parse: (part) => {
seen.push(part)
},
convert: (part) => part,
tail: async () => undefined,
metric: (values) => logs.push(values),
})
const text = await read(stream)
expect(text).toContain("response.created")
expect(text).toContain("response.completed")
expect(seen).toHaveLength(2)
expect(logs.at(-1)?.["stream.event"]).toBe("finished")
expect(logs.at(-1)?.["stream.saw_completed"]).toBe(true)
})
test("keeps reading when binary decoder needs another chunk", async () => {
let calls = 0
const logs: Array<Record<string, unknown>> = []
const stream = relay({
body: body(["a", "b"]),
separator: "\n\n",
signal: new AbortController().signal,
start: Date.now(),
same: true,
binary: (chunk) => {
calls += 1
if (calls === 1) return
return chunk
},
parse: () => undefined,
convert: (part) => part,
tail: async () => undefined,
metric: (values) => logs.push(values),
})
const text = await read(stream)
expect(text).toBe("b")
expect(logs.at(-1)?.["stream.event"]).toBe("finished")
})
test("flushes a final unterminated event at EOF", async () => {
const seen: string[] = []
const logs: Array<Record<string, unknown>> = []
const stream = relay({
body: body(['event: response.completed\ndata: {"response":{"usage":{"input_tokens":1}}}']),
separator: "\n\n",
signal: new AbortController().signal,
start: Date.now(),
same: true,
parse: (part) => {
seen.push(part)
},
convert: (part) => part,
tail: async () => undefined,
metric: (values) => logs.push(values),
})
const text = await read(stream)
expect(text).toContain("response.completed")
expect(seen).toHaveLength(1)
expect(logs.at(-1)?.["stream.saw_completed"]).toBe(true)
})
test("closes cleanly when upstream body is missing", async () => {
const logs: Array<Record<string, unknown>> = []
const stream = relay({
body: null,
separator: "\n\n",
signal: new AbortController().signal,
start: Date.now(),
same: true,
parse: () => undefined,
convert: (part) => part,
tail: async () => undefined,
metric: (values) => logs.push(values),
})
expect(await read(stream)).toBe("")
expect(logs.at(-1)?.["stream.event"]).toBe("missing_body")
})
test("surfaces postprocess failures with stream metrics", async () => {
const logs: Array<Record<string, unknown>> = []
const stream = relay({
body: body(["event: response.created\ndata: {}\n\n"]),
separator: "\n\n",
signal: new AbortController().signal,
start: Date.now(),
same: true,
parse: () => undefined,
convert: (part) => part,
tail: async () => {
throw new Error("boom")
},
metric: (values) => logs.push(values),
})
await expect(read(stream)).rejects.toThrow("boom")
expect(logs.at(-1)?.["stream.event"]).toBe("error")
expect(logs.at(-1)?.["stream.phase"]).toBe("tail")
})
})
@@ -0,0 +1,16 @@
import { cmd } from "./cmd"
import { withNetworkOptions, resolveNetworkOptions } from "../network"
import { WorkspaceServer } from "../../control-plane/workspace-server/server"
export const WorkspaceServeCommand = cmd({
command: "workspace-serve",
builder: (yargs) => withNetworkOptions(yargs),
describe: "starts a remote workspace event server",
handler: async (args) => {
const opts = await resolveNetworkOptions(args)
const server = WorkspaceServer.Listen(opts)
console.log(`workspace event server listening on http://${server.hostname}:${server.port}/event`)
await new Promise(() => {})
await server.stop()
},
})
@@ -2,8 +2,6 @@ import z from "zod"
import { Worktree } from "@/worktree"
import { type Adaptor, WorkspaceInfo } from "../types"
import { Server } from "../../server/server"
const Config = WorkspaceInfo.extend({
name: WorkspaceInfo.shape.name.unwrap(),
branch: WorkspaceInfo.shape.branch.unwrap(),
@@ -36,11 +34,12 @@ export const WorktreeAdaptor: Adaptor = {
},
async fetch(info, input: RequestInfo | URL, init?: RequestInit) {
const config = Config.parse(info)
const { WorkspaceServer } = await import("../workspace-server/server")
const url = input instanceof Request || input instanceof URL ? input : new URL(input, "http://opencode.internal")
const headers = new Headers(init?.headers ?? (input instanceof Request ? input.headers : undefined))
headers.set("x-opencode-directory", config.directory)
const request = new Request(url, { ...init, headers })
return Server.Default().fetch(request)
return WorkspaceServer.App().fetch(request)
},
}
@@ -0,0 +1,33 @@
import { GlobalBus } from "../../bus/global"
import { Hono } from "hono"
import { streamSSE } from "hono/streaming"
export function WorkspaceServerRoutes() {
return new Hono().get("/event", async (c) => {
c.header("X-Accel-Buffering", "no")
c.header("X-Content-Type-Options", "nosniff")
return streamSSE(c, async (stream) => {
const send = async (event: unknown) => {
await stream.writeSSE({
data: JSON.stringify(event),
})
}
const handler = async (event: { directory?: string; payload: unknown }) => {
await send(event.payload)
}
GlobalBus.on("event", handler)
await send({ type: "server.connected", properties: {} })
const heartbeat = setInterval(() => {
void send({ type: "server.heartbeat", properties: {} })
}, 10_000)
await new Promise<void>((resolve) => {
stream.onAbort(() => {
clearInterval(heartbeat)
GlobalBus.off("event", handler)
resolve()
})
})
})
})
}
@@ -0,0 +1,65 @@
import { Hono } from "hono"
import { Instance } from "../../project/instance"
import { InstanceBootstrap } from "../../project/bootstrap"
import { SessionRoutes } from "../../server/routes/session"
import { WorkspaceServerRoutes } from "./routes"
import { WorkspaceContext } from "../workspace-context"
import { WorkspaceID } from "../schema"
export namespace WorkspaceServer {
export function App() {
const session = new Hono()
.use(async (c, next) => {
// Right now, we need handle all requests because we don't
// have syncing. In the future all GET requests will handled
// by the control plane
//
// if (c.req.method === "GET") return c.notFound()
await next()
})
.route("/", SessionRoutes())
return new Hono()
.use(async (c, next) => {
const rawWorkspaceID = c.req.query("workspace") || c.req.header("x-opencode-workspace")
const raw = c.req.query("directory") || c.req.header("x-opencode-directory")
if (rawWorkspaceID == null) {
throw new Error("workspaceID parameter is required")
}
if (raw == null) {
throw new Error("directory parameter is required")
}
const directory = (() => {
try {
return decodeURIComponent(raw)
} catch {
return raw
}
})()
return WorkspaceContext.provide({
workspaceID: WorkspaceID.make(rawWorkspaceID),
async fn() {
return Instance.provide({
directory,
init: InstanceBootstrap,
async fn() {
return next()
},
})
},
})
})
.route("/session", session)
.route("/", WorkspaceServerRoutes())
}
export function Listen(opts: { hostname: string; port: number }) {
return Bun.serve({
hostname: opts.hostname,
port: opts.port,
fetch: App().fetch,
})
}
}
+8 -1
View File
@@ -14,6 +14,7 @@ import { Installation } from "./installation"
import { NamedError } from "@opencode-ai/util/error"
import { FormatError } from "./cli/error"
import { ServeCommand } from "./cli/cmd/serve"
import { WorkspaceServeCommand } from "./cli/cmd/workspace-serve"
import { Filesystem } from "./util/filesystem"
import { DebugCommand } from "./cli/cmd/debug"
import { StatsCommand } from "./cli/cmd/stats"
@@ -46,7 +47,7 @@ process.on("uncaughtException", (e) => {
})
})
const cli = yargs(hideBin(process.argv))
let cli = yargs(hideBin(process.argv))
.parserConfiguration({ "populate--": true })
.scriptName("opencode")
.wrap(100)
@@ -144,6 +145,12 @@ const cli = yargs(hideBin(process.argv))
.command(PrCommand)
.command(SessionCommand)
.command(DbCommand)
if (Installation.isLocal()) {
cli = cli.command(WorkspaceServeCommand)
}
cli = cli
.fail((msg, err) => {
if (
msg?.startsWith("Unknown argument") ||
@@ -0,0 +1,70 @@
import { afterEach, describe, expect, test } from "bun:test"
import { Log } from "../../src/util/log"
import { WorkspaceServer } from "../../src/control-plane/workspace-server/server"
import { parseSSE } from "../../src/control-plane/sse"
import { GlobalBus } from "../../src/bus/global"
import { resetDatabase } from "../fixture/db"
import { tmpdir } from "../fixture/fixture"
afterEach(async () => {
await resetDatabase()
})
Log.init({ print: false })
describe("control-plane/workspace-server SSE", () => {
test("streams GlobalBus events and parseSSE reads them", async () => {
await using tmp = await tmpdir({ git: true })
const app = WorkspaceServer.App()
const stop = new AbortController()
const seen: unknown[] = []
try {
const response = await app.request("/event", {
signal: stop.signal,
headers: {
"x-opencode-workspace": "wrk_test_workspace",
"x-opencode-directory": tmp.path,
},
})
expect(response.status).toBe(200)
expect(response.body).toBeDefined()
const done = new Promise<void>((resolve, reject) => {
const timeout = setTimeout(() => {
reject(new Error("timed out waiting for workspace.test event"))
}, 3000)
void parseSSE(response.body!, stop.signal, (event) => {
seen.push(event)
const next = event as { type?: string }
if (next.type === "server.connected") {
GlobalBus.emit("event", {
payload: {
type: "workspace.test",
properties: { ok: true },
},
})
return
}
if (next.type !== "workspace.test") return
clearTimeout(timeout)
resolve()
}).catch((error) => {
clearTimeout(timeout)
reject(error)
})
})
await done
expect(seen.some((event) => (event as { type?: string }).type === "server.connected")).toBe(true)
expect(seen).toContainEqual({
type: "workspace.test",
properties: { ok: true },
})
} finally {
stop.abort()
}
})
})