Merge origin/master (v0.5.81) into gitea/new_feature
Resolve conflicts: - streamingHandler.js: merge buildStreamErrorBytes onAbortTerminal + local shouldPersistRequestDetail & streamStatusForContent - capabilities.js: preserve server-injected user-asserted caps and models.dev catalog lookup; wire CommandCode /alpha/generate caps inside resolve() - commandcode.js (services/usage): adopt upstream whoami + billing credits/subscriptions with 5h/weekly rate windows and plan caps - openai-to-commandcode.js: merge toNativeImageBlock (data-URI & http(s) support) and assistant reasoning_content preservation - commandcode-to-openai.js: adopt upstream mid-stream error throw for clean retry and abortion - tests: sync commandcode test suite and exclude .next from vitest config
This commit is contained in:
@@ -95,6 +95,9 @@ export function createStreamController({ onDisconnect, onError, log, provider, m
|
||||
* activity), not here — output of the transform stream may be silent
|
||||
* for long periods while raw bytes still flow (e.g. Kiro EventStream
|
||||
* binary frames buffering, Claude reasoning streams).
|
||||
*
|
||||
* @param {function} [onAbortTerminal] - Receives a human-readable abort
|
||||
* message and returns terminal SSE bytes to emit downstream.
|
||||
*/
|
||||
export function createDisconnectAwareStream(transformStream, streamController, onAbortTerminal = null) {
|
||||
const reader = transformStream.readable.getReader();
|
||||
@@ -194,6 +197,7 @@ export function pipeWithDisconnect(providerResponse, transformStream, streamCont
|
||||
let chunkCount = 0;
|
||||
let totalBytes = 0;
|
||||
let lastChunkAt = Date.now();
|
||||
let abortMessage = "upstream connection lost";
|
||||
const t0 = Date.now();
|
||||
const tag = "STREAM";
|
||||
const clearStall = () => {
|
||||
@@ -203,6 +207,7 @@ export function pipeWithDisconnect(providerResponse, transformStream, streamCont
|
||||
clearStall();
|
||||
stallTimer = setTimeout(() => {
|
||||
stallTimer = null;
|
||||
abortMessage = "stream stall timeout";
|
||||
dbg(tag, `STALL TIMEOUT ${stallTimeoutMs}ms | chunks=${chunkCount} | bytes=${totalBytes} | sinceLast=${Date.now() - lastChunkAt}ms`);
|
||||
streamController.handleError?.(new Error("stream stall timeout"));
|
||||
streamController.abort?.();
|
||||
@@ -249,7 +254,7 @@ export function pipeWithDisconnect(providerResponse, transformStream, streamCont
|
||||
return createDisconnectAwareStream(
|
||||
{ readable: transformedBody, writable: { getWriter: () => ({ abort: () => Promise.resolve() }) } },
|
||||
wrappedController,
|
||||
onAbortTerminal
|
||||
onAbortTerminal ? () => onAbortTerminal(abortMessage) : null
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
@@ -1,4 +1,8 @@
|
||||
import { FORMATS } from "../translator/formats.js";
|
||||
import { buildErrorBody } from "./error.js";
|
||||
import { SSE_DONE } from "./sseConstants.js";
|
||||
|
||||
const sharedEncoder = new TextEncoder();
|
||||
|
||||
// Parse SSE data line
|
||||
export function parseSSELine(line, format = null) {
|
||||
@@ -120,3 +124,24 @@ export function formatSSE(data, sourceFormat) {
|
||||
|
||||
return `data: ${JSON.stringify(data)}\n\n`;
|
||||
}
|
||||
|
||||
// Terminal frames for a stream that aborted after HTTP 200 was already sent, so
|
||||
// the status code can no longer change. OpenAI-compatible clients (openai-python
|
||||
// raises APIError on any `data:` payload carrying an `error` key, checked before
|
||||
// [DONE]) need the error frame first, then [DONE]; Anthropic clients need
|
||||
// `event: error`. Never fabricate a successful finish_reason instead.
|
||||
//
|
||||
// Returns encoded bytes: onAbortTerminal callbacks are enqueued verbatim, same
|
||||
// as buildAbortedResponsesTerminalBytes.
|
||||
//
|
||||
// NOTE: non-SSE client formats (Ollama NDJSON) get an SSE frame here — dead in
|
||||
// practice because detectFormatByEndpoint never resolves to OLLAMA.
|
||||
export function buildStreamErrorBytes(statusCode, message, clientFormat) {
|
||||
const { error } = buildErrorBody(statusCode, message);
|
||||
|
||||
const sse = clientFormat === FORMATS.CLAUDE
|
||||
? formatSSE({ type: "error", error }, FORMATS.CLAUDE)
|
||||
: formatSSE({ error }, clientFormat) + SSE_DONE;
|
||||
|
||||
return sharedEncoder.encode(sse);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user