Merge pull request #209 from pikasTech/fix/d601-code-agent-provider-502-v2

fix: stream Code Agent OpenAI responses
This commit is contained in:
Lyon
2026-05-23 08:34:21 +08:00
committed by GitHub
2 changed files with 209 additions and 5 deletions
+101 -5
View File
@@ -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 {
+108
View File
@@ -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: {