Merge origin/master (v0.5.69) into gitea/new_feature
This commit is contained in:
@@ -46,203 +46,242 @@ export class CommandCodeExecutor extends BaseExecutor {
|
||||
return headers;
|
||||
}
|
||||
|
||||
async execute(opts) {
|
||||
const result = await super.execute(opts);
|
||||
if (!result?.response?.ok || !result.response.body) return result;
|
||||
result.response = await peekForUpstreamError(result.response, opts.model, {
|
||||
signal: opts.signal,
|
||||
});
|
||||
return result;
|
||||
}
|
||||
async execute(opts) {
|
||||
const result = await super.execute(opts);
|
||||
if (!result?.response?.ok || !result.response.body) return result;
|
||||
result.response = await inspectAndWrapCommandCodeResponse(result.response, opts.model);
|
||||
return result;
|
||||
}
|
||||
|
||||
parseError(response, bodyText) {
|
||||
let parsed = null;
|
||||
try {
|
||||
parsed = JSON.parse(bodyText || "{}");
|
||||
} catch {
|
||||
parsed = null;
|
||||
}
|
||||
const errObj = parsed?.error || parsed;
|
||||
const msg = errObj?.message || parsed?.message || bodyText || response.statusText;
|
||||
const status = Number(errObj?.code || errObj?.statusCode || response.status) || response.status;
|
||||
return {
|
||||
status,
|
||||
message: msg || `CommandCode upstream error: ${response.status}`,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
// How long to hold the response open while peeking the first upstream events.
|
||||
// An upstream error event ("Network connection lost") is emitted at stream
|
||||
// start, so the peek is fast; the bound just prevents a slow-started stream
|
||||
// from being held hostage. Env: COMMANDCODE_PEEK_TIMEOUT_MS.
|
||||
const PEEK_TIMEOUT_MS = (() => {
|
||||
const raw = process.env.COMMANDCODE_PEEK_TIMEOUT_MS;
|
||||
const n = raw ? parseInt(raw, 10) : NaN;
|
||||
return Number.isFinite(n) && n > 0 ? n : 10 * 1000;
|
||||
})();
|
||||
export function parseCommandCodeError(event) {
|
||||
if (!event || typeof event !== "object") {
|
||||
return {
|
||||
statusCode: 503,
|
||||
message: "CommandCode upstream error",
|
||||
type: "server_error",
|
||||
};
|
||||
}
|
||||
|
||||
// Event types that count as "the stream has started producing". Everything
|
||||
// else (start, start-step, reasoning-start, text-start, ...) is metadata and
|
||||
// does not end the peek.
|
||||
const MEANINGFUL_EVENT_TYPES = new Set([
|
||||
"text-delta",
|
||||
"reasoning-delta",
|
||||
"tool-input-start",
|
||||
"tool-input-delta",
|
||||
"tool-input-end",
|
||||
"tool-call",
|
||||
"finish-step",
|
||||
"finish",
|
||||
]);
|
||||
const errVal = event.error ?? event.message ?? "unknown";
|
||||
let message = "";
|
||||
let statusCode = null;
|
||||
let type = "server_error";
|
||||
|
||||
function makeAbortError(reason) {
|
||||
const error = new Error(reason?.message || reason || "Request aborted");
|
||||
error.name = "AbortError";
|
||||
return error;
|
||||
if (typeof errVal === "object" && errVal !== null) {
|
||||
message = errVal.message || errVal.error || JSON.stringify(errVal);
|
||||
if (errVal.statusCode && Number.isInteger(Number(errVal.statusCode))) {
|
||||
statusCode = Number(errVal.statusCode);
|
||||
} else if (errVal.status && Number.isInteger(Number(errVal.status))) {
|
||||
statusCode = Number(errVal.status);
|
||||
}
|
||||
if (errVal.type) type = errVal.type;
|
||||
} else if (typeof errVal === "string") {
|
||||
message = errVal;
|
||||
} else {
|
||||
message = JSON.stringify(errVal);
|
||||
}
|
||||
|
||||
if (event.statusCode && Number.isInteger(Number(event.statusCode))) {
|
||||
statusCode = Number(event.statusCode);
|
||||
}
|
||||
|
||||
if (!statusCode || statusCode < 400 || statusCode > 599) {
|
||||
const lower = message.toLowerCase();
|
||||
if (lower.includes("rate limit") || lower.includes("too many requests")) {
|
||||
statusCode = 429;
|
||||
type = "rate_limit_error";
|
||||
} else if (lower.includes("unauthorized") || lower.includes("invalid api key") || lower.includes("authentication")) {
|
||||
statusCode = 401;
|
||||
type = "authentication_error";
|
||||
} else if (lower.includes("payment required") || lower.includes("billing")) {
|
||||
statusCode = 402;
|
||||
type = "billing_error";
|
||||
} else if (lower.includes("quota") || lower.includes("forbidden") || lower.includes("permission")) {
|
||||
statusCode = 403;
|
||||
type = "permission_error";
|
||||
} else if (lower.includes("not found")) {
|
||||
statusCode = 404;
|
||||
type = "invalid_request_error";
|
||||
} else if (lower.includes("unavailable") || lower.includes("overloaded") || lower.includes("server error")) {
|
||||
statusCode = 503;
|
||||
type = "server_error";
|
||||
} else {
|
||||
statusCode = 503;
|
||||
}
|
||||
}
|
||||
|
||||
return { statusCode, message, type };
|
||||
}
|
||||
|
||||
function tryParseEvent(line) {
|
||||
const trimmed = line.trim();
|
||||
if (!trimmed) return null;
|
||||
const json = trimmed.startsWith("data:") ? trimmed.slice(5).trim() : trimmed;
|
||||
if (!json || json === "[DONE]") return null;
|
||||
try {
|
||||
return JSON.parse(json);
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
export async function inspectAndWrapCommandCodeResponse(originalResponse, model) {
|
||||
const reader = originalResponse.body.getReader();
|
||||
const decoder = new TextDecoder();
|
||||
let buffer = "";
|
||||
const bufferedLines = [];
|
||||
let detectedError = null;
|
||||
|
||||
try {
|
||||
while (true) {
|
||||
const { value, done } = await reader.read();
|
||||
if (done) {
|
||||
const trimmed = buffer.trim();
|
||||
if (trimmed) {
|
||||
try {
|
||||
const jsonStr = trimmed.startsWith("data:") ? trimmed.slice(5).trim() : trimmed;
|
||||
const parsed = JSON.parse(jsonStr);
|
||||
if (parsed?.type === "error") {
|
||||
detectedError = parsed;
|
||||
} else {
|
||||
bufferedLines.push(trimmed);
|
||||
}
|
||||
} catch {
|
||||
bufferedLines.push(trimmed);
|
||||
}
|
||||
}
|
||||
break;
|
||||
}
|
||||
|
||||
buffer += decoder.decode(value, { stream: true });
|
||||
const lines = buffer.split("\n");
|
||||
buffer = lines.pop() || "";
|
||||
|
||||
let stopLoop = false;
|
||||
for (const line of lines) {
|
||||
const trimmed = line.trim();
|
||||
if (!trimmed) continue;
|
||||
const jsonStr = trimmed.startsWith("data:") ? trimmed.slice(5).trim() : trimmed;
|
||||
if (!jsonStr || jsonStr === "[DONE]") {
|
||||
bufferedLines.push(trimmed);
|
||||
stopLoop = true;
|
||||
break;
|
||||
}
|
||||
|
||||
let event;
|
||||
try {
|
||||
event = JSON.parse(jsonStr);
|
||||
} catch {
|
||||
bufferedLines.push(trimmed);
|
||||
continue;
|
||||
}
|
||||
|
||||
if (event?.type === "error") {
|
||||
detectedError = event;
|
||||
stopLoop = true;
|
||||
break;
|
||||
}
|
||||
|
||||
bufferedLines.push(trimmed);
|
||||
|
||||
if (
|
||||
event?.type === "text-delta" ||
|
||||
event?.type === "reasoning-delta" ||
|
||||
event?.type === "tool-input-start" ||
|
||||
event?.type === "tool-call" ||
|
||||
event?.type === "finish" ||
|
||||
event?.type === "finish-step"
|
||||
) {
|
||||
stopLoop = true;
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
if (stopLoop) break;
|
||||
}
|
||||
} catch {
|
||||
try { reader.releaseLock(); } catch { /* ignore */ }
|
||||
return originalResponse;
|
||||
}
|
||||
|
||||
if (detectedError) {
|
||||
try { await reader.cancel(); } catch { /* ignore */ }
|
||||
const { statusCode, message, type } = parseCommandCodeError(detectedError);
|
||||
return new Response(
|
||||
JSON.stringify({
|
||||
error: {
|
||||
message: `[CommandCode error: ${message}]`,
|
||||
type,
|
||||
code: statusCode,
|
||||
},
|
||||
}),
|
||||
{
|
||||
status: statusCode,
|
||||
statusText: statusCode === 503 ? "Service Unavailable" : (statusCode === 429 ? "Too Many Requests" : "Bad Gateway"),
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
"Access-Control-Allow-Origin": "*",
|
||||
},
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
const combinedStream = createReplayedStream(bufferedLines, buffer, reader);
|
||||
return wrapNdjsonAsOpenAISse(combinedStream, model, originalResponse);
|
||||
}
|
||||
|
||||
function formatErrorValue(errVal) {
|
||||
const errStr =
|
||||
typeof errVal === "string"
|
||||
? errVal
|
||||
: typeof errVal?.message === "string"
|
||||
? errVal.message
|
||||
: JSON.stringify(errVal);
|
||||
const errType =
|
||||
typeof errVal === "string"
|
||||
? "upstream_error"
|
||||
: errVal?.type || "upstream_error";
|
||||
return { message: errStr, type: errType };
|
||||
function createReplayedStream(bufferedLines, remainingBuffer, reader) {
|
||||
const encoder = new TextEncoder();
|
||||
let replayed = false;
|
||||
|
||||
return new ReadableStream({
|
||||
async pull(controller) {
|
||||
if (!replayed) {
|
||||
replayed = true;
|
||||
let prefix = bufferedLines.join("\n");
|
||||
if (prefix && remainingBuffer) {
|
||||
prefix += "\n" + remainingBuffer;
|
||||
} else if (remainingBuffer) {
|
||||
prefix = remainingBuffer;
|
||||
} else if (prefix) {
|
||||
prefix += "\n";
|
||||
}
|
||||
if (prefix) {
|
||||
controller.enqueue(encoder.encode(prefix));
|
||||
}
|
||||
}
|
||||
|
||||
try {
|
||||
const { value, done } = await reader.read();
|
||||
if (done) {
|
||||
controller.close();
|
||||
} else {
|
||||
controller.enqueue(value);
|
||||
}
|
||||
} catch (err) {
|
||||
controller.error(err);
|
||||
}
|
||||
},
|
||||
async cancel(reason) {
|
||||
try {
|
||||
await reader.cancel(reason);
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Read the first upstream events before committing the response.
|
||||
*
|
||||
* - `{"type":"error"}` as the first meaningful event → return a 502 Response so
|
||||
* chatCore's `!response.ok` path parses the error and triggers fallback.
|
||||
* - Otherwise → re-emit the buffered bytes + the rest of the stream through the
|
||||
* normal NDJSON → OpenAI SSE wrapper and return it untouched in spirit.
|
||||
*
|
||||
* Bounded by `timeoutMs` (default PEEK_TIMEOUT_MS): if no meaningful event
|
||||
* arrives in time, or the request signal aborts, we commit whatever we have and
|
||||
* let the regular stream pipeline (stall detection, abort handling) take over.
|
||||
*/
|
||||
export async function peekForUpstreamError(
|
||||
originalResponse,
|
||||
model,
|
||||
{ signal = null, timeoutMs = PEEK_TIMEOUT_MS } = {},
|
||||
) {
|
||||
const reader = originalResponse.body.getReader();
|
||||
const decoder = new TextDecoder();
|
||||
const abortController = new AbortController();
|
||||
const forwardAbort = () => abortController.abort(signal?.reason);
|
||||
if (signal?.aborted) abortController.abort(signal?.reason);
|
||||
else if (signal)
|
||||
signal.addEventListener("abort", forwardAbort, { once: true });
|
||||
|
||||
// Raw bytes for lossless re-emission; decoded text is only used for line
|
||||
// parsing / error detection. Never re-encode decoded text: TextDecoder
|
||||
// holds a split multi-byte char internally and flush() would replace it
|
||||
// with U+FFFD, corrupting the stream.
|
||||
const rawChunks = [];
|
||||
let peeked = "";
|
||||
let errorEvent = null;
|
||||
let committed = false;
|
||||
|
||||
const readWithTimeout = (ms) => {
|
||||
if (abortController.signal.aborted) {
|
||||
return Promise.reject(makeAbortError(abortController.signal.reason));
|
||||
}
|
||||
const timeoutPromise = new Promise((_, reject) => {
|
||||
const t = setTimeout(() => reject(new Error("peek timeout")), ms);
|
||||
t.unref?.();
|
||||
});
|
||||
const abortPromise = new Promise((_, reject) => {
|
||||
abortController.signal.addEventListener(
|
||||
"abort",
|
||||
() => reject(makeAbortError(abortController.signal.reason)),
|
||||
{ once: true },
|
||||
);
|
||||
});
|
||||
return Promise.race([reader.read(), timeoutPromise, abortPromise]);
|
||||
};
|
||||
|
||||
try {
|
||||
const deadline = Date.now() + timeoutMs;
|
||||
while (!errorEvent && !committed && Date.now() < deadline) {
|
||||
const { done, value } = await readWithTimeout(
|
||||
Math.max(deadline - Date.now(), 1),
|
||||
);
|
||||
if (done) break;
|
||||
rawChunks.push(value);
|
||||
peeked += decoder.decode(value, { stream: true });
|
||||
const lines = peeked.split("\n");
|
||||
// The last segment may be a partial line — only parse complete ones.
|
||||
for (const line of lines.slice(0, -1)) {
|
||||
const event = tryParseEvent(line);
|
||||
if (!event?.type) continue;
|
||||
if (event.type === "error") {
|
||||
errorEvent = event;
|
||||
break;
|
||||
}
|
||||
if (MEANINGFUL_EVENT_TYPES.has(event.type)) {
|
||||
committed = true;
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch {
|
||||
// timeout / abort / read failure during the peek → commit whatever we have;
|
||||
// the downstream stream pipeline (stall detection, abort handling) takes over.
|
||||
}
|
||||
|
||||
if (signal) signal.removeEventListener("abort", forwardAbort);
|
||||
|
||||
if (errorEvent) {
|
||||
await reader.cancel("commandcode early error detected").catch(() => {});
|
||||
const { message, type } = formatErrorValue(
|
||||
errorEvent.error ?? errorEvent.message ?? "unknown",
|
||||
);
|
||||
return new Response(JSON.stringify({ error: { message, type } }), {
|
||||
status: HTTP_STATUS.BAD_GATEWAY,
|
||||
statusText: message.slice(0, 200),
|
||||
headers: { "Content-Type": "application/json" },
|
||||
});
|
||||
}
|
||||
|
||||
const remaining = new ReadableStream({
|
||||
start(controller) {
|
||||
(async () => {
|
||||
try {
|
||||
// Re-emit RAW bytes (never re-encoded decoded text) so split
|
||||
// multi-byte UTF-8 sequences survive the peek untouched.
|
||||
for (const c of rawChunks) controller.enqueue(c);
|
||||
while (true) {
|
||||
const { done, value } = await reader.read();
|
||||
if (done) break;
|
||||
controller.enqueue(value);
|
||||
}
|
||||
controller.close();
|
||||
} catch (err) {
|
||||
controller.error(err);
|
||||
}
|
||||
})();
|
||||
},
|
||||
cancel() {
|
||||
reader.cancel("commandcode stream cancelled").catch(() => {});
|
||||
},
|
||||
});
|
||||
|
||||
const combined = new Response(remaining, {
|
||||
status: originalResponse.status,
|
||||
statusText: originalResponse.statusText,
|
||||
headers: originalResponse.headers,
|
||||
});
|
||||
return wrapNdjsonAsOpenAISse(combined, model);
|
||||
}
|
||||
|
||||
function wrapNdjsonAsOpenAISse(originalResponse, model) {
|
||||
const decoder = new TextDecoder();
|
||||
const encoder = new TextEncoder();
|
||||
let buffer = "";
|
||||
const state = { model };
|
||||
function wrapNdjsonAsOpenAISse(streamBody, model, originalResponse = null) {
|
||||
const decoder = new TextDecoder();
|
||||
const encoder = new TextEncoder();
|
||||
let buffer = "";
|
||||
const state = { model };
|
||||
|
||||
const emitChunks = (chunks, controller) => {
|
||||
if (!chunks) return;
|
||||
@@ -253,33 +292,38 @@ function wrapNdjsonAsOpenAISse(originalResponse, model) {
|
||||
}
|
||||
};
|
||||
|
||||
const transform = new TransformStream({
|
||||
transform(chunk, controller) {
|
||||
buffer += decoder.decode(chunk, { stream: true });
|
||||
const lines = buffer.split("\n");
|
||||
buffer = lines.pop() || "";
|
||||
for (const line of lines) {
|
||||
const trimmed = line.trim();
|
||||
if (!trimmed) continue;
|
||||
// Translate AI SDK v5 NDJSON line to one or more OpenAI chunks
|
||||
emitChunks(commandCodeToOpenAIResponse(trimmed, state), controller);
|
||||
}
|
||||
},
|
||||
flush(controller) {
|
||||
const trimmed = buffer.trim();
|
||||
if (trimmed) {
|
||||
emitChunks(commandCodeToOpenAIResponse(trimmed, state), controller);
|
||||
}
|
||||
controller.enqueue(encoder.encode(SSE_DONE));
|
||||
},
|
||||
});
|
||||
const transform = new TransformStream({
|
||||
transform(chunk, controller) {
|
||||
buffer += decoder.decode(chunk, { stream: true });
|
||||
const lines = buffer.split("\n");
|
||||
buffer = lines.pop() || "";
|
||||
for (const line of lines) {
|
||||
const trimmed = line.trim();
|
||||
if (!trimmed) continue;
|
||||
emitChunks(commandCodeToOpenAIResponse(trimmed, state), controller);
|
||||
}
|
||||
},
|
||||
flush(controller) {
|
||||
const trimmed = buffer.trim();
|
||||
if (trimmed) {
|
||||
emitChunks(commandCodeToOpenAIResponse(trimmed, state), controller);
|
||||
}
|
||||
controller.enqueue(encoder.encode(SSE_DONE));
|
||||
},
|
||||
});
|
||||
|
||||
const newBody = originalResponse.body.pipeThrough(transform);
|
||||
return new Response(newBody, {
|
||||
status: originalResponse.status,
|
||||
statusText: originalResponse.statusText,
|
||||
headers: originalResponse.headers,
|
||||
});
|
||||
const newBody = streamBody.pipeThrough(transform);
|
||||
return new Response(newBody, {
|
||||
status: originalResponse?.status || 200,
|
||||
statusText: originalResponse?.statusText || "OK",
|
||||
headers: {
|
||||
"Content-Type": "text/event-stream",
|
||||
"Cache-Control": "no-cache",
|
||||
"Connection": "keep-alive",
|
||||
...(originalResponse?.headers ? Object.fromEntries(originalResponse.headers.entries()) : {}),
|
||||
"content-type": "text/event-stream",
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
export default CommandCodeExecutor;
|
||||
|
||||
Reference in New Issue
Block a user