a lot
This commit is contained in:
@@ -140,6 +140,7 @@ export function translateAnthropicModels(payload) {
|
||||
object: 'list',
|
||||
data: (payload.data || []).map((model) => ({
|
||||
id: model.id, display_name: model.name || model.id, created_at: new Date((model.created || 0) * 1000 || Date.now()).toISOString(), type: 'model',
|
||||
context_length: model.context_length, context_window: model.context_window, max_context_tokens: model.max_context_tokens,
|
||||
})),
|
||||
has_more: false, first_id: payload.data?.[0]?.id, last_id: payload.data?.at(-1)?.id,
|
||||
};
|
||||
|
||||
@@ -18,6 +18,26 @@ export async function loadKeys(path = new URL('../keys.txt', import.meta.url)) {
|
||||
return keys;
|
||||
}
|
||||
|
||||
export function parseKeysJson(value) {
|
||||
let parsed;
|
||||
try {
|
||||
parsed = JSON.parse(value);
|
||||
} catch (error) {
|
||||
throw new Error(`OPENCODE_API_KEYS must be valid JSON: ${error.message}`);
|
||||
}
|
||||
if (!Array.isArray(parsed) || parsed.some((key) => typeof key !== 'string' || !key.trim())) {
|
||||
throw new Error('OPENCODE_API_KEYS must be a non-empty JSON array of strings');
|
||||
}
|
||||
const keys = parsed.map((key) => key.trim()).filter(Boolean);
|
||||
if (keys.length === 0) throw new Error('OPENCODE_API_KEYS must contain at least one key');
|
||||
return keys;
|
||||
}
|
||||
|
||||
export async function loadConfiguredKeys({ env = process.env, path } = {}) {
|
||||
if (env.OPENCODE_API_KEYS !== undefined) return parseKeysJson(env.OPENCODE_API_KEYS);
|
||||
return loadKeys(path);
|
||||
}
|
||||
|
||||
export function createKeyRotator(keys) {
|
||||
if (!Array.isArray(keys) || keys.length === 0) {
|
||||
throw new Error('At least one OpenCode Go API key is required');
|
||||
|
||||
@@ -0,0 +1,7 @@
|
||||
import fs from 'node:fs/promises';
|
||||
import { parseKeys } from './key-store.js';
|
||||
|
||||
const contents = await fs.readFile(new URL('../keys.txt', import.meta.url), 'utf8');
|
||||
const keys = parseKeys(contents);
|
||||
if (keys.length === 0) throw new Error('No usable keys found in keys.txt');
|
||||
console.log(JSON.stringify(keys));
|
||||
+82
-9
@@ -1,6 +1,7 @@
|
||||
import { Readable } from 'node:stream';
|
||||
import express from 'express';
|
||||
import { createKeyRotator } from './key-store.js';
|
||||
import { clientIp, createRateLimiter, numericCost, rateLimitError } from './rate-limit.js';
|
||||
import {
|
||||
anthropicError,
|
||||
translateAnthropicModels,
|
||||
@@ -24,7 +25,36 @@ async function forwardError(upstream, response) {
|
||||
response.status(upstream.status).send(body || upstream.statusText);
|
||||
}
|
||||
|
||||
export function createProxyApp({ keys, fetchImpl = globalThis.fetch, upstreamBaseUrl = DEFAULT_UPSTREAM_BASE_URL } = {}) {
|
||||
export const DEFAULT_MAX_CONTEXT_TOKENS = 262144;
|
||||
|
||||
function maxContextFromEnv(value = process.env.MAX_CONTEXT_TOKENS) {
|
||||
if (value === undefined) return DEFAULT_MAX_CONTEXT_TOKENS;
|
||||
const parsed = Number(value);
|
||||
if (!Number.isSafeInteger(parsed) || parsed <= 0) throw new Error('MAX_CONTEXT_TOKENS must be a positive integer');
|
||||
return parsed;
|
||||
}
|
||||
|
||||
function validateContext(body, maxContextTokens) {
|
||||
if (body?.max_tokens !== undefined && body.max_tokens > maxContextTokens) {
|
||||
const error = new Error(`max_tokens exceeds the global context limit of ${maxContextTokens} tokens`);
|
||||
error.status = 400;
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
function withContextLimit(payload, maxContextTokens) {
|
||||
return {
|
||||
...payload,
|
||||
data: (payload.data || []).map((model) => ({
|
||||
...model,
|
||||
context_length: maxContextTokens,
|
||||
context_window: maxContextTokens,
|
||||
max_context_tokens: maxContextTokens,
|
||||
})),
|
||||
};
|
||||
}
|
||||
|
||||
export function createProxyApp({ keys, fetchImpl = globalThis.fetch, upstreamBaseUrl = DEFAULT_UPSTREAM_BASE_URL, maxContextTokens = maxContextFromEnv(), rateLimiter = createRateLimiter() } = {}) {
|
||||
const rotator = createKeyRotator(keys);
|
||||
const app = express();
|
||||
app.use(express.json({ limit: '10mb' }));
|
||||
@@ -59,8 +89,8 @@ export function createProxyApp({ keys, fetchImpl = globalThis.fetch, upstreamBas
|
||||
headers: { accept: 'application/json' },
|
||||
});
|
||||
if (!upstream.ok) return forwardError(upstream, response);
|
||||
copyResponseHeaders(upstream, response);
|
||||
response.status(upstream.status).send(await upstream.text());
|
||||
const payload = withContextLimit(JSON.parse(await upstream.text()), maxContextTokens);
|
||||
response.type('application/json').status(upstream.status).send(JSON.stringify(payload));
|
||||
} catch (error) {
|
||||
response.status(502).json({ error: { message: `Unable to reach OpenCode Go: ${error.message}`, type: 'upstream_error' } });
|
||||
}
|
||||
@@ -70,7 +100,7 @@ export function createProxyApp({ keys, fetchImpl = globalThis.fetch, upstreamBas
|
||||
try {
|
||||
const upstream = await fetchWithKeyRetries(`${upstreamBaseUrl}/models`, { headers: { accept: 'application/json' } });
|
||||
if (!upstream.ok) return forwardError(upstream, response);
|
||||
const payload = JSON.parse(await upstream.text());
|
||||
const payload = withContextLimit(JSON.parse(await upstream.text()), maxContextTokens);
|
||||
response.type('application/json').send(JSON.stringify(translateAnthropicModels(payload)));
|
||||
} catch (error) {
|
||||
response.status(502).json({ type: 'error', error: { type: 'upstream_error', message: `Unable to reach OpenCode Go: ${error.message}` } });
|
||||
@@ -79,6 +109,10 @@ export function createProxyApp({ keys, fetchImpl = globalThis.fetch, upstreamBas
|
||||
|
||||
app.post('/oai/v1/chat/completions', async (request, response) => {
|
||||
try {
|
||||
validateContext(request.body, maxContextTokens);
|
||||
const ip = clientIp(request);
|
||||
const limit = rateLimiter.check(ip);
|
||||
if (!limit.allowed) return rateLimitError(response, false, limit);
|
||||
const upstream = await fetchWithKeyRetries(`${upstreamBaseUrl}/chat/completions`, {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
@@ -93,14 +127,40 @@ export function createProxyApp({ keys, fetchImpl = globalThis.fetch, upstreamBas
|
||||
|
||||
if (request.body?.stream && upstream.body) {
|
||||
response.status(upstream.status);
|
||||
Readable.fromWeb(upstream.body).pipe(response);
|
||||
let streamBuffer = '';
|
||||
let recordedCost = false;
|
||||
for await (const chunk of Readable.fromWeb(upstream.body)) {
|
||||
const text = chunk.toString();
|
||||
response.write(text);
|
||||
streamBuffer += text;
|
||||
const events = streamBuffer.split(/\n\n/);
|
||||
streamBuffer = events.pop() || '';
|
||||
for (const eventText of events) {
|
||||
const data = eventText.split('\n').find((line) => line.startsWith('data:'))?.slice(5).trim();
|
||||
if (!recordedCost && data && data !== '[DONE]') {
|
||||
try {
|
||||
const cost = numericCost(JSON.parse(data));
|
||||
if (cost !== null) { rateLimiter.recordCost(ip, cost); recordedCost = true; }
|
||||
} catch { /* ignore malformed stream events */ }
|
||||
}
|
||||
}
|
||||
}
|
||||
if (!recordedCost) {
|
||||
const data = streamBuffer.split('\n').find((line) => line.startsWith('data:'))?.slice(5).trim();
|
||||
try { const cost = numericCost(JSON.parse(data)); if (cost !== null) rateLimiter.recordCost(ip, cost); } catch { /* ignore */ }
|
||||
}
|
||||
response.end();
|
||||
return;
|
||||
}
|
||||
|
||||
response.status(upstream.status).send(await upstream.text());
|
||||
const body = await upstream.text();
|
||||
try { rateLimiter.recordCost(ip, numericCost(JSON.parse(body))); } catch { /* non-JSON upstream response */ }
|
||||
response.status(upstream.status).send(body);
|
||||
} catch (error) {
|
||||
if (!response.headersSent) {
|
||||
response.status(502).json({ error: { message: `Unable to reach OpenCode Go: ${error.message}`, type: 'upstream_error' } });
|
||||
response.status(error.status || 502).json(error.status === 400
|
||||
? { error: { message: error.message, type: 'invalid_request_error' } }
|
||||
: { error: { message: `Unable to reach OpenCode Go: ${error.message}`, type: 'upstream_error' } });
|
||||
}
|
||||
}
|
||||
});
|
||||
@@ -115,6 +175,10 @@ export function createProxyApp({ keys, fetchImpl = globalThis.fetch, upstreamBas
|
||||
}
|
||||
|
||||
try {
|
||||
validateContext(translated, maxContextTokens);
|
||||
const ip = clientIp(request);
|
||||
const limit = rateLimiter.check(ip);
|
||||
if (!limit.allowed) return rateLimitError(response, true, limit);
|
||||
const upstream = await fetchWithKeyRetries(`${upstreamBaseUrl}/chat/completions`, {
|
||||
method: 'POST',
|
||||
headers: { 'content-type': 'application/json', accept: translated.stream ? 'text/event-stream' : 'application/json' },
|
||||
@@ -123,6 +187,7 @@ export function createProxyApp({ keys, fetchImpl = globalThis.fetch, upstreamBas
|
||||
if (!upstream.ok) return forwardError(upstream, response);
|
||||
if (!translated.stream) {
|
||||
const payload = JSON.parse(await upstream.text());
|
||||
rateLimiter.recordCost(ip, numericCost(payload));
|
||||
return response.type('application/json').send(JSON.stringify(translateAnthropicResponse(payload, translated.model)));
|
||||
}
|
||||
if (!upstream.body) return response.status(502).json({ type: 'error', error: { type: 'upstream_error', message: 'Upstream returned no streaming body' } });
|
||||
@@ -132,12 +197,18 @@ export function createProxyApp({ keys, fetchImpl = globalThis.fetch, upstreamBas
|
||||
const decoder = new TextDecoder();
|
||||
let buffer = '';
|
||||
const state = { id: `msg_${Date.now()}`, model: translated.model, nextBlockIndex: 0 };
|
||||
let recordedCost = false;
|
||||
const writeEvents = (text) => {
|
||||
for (const eventText of text.split('\n\n')) {
|
||||
if (!eventText.trim()) continue;
|
||||
const data = eventText.split('\n').find((line) => line.startsWith('data:'))?.slice(5).trim();
|
||||
if (!data || data === '[DONE]') continue;
|
||||
try { response.write(translateOpenAIChunk(JSON.parse(data), state)); } catch { /* ignore malformed upstream event */ }
|
||||
try {
|
||||
const payload = JSON.parse(data);
|
||||
const cost = numericCost(payload);
|
||||
if (!recordedCost && cost !== null) { rateLimiter.recordCost(ip, cost); recordedCost = true; }
|
||||
response.write(translateOpenAIChunk(payload, state));
|
||||
} catch { /* ignore malformed upstream event */ }
|
||||
}
|
||||
};
|
||||
while (true) {
|
||||
@@ -151,7 +222,9 @@ export function createProxyApp({ keys, fetchImpl = globalThis.fetch, upstreamBas
|
||||
writeEvents(buffer);
|
||||
response.end();
|
||||
} catch (error) {
|
||||
if (!response.headersSent) response.status(502).json({ type: 'error', error: { type: 'upstream_error', message: `Unable to reach OpenCode Go: ${error.message}` } });
|
||||
if (!response.headersSent) response.status(error.status || 502).json(error.status === 400
|
||||
? { type: 'error', error: { type: 'invalid_request_error', message: error.message } }
|
||||
: { type: 'error', error: { type: 'upstream_error', message: `Unable to reach OpenCode Go: ${error.message}` } });
|
||||
else response.end();
|
||||
}
|
||||
});
|
||||
|
||||
@@ -0,0 +1,73 @@
|
||||
const SECOND = 1_000;
|
||||
const FIVE_HOURS = 5 * 60 * 60 * 1_000;
|
||||
|
||||
export const RATE_LIMITS = {
|
||||
requestsPerSecond: 2,
|
||||
requestsPerFiveHours: 100,
|
||||
spendPerFiveHours: 10,
|
||||
};
|
||||
|
||||
function prune(timestamps, now, window) {
|
||||
return timestamps.filter((timestamp) => timestamp > now - window);
|
||||
}
|
||||
|
||||
export function clientIp(request) {
|
||||
const railwayIp = request.get('x-real-ip');
|
||||
return railwayIp?.split(',')[0].trim() || request.ip || request.socket.remoteAddress || 'unknown';
|
||||
}
|
||||
|
||||
export function createRateLimiter({ now = () => Date.now() } = {}) {
|
||||
const clients = new Map();
|
||||
|
||||
function stateFor(ip, timestamp) {
|
||||
let state = clients.get(ip);
|
||||
if (!state) {
|
||||
state = { requests: [], spend: [] };
|
||||
clients.set(ip, state);
|
||||
}
|
||||
state.requests = prune(state.requests, timestamp, FIVE_HOURS);
|
||||
state.spend = state.spend.filter(({ timestamp: spentAt }) => spentAt > timestamp - FIVE_HOURS);
|
||||
return state;
|
||||
}
|
||||
|
||||
return {
|
||||
check(ip) {
|
||||
const timestamp = now();
|
||||
const state = stateFor(ip, timestamp);
|
||||
const recentRequests = state.requests.filter((requestAt) => requestAt > timestamp - SECOND);
|
||||
const spend = state.spend.reduce((total, entry) => total + entry.amount, 0);
|
||||
if (recentRequests.length >= RATE_LIMITS.requestsPerSecond) {
|
||||
return { allowed: false, retryAfter: Math.max(1, Math.ceil((recentRequests[0] + SECOND - timestamp) / 1_000)), reason: 'rate' };
|
||||
}
|
||||
if (state.requests.length >= RATE_LIMITS.requestsPerFiveHours) {
|
||||
return { allowed: false, retryAfter: Math.max(1, Math.ceil((state.requests[0] + FIVE_HOURS - timestamp) / 1_000)), reason: 'rate' };
|
||||
}
|
||||
if (spend >= RATE_LIMITS.spendPerFiveHours) {
|
||||
const firstSpend = state.spend[0];
|
||||
return { allowed: false, retryAfter: Math.max(1, Math.ceil((firstSpend.timestamp + FIVE_HOURS - timestamp) / 1_000)), reason: 'spend' };
|
||||
}
|
||||
state.requests.push(timestamp);
|
||||
return { allowed: true };
|
||||
},
|
||||
|
||||
recordCost(ip, amount) {
|
||||
if (!Number.isFinite(amount) || amount < 0) return;
|
||||
const timestamp = now();
|
||||
const state = stateFor(ip, timestamp);
|
||||
state.spend.push({ timestamp, amount });
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
export function numericCost(payload) {
|
||||
const cost = payload?.cost ?? payload?.usage?.cost;
|
||||
return typeof cost === 'number' && Number.isFinite(cost) ? cost : null;
|
||||
}
|
||||
|
||||
export function rateLimitError(response, anthropic = false, result) {
|
||||
response.set('retry-after', String(result.retryAfter));
|
||||
if (anthropic) {
|
||||
return response.status(429).json({ type: 'error', error: { type: 'rate_limit_error', message: `Rate limit exceeded (${result.reason === 'spend' ? 'spend' : 'request'} limit)` } });
|
||||
}
|
||||
return response.status(429).json({ error: { message: `Rate limit exceeded (${result.reason === 'spend' ? 'spend' : 'request'} limit)`, type: 'rate_limit_error' } });
|
||||
}
|
||||
+3
-3
@@ -1,13 +1,13 @@
|
||||
import path from 'node:path';
|
||||
import { fileURLToPath } from 'node:url';
|
||||
import { createProxyApp } from './proxy.js';
|
||||
import { loadKeys } from './key-store.js';
|
||||
import { loadConfiguredKeys } from './key-store.js';
|
||||
|
||||
const projectRoot = path.dirname(path.dirname(fileURLToPath(import.meta.url)));
|
||||
const keys = await loadKeys(path.join(projectRoot, 'keys.txt'));
|
||||
const keys = await loadConfiguredKeys({ path: path.join(projectRoot, 'keys.txt') });
|
||||
const app = createProxyApp({ keys });
|
||||
const port = Number(process.env.PORT || 4005);
|
||||
|
||||
app.listen(port, '127.0.0.1', () => {
|
||||
app.listen(port, '0.0.0.0', () => {
|
||||
console.log(`OpenCode Go proxy listening at http://localhost:${port}/oai/v1`);
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user