From 4e7e86a3bef6131a238df8c12894db1f9b3a7a36 Mon Sep 17 00:00:00 2001 From: Owen Qwen Date: Tue, 21 Jul 2026 01:31:53 -0500 Subject: [PATCH] a lot --- README.md | 16 ++++++++ package.json | 1 + src/anthropic.js | 1 + src/key-store.js | 20 +++++++++ src/keys-env.js | 7 ++++ src/proxy.js | 91 +++++++++++++++++++++++++++++++++++++---- src/rate-limit.js | 73 +++++++++++++++++++++++++++++++++ src/server.js | 6 +-- test/key-store.test.js | 8 +++- test/proxy.test.js | 18 ++++++++ test/rate-limit.test.js | 26 ++++++++++++ 11 files changed, 254 insertions(+), 13 deletions(-) create mode 100644 src/keys-env.js create mode 100644 src/rate-limit.js create mode 100644 test/rate-limit.test.js diff --git a/README.md b/README.md index c81d068..c5bfde8 100644 --- a/README.md +++ b/README.md @@ -18,6 +18,22 @@ go-second-key # optional comment # disabled-key ``` +For deployments such as Railway, keys can instead be supplied as a JSON environment variable. When set, `OPENCODE_API_KEYS` takes precedence over `keys.txt`: + +```sh +OPENCODE_API_KEYS='["go-first-key","go-second-key"]' +``` + +To convert `keys.txt` into a copyable JSON value: + +```sh +npm run keys:env +``` + +The global context limit defaults to 256k tokens and can be changed with `MAX_CONTEXT_TOKENS`. + +Inference requests are limited per client IP to 2 requests per second, 100 requests per five hours, and $10 of reported upstream cost per five hours. Railway's `X-Real-IP` header is used to identify clients. + Start the proxy: ```sh diff --git a/package.json b/package.json index d3a63ab..6503775 100644 --- a/package.json +++ b/package.json @@ -7,6 +7,7 @@ "main": "src/server.js", "scripts": { "start": "node src/server.js", + "keys:env": "node src/keys-env.js", "test": "vitest run" }, "engines": { diff --git a/src/anthropic.js b/src/anthropic.js index cff0775..d0ab1d2 100644 --- a/src/anthropic.js +++ b/src/anthropic.js @@ -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, }; diff --git a/src/key-store.js b/src/key-store.js index 608bc4d..9ea96b3 100644 --- a/src/key-store.js +++ b/src/key-store.js @@ -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'); diff --git a/src/keys-env.js b/src/keys-env.js new file mode 100644 index 0000000..e51afea --- /dev/null +++ b/src/keys-env.js @@ -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)); diff --git a/src/proxy.js b/src/proxy.js index b790462..f75aa6d 100644 --- a/src/proxy.js +++ b/src/proxy.js @@ -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(); } }); diff --git a/src/rate-limit.js b/src/rate-limit.js new file mode 100644 index 0000000..4f4a652 --- /dev/null +++ b/src/rate-limit.js @@ -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' } }); +} diff --git a/src/server.js b/src/server.js index b9fbe2f..dca8036 100644 --- a/src/server.js +++ b/src/server.js @@ -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`); }); diff --git a/test/key-store.test.js b/test/key-store.test.js index 429c8ab..e4ca0a5 100644 --- a/test/key-store.test.js +++ b/test/key-store.test.js @@ -1,7 +1,13 @@ import { describe, expect, it } from 'vitest'; -import { createKeyRotator, parseKeys } from '../src/key-store.js'; +import { createKeyRotator, loadConfiguredKeys, parseKeys, parseKeysJson } from '../src/key-store.js'; describe('key store', () => { + it('parses JSON keys and prefers the environment variable', async () => { + expect(parseKeysJson('[" a ","b"]')).toEqual(['a', 'b']); + expect(() => parseKeysJson('{"key":"a"}')).toThrow(); + await expect(loadConfiguredKeys({ env: { OPENCODE_API_KEYS: '["env-key"]' }, path: '/does-not-exist' })).resolves.toEqual(['env-key']); + }); + it('parses keys, blank lines, and comments', () => { expect(parseKeys(' first \n\nsecond # note\n# ignored\n')).toEqual(['first', 'second']); }); diff --git a/test/proxy.test.js b/test/proxy.test.js index a77cf84..958f5ac 100644 --- a/test/proxy.test.js +++ b/test/proxy.test.js @@ -58,6 +58,24 @@ describe('proxy', () => { expect(JSON.parse(fetchImpl.mock.calls[0][1].body)).toEqual(body); }); + it('enforces the configured context and IP limits and records cost', async () => { + const fetchImpl = vi.fn() + .mockResolvedValueOnce(response(JSON.stringify({ id: 'c', cost: 0 }))) + .mockResolvedValue(response(JSON.stringify({ id: 'c', cost: 10 }))); + const app = createProxyApp({ keys: ['key'], fetchImpl, maxContextTokens: 100 }); + await request(app).post('/oai/v1/chat/completions').set('X-Real-IP', '198.51.100.1').send({ model: 'm', max_tokens: 101 }).expect(400); + await request(app).post('/oai/v1/chat/completions').set('X-Real-IP', '198.51.100.1').send({ model: 'm', max_tokens: 1 }).expect(200); + await request(app).post('/oai/v1/chat/completions').set('X-Real-IP', '198.51.100.1').send({ model: 'm', max_tokens: 1 }).expect(200); + await request(app).post('/oai/v1/chat/completions').set('X-Real-IP', '198.51.100.1').send({ model: 'm', max_tokens: 1 }).expect(429); + }); + + it('adds the global context limit to model listings', async () => { + const fetchImpl = vi.fn().mockResolvedValue(response('{"data":[{"id":"m"}]}')); + const app = createProxyApp({ keys: ['key'], fetchImpl, maxContextTokens: 123 }); + const result = await request(app).get('/oai/v1/models').expect(200); + expect(result.body.data[0]).toMatchObject({ context_length: 123, context_window: 123, max_context_tokens: 123 }); + }); + it('passes through streaming responses', async () => { const fetchImpl = vi.fn().mockResolvedValue(response('data: {"delta":"Hi"}\n\ndata: [DONE]\n\n', { headers: { 'content-type': 'text/event-stream' }, diff --git a/test/rate-limit.test.js b/test/rate-limit.test.js new file mode 100644 index 0000000..6267557 --- /dev/null +++ b/test/rate-limit.test.js @@ -0,0 +1,26 @@ +import { describe, expect, it } from 'vitest'; +import { createRateLimiter } from '../src/rate-limit.js'; + +describe('rate limiter', () => { + it('enforces request and spend windows', () => { + let time = 0; + const limiter = createRateLimiter({ now: () => time }); + expect(limiter.check('ip').allowed).toBe(true); + expect(limiter.check('ip').allowed).toBe(true); + expect(limiter.check('ip').allowed).toBe(false); + time = 1_001; + expect(limiter.check('ip').allowed).toBe(true); + limiter.recordCost('ip', 10); + expect(limiter.check('ip').allowed).toBe(false); + time = 5 * 60 * 60 * 1_000 + 1_002; + expect(limiter.check('ip').allowed).toBe(true); + }); + + it('keeps clients independent', () => { + const limiter = createRateLimiter(); + expect(limiter.check('a').allowed).toBe(true); + expect(limiter.check('a').allowed).toBe(true); + expect(limiter.check('a').allowed).toBe(false); + expect(limiter.check('b').allowed).toBe(true); + }); +});