File size: 7,438 Bytes
88c4c60 2dd17bb 88c4c60 2dd17bb 50f198f 2dd17bb 88c4c60 2dd17bb 88c4c60 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 | import { HTTP_STATUS, RETRY_CONFIG, DEFAULT_RETRY_CONFIG, resolveRetryEntry, FETCH_CONNECT_TIMEOUT_MS } from "../config/runtimeConfig.js";
import { shouldRefreshCredentials } from "../services/oauthCredentialManager.js";
import { proxyAwareFetch } from "../utils/proxyFetch.js";
import { dbg } from "../utils/debugLog.js";
import { getSettings } from "@/lib/localDb";
/**
* BaseExecutor - Base class for provider executors
*/
export class BaseExecutor {
constructor(provider, config) {
this.provider = provider;
this.config = config;
this.noAuth = config?.noAuth || false;
}
getProvider() {
return this.provider;
}
getBaseUrls() {
return this.config.baseUrls || (this.config.baseUrl ? [this.config.baseUrl] : []);
}
getFallbackCount() {
return this.getBaseUrls().length || 1;
}
buildUrl(model, stream, urlIndex = 0, credentials = null) {
if (this.provider?.startsWith?.("openai-compatible-")) {
const baseUrl = credentials?.providerSpecificData?.baseUrl || "https://api.openai.com/v1";
const normalized = baseUrl.replace(/\/$/, "");
const path = this.provider.includes("responses") ? "/responses" : "/chat/completions";
return `${normalized}${path}`;
}
if (this.provider?.startsWith?.("anthropic-compatible-")) {
const baseUrl = credentials?.providerSpecificData?.baseUrl || "https://api.anthropic.com/v1";
const normalized = baseUrl.replace(/\/$/, "");
return `${normalized}/messages`;
}
const baseUrls = this.getBaseUrls();
return baseUrls[urlIndex] || baseUrls[0] || this.config.baseUrl;
}
buildHeaders(credentials, stream = true) {
const headers = {
"Content-Type": "application/json",
...this.config.headers
};
if (this.provider?.startsWith?.("anthropic-compatible-")) {
// Anthropic-compatible providers use x-api-key header
if (credentials.apiKey) {
headers["x-api-key"] = credentials.apiKey;
} else if (credentials.accessToken) {
headers["Authorization"] = `Bearer ${credentials.accessToken}`;
}
if (!headers["anthropic-version"]) {
headers["anthropic-version"] = "2023-06-01";
}
} else {
// Standard Bearer token auth for other providers
if (credentials.accessToken) {
headers["Authorization"] = `Bearer ${credentials.accessToken}`;
} else if (credentials.apiKey) {
headers["Authorization"] = `Bearer ${credentials.apiKey}`;
}
}
if (stream) {
headers["Accept"] = "text/event-stream";
}
return headers;
}
// Override in subclass for provider-specific transformations
transformRequest(model, body, stream, credentials) {
return body;
}
shouldRetry(status, urlIndex) {
return status === HTTP_STATUS.RATE_LIMITED && urlIndex + 1 < this.getFallbackCount();
}
// Override in subclass for provider-specific refresh
async refreshCredentials(credentials, log, proxyOptions = null) {
return null;
}
needsRefresh(credentials) {
return shouldRefreshCredentials(this.provider, credentials);
}
parseError(response, bodyText) {
return { status: response.status, message: bodyText || `HTTP ${response.status}` };
}
async execute({ model, body, stream, credentials, signal, log, proxyOptions = null }) {
const fallbackCount = this.getFallbackCount();
let lastError = null;
let lastStatus = 0;
const retryAttemptsByUrl = {};
// Merge default retry config with provider-specific config
const retryConfig = { ...DEFAULT_RETRY_CONFIG, ...this.config.retry };
// Resolve connect timeout: settings DB > env var > hardcoded default
let connectTimeoutMs = FETCH_CONNECT_TIMEOUT_MS;
try {
const settings = await getSettings();
if (settings.fetchConnectTimeoutMs > 0) {
connectTimeoutMs = settings.fetchConnectTimeoutMs * 1000;
}
} catch (_) { /* fallback to env/hardcoded */ }
// Schedule retry via retryConfig[statusKey]. Returns true when caller should `urlIndex--; continue`
const tryRetry = async (urlIndex, statusKey, reason) => {
const { attempts, delayMs } = resolveRetryEntry(retryConfig[statusKey]);
if (attempts <= 0 || retryAttemptsByUrl[urlIndex] >= attempts) return false;
retryAttemptsByUrl[urlIndex]++;
log?.debug?.("RETRY", `${reason} retry ${retryAttemptsByUrl[urlIndex]}/${attempts} after ${delayMs / 1000}s`);
await new Promise(resolve => setTimeout(resolve, delayMs));
return true;
};
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);
if (!retryAttemptsByUrl[urlIndex]) retryAttemptsByUrl[urlIndex] = 0;
// Abort if upstream doesn't return response headers within connection timeout
const connectCtrl = new AbortController();
const timeoutMs = this.config?.timeoutMs || connectTimeoutMs;
const connectTimer = setTimeout(() => connectCtrl.abort(new Error("fetch connect timeout")), timeoutMs);
const mergedSignal = signal ? AbortSignal.any([signal, connectCtrl.signal]) : connectCtrl.signal;
try {
const bodyStr = JSON.stringify(transformedBody);
const fetchT0 = Date.now();
dbg("FETCH", `${this.provider.toUpperCase()} → ${url} | body=${bodyStr.length}B | connectTimeout=${timeoutMs}ms`);
const response = await proxyAwareFetch(url, {
method: "POST",
headers,
body: bodyStr,
signal: mergedSignal
}, proxyOptions);
clearTimeout(connectTimer);
const ct = response.headers?.get?.("content-type") || "";
const cl = response.headers?.get?.("content-length") || "?";
dbg("FETCH", `${this.provider.toUpperCase()} ← ${response.status} | ttft=${Date.now() - fetchT0}ms | ct=${ct} | cl=${cl}`);
if (await tryRetry(urlIndex, response.status, `status ${response.status}`)) { urlIndex--; continue; }
if (this.shouldRetry(response.status, urlIndex)) {
log?.debug?.("RETRY", `${response.status} on ${url}, trying fallback ${urlIndex + 1}`);
lastStatus = response.status;
continue;
}
return { response, url, headers, transformedBody };
} catch (error) {
clearTimeout(connectTimer);
lastError = error;
const isConnectTimeout = connectCtrl.signal.aborted && error.name === "AbortError";
dbg("FETCH", `${this.provider.toUpperCase()} ✖ ${error.name}: ${error.message}${isConnectTimeout ? " (connect timeout)" : ""}`);
// Connect timeout is internal — convert to retryable network error, don't propagate AbortError
if (error.name === "AbortError" && !isConnectTimeout) throw error;
// Map network/fetch exceptions to 502 retry config
if (await tryRetry(urlIndex, HTTP_STATUS.BAD_GATEWAY, `network "${error.message}"`)) { urlIndex--; continue; }
if (urlIndex + 1 < fallbackCount) {
log?.debug?.("RETRY", `Error on ${url}, trying fallback ${urlIndex + 1}`);
continue;
}
throw error;
}
}
throw lastError || new Error(`All ${fallbackCount} URLs failed with status ${lastStatus}`);
}
}
export default BaseExecutor;
|