From de7ed3639d04badee5268c359a662456cad1459b Mon Sep 17 00:00:00 2001 From: Code Queue Review Date: Sat, 23 May 2026 00:27:25 +0000 Subject: [PATCH] fix: stream code agent OpenAI responses --- internal/cloud/code-agent-chat.mjs | 106 ++++++++++++++++++++++++++-- internal/cloud/server.test.mjs | 108 +++++++++++++++++++++++++++++ 2 files changed, 209 insertions(+), 5 deletions(-) diff --git a/internal/cloud/code-agent-chat.mjs b/internal/cloud/code-agent-chat.mjs index 3fd75d2f..45e7ce9b 100644 --- a/internal/cloud/code-agent-chat.mjs +++ b/internal/cloud/code-agent-chat.mjs @@ -258,6 +258,7 @@ async function callOpenAiResponses({ providerPlan, message, conversationId, trac method: "POST", headers: { "content-type": "application/json", + accept: "text/event-stream", authorization: `Bearer ${env.OPENAI_API_KEY}` }, body: JSON.stringify({ @@ -274,11 +275,12 @@ async function callOpenAiResponses({ providerPlan, message, conversationId, trac ] } ], - store: false + store: false, + stream: true }), signal: controller.signal }); - payload = await response.json().catch(() => null); + payload = parseOpenAiResponseBody(await response.text()); } catch (error) { if (error.name === "AbortError") { throw providerUnavailable(`OpenAI Responses request timed out after ${effectiveTimeout(timeoutMs)}ms`, { @@ -296,8 +298,9 @@ async function callOpenAiResponses({ providerPlan, message, conversationId, trac clearTimeout(timer); } + const providerError = payload.error ?? payload.raw?.error ?? null; if (!response.ok) { - throw providerUnavailable(`OpenAI Responses returned HTTP ${response.status}: ${payload?.error?.message || "request rejected"}`, { + throw providerUnavailable(`OpenAI Responses returned HTTP ${response.status}: ${providerError?.message || "request rejected"}`, { provider: providerPlan.provider, model: providerPlan.model, backend: providerPlan.backend, @@ -305,7 +308,15 @@ async function callOpenAiResponses({ providerPlan, message, conversationId, trac }); } - const content = extractOpenAiOutputText(payload); + if (providerError) { + throw providerUnavailable(`OpenAI Responses stream error: ${providerError.message || "request rejected"}`, { + provider: providerPlan.provider, + model: providerPlan.model, + backend: providerPlan.backend + }); + } + + const content = payload.content || extractOpenAiOutputText(payload.raw); if (!content) { throw providerUnavailable("OpenAI Responses returned no assistant text", { provider: providerPlan.provider, @@ -321,7 +332,7 @@ async function callOpenAiResponses({ providerPlan, message, conversationId, trac content, usage: payload.usage ?? null, providerTrace: { - responseId: payload.id ?? null + responseId: payload.responseId ?? null } }; } @@ -477,6 +488,91 @@ function extractOpenAiOutputText(payload) { return chunks.join("\n").trim(); } +function parseOpenAiResponseBody(text) { + const raw = parseJsonOrNull(text); + if (raw) { + return { + raw, + content: extractOpenAiOutputText(raw), + responseId: raw.id ?? null, + model: raw.model ?? null, + usage: raw.usage ?? null, + error: raw.error ?? null + }; + } + + const stream = parseOpenAiResponsesSse(text); + return { + raw: stream.finalResponse, + content: stream.content, + responseId: stream.responseId, + model: stream.model, + usage: stream.usage, + error: stream.error + }; +} + +function parseOpenAiResponsesSse(text) { + const deltas = []; + let finalResponse = null; + let responseId = null; + let model = null; + let usage = null; + let error = null; + + for (const event of parseSseDataMessages(text)) { + const payload = parseJsonOrNull(event); + if (!payload) continue; + const response = payload.response && typeof payload.response === "object" ? payload.response : null; + if (response) { + finalResponse = response; + responseId = response.id ?? responseId; + model = response.model ?? model; + usage = response.usage ?? usage; + if (response.error) error = response.error; + } + if (payload.response_id) responseId = payload.response_id; + if (payload.item?.id) responseId = responseId ?? payload.item.id; + if (payload.error) error = payload.error; + if (payload.type === "error" && !payload.error) error = payload; + if (typeof payload.delta === "string") deltas.push(payload.delta); + else if (typeof payload.text === "string" && payload.type === "response.output_text.done" && deltas.length === 0) { + deltas.push(payload.text); + } + } + + const content = deltas.join("").trim() || extractOpenAiOutputText(finalResponse); + return { + finalResponse, + content, + responseId, + model, + usage, + error + }; +} + +function parseSseDataMessages(text) { + const messages = []; + for (const block of String(text ?? "").split(/\r?\n\r?\n/u)) { + const data = []; + for (const line of block.split(/\r?\n/u)) { + if (line.startsWith("data:")) data.push(line.slice(5).trimStart()); + } + const message = data.join("\n").trim(); + if (message && message !== "[DONE]") messages.push(message); + } + return messages; +} + +function parseJsonOrNull(value) { + try { + return JSON.parse(value); + } catch { + return null; + } +} + async function commandExists(command, env) { if (command.includes("/") || command.includes("\\")) { try { diff --git a/internal/cloud/server.test.mjs b/internal/cloud/server.test.mjs index 0b7d2112..d28e8bb4 100644 --- a/internal/cloud/server.test.mjs +++ b/internal/cloud/server.test.mjs @@ -1,4 +1,5 @@ import assert from "node:assert/strict"; +import { createServer as createHttpServer } from "node:http"; import { createServer as createTcpServer } from "node:net"; import test from "node:test"; @@ -442,6 +443,113 @@ test("cloud api /v1/agent/chat returns structured completed Code Agent payload", } }); +test("cloud api /v1/agent/chat uses streaming OpenAI Responses and parses SSE text", async () => { + const providerRequests = []; + const providerServer = createHttpServer((request, response) => { + const chunks = []; + request.on("data", (chunk) => chunks.push(chunk)); + request.on("end", () => { + const bodyText = Buffer.concat(chunks).toString("utf8"); + const body = JSON.parse(bodyText); + providerRequests.push({ + method: request.method, + url: request.url, + authorizationPresent: Boolean(request.headers.authorization), + accept: request.headers.accept, + contentType: request.headers["content-type"], + body + }); + response.writeHead(200, { + "content-type": "text/event-stream" + }); + response.end([ + `data: ${JSON.stringify({ + type: "response.created", + response: { + id: "resp_server_test_stream", + model: body.model, + usage: null + } + })}`, + `data: ${JSON.stringify({ + type: "response.output_text.delta", + response_id: "resp_server_test_stream", + delta: "HWLAB Code Agent " + })}`, + `data: ${JSON.stringify({ + type: "response.output_text.delta", + response_id: "resp_server_test_stream", + delta: "streaming ready." + })}`, + `data: ${JSON.stringify({ + type: "response.completed", + response: { + id: "resp_server_test_stream", + model: body.model, + usage: null + } + })}`, + "data: [DONE]", + "" + ].join("\n\n")); + }); + }); + await new Promise((resolve) => providerServer.listen(0, "127.0.0.1", resolve)); + const providerPort = providerServer.address().port; + + const server = createCloudApiServer({ + env: { + OPENAI_API_KEY: "test-openai-key-material", + HWLAB_CODE_AGENT_PROVIDER: "openai", + HWLAB_CODE_AGENT_MODEL: "gpt-test", + HWLAB_CODE_AGENT_OPENAI_BASE_URL: `http://127.0.0.1:${providerPort}/v1/responses` + } + }); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + + try { + const { port } = server.address(); + const response = await fetch(`http://127.0.0.1:${port}/v1/agent/chat`, { + method: "POST", + headers: { + "content-type": "application/json", + "x-trace-id": "trc_server-test-agent-chat-stream" + }, + body: JSON.stringify({ + conversationId: "cnv_server-test-agent-chat-stream", + message: "请用一句话说明当前 HWLAB 工作台可以做什么。" + }) + }); + assert.equal(response.status, 200); + const payload = await response.json(); + assert.equal(payload.status, "completed"); + assert.equal(payload.conversationId, "cnv_server-test-agent-chat-stream"); + assert.equal(payload.traceId, "trc_server-test-agent-chat-stream"); + assert.equal(payload.provider, "openai-responses"); + assert.equal(payload.model, "gpt-test"); + assert.equal(payload.providerTrace.responseId, "resp_server_test_stream"); + assert.equal(payload.reply.content, "HWLAB Code Agent streaming ready."); + assert.equal(JSON.stringify(payload).includes("test-openai-key-material"), false); + + assert.equal(providerRequests.length, 1); + assert.equal(providerRequests[0].method, "POST"); + assert.equal(providerRequests[0].url, "/v1/responses"); + assert.equal(providerRequests[0].authorizationPresent, true); + assert.match(providerRequests[0].accept, /text\/event-stream/u); + assert.equal(providerRequests[0].contentType, "application/json"); + assert.equal(providerRequests[0].body.model, "gpt-test"); + assert.equal(providerRequests[0].body.store, false); + assert.equal(providerRequests[0].body.stream, true); + } finally { + await new Promise((resolve, reject) => { + server.close((error) => (error ? reject(error) : resolve())); + }); + await new Promise((resolve, reject) => { + providerServer.close((error) => (error ? reject(error) : resolve())); + }); + } +}); + test("cloud api /v1/agent/chat reports provider gaps without faking a reply", async () => { const server = createCloudApiServer({ env: {