Merge remote-tracking branch 'origin/master' into gitea/new_feature
# Conflicts: # open-sse/handlers/chatCore.js # open-sse/services/combo.js # src/app/(dashboard)/dashboard/profile/page.js # src/app/api/v1/models/route.js # src/lib/db/repos/settingsRepo.js
This commit is contained in:
@@ -245,6 +245,18 @@ export class AntigravityExecutor extends BaseExecutor {
|
||||
// Strip tools/toolConfig (handled separately) and blacklisted fields that Google rejects
|
||||
const { tools: _originalTools, toolConfig: _originalToolConfig, ...requestWithoutTools } = body.request || {};
|
||||
stripBlacklisted(requestWithoutTools);
|
||||
|
||||
// Rewrite competitive system prompts (e.g. Zed IDE's Claude prompt) to prevent Antigravity from
|
||||
// flagging the request and immediately blocking it with a 429 Quota Exhausted response.
|
||||
if (requestWithoutTools.systemInstruction?.parts) {
|
||||
const oldText = "You are a Claude agent, built on Anthropic's Claude Agent SDK.";
|
||||
for (const part of requestWithoutTools.systemInstruction.parts) {
|
||||
if (typeof part.text === "string" && part.text.includes(oldText)) {
|
||||
part.text = part.text.split(oldText).join("");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const generationConfig = { ...(requestWithoutTools.generationConfig || {}) };
|
||||
if (generationConfig.maxOutputTokens > MAX_ANTIGRAVITY_OUTPUT_TOKENS) {
|
||||
generationConfig.maxOutputTokens = MAX_ANTIGRAVITY_OUTPUT_TOKENS;
|
||||
|
||||
@@ -10,7 +10,6 @@ import { CodexExecutor } from "./codex.js";
|
||||
import { CursorExecutor } from "./cursor.js";
|
||||
import { VertexExecutor } from "./vertex.js";
|
||||
import { OpenCodeExecutor } from "./opencode.js";
|
||||
import { OpenCodeGoExecutor } from "./opencode-go.js";
|
||||
import { GrokWebExecutor } from "./grok-web.js";
|
||||
import { GrokCliExecutor } from "./grok-cli.js";
|
||||
import { PerplexityWebExecutor } from "./perplexity-web.js";
|
||||
@@ -41,7 +40,6 @@ const executors = {
|
||||
vertex: new VertexExecutor("vertex"),
|
||||
"vertex-partner": new VertexExecutor("vertex-partner"),
|
||||
opencode: new OpenCodeExecutor(),
|
||||
"opencode-go": new OpenCodeGoExecutor(),
|
||||
"grok-web": new GrokWebExecutor(),
|
||||
"grok-cli": new GrokCliExecutor(),
|
||||
gcli: new GrokCliExecutor(), // Alias
|
||||
@@ -86,7 +84,6 @@ export { CursorExecutor } from "./cursor.js";
|
||||
export { VertexExecutor } from "./vertex.js";
|
||||
export { DefaultExecutor } from "./default.js";
|
||||
export { OpenCodeExecutor } from "./opencode.js";
|
||||
export { OpenCodeGoExecutor } from "./opencode-go.js";
|
||||
export { GrokWebExecutor } from "./grok-web.js";
|
||||
export { GrokCliExecutor } from "./grok-cli.js";
|
||||
export { PerplexityWebExecutor } from "./perplexity-web.js";
|
||||
|
||||
@@ -144,6 +144,12 @@ function normalizeStopReason(value) {
|
||||
return reason || null;
|
||||
}
|
||||
|
||||
// Of the reasons stopDisposition() folds into "terminal_incomplete", only these
|
||||
// mean "usable as far as it got, then the budget ran out" -- the case
|
||||
// finish_reason "length" exists for. cancelled / pause_turn are abandoned turns
|
||||
// whose partial content must stay private, so they are deliberately absent.
|
||||
const KIRO_TRUNCATION_STOP_REASONS = new Set(["model_context_window_exceeded", "max_tokens"]);
|
||||
|
||||
function stopDisposition(stopReason, hasToolCalls) {
|
||||
if (["malformed_model_output", "invalid_model_output"].includes(stopReason)) return "retryable_protocol_failure";
|
||||
if (["cancelled", "pause_turn", "model_context_window_exceeded"].includes(stopReason)) return "terminal_incomplete";
|
||||
@@ -711,14 +717,25 @@ export class KiroExecutor extends BaseExecutor {
|
||||
};
|
||||
const emitTools = (controller) => {
|
||||
for (const tool of state.tools.values()) {
|
||||
const input = parsedToolInput(tool);
|
||||
if (tool.name === "tool_call") {
|
||||
if (typeof input.name !== "string" || !input.name.trim()) {
|
||||
throw new Error("Invalid Kiro tool_call payload: missing nested MCP tool name");
|
||||
}
|
||||
if (!Object.prototype.hasOwnProperty.call(input, "arguments")) {
|
||||
throw new Error("Invalid Kiro tool_call payload: missing nested MCP tool arguments");
|
||||
// Validate per tool, not per turn: one unusable fragment used to throw out
|
||||
// of emitTools and take every other complete tool call in the same turn
|
||||
// with it, which the client saw as a turn that answered nothing.
|
||||
let input;
|
||||
try {
|
||||
input = parsedToolInput(tool);
|
||||
if (tool.name === "tool_call") {
|
||||
if (typeof input.name !== "string" || !input.name.trim()) {
|
||||
throw new Error("Invalid Kiro tool_call payload: missing nested MCP tool name");
|
||||
}
|
||||
if (!Object.prototype.hasOwnProperty.call(input, "arguments")) {
|
||||
throw new Error("Invalid Kiro tool_call payload: missing nested MCP tool arguments");
|
||||
}
|
||||
}
|
||||
} catch (error) {
|
||||
state.droppedTools = (state.droppedTools || 0) + 1;
|
||||
state.toolValidationError ||= error.message;
|
||||
console.error(`[Kiro] dropping unusable tool call ${tool.id} (${tool.name}): ${error.message}`);
|
||||
continue;
|
||||
}
|
||||
const index = state.toolCounter++;
|
||||
emitDelta(controller, {
|
||||
@@ -729,14 +746,26 @@ export class KiroExecutor extends BaseExecutor {
|
||||
function: { name: tool.name, arguments: "" }
|
||||
}]
|
||||
});
|
||||
const serializedInput = JSON.stringify(input);
|
||||
emitDelta(controller, {
|
||||
tool_calls: [{ index, function: { arguments: JSON.stringify(input) } }]
|
||||
tool_calls: [{ index, function: { arguments: serializedInput } }]
|
||||
});
|
||||
// Tool arguments are billed output like any other completion bytes. They
|
||||
// were never added to totalContentLength, so the /4 estimator in finish()
|
||||
// reported OUT 0 -- or the Math.max floor of 1 -- for every turn whose
|
||||
// entire answer was a tool call.
|
||||
state.totalContentLength += tool.name.length + serializedInput.length;
|
||||
state.hasToolCalls = true;
|
||||
}
|
||||
state.tools.clear();
|
||||
state.bufferedToolBytes = 0;
|
||||
if (state.stopReason === "tool_use" && !state.hasToolCalls) {
|
||||
// A declared tool turn that emitted no usable call is only fatal when the
|
||||
// turn produced nothing else. Throwing unconditionally here escaped
|
||||
// emitTools() with provenance "invalid_tool_call", which the integrity gate
|
||||
// re-derived into a repair retry -- discarding text the client had already
|
||||
// been promised.
|
||||
if (state.stopReason === "tool_use" && !state.hasToolCalls &&
|
||||
!state.hasText && !state.hasReasoning && !state.hasCode) {
|
||||
throw new Error("Kiro tool_use stop reason did not include a complete tool call");
|
||||
}
|
||||
};
|
||||
@@ -796,7 +825,6 @@ export class KiroExecutor extends BaseExecutor {
|
||||
emitDelta(controller, { content: event.payload.content });
|
||||
} else if (eventType === "toolUseEvent") {
|
||||
state.sawToolUse = true;
|
||||
if (state.toolValidationError) return true;
|
||||
const values = Array.isArray(event.payload) ? event.payload : [event.payload];
|
||||
if (!values[0]) throw new Error("Kiro toolUseEvent is empty");
|
||||
for (const value of values) {
|
||||
@@ -924,9 +952,10 @@ export class KiroExecutor extends BaseExecutor {
|
||||
} catch (error) {
|
||||
const bufferExceeded = error.code === "KIRO_BUFFER_EXCEEDED";
|
||||
if (!bufferExceeded) {
|
||||
// Keep whatever is already buffered: the rejected fragment belongs to
|
||||
// one tool, and clearing the map dropped the complete calls too.
|
||||
state.toolValidationError ||= error.message;
|
||||
state.tools.clear();
|
||||
state.bufferedToolBytes = 0;
|
||||
console.error(`[Kiro] tool fragment rejected, keeping ${state.tools.size} buffered tool(s): ${error.message}`);
|
||||
continue;
|
||||
}
|
||||
fail(
|
||||
@@ -958,7 +987,16 @@ export class KiroExecutor extends BaseExecutor {
|
||||
}
|
||||
state.transportState = "clean_eof";
|
||||
const declaredDisposition = stopDisposition(state.stopReason, state.sawToolUse);
|
||||
if (["retryable_protocol_failure", "terminal_incomplete", "terminal_refusal", "unknown_failure"].includes(declaredDisposition)) {
|
||||
// model_context_window_exceeded / max_tokens map to terminal_incomplete. When
|
||||
// they arrive after the model already streamed content, fail() threw away a
|
||||
// complete-enough answer; a truncated turn is what finish_reason "length" is
|
||||
// for. chunkIndex > 0 means at least one delta already reached the client.
|
||||
const declaredTruncatedAfterOutput = declaredDisposition === "terminal_incomplete" &&
|
||||
KIRO_TRUNCATION_STOP_REASONS.has(state.stopReason) && state.chunkIndex > 0;
|
||||
if (declaredTruncatedAfterOutput) {
|
||||
console.error(`[Kiro] truncated after ${state.chunkIndex} chunk(s) (stop_reason=${state.stopReason}); keeping output`);
|
||||
}
|
||||
if (!declaredTruncatedAfterOutput && ["retryable_protocol_failure", "terminal_incomplete", "terminal_refusal", "unknown_failure"].includes(declaredDisposition)) {
|
||||
const code = declaredDisposition === "retryable_protocol_failure"
|
||||
? "kiro_retryable_protocol_failure"
|
||||
: declaredDisposition === "terminal_refusal"
|
||||
@@ -975,16 +1013,6 @@ export class KiroExecutor extends BaseExecutor {
|
||||
);
|
||||
return;
|
||||
}
|
||||
if (state.toolValidationError) {
|
||||
fail(
|
||||
controller,
|
||||
"invalid_tool_call",
|
||||
"invalid_kiro_tool_call",
|
||||
state.toolValidationError,
|
||||
{ transport_state: state.transportState, stop_disposition: "retryable_protocol_failure" }
|
||||
);
|
||||
return;
|
||||
}
|
||||
try {
|
||||
emitTools(controller);
|
||||
} catch (error) {
|
||||
@@ -997,6 +1025,22 @@ export class KiroExecutor extends BaseExecutor {
|
||||
);
|
||||
return;
|
||||
}
|
||||
// Fail only when the turn has nothing usable left. emitTools() validates
|
||||
// per tool and drops just the unusable ones, so this has to run AFTER it:
|
||||
// before, the rejected tool was still buffered and tools.size was never 0.
|
||||
// A turn that also produced text keeps that text -- the dropped call is
|
||||
// logged, not fatal.
|
||||
if (state.toolValidationError && !state.hasToolCalls &&
|
||||
!state.hasText && !state.hasReasoning && !state.hasCode) {
|
||||
fail(
|
||||
controller,
|
||||
"invalid_tool_call",
|
||||
"invalid_kiro_tool_call",
|
||||
state.toolValidationError,
|
||||
{ transport_state: state.transportState, stop_disposition: "retryable_protocol_failure" }
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
const hasOutput = state.hasText || state.hasReasoning || state.hasCode || state.hasToolCalls;
|
||||
if (!hasOutput && !state.explicitStop) {
|
||||
@@ -1011,7 +1055,13 @@ export class KiroExecutor extends BaseExecutor {
|
||||
}
|
||||
|
||||
const disposition = stopDisposition(state.stopReason, state.hasToolCalls);
|
||||
if (["retryable_protocol_failure", "terminal_incomplete", "terminal_refusal", "unknown_failure"].includes(disposition)) {
|
||||
// Same reasoning as declaredTruncatedAfterOutput above.
|
||||
const truncatedAfterOutput = disposition === "terminal_incomplete" &&
|
||||
KIRO_TRUNCATION_STOP_REASONS.has(state.stopReason) && state.chunkIndex > 0;
|
||||
if (truncatedAfterOutput) {
|
||||
console.error(`[Kiro] truncated after ${state.chunkIndex} chunk(s) (stop_reason=${state.stopReason}); closing as length`);
|
||||
}
|
||||
if (!truncatedAfterOutput && ["retryable_protocol_failure", "terminal_incomplete", "terminal_refusal", "unknown_failure"].includes(disposition)) {
|
||||
const code = disposition === "retryable_protocol_failure"
|
||||
? "kiro_retryable_protocol_failure"
|
||||
: disposition === "terminal_refusal"
|
||||
@@ -1041,18 +1091,24 @@ export class KiroExecutor extends BaseExecutor {
|
||||
total_tokens: prompt + completion
|
||||
};
|
||||
}
|
||||
const finishReason = state.hasToolCalls
|
||||
? "tool_calls"
|
||||
: disposition === "length"
|
||||
? "length"
|
||||
: "stop";
|
||||
const finishReason = truncatedAfterOutput
|
||||
? "length"
|
||||
: state.hasToolCalls
|
||||
? "tool_calls"
|
||||
: disposition === "length"
|
||||
? "length"
|
||||
: "stop";
|
||||
controller.enqueue(sseChunk({}, finishReason, state.usage));
|
||||
controller.enqueue(encoder.encode(SSE_DONE));
|
||||
state.finished = true;
|
||||
options.onTerminalState?.(diagnostics({
|
||||
terminal_provenance: state.terminalProvenance || "clean_eventstream_eof",
|
||||
transport_state: state.transportState,
|
||||
stop_disposition: disposition
|
||||
// Report what this exit actually did, not the raw disposition. The
|
||||
// integrity gate re-derives its verdict from stop_disposition, so
|
||||
// reporting "terminal_incomplete" for a turn we deliberately kept made
|
||||
// it discard the very bytes we just released to the client.
|
||||
stop_disposition: truncatedAfterOutput ? "length" : disposition
|
||||
}));
|
||||
};
|
||||
|
||||
|
||||
@@ -1,49 +0,0 @@
|
||||
import { BaseExecutor } from "./base.js";
|
||||
import { PROVIDERS } from "../config/providers.js";
|
||||
import { injectReasoningContent } from "../utils/reasoningContentInjector.js";
|
||||
import { ANTHROPIC_API_VERSION } from "../providers/shared.js";
|
||||
|
||||
// Models that use /zen/go/v1/messages (Anthropic/Claude format + x-api-key auth)
|
||||
const MESSAGES_FORMAT_MODELS = new Set([
|
||||
"minimax-m3",
|
||||
"minimax-m2.7",
|
||||
"minimax-m2.5",
|
||||
"qwen3.7-max",
|
||||
"qwen3.7-plus",
|
||||
"qwen3.6-plus",
|
||||
]);
|
||||
|
||||
const BASE = "https://opencode.ai/zen/go/v1";
|
||||
|
||||
export class OpenCodeGoExecutor extends BaseExecutor {
|
||||
constructor() {
|
||||
super("opencode-go", PROVIDERS["opencode-go"]);
|
||||
}
|
||||
|
||||
// buildUrl runs before buildHeaders in BaseExecutor.execute, cache model here
|
||||
buildUrl(model) {
|
||||
this._lastModel = model;
|
||||
return MESSAGES_FORMAT_MODELS.has(model)
|
||||
? `${BASE}/messages`
|
||||
: `${BASE}/chat/completions`;
|
||||
}
|
||||
|
||||
buildHeaders(credentials, stream = true) {
|
||||
const key = credentials?.apiKey || credentials?.accessToken;
|
||||
const headers = { "Content-Type": "application/json" };
|
||||
|
||||
if (MESSAGES_FORMAT_MODELS.has(this._lastModel)) {
|
||||
headers["x-api-key"] = key;
|
||||
headers["anthropic-version"] = ANTHROPIC_API_VERSION;
|
||||
} else {
|
||||
headers["Authorization"] = `Bearer ${key}`;
|
||||
}
|
||||
|
||||
if (stream) headers["Accept"] = "text/event-stream";
|
||||
return headers;
|
||||
}
|
||||
|
||||
transformRequest(model, body) {
|
||||
return injectReasoningContent({ provider: this.provider, model, body });
|
||||
}
|
||||
}
|
||||
@@ -1,16 +1,43 @@
|
||||
import crypto from "crypto";
|
||||
import { BaseExecutor } from "./base.js";
|
||||
import { PROVIDERS } from "../config/providers.js";
|
||||
import { injectReasoningContent } from "../utils/reasoningContentInjector.js";
|
||||
import { resolveSessionId } from "../utils/sessionManager.js";
|
||||
|
||||
// Models that use /zen/v1/messages (claude format)
|
||||
const OPENCODE_UA = "opencode";
|
||||
const MESSAGES_MODELS = new Set();
|
||||
|
||||
function generateRequestId() {
|
||||
return `msg_${crypto.randomUUID().replace(/-/g, "")}`;
|
||||
}
|
||||
|
||||
function generateSessionId() {
|
||||
return `ses_${crypto.randomUUID().replace(/-/g, "")}`;
|
||||
}
|
||||
|
||||
// Normalize any resolved id into opencode's ses_ format (stable per-conversation)
|
||||
function toOpencodeSession(id) {
|
||||
const stripped = String(id || "").replace(/^ses_/, "").replace(/-/g, "");
|
||||
return stripped ? `ses_${stripped}` : null;
|
||||
}
|
||||
|
||||
function resolveOpencodeSession(body, credentials) {
|
||||
return toOpencodeSession(resolveSessionId({
|
||||
headers: credentials?.rawHeaders,
|
||||
body,
|
||||
connectionId: credentials?.connectionId,
|
||||
scope: "opencode",
|
||||
}));
|
||||
}
|
||||
|
||||
export class OpenCodeExecutor extends BaseExecutor {
|
||||
constructor() {
|
||||
super("opencode", PROVIDERS.opencode);
|
||||
this._currentSessionId = null;
|
||||
}
|
||||
|
||||
transformRequest(model, body) {
|
||||
transformRequest(model, body, stream, credentials) {
|
||||
this._currentSessionId = resolveOpencodeSession(body, credentials);
|
||||
return injectReasoningContent({ provider: this.provider, model, body });
|
||||
}
|
||||
|
||||
@@ -21,12 +48,23 @@ export class OpenCodeExecutor extends BaseExecutor {
|
||||
: `${base}/zen/v1/chat/completions`;
|
||||
}
|
||||
|
||||
buildHeaders() {
|
||||
buildHeaders(credentials, stream = true) {
|
||||
const raw = credentials?.rawHeaders || {};
|
||||
const lower = {};
|
||||
for (const [k, v] of Object.entries(raw)) lower[k.toLowerCase()] = v;
|
||||
|
||||
const downstreamUa = lower["user-agent"] || "";
|
||||
const isOpencodeDownstream = downstreamUa.toLowerCase().includes("opencode");
|
||||
|
||||
return {
|
||||
"Content-Type": "application/json",
|
||||
"Authorization": "Bearer public",
|
||||
"x-opencode-client": "desktop",
|
||||
"Accept": "text/event-stream"
|
||||
"User-Agent": isOpencodeDownstream ? downstreamUa : OPENCODE_UA,
|
||||
"x-opencode-client": lower["x-opencode-client"] || "desktop",
|
||||
"x-opencode-session": lower["x-opencode-session"] || this._currentSessionId || generateSessionId(),
|
||||
"x-opencode-request": lower["x-opencode-request"] || generateRequestId(),
|
||||
"x-opencode-project": lower["x-opencode-project"] || "global",
|
||||
"Accept": stream ? "text/event-stream" : "*/*",
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
@@ -216,6 +216,52 @@ async function buildQoderRequestBody({ model, body, credentials, log, proxyOptio
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Check if a qoder error message indicates a billing/quota block.
|
||||
* Signatures: code 112 (quota exhausted), code 10605 (queue throttle), pricingUrl field.
|
||||
*/
|
||||
function isBillingBlock(inner) {
|
||||
if (!inner || typeof inner !== "string") return false;
|
||||
const lowerMsg = inner.toLowerCase();
|
||||
// Match: {"code":"112",...}, {"code":"10605",...}, or pricingUrl field
|
||||
return /\"code\"\s*:\s*\"(112|10605)\"/.test(inner) || lowerMsg.includes("pricingurl");
|
||||
}
|
||||
|
||||
/**
|
||||
* Peek the first SSE frame to detect billing errors before piping.
|
||||
* Returns { isBilling, statusVal, message, consumed } — `consumed` is every
|
||||
* byte read so far (including the peeked line) so the caller can re-process
|
||||
* it and nothing is dropped from the stream.
|
||||
*/
|
||||
async function peekFirstQoderFrame(reader, decoder) {
|
||||
let consumed = "";
|
||||
while (true) {
|
||||
const { done, value } = await reader.read();
|
||||
if (done) return { isBilling: false, consumed, upstreamDone: true };
|
||||
|
||||
consumed += decoder.decode(value, { stream: true });
|
||||
const nl = consumed.indexOf("\n");
|
||||
if (nl === -1) continue; // need a full line first
|
||||
|
||||
const line = consumed.slice(0, nl).replace(/\r$/, "").trim();
|
||||
if (!line.startsWith("data:")) continue;
|
||||
|
||||
const data = line.slice(5).trimStart();
|
||||
if (data === "[DONE]") return { isBilling: false, consumed };
|
||||
|
||||
let envelope;
|
||||
try { envelope = JSON.parse(data); } catch { return { isBilling: false, consumed }; }
|
||||
|
||||
const statusVal = typeof envelope.statusCodeValue === "number" ? envelope.statusCodeValue : 200;
|
||||
const inner = typeof envelope.body === "string" ? envelope.body : "";
|
||||
|
||||
if (statusVal !== 200 && isBillingBlock(inner)) {
|
||||
return { isBilling: true, statusVal, message: inner || `qoder billing block (${statusVal})` };
|
||||
}
|
||||
return { isBilling: false, consumed };
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Wrap the upstream's `{statusCodeValue, body}` SSE envelope into plain
|
||||
* OpenAI SSE chunks the rest of the chatCore pipeline understands.
|
||||
@@ -230,16 +276,34 @@ async function buildQoderRequestBody({ model, body, credentials, log, proxyOptio
|
||||
* [DONE]/error frame (agent keepalive). Non-streaming clients drain via
|
||||
* response.text() which hangs until the socket closes — so on terminal
|
||||
* events we cancel the upstream reader and close our stream immediately.
|
||||
*
|
||||
* NEW: Peek first frame to detect billing blocks (code 112/10605/pricingUrl).
|
||||
* If detected, return 403 response so chatCore marks connection unavailable
|
||||
* and triggers combo fallback instead of leaking error text into chat.
|
||||
*/
|
||||
function wrapQoderSSE(response, model) {
|
||||
async function wrapQoderSSE(response, model) {
|
||||
if (!response.ok || !response.body) return response;
|
||||
|
||||
const decoder = new TextDecoder();
|
||||
const encoder = new TextEncoder();
|
||||
let buffer = "";
|
||||
let doneEmitted = false;
|
||||
const reader = response.body.getReader();
|
||||
|
||||
// Peek first frame to detect billing block
|
||||
const peek = await peekFirstQoderFrame(reader, decoder);
|
||||
if (peek?.isBilling) {
|
||||
// Billing block detected — return 403 so chatCore fails this connection
|
||||
await reader.cancel().catch(() => {});
|
||||
return new Response(
|
||||
JSON.stringify({ error: { message: peek.message, code: peek.statusVal } }),
|
||||
{ status: 403, headers: { "Content-Type": "application/json" } }
|
||||
);
|
||||
}
|
||||
|
||||
// Normal flow: re-process every byte the peek consumed, then continue.
|
||||
let buffer = peek.consumed || "";
|
||||
const upstreamDrained = peek.upstreamDone === true;
|
||||
const encoder = new TextEncoder();
|
||||
let doneEmitted = false;
|
||||
|
||||
// Process one already-extracted SSE line (no trailing newline).
|
||||
const processLine = (line, controller) => {
|
||||
const trimmed = line.replace(/\r$/, "").trim();
|
||||
@@ -288,7 +352,28 @@ function wrapQoderSSE(response, model) {
|
||||
// enqueueing would never be re-invoked, hanging consumers like .text().
|
||||
async start(controller) {
|
||||
try {
|
||||
while (!doneEmitted) {
|
||||
// Drain whatever the peek already pulled off the socket first.
|
||||
let nlSeed;
|
||||
while ((nlSeed = buffer.indexOf("\n")) !== -1) {
|
||||
const line = buffer.slice(0, nlSeed);
|
||||
buffer = buffer.slice(nlSeed + 1);
|
||||
processLine(line, controller);
|
||||
if (doneEmitted) {
|
||||
await reader.cancel().catch(() => {});
|
||||
controller.close();
|
||||
return;
|
||||
}
|
||||
}
|
||||
if (upstreamDrained) {
|
||||
// Peek hit end-of-stream: flush any trailing partial line.
|
||||
buffer += decoder.decode();
|
||||
if (buffer.length > 0) {
|
||||
processLine(buffer, controller);
|
||||
buffer = "";
|
||||
}
|
||||
}
|
||||
|
||||
while (!doneEmitted && !upstreamDrained) {
|
||||
const { done, value } = await reader.read();
|
||||
if (done) {
|
||||
buffer += decoder.decode();
|
||||
@@ -473,7 +558,7 @@ export class QoderExecutor extends BaseExecutor {
|
||||
return { response, url, headers, transformedBody: payload };
|
||||
}
|
||||
|
||||
const wrapped = wrapQoderSSE(response, `qoder/${qoderKey}`);
|
||||
const wrapped = await wrapQoderSSE(response, `qoder/${qoderKey}`);
|
||||
return { response: wrapped, url, headers, transformedBody: payload };
|
||||
}
|
||||
|
||||
@@ -497,4 +582,5 @@ export const __test__ = {
|
||||
normalizeMessages,
|
||||
wrapQoderSSE,
|
||||
buildQoderRequestBody,
|
||||
isBillingBlock,
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user