merge: sync upstream v0.5.81 into MIBP fork
# Conflicts: # .gitignore # Dockerfile # open-sse/handlers/chatCore.js # open-sse/providers/registry/cline.js # open-sse/providers/registry/index.js # open-sse/services/usage.js # open-sse/utils/streamHandler.js # package.json # src/app/(dashboard)/dashboard/profile/page.js # src/app/(dashboard)/dashboard/providers/[id]/page.js
This commit is contained in:
@@ -0,0 +1,80 @@
|
||||
// Codex-specific tool JSON Schema compatibility.
|
||||
//
|
||||
// `https://chatgpt.com/backend-api/codex/responses` validates every function
|
||||
// tool's `parameters` with a regex engine that does not implement Unicode
|
||||
// property escapes. A `pattern` such as
|
||||
//
|
||||
// "^(?!__.*__$)[^\\p{Cc}\\p{Cf}\\p{Zl}\\p{Zp}\"\\\\./\\[\\]]{1,200}$"
|
||||
//
|
||||
// is a perfectly valid ECMAScript `u`-mode regex, but Codex answers
|
||||
//
|
||||
// 400 Invalid schema for function 'Artifact': '^\p{Cc}...' is not a 'regex'
|
||||
// param: tools[0].parameters
|
||||
//
|
||||
// The request is deterministically malformed for this provider, so every
|
||||
// account fails identically and the combo pays a full failover before landing
|
||||
// somewhere that accepts it (#3922).
|
||||
//
|
||||
// Scope guardrail (#3667): this is NOT a global schema sanitizer. Providers
|
||||
// that do support `\p{...}` keep the constraint untouched — the strip runs only
|
||||
// on the Codex dispatch path, and only on `pattern` strings that actually
|
||||
// contain a property escape. Everything else in the schema (including valid
|
||||
// patterns) passes through byte-identical.
|
||||
|
||||
// `\p{...}` / `\P{...}` with an odd number of preceding backslashes — an even
|
||||
// count means the backslash itself is escaped, so `\\p{Cc}` is a literal "p".
|
||||
const UNICODE_PROPERTY_ESCAPE = /(^|[^\\])(\\\\)*\\[pP]\{/;
|
||||
|
||||
export function hasUnicodePropertyEscape(pattern) {
|
||||
return typeof pattern === "string" && UNICODE_PROPERTY_ESCAPE.test(pattern);
|
||||
}
|
||||
|
||||
// Copy-on-write walk: returns the original reference when nothing changed, so
|
||||
// untouched schemas keep object identity and callers can cheaply detect a no-op.
|
||||
// `properties` is special-cased because its keys are arbitrary property *names*
|
||||
// (which may themselves be "pattern" or "properties") and must never be read as
|
||||
// schema keywords; every other key recurses as an ordinary schema node.
|
||||
function stripNode(node, stats) {
|
||||
if (Array.isArray(node)) {
|
||||
let changed = false;
|
||||
const next = node.map((item) => {
|
||||
const cleaned = stripNode(item, stats);
|
||||
if (cleaned !== item) changed = true;
|
||||
return cleaned;
|
||||
});
|
||||
return changed ? next : node;
|
||||
}
|
||||
if (!node || typeof node !== "object") return node;
|
||||
|
||||
let changed = false;
|
||||
const next = {};
|
||||
for (const [key, value] of Object.entries(node)) {
|
||||
if (key === "pattern" && hasUnicodePropertyEscape(value)) {
|
||||
stats.removed++;
|
||||
changed = true;
|
||||
continue;
|
||||
}
|
||||
if (key === "properties" && value && typeof value === "object" && !Array.isArray(value)) {
|
||||
let propsChanged = false;
|
||||
const props = {};
|
||||
for (const [propName, propSchema] of Object.entries(value)) {
|
||||
const cleaned = stripNode(propSchema, stats);
|
||||
if (cleaned !== propSchema) propsChanged = true;
|
||||
props[propName] = cleaned;
|
||||
}
|
||||
if (propsChanged) changed = true;
|
||||
next[key] = propsChanged ? props : value;
|
||||
continue;
|
||||
}
|
||||
const cleaned = stripNode(value, stats);
|
||||
if (cleaned !== value) changed = true;
|
||||
next[key] = cleaned;
|
||||
}
|
||||
return changed ? next : node;
|
||||
}
|
||||
|
||||
// Remove only the `pattern` constraints Codex's validator rejects.
|
||||
// Returns the same reference when the schema is already compatible.
|
||||
export function stripCodexUnsupportedPatterns(schema, stats = { removed: 0 }) {
|
||||
return stripNode(schema, stats);
|
||||
}
|
||||
@@ -49,7 +49,8 @@ export function createSSEStream(options = {}) {
|
||||
connectionId = null,
|
||||
body = null,
|
||||
onStreamComplete = null,
|
||||
apiKey = null
|
||||
apiKey = null,
|
||||
credentials = null
|
||||
} = options;
|
||||
|
||||
let buffer = "";
|
||||
@@ -59,7 +60,7 @@ export function createSSEStream(options = {}) {
|
||||
const decoder = new TextDecoder("utf-8", { fatal: false });
|
||||
|
||||
const state = mode === STREAM_MODE.TRANSLATE
|
||||
? { ...initState(sourceFormat), provider, toolNameMap, customToolNames: new Set(customToolNames || []), model }
|
||||
? { ...initState(sourceFormat), provider, toolNameMap, customToolNames: new Set(customToolNames || []), model, sessionId: credentials?._clientSessionId || null }
|
||||
: null;
|
||||
|
||||
let totalContentLength = 0;
|
||||
@@ -485,7 +486,7 @@ export function createSSEStream(options = {}) {
|
||||
});
|
||||
}
|
||||
|
||||
export function createSSETransformStreamWithLogger(targetFormat, sourceFormat, provider = null, reqLogger = null, toolNameMap = null, model = null, connectionId = null, body = null, onStreamComplete = null, apiKey = null, customToolNames = null) {
|
||||
export function createSSETransformStreamWithLogger(targetFormat, sourceFormat, provider = null, reqLogger = null, toolNameMap = null, model = null, connectionId = null, body = null, onStreamComplete = null, apiKey = null, customToolNames = null, credentials = null) {
|
||||
return createSSEStream({
|
||||
mode: STREAM_MODE.TRANSLATE,
|
||||
targetFormat,
|
||||
@@ -498,7 +499,8 @@ export function createSSETransformStreamWithLogger(targetFormat, sourceFormat, p
|
||||
connectionId,
|
||||
body,
|
||||
onStreamComplete,
|
||||
apiKey
|
||||
apiKey,
|
||||
credentials
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -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();
|
||||
@@ -191,6 +194,7 @@ export function createDisconnectAwareStream(transformStream, streamController, o
|
||||
*/
|
||||
export function pipeWithDisconnect(providerResponse, transformStream, streamController, onAbortTerminal = null, stallTimeoutMs = STREAM_STALL_TIMEOUT_MS) {
|
||||
let stallTimer = null;
|
||||
let abortMessage = "upstream connection lost";
|
||||
const clearStall = () => {
|
||||
if (stallTimer) { clearTimeout(stallTimer); stallTimer = null; }
|
||||
};
|
||||
@@ -198,6 +202,7 @@ export function pipeWithDisconnect(providerResponse, transformStream, streamCont
|
||||
clearStall();
|
||||
stallTimer = setTimeout(() => {
|
||||
stallTimer = null;
|
||||
abortMessage = "stream stall timeout";
|
||||
streamController.handleError?.(new Error("stream stall timeout"));
|
||||
streamController.abort?.();
|
||||
}, stallTimeoutMs);
|
||||
@@ -233,7 +238,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