From 98412aa0bb83cfd7b24485e3d84a3dbb23b9f2fe Mon Sep 17 00:00:00 2001 From: luulam Date: Mon, 22 Jun 2026 11:26:25 +0700 Subject: [PATCH] update: sync local changes with latest features --- .commandcode/taste/taste.md | 4 + docker-compose.yml | 23 ++++ open-sse/executors/base.js | 20 ++- open-sse/executors/qoder.js | 76 ++++++++++- open-sse/handlers/chatCore.js | 11 +- .../handlers/chatCore/nonStreamingHandler.js | 17 ++- .../handlers/chatCore/streamingHandler.js | 15 ++- open-sse/services/accountFallback.js | 12 +- package.json | 9 +- .../providers/[id]/AddApiKeyModal.js | 75 +++++++++-- .../dashboard/providers/[id]/page.js | 118 ++++++++++++++++++ src/lib/db/repos/settingsRepo.js | 2 + src/shared/components/EditConnectionModal.js | 79 +++++++++++- src/sse/handlers/chat.js | 6 + src/sse/services/auth.js | 11 +- tests/unit/qoder.test.js | 46 +++++++ 16 files changed, 491 insertions(+), 33 deletions(-) create mode 100644 .commandcode/taste/taste.md create mode 100644 docker-compose.yml diff --git a/.commandcode/taste/taste.md b/.commandcode/taste/taste.md new file mode 100644 index 00000000..f562cacd --- /dev/null +++ b/.commandcode/taste/taste.md @@ -0,0 +1,4 @@ +# Taste (Continuously Learned by [CommandCode][cmd]) + +[cmd]: https://commandcode.ai/ + diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 00000000..5df1dce9 --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,23 @@ +services: + 9router: + build: + context: . + image: 9router:local + container_name: 9router + restart: unless-stopped + env_file: + - .env + environment: + NODE_ENV: production + PORT: 20128 + HOSTNAME: 0.0.0.0 + DATA_DIR: /app/data + BASE_URL: http://localhost:20128 + NEXT_PUBLIC_BASE_URL: http://localhost:20128 + ports: + - "20128:20128" + volumes: + - 9router-data:/app/data + +volumes: + 9router-data: diff --git a/open-sse/executors/base.js b/open-sse/executors/base.js index aaf6ede8..33b16cbf 100644 --- a/open-sse/executors/base.js +++ b/open-sse/executors/base.js @@ -116,8 +116,20 @@ export class BaseExecutor { for (let urlIndex = 0; urlIndex < fallbackCount; urlIndex++) { const url = this.buildUrl(model, stream, urlIndex, credentials); - const transformedBody = this.transformRequest(model, body, stream, credentials); - const headers = this.buildHeaders(credentials, stream); + + // Merge extra body params from provider-specific data (request body takes priority) + let mergedBody = this.transformRequest(model, body, stream, credentials); + const extraBodyParams = credentials?.providerSpecificData?.bodyParams; + if (extraBodyParams && typeof extraBodyParams === "object" && !Array.isArray(extraBodyParams)) { + mergedBody = { ...extraBodyParams, ...mergedBody }; + } + + // Merge extra header params from provider-specific data (request headers take priority) + let headers = this.buildHeaders(credentials, stream); + const extraHeaderParams = credentials?.providerSpecificData?.headerParams; + if (extraHeaderParams && typeof extraHeaderParams === "object" && !Array.isArray(extraHeaderParams)) { + headers = { ...extraHeaderParams, ...headers }; + } if (!retryAttemptsByUrl[urlIndex]) retryAttemptsByUrl[urlIndex] = 0; @@ -128,7 +140,7 @@ export class BaseExecutor { const mergedSignal = signal ? AbortSignal.any([signal, connectCtrl.signal]) : connectCtrl.signal; try { - const bodyStr = JSON.stringify(transformedBody); + const bodyStr = JSON.stringify(mergedBody); const fetchT0 = Date.now(); dbg("FETCH", `${this.provider.toUpperCase()} → ${url} | body=${bodyStr.length}B | connectTimeout=${timeoutMs}ms`); const response = await proxyAwareFetch(url, { @@ -150,7 +162,7 @@ export class BaseExecutor { continue; } - return { response, url, headers, transformedBody }; + return { response, url, headers, transformedBody: mergedBody }; } catch (error) { clearTimeout(connectTimer); lastError = error; diff --git a/open-sse/executors/qoder.js b/open-sse/executors/qoder.js index 8e2f263a..4bf2afac 100644 --- a/open-sse/executors/qoder.js +++ b/open-sse/executors/qoder.js @@ -121,6 +121,68 @@ function truncate(s, n) { return s && s.length > n ? `${s.slice(0, n)}...` : s || ""; } +/** + * Parse Qoder-specific mid-stream error codes into user-friendly messages. + * Qoder embeds errors inside SSE envelope body as JSON: {"code":"112","message":"{...}"} + * + * Known codes: + * 112 → quota/billing exceeded (message contains pricingUrl) + * 113 → model not available for current plan + * + * @param {number} statusVal - Upstream status code from envelope + * @param {string} bodyStr - Raw body string from envelope + * @returns {{ statusCode: number, message: string, errorCode: string|null }} + */ +function parseQoderStreamError(statusVal, bodyStr) { + let code = null; + let innerMessage = ""; + + try { + const inner = JSON.parse(bodyStr); + code = String(inner.code || ""); + innerMessage = inner.message || bodyStr; + + // Code 112: quota exceeded — inner message is JSON with pricingUrl + if (code === "112") { + let pricingUrl = ""; + try { + const msgObj = JSON.parse(innerMessage); + pricingUrl = msgObj.pricingUrl || ""; + } catch { /* innerMessage is not JSON */ } + return { + statusCode: statusVal, + message: pricingUrl + ? `Qoder quota exceeded. Upgrade your plan at: ${pricingUrl}` + : "Qoder quota exceeded. Please check your plan limits.", + errorCode: code, + }; + } + + // Code 113: model not available + if (code === "113") { + return { + statusCode: statusVal, + message: `Qoder model not available for your current plan. ${innerMessage}`, + errorCode: code, + }; + } + } catch { + // bodyStr is not valid JSON — fall through to generic message + } + + // Generic fallback + const statusLabel = statusVal >= 500 ? "server error" + : statusVal === 429 ? "rate limit exceeded" + : statusVal === 403 ? "permission error" + : statusVal === 401 ? "authentication error" + : "request error"; + return { + statusCode: statusVal, + message: `Qoder ${statusLabel}: ${truncate(innerMessage || bodyStr, 200)}`, + errorCode: code, + }; +} + /** * Map the OpenAI-style request body into the exact shape Qoder expects. */ @@ -222,7 +284,7 @@ async function buildQoderRequestBody({ model, body, credentials, log, proxyOptio * and re-emit as `data: \n\n`. Errors become `data: [DONE]\n\n` plus * a synthetic OpenAI error chunk. */ -async function wrapQoderSSE(response, model) { +async function wrapQoderSSE(response, model, midStreamError = {}) { if (!response.ok || !response.body) return response; const decoder = new TextDecoder(); @@ -304,13 +366,15 @@ async function wrapQoderSSE(response, model) { const statusVal = typeof envelope.statusCodeValue === "number" ? envelope.statusCodeValue : 200; const inner = typeof envelope.body === "string" ? envelope.body : ""; if (statusVal !== 200) { - const msg = inner || `upstream status ${statusVal}`; + const parsed = parseQoderStreamError(statusVal, inner); + // Store error in shared state so onStreamComplete can trigger cooldown + midStreamError.error = { status: parsed.statusCode, message: parsed.message, errorCode: parsed.errorCode }; const errChunk = JSON.stringify({ id: `qoder-error-${Date.now()}`, object: "chat.completion.chunk", created: Math.floor(Date.now() / 1000), model, - choices: [{ index: 0, delta: { content: `\n[qoder error ${statusVal}: ${truncate(msg, 200)}]` }, finish_reason: "stop" }], + choices: [{ index: 0, delta: { content: `\n\n${parsed.message}` }, finish_reason: "stop" }], }); controller.enqueue(encoder.encode(`data: ${errChunk}\n\n`)); controller.enqueue(encoder.encode("data: [DONE]\n\n")); @@ -489,8 +553,9 @@ export class QoderExecutor extends BaseExecutor { return { response, url, headers, transformedBody: payload }; } - const wrapped = await wrapQoderSSE(response, `qoder/${qoderKey}`); - return { response: wrapped, url, headers, transformedBody: payload }; + const midStreamError = {}; + const wrapped = await wrapQoderSSE(response, `qoder/${qoderKey}`, midStreamError); + return { response: wrapped, url, headers, transformedBody: payload, midStreamError }; } // Qoder device tokens don't refresh through OAuth — the upstream returns @@ -511,6 +576,7 @@ export default QoderExecutor; // should import QoderExecutor and use its public methods. export const __test__ = { normalizeMessages, + parseQoderStreamError, wrapQoderSSE, buildQoderRequestBody, }; diff --git a/open-sse/handlers/chatCore.js b/open-sse/handlers/chatCore.js index 65281cfa..7764b828 100644 --- a/open-sse/handlers/chatCore.js +++ b/open-sse/handlers/chatCore.js @@ -28,7 +28,7 @@ import { compressMessages, formatRtkLog } from "../rtk/index.js"; * @param {object} options.credentials - Provider credentials * @param {string} options.sourceFormatOverride - Override detected source format (e.g. "openai-responses") */ -export async function handleChatCore({ body, modelInfo, credentials, log, onCredentialsRefreshed, onRequestSuccess, onDisconnect, clientRawRequest, connectionId, userAgent, apiKey, ccFilterNaming, rtkEnabled, cavemanEnabled, cavemanLevel, sourceFormatOverride, providerThinking }) { +export async function handleChatCore({ body, modelInfo, credentials, log, onCredentialsRefreshed, onRequestSuccess, onDisconnect, onMidStreamError, clientRawRequest, connectionId, userAgent, apiKey, ccFilterNaming, rtkEnabled, cavemanEnabled, cavemanLevel, sourceFormatOverride, providerThinking }) { const { provider, model } = modelInfo; const requestStartTime = Date.now(); @@ -185,13 +185,14 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred } // Execute request - let providerResponse, providerUrl, providerHeaders, finalBody; + let providerResponse, providerUrl, providerHeaders, finalBody, midStreamError; try { const result = await executor.execute({ model, body: translatedBody, stream, credentials, signal: streamController.signal, log, proxyOptions }); providerResponse = result.response; providerUrl = result.url; providerHeaders = result.headers; finalBody = result.transformedBody; + midStreamError = result.midStreamError; reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); } catch (error) { trackPendingRequest(model, provider, connectionId, false, true); @@ -258,7 +259,7 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred return createErrorResult(statusCode, errMsg, resetsAtMs); } - const sharedCtx = { provider, model, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess }; + const sharedCtx = { provider, model, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, midStreamError, onMidStreamError }; const appendLog = (extra) => appendRequestLog({ model, provider, connectionId, ...extra }).catch(() => { }); const trackDone = () => trackPendingRequest(model, provider, connectionId, false); @@ -275,8 +276,8 @@ export async function handleChatCore({ body, modelInfo, credentials, log, onCred return result; } - // Streaming response - const { onStreamComplete } = buildOnStreamComplete({ ...sharedCtx }); + // Streaming response (midStreamError + onMidStreamError already in sharedCtx) + const { onStreamComplete } = buildOnStreamComplete(sharedCtx); return handleStreamingResponse({ ...sharedCtx, providerResponse, sourceFormat, targetFormat, userAgent, reqLogger, toolNameMap, streamController, onStreamComplete }); } diff --git a/open-sse/handlers/chatCore/nonStreamingHandler.js b/open-sse/handlers/chatCore/nonStreamingHandler.js index 3c31e7a1..e58b6249 100644 --- a/open-sse/handlers/chatCore/nonStreamingHandler.js +++ b/open-sse/handlers/chatCore/nonStreamingHandler.js @@ -132,7 +132,7 @@ export function translateNonStreamingResponse(responseBody, targetFormat, source /** * Handle non-streaming response from provider. */ -export async function handleNonStreamingResponse({ providerResponse, provider, model, sourceFormat, targetFormat, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, reqLogger, toolNameMap, trackDone, appendLog }) { +export async function handleNonStreamingResponse({ providerResponse, provider, model, sourceFormat, targetFormat, body, stream, translatedBody, finalBody, requestStartTime, connectionId, apiKey, clientRawRequest, onRequestSuccess, reqLogger, toolNameMap, trackDone, appendLog, midStreamError, onMidStreamError }) { trackDone(); const contentType = providerResponse.headers.get("content-type") || ""; let responseBody; @@ -144,6 +144,21 @@ export async function handleNonStreamingResponse({ providerResponse, provider, m appendLog({ status: `FAILED ${HTTP_STATUS.BAD_GATEWAY}` }); return createErrorResult(HTTP_STATUS.BAD_GATEWAY, "Invalid SSE response for non-streaming request"); } + + // Check for mid-stream errors detected during SSE parsing (e.g. Qoder quota exceeded). + // The upstream returned HTTP 200 but the SSE envelope contained a non-200 statusCodeValue. + // Without this check, the error message would be silently served as normal content with HTTP 200. + if (midStreamError?.error) { + const err = midStreamError.error; + if (typeof onMidStreamError === "function") { + onMidStreamError(err).catch(e => { + console.error("[MidStreamError] Failed to apply cooldown:", e.message); + }); + } + appendLog({ status: `FAILED ${err.status || 429}` }); + return createErrorResult(err.status || 429, err.message || "Upstream error detected in SSE stream"); + } + responseBody = parsed; } else { try { diff --git a/open-sse/handlers/chatCore/streamingHandler.js b/open-sse/handlers/chatCore/streamingHandler.js index 757bb5b4..5c5a5ad9 100644 --- a/open-sse/handlers/chatCore/streamingHandler.js +++ b/open-sse/handlers/chatCore/streamingHandler.js @@ -75,8 +75,10 @@ export function handleStreamingResponse({ providerResponse, provider, model, sou /** * Build onStreamComplete callback for streaming usage tracking. + * @param {object} options.midStreamError - Shared state object from executor (filled during streaming) + * @param {function} options.onMidStreamError - Callback to invoke when mid-stream error is detected */ -export function buildOnStreamComplete({ provider, model, connectionId, apiKey, requestStartTime, body, stream, finalBody, translatedBody, clientRawRequest }) { +export function buildOnStreamComplete({ provider, model, connectionId, apiKey, requestStartTime, body, stream, finalBody, translatedBody, clientRawRequest, midStreamError, onMidStreamError }) { const streamDetailId = `${Date.now()}-${Math.random().toString(36).slice(2, 11)}`; const onStreamComplete = (contentObj, usage, ttftAt) => { @@ -87,6 +89,15 @@ export function buildOnStreamComplete({ provider, model, connectionId, apiKey, r const safeContent = contentObj?.content || "[Empty streaming response]"; const safeThinking = contentObj?.thinking || null; + // Check if a mid-stream error was detected during streaming (e.g. Qoder quota exceeded) + // The stream already returned HTTP 200, so the pre-stream error path in chatCore didn't fire. + // Apply cooldown here so the account gets locked for subsequent requests. + if (midStreamError?.error && typeof onMidStreamError === "function") { + onMidStreamError(midStreamError.error).catch(err => { + console.error("[MidStreamError] Failed to apply cooldown:", err.message); + }); + } + saveRequestDetail(buildRequestDetail({ provider, model, connectionId, latency, @@ -95,7 +106,7 @@ export function buildOnStreamComplete({ provider, model, connectionId, apiKey, r providerRequest: finalBody || translatedBody || null, providerResponse: safeContent, response: { content: safeContent, thinking: safeThinking, type: "streaming" }, - status: "success" + status: midStreamError?.error ? "error" : "success" }, { id: streamDetailId })).catch(err => { console.error("[RequestDetail] Failed to update streaming content:", err.message); }); diff --git a/open-sse/services/accountFallback.js b/open-sse/services/accountFallback.js index 8d280da4..42189b29 100644 --- a/open-sse/services/accountFallback.js +++ b/open-sse/services/accountFallback.js @@ -18,9 +18,10 @@ export function getQuotaCooldown(backoffLevel = 0) { * @param {number} status - HTTP status code * @param {string} errorText - Error message text * @param {number} backoffLevel - Current backoff level for exponential backoff + * @param {number} [fixedCooldownMs=0] - When >0, override all cooldowns with this fixed value (ms) * @returns {{ shouldFallback: boolean, cooldownMs: number, newBackoffLevel?: number }} */ -export function checkFallbackError(status, errorText, backoffLevel = 0) { +export function checkFallbackError(status, errorText, backoffLevel = 0, fixedCooldownMs = 0) { const lowerError = errorText ? (typeof errorText === "string" ? errorText : JSON.stringify(errorText)).toLowerCase() : ""; @@ -28,6 +29,9 @@ export function checkFallbackError(status, errorText, backoffLevel = 0) { for (const rule of ERROR_RULES) { // Text-based rule: match substring in error message if (rule.text && lowerError && lowerError.includes(rule.text)) { + if (fixedCooldownMs > 0) { + return { shouldFallback: true, cooldownMs: fixedCooldownMs, newBackoffLevel: 0 }; + } if (rule.backoff) { const newLevel = Math.min(backoffLevel + 1, BACKOFF_CONFIG.maxLevel); return { shouldFallback: true, cooldownMs: getQuotaCooldown(newLevel), newBackoffLevel: newLevel }; @@ -37,6 +41,9 @@ export function checkFallbackError(status, errorText, backoffLevel = 0) { // Status-based rule: match HTTP status code if (rule.status && rule.status === status) { + if (fixedCooldownMs > 0) { + return { shouldFallback: true, cooldownMs: fixedCooldownMs, newBackoffLevel: 0 }; + } if (rule.backoff) { const newLevel = Math.min(backoffLevel + 1, BACKOFF_CONFIG.maxLevel); return { shouldFallback: true, cooldownMs: getQuotaCooldown(newLevel), newBackoffLevel: newLevel }; @@ -46,7 +53,8 @@ export function checkFallbackError(status, errorText, backoffLevel = 0) { } // Default: transient cooldown for any unmatched error - return { shouldFallback: true, cooldownMs: TRANSIENT_COOLDOWN_MS }; + const defaultCooldown = fixedCooldownMs > 0 ? fixedCooldownMs : TRANSIENT_COOLDOWN_MS; + return { shouldFallback: true, cooldownMs: defaultCooldown, newBackoffLevel: 0 }; } /** diff --git a/package.json b/package.json index 91dd2f1d..ed78ff38 100644 --- a/package.json +++ b/package.json @@ -19,6 +19,7 @@ "@dnd-kit/sortable": "^10.0.0", "@dnd-kit/utilities": "^3.2.2", "@monaco-editor/react": "^4.7.0", + "@tailwindcss/postcss": "^4.3.1", "@xyflow/react": "^12.10.1", "bcryptjs": "^3.0.3", "confbox": "^0.2.4", @@ -34,6 +35,8 @@ "node-machine-id": "^1.1.12", "open": "^11.0.0", "ora": "^9.1.0", + "postcss": "^8.5.6", + "prop-types": "^15.8.1", "react": "19.2.4", "react-dom": "19.2.4", "react-is": "^16.13.1", @@ -41,6 +44,7 @@ "selfsigned": "^5.5.0", "socks-proxy-agent": "^8.0.5", "sql.js": "^1.14.1", + "tailwindcss": "^4", "undici": "^7.19.2", "uuid": "^13.0.0", "zustand": "^5.0.10" @@ -50,10 +54,7 @@ }, "comment_better_sqlite3": "kept in optionalDependencies so npm install doesn't fail on systems without build tools — sql.js is used as fallback at runtime", "devDependencies": { - "@tailwindcss/postcss": "^4.1.18", "eslint": "^9", - "eslint-config-next": "16.1.6", - "postcss": "^8.5.6", - "tailwindcss": "^4" + "eslint-config-next": "16.1.6" } } diff --git a/src/app/(dashboard)/dashboard/providers/[id]/AddApiKeyModal.js b/src/app/(dashboard)/dashboard/providers/[id]/AddApiKeyModal.js index 77fd9c07..1f7f45ff 100644 --- a/src/app/(dashboard)/dashboard/providers/[id]/AddApiKeyModal.js +++ b/src/app/(dashboard)/dashboard/providers/[id]/AddApiKeyModal.js @@ -1,6 +1,6 @@ "use client"; -import { useState } from "react"; +import { useState, useCallback } from "react"; import PropTypes from "prop-types"; import { Button, Badge, Input, Modal, Select } from "@/shared/components"; import { AI_PROVIDERS } from "@/shared/constants/providers"; @@ -38,6 +38,9 @@ export default function AddApiKeyModal({ isOpen, provider, providerName, isCompa }); const [cloudflareData, setCloudflareData] = useState({ accountId: "" }); const [region, setRegion] = useState(defaultRegion); + const [extraBodyParams, setExtraBodyParams] = useState(""); + const [extraHeaderParams, setExtraHeaderParams] = useState(""); + const [jsonError, setJsonError] = useState(""); const [validating, setValidating] = useState(false); const [validationResult, setValidationResult] = useState(null); const [saving, setSaving] = useState(false); @@ -45,25 +48,54 @@ export default function AddApiKeyModal({ isOpen, provider, providerName, isCompa const [bulkText, setBulkText] = useState(""); const [bulkResult, setBulkResult] = useState(null); // { success, failed } + const handleJsonChange = useCallback((value, setter) => { + setter(value); + if (value.trim()) { + try { + JSON.parse(value); + setJsonError(""); + } catch { + setJsonError("Invalid JSON"); + } + } else { + setJsonError(""); + } + }, []); + + const parseJsonSafe = (str) => { + if (!str.trim()) return undefined; + try { + return JSON.parse(str); + } catch { + return undefined; + } + }; + const buildProviderSpecificData = () => { + const data = {}; if (isOllamaLocal && formData.ollamaHostUrl.trim()) { - return { baseUrl: formData.ollamaHostUrl.trim() }; + data.baseUrl = formData.ollamaHostUrl.trim(); } if (isAzure) { - return { + Object.assign(data, { azureEndpoint: azureData.azureEndpoint, apiVersion: azureData.apiVersion, deployment: azureData.deployment, organization: azureData.organization, - }; + }); } if (isCloudflareAi) { - return { accountId: cloudflareData.accountId }; + data.accountId = cloudflareData.accountId; } if (providerRegions && region) { - return { region }; + data.region = region; } - return undefined; + // Extra params + const parsedBody = parseJsonSafe(extraBodyParams); + const parsedHeaders = parseJsonSafe(extraHeaderParams); + if (parsedBody) data.bodyParams = parsedBody; + if (parsedHeaders) data.headerParams = parsedHeaders; + return Object.keys(data).length > 0 ? data : undefined; }; const handleValidate = async () => { @@ -295,6 +327,35 @@ export default function AddApiKeyModal({ isOpen, provider, providerName, isCompa

)} + {/* Extra Request Parameters */} +
+ + chevron_right + Extra Request Parameters + +
+
+ +