Merge branch 'master' of https://github.com/decolua/9router
Some checks failed
Deploy GitBook to 9router.github.io / build-deploy (push) Has been cancelled
Some checks failed
Deploy GitBook to 9router.github.io / build-deploy (push) Has been cancelled
# Conflicts: # open-sse/utils/usageTracking.js
This commit is contained in:
@@ -13,11 +13,18 @@ function generateBillingHeader(payload) {
|
||||
return `x-anthropic-billing-header: cc_version=${CLAUDE_VERSION}.${buildHash}; cc_entrypoint=${CC_ENTRYPOINT}; cch=${cch};`;
|
||||
}
|
||||
|
||||
// Derive a deterministic UUID-v4-shaped string from a seed (stable per account)
|
||||
function deriveUuid(seed) {
|
||||
const h = createHash("sha256").update(seed).digest("hex");
|
||||
return `${h.slice(0, 8)}-${h.slice(8, 12)}-4${h.slice(13, 16)}-${((parseInt(h[16], 16) & 0x3) | 0x8).toString(16)}${h.slice(17, 20)}-${h.slice(20, 32)}`;
|
||||
}
|
||||
|
||||
// Generate fake user ID in Claude Code 2.1.92+ JSON format:
|
||||
// {"device_id":"<64hex>","account_uuid":"<uuid>","session_id":"<uuid>"}
|
||||
function generateFakeUserID(sessionId) {
|
||||
const deviceId = randomBytes(32).toString("hex");
|
||||
const accountUuid = randomUUID();
|
||||
// device_id/account_uuid derive from apiKey (stable per account), session_id per-conversation
|
||||
function generateFakeUserID(sessionId, apiKey) {
|
||||
const deviceId = apiKey ? createHash("sha256").update(`device:${apiKey}`).digest("hex") : randomBytes(32).toString("hex");
|
||||
const accountUuid = apiKey ? deriveUuid(`account:${apiKey}`) : randomUUID();
|
||||
const sessionUuid = sessionId || randomUUID();
|
||||
return `{"device_id":"${deviceId}","account_uuid":"${accountUuid}","session_id":"${sessionUuid}"}`;
|
||||
}
|
||||
@@ -40,8 +47,11 @@ export function cloakClaudeTools(body) {
|
||||
const clientToolNames = new Set();
|
||||
const clientDeclarations = [];
|
||||
|
||||
// All client tools get renamed with suffix
|
||||
// All client tools get renamed with suffix.
|
||||
// Built-in server tools (web_search_20250305, etc.) carry a `type` and require
|
||||
// an exact reserved `name` — never suffix those or Claude rejects the request.
|
||||
for (const tool of tools) {
|
||||
if (tool.type) { clientDeclarations.push(tool); continue; }
|
||||
const suffixed = suffix(tool.name);
|
||||
toolNameMap.set(suffixed, tool.name);
|
||||
clientToolNames.add(tool.name);
|
||||
@@ -148,7 +158,7 @@ export function applyCloaking(body, apiKey, sessionId) {
|
||||
// Inject fake user ID into metadata (session_id must match X-Claude-Code-Session-Id)
|
||||
const existingUserId = result.metadata?.user_id;
|
||||
if (!existingUserId) {
|
||||
result.metadata = { ...result.metadata, user_id: generateFakeUserID(sessionId) };
|
||||
result.metadata = { ...result.metadata, user_id: generateFakeUserID(sessionId, apiKey) };
|
||||
}
|
||||
|
||||
return result;
|
||||
|
||||
41
open-sse/utils/claudeSignature.js
Normal file
41
open-sse/utils/claudeSignature.js
Normal file
@@ -0,0 +1,41 @@
|
||||
// Claude thinking signature validation (ported from CLIProxyAPI internal/signature).
|
||||
// E-form: single-layer base64, decoded[0] == 0x12 (Claude marker).
|
||||
// R-form: double-layer base64, outer decoded[0] == 'E', inner decoded[0] == 0x12.
|
||||
// Cache prefix "...#sig" stripped before validation.
|
||||
|
||||
const MAX_CLAUDE_SIGNATURE_LEN = 32 * 1024 * 1024;
|
||||
const CLAUDE_SIGNATURE_MARKER = 0x12;
|
||||
|
||||
function stripCachePrefix(rawSignature) {
|
||||
const sig = (rawSignature || "").trim();
|
||||
if (!sig) return "";
|
||||
const idx = sig.indexOf("#");
|
||||
return idx >= 0 ? sig.slice(idx + 1).trim() : sig;
|
||||
}
|
||||
|
||||
export function hasClaudeSignaturePrefix(rawSignature) {
|
||||
const sig = stripCachePrefix(rawSignature);
|
||||
return sig.length > 0 && (sig[0] === "E" || sig[0] === "R");
|
||||
}
|
||||
|
||||
// Strict-ish: validates base64 layers + Claude marker byte.
|
||||
export function isValidClaudeSignature(rawSignature) {
|
||||
const sig = stripCachePrefix(rawSignature);
|
||||
if (!sig || sig.length > MAX_CLAUDE_SIGNATURE_LEN) return false;
|
||||
|
||||
try {
|
||||
if (sig[0] === "E") {
|
||||
const decoded = Buffer.from(sig, "base64");
|
||||
return decoded.length > 0 && decoded[0] === CLAUDE_SIGNATURE_MARKER;
|
||||
}
|
||||
if (sig[0] === "R") {
|
||||
const outer = Buffer.from(sig, "base64");
|
||||
if (!outer.length || outer[0] !== 0x45) return false; // 'E'
|
||||
const inner = Buffer.from(outer.toString(), "base64");
|
||||
return inner.length > 0 && inner[0] === CLAUDE_SIGNATURE_MARKER;
|
||||
}
|
||||
return false;
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
@@ -1,15 +1,12 @@
|
||||
// Some thinking-mode providers (DeepSeek, Kimi, MiniMax, ...) require reasoning_content
|
||||
// to be echoed back on assistant messages. Clients in OpenAI format don't send it,
|
||||
// so we inject a non-empty placeholder to satisfy upstream validation.
|
||||
import { PROVIDERS } from "../config/providers.js";
|
||||
|
||||
const PLACEHOLDER = " ";
|
||||
|
||||
// Provider-level rules: keyed by executor.provider
|
||||
const PROVIDER_RULES = {
|
||||
deepseek: { scope: "all" },
|
||||
minimax: { scope: "all" },
|
||||
"minimax-cn": { scope: "all" }
|
||||
};
|
||||
// Provider-level rules derive from registry transport.reasoningInject (single source)
|
||||
const providerRuleFor = (provider) => PROVIDERS[provider]?.reasoningInject;
|
||||
|
||||
// Model-level rules: matched by predicate against model id
|
||||
const MODEL_RULES = [
|
||||
@@ -71,7 +68,7 @@ function applyDeepSeekV4ProAlias({ provider, model, body }) {
|
||||
}
|
||||
|
||||
export function injectReasoningContent({ provider, model, body }) {
|
||||
const providerRule = PROVIDER_RULES[provider];
|
||||
const providerRule = providerRuleFor(provider);
|
||||
const modelRule = MODEL_RULES.find(r => r.match(model));
|
||||
const rule = providerRule || modelRule;
|
||||
const nextBody = applyDeepSeekV4ProAlias({ provider, model, body });
|
||||
|
||||
@@ -79,4 +79,153 @@ export function generateBinaryStyleId() {
|
||||
*/
|
||||
export function clearSessionStore() {
|
||||
runtimeSessionStore.clear();
|
||||
assistantSessionStore.clear();
|
||||
}
|
||||
|
||||
// Conversation-stable session store: Key = hash(scope+assistant text), Value = { sessionId, lastUsed }
|
||||
const assistantSessionStore = new Map();
|
||||
const ASSISTANT_MIN_LEN = 50;
|
||||
const ASSISTANT_CAP_LEN = 50;
|
||||
const MAX_ASSISTANT_SESSIONS = 5000;
|
||||
|
||||
// Client headers/body fields that carry an upstream session id (priority order)
|
||||
const SESSION_HEADER_KEYS = ["x-session-id", "session-id", "session_id", "x-amp-thread-id", "x-client-request-id"];
|
||||
const CLAUDE_CODE_SESSION_RE = /_session_([a-f0-9-]+)$/;
|
||||
|
||||
function sha16(text) {
|
||||
return crypto.createHash("sha256").update(text).digest("hex").slice(0, 16);
|
||||
}
|
||||
|
||||
// Normalize a session id candidate (trim, length cap)
|
||||
function normalizeSessionId(value) {
|
||||
if (typeof value !== "string") return null;
|
||||
const v = value.trim();
|
||||
if (!v || v.length > 256) return null;
|
||||
return v;
|
||||
}
|
||||
|
||||
// Extract Claude Code session id from metadata.user_id (_session_{uuid} | JSON {session_id})
|
||||
function extractClaudeCodeSession(userId) {
|
||||
if (typeof userId !== "string" || !userId) return null;
|
||||
const m = userId.match(CLAUDE_CODE_SESSION_RE);
|
||||
if (m) return m[1];
|
||||
if (userId[0] === "{") {
|
||||
try { return normalizeSessionId(JSON.parse(userId)?.session_id); } catch { /* noop */ }
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
// Lowercase-key lookup for raw client headers
|
||||
function headerValue(headers, key) {
|
||||
if (!headers || typeof headers !== "object") return null;
|
||||
return normalizeSessionId(headers[key] ?? headers[key.toLowerCase()]);
|
||||
}
|
||||
|
||||
// Read client-provided session id from headers/body (no generation)
|
||||
// Antigravity envelope carries session in request.sessionId; requestId embeds conversation uuid
|
||||
const ANTIGRAVITY_CONV_RE = /^[a-z]+\/([0-9a-f-]{36})\//i;
|
||||
function extractAntigravitySession(body) {
|
||||
const sid = body?.request?.sessionId;
|
||||
if (sid != null && sid !== "") return normalizeSessionId(String(sid));
|
||||
const m = typeof body?.requestId === "string" ? body.requestId.match(ANTIGRAVITY_CONV_RE) : null;
|
||||
return m ? normalizeSessionId(m[1]) : null;
|
||||
}
|
||||
|
||||
function extractClientSessionId(headers, body) {
|
||||
const claude = extractClaudeCodeSession(body?.metadata?.user_id);
|
||||
if (claude) return `claude:${claude}`;
|
||||
const antigravity = extractAntigravitySession(body);
|
||||
if (antigravity) return `antigravity:${antigravity}`;
|
||||
for (const key of SESSION_HEADER_KEYS) {
|
||||
const v = headerValue(headers, key);
|
||||
if (v) return v;
|
||||
}
|
||||
const fromBody =
|
||||
normalizeSessionId(body?.prompt_cache_key) ||
|
||||
normalizeSessionId(body?.session_id) ||
|
||||
normalizeSessionId(body?.conversation_id) ||
|
||||
normalizeSessionId(body?.metadata?.user_id);
|
||||
return fromBody || null;
|
||||
}
|
||||
|
||||
// Accumulate assistant text from OpenAI/Responses-style input/messages (cap-limited)
|
||||
function accumulateAssistantText(body) {
|
||||
const items = Array.isArray(body?.input) ? body.input
|
||||
: Array.isArray(body?.messages) ? body.messages : null;
|
||||
if (!items) return "";
|
||||
let text = "";
|
||||
for (const item of items) {
|
||||
if (item?.role !== "assistant") continue;
|
||||
if (typeof item.content === "string") text += item.content;
|
||||
else if (Array.isArray(item.content)) {
|
||||
for (const c of item.content) text += c?.text || c?.output || "";
|
||||
}
|
||||
if (text.length >= ASSISTANT_CAP_LEN) break;
|
||||
}
|
||||
return text;
|
||||
}
|
||||
|
||||
// Stable session id keyed on accumulated assistant text (avoids collision on identical first user prompt)
|
||||
function assistantTextSessionId(scope, body) {
|
||||
const text = accumulateAssistantText(body);
|
||||
if (text.length < ASSISTANT_MIN_LEN) return null;
|
||||
const hash = sha16(`${scope}:${text.slice(0, ASSISTANT_CAP_LEN)}`);
|
||||
const existing = assistantSessionStore.get(hash);
|
||||
if (existing) {
|
||||
existing.lastUsed = Date.now();
|
||||
return existing.sessionId;
|
||||
}
|
||||
if (assistantSessionStore.size >= MAX_ASSISTANT_SESSIONS) {
|
||||
assistantSessionStore.delete(assistantSessionStore.keys().next().value);
|
||||
}
|
||||
const sessionId = generateBinaryStyleId();
|
||||
assistantSessionStore.set(hash, { sessionId, lastUsed: Date.now() });
|
||||
return sessionId;
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve a conversation-stable session id (generalizes Codex resolveCacheSessionId).
|
||||
* Priority: client session → accumulated-assistant-text hash → workspaceId → per-connection.
|
||||
*
|
||||
* @param {object} opts
|
||||
* @param {object} [opts.headers] - Raw client request headers (lowercase keys)
|
||||
* @param {object} [opts.body] - Parsed request body
|
||||
* @param {string} [opts.connectionId] - Connection identifier (fallback scope)
|
||||
* @param {string} [opts.workspaceId] - Provider workspace id (account-wide fallback)
|
||||
* @param {string} [opts.scope] - Provider scope to isolate cache keys across providers
|
||||
* @returns {string} A stable session id
|
||||
*/
|
||||
export function resolveSessionId({ headers, body, connectionId, workspaceId, scope = "" } = {}) {
|
||||
const client = extractClientSessionId(headers, body);
|
||||
if (client) return client;
|
||||
const fromAssistant = assistantTextSessionId(`${scope}:${connectionId || ""}`, body);
|
||||
if (fromAssistant) return fromAssistant;
|
||||
const ws = normalizeSessionId(workspaceId);
|
||||
if (ws) return ws;
|
||||
return deriveSessionId(connectionId);
|
||||
}
|
||||
|
||||
// Capture session id from request body + credentials (envelope still intact here)
|
||||
export function captureSessionId(body, credentials, connectionId, scope = "") {
|
||||
return resolveSessionId({ headers: credentials?.rawHeaders, body, connectionId, scope });
|
||||
}
|
||||
|
||||
// Convert any session id to Antigravity numeric format "-<int64>" (matches real AG / CLIProxyAPI).
|
||||
// Already-numeric ids (native AG sessionId) pass through unchanged.
|
||||
export function toNumericSessionId(sessionId) {
|
||||
const v = normalizeSessionId(sessionId);
|
||||
if (!v) return null;
|
||||
if (/^-?\d+$/.test(v)) return v;
|
||||
const h = crypto.createHash("sha256").update(v).digest();
|
||||
const n = h.readBigUInt64BE(0) & 0x7fffffffffffffffn;
|
||||
return `-${n.toString()}`;
|
||||
}
|
||||
|
||||
// Cleanup expired assistant-session entries
|
||||
const assistantCleanup = setInterval(() => {
|
||||
const now = Date.now();
|
||||
for (const [key, entry] of assistantSessionStore) {
|
||||
if (now - entry.lastUsed > MEMORY_CONFIG.sessionTtlMs) assistantSessionStore.delete(key);
|
||||
}
|
||||
}, MEMORY_CONFIG.sessionCleanupIntervalMs);
|
||||
if (assistantCleanup.unref) assistantCleanup.unref();
|
||||
|
||||
14
open-sse/utils/sse.js
Normal file
14
open-sse/utils/sse.js
Normal file
@@ -0,0 +1,14 @@
|
||||
export function sseChunk(data) {
|
||||
return `data: ${JSON.stringify(data)}\n\n`;
|
||||
}
|
||||
|
||||
// Build OpenAI chat.completion.chunk SSE frame. Key order: id, object, created, model, choices.
|
||||
export function chatChunkSse({ id, created, model, delta, finishReason = null }) {
|
||||
return sseChunk({
|
||||
id,
|
||||
object: "chat.completion.chunk",
|
||||
created,
|
||||
model,
|
||||
choices: [{ index: 0, delta, finish_reason: finishReason }],
|
||||
});
|
||||
}
|
||||
23
open-sse/utils/sseConstants.js
Normal file
23
open-sse/utils/sseConstants.js
Normal file
@@ -0,0 +1,23 @@
|
||||
// Shared SSE primitives (no imports → safe for executors + stream.js)
|
||||
export const SSE_DONE = "data: [DONE]\n\n";
|
||||
|
||||
export const SSE_HEADERS = {
|
||||
"Content-Type": "text/event-stream",
|
||||
"Cache-Control": "no-cache",
|
||||
"Connection": "keep-alive"
|
||||
};
|
||||
|
||||
// Variant for web-cookie executors behind nginx (disable proxy buffering)
|
||||
export const SSE_HEADERS_NO_BUFFER = {
|
||||
"Content-Type": "text/event-stream",
|
||||
"Cache-Control": "no-cache",
|
||||
"X-Accel-Buffering": "no"
|
||||
};
|
||||
|
||||
// Variant for client-facing SSE responses (adds permissive CORS)
|
||||
export const SSE_HEADERS_CORS = {
|
||||
"Content-Type": "text/event-stream",
|
||||
"Cache-Control": "no-cache",
|
||||
"Connection": "keep-alive",
|
||||
"Access-Control-Allow-Origin": "*"
|
||||
};
|
||||
@@ -7,7 +7,10 @@ import { parseSSELine, hasValuableContent, fixInvalidId, formatSSE } from "./str
|
||||
import { getOpenAIResponsesEventName, isOpenAIResponsesTerminalEvent, formatIncompleteOpenAIResponsesStreamFailure } from "./responsesStreamHelpers.js";
|
||||
import { dbg, isDebugEnabled } from "./debugLog.js";
|
||||
|
||||
import { SSE_DONE, SSE_HEADERS, SSE_HEADERS_NO_BUFFER } from "./sseConstants.js";
|
||||
|
||||
export { COLORS, formatSSE };
|
||||
export { SSE_DONE, SSE_HEADERS, SSE_HEADERS_NO_BUFFER };
|
||||
|
||||
// sharedEncoder is stateless — safe to share across streams
|
||||
const sharedEncoder = new TextEncoder();
|
||||
@@ -69,6 +72,7 @@ export function createSSEStream(options = {}) {
|
||||
let currentOpenAIResponsesEvent = null;
|
||||
let openAIResponsesTerminalSeen = false;
|
||||
let openAIResponsesDoneSent = false;
|
||||
let streamDoneSent = false; // track duplicate [DONE] across transform + flush
|
||||
|
||||
return new TransformStream({
|
||||
transform(chunk, controller) {
|
||||
@@ -164,7 +168,12 @@ export function createSSEStream(options = {}) {
|
||||
output = `data: ${JSON.stringify(parsed)}\n`;
|
||||
injectedUsage = true;
|
||||
}
|
||||
} catch { }
|
||||
} catch {
|
||||
// Skip non-JSON data lines silently — don't forward garbage to clients.
|
||||
// Upstream providers sometimes return plain-text errors (HTML, rate-limit
|
||||
// messages) in the SSE stream that would break downstream JSON decoders.
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
||||
if (!injectedUsage) {
|
||||
@@ -209,9 +218,10 @@ export function createSSEStream(options = {}) {
|
||||
sseEmittedCount++;
|
||||
}
|
||||
|
||||
const output = "data: [DONE]\n\n";
|
||||
reqLogger?.appendConvertedChunk?.(output);
|
||||
controller.enqueue(sharedEncoder.encode(output));
|
||||
// [DONE] not emitted in translate mode — some clients' SSE decoders
|
||||
// fail to parse the OpenAI sentinel on Claude-format translated streams.
|
||||
// message_stop already signals end-of-response; stream close handles it.
|
||||
streamDoneSent = true;
|
||||
if (keepsOpenAIResponsesFormat) openAIResponsesDoneSent = true;
|
||||
continue;
|
||||
}
|
||||
@@ -341,9 +351,11 @@ export function createSSEStream(options = {}) {
|
||||
// Some clients (e.g. OpenClaw) expect the OpenAI-style sentinel:
|
||||
// data: [DONE]\n\n
|
||||
// Without it they can hang until timeout and trigger failover.
|
||||
const doneOutput = "data: [DONE]\n\n";
|
||||
reqLogger?.appendConvertedChunk?.(doneOutput);
|
||||
controller.enqueue(sharedEncoder.encode(doneOutput));
|
||||
if (!streamDoneSent) {
|
||||
const doneOutput = "data: [DONE]\n\n";
|
||||
reqLogger?.appendConvertedChunk?.(doneOutput);
|
||||
controller.enqueue(sharedEncoder.encode(doneOutput));
|
||||
}
|
||||
|
||||
if (onStreamComplete) {
|
||||
onStreamComplete({
|
||||
@@ -404,11 +416,8 @@ export function createSSEStream(options = {}) {
|
||||
openAIResponsesTerminalSeen = true;
|
||||
}
|
||||
|
||||
if (!keepsOpenAIResponsesFormat || !openAIResponsesDoneSent) {
|
||||
const doneOutput = "data: [DONE]\n\n";
|
||||
reqLogger?.appendConvertedChunk?.(doneOutput);
|
||||
controller.enqueue(sharedEncoder.encode(doneOutput));
|
||||
}
|
||||
// [DONE] not emitted in translate mode — see comment above.
|
||||
// Passthrough mode still emits it for standard OpenAI clients.
|
||||
|
||||
if (!hasValidUsage(state?.usage) && totalContentLength > 0) {
|
||||
state.usage = estimateUsage(body, totalContentLength, sourceFormat);
|
||||
|
||||
@@ -128,7 +128,10 @@ export function createDisconnectAwareStream(transformStream, streamController, o
|
||||
controller.enqueue(value);
|
||||
} catch (error) {
|
||||
const wasConnected = streamController.isConnected();
|
||||
streamController.handleError(error);
|
||||
// Controller already closed = downstream ended; not an upstream error, skip noisy log.
|
||||
const msg0 = error?.message || "";
|
||||
const isControllerClosed = msg0.includes("already closed") || msg0.includes("Invalid state");
|
||||
if (!isControllerClosed) streamController.handleError(error);
|
||||
reader.cancel().catch(() => {});
|
||||
writer.abort().catch(() => {});
|
||||
|
||||
|
||||
@@ -298,3 +298,37 @@ export function estimateUsage(body, contentLength, targetFormat = FORMATS.OPENAI
|
||||
targetFormat
|
||||
);
|
||||
}
|
||||
/**
|
||||
* Log usage with cache info (green color)
|
||||
*/
|
||||
export function logUsage(provider, usage, model = null, connectionId = null, apiKey = null) {
|
||||
if (!usage || typeof usage !== "object") return;
|
||||
|
||||
const p = provider?.toUpperCase() || "UNKNOWN";
|
||||
|
||||
// Support both formats:
|
||||
// - OpenAI: prompt_tokens, completion_tokens
|
||||
// - Claude: input_tokens, output_tokens
|
||||
const inTokens = usage?.prompt_tokens || usage?.input_tokens || 0;
|
||||
const outTokens = usage?.completion_tokens || usage?.output_tokens || 0;
|
||||
const accountPrefix = connectionId ? connectionId.slice(0, 8) + "..." : "unknown";
|
||||
|
||||
let msg = `[${getTimeString()}] 📊 ${COLORS.green}[USAGE] ${p} | in=${inTokens} | out=${outTokens} | account=${accountPrefix}${COLORS.reset}`;
|
||||
|
||||
// Add estimated flag if present
|
||||
if (usage.estimated) {
|
||||
msg += ` ${COLORS.yellow}(estimated)${COLORS.reset}`;
|
||||
}
|
||||
|
||||
// Add cache info if present (unified from different formats)
|
||||
const cacheRead = usage.cache_read_input_tokens || usage.cached_tokens || usage.prompt_tokens_details?.cached_tokens;
|
||||
if (cacheRead) msg += ` | cache_read=${cacheRead}`;
|
||||
|
||||
const cacheCreation = usage.cache_creation_input_tokens;
|
||||
if (cacheCreation) msg += ` | cache_create=${cacheCreation}`;
|
||||
|
||||
const reasoning = usage.reasoning_tokens;
|
||||
if (reasoning) msg += ` | reasoning=${reasoning}`;
|
||||
|
||||
console.log(msg);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user