diff --git a/docker-compose.yml b/docker-compose.yml index 146b752..9f216a4 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -35,6 +35,19 @@ services: - BOTFIGHTS_CREATOR_PUBKEYS=${BOTFIGHTS_CREATOR_PUBKEYS:-da5e0c1b646bdb13c2300f805b0ca3e5afe5b052c594ce78bac8978d21c3fa39} # SQLite database path (defaults to /app/server/data/botfights.db) # - DB_PATH=/app/server/data/botfights.db + # ── Arena federation (BOT-03) ── + # Set on a NODE instance to make it a thin client of a shared canonical + # arena: every /api/* request is proxied there instead of touching this + # instance's own local SQLite DB. Leave UNSET on the canonical arena + # itself (it stays standalone). Any BotFights instance can be a + # canonical arena for others — this is not hardcoded to one host; the + # Foundation's VPS2 instance is only the well-known default. + # - ARENA_UPSTREAM_URL=http://146.59.87.168:9100 + # Set to 1 ONLY on the canonical arena instance when it sits behind a + # reverse proxy (e.g. nginx-proxy-manager) — makes the arena trust + # cf-connecting-ip/x-real-ip/x-forwarded-for from the proxy for + # per-IP rate limiting. Never set on a node's own proxying instance. + # - TRUSTED_PROXY=1 volumes: botfights-data: diff --git a/server/src/middleware/arena-proxy.test.ts b/server/src/middleware/arena-proxy.test.ts index 1c880b3..c7a8b56 100644 --- a/server/src/middleware/arena-proxy.test.ts +++ b/server/src/middleware/arena-proxy.test.ts @@ -40,6 +40,24 @@ upstream.post('/api/echo', async (c) => { }) }) +upstream.get('/api/sse', (c) => { + const encoder = new TextEncoder() + const stream = new ReadableStream({ + async start(controller) { + controller.enqueue(encoder.encode('event: frame\ndata: {"n":1}\n\n')) + await new Promise((r) => setTimeout(r, 60)) + controller.enqueue(encoder.encode('event: frame\ndata: {"n":2}\n\n')) + await new Promise((r) => setTimeout(r, 60)) + controller.enqueue(encoder.encode('event: frame\ndata: {"n":3}\n\n')) + controller.close() + }, + }) + return new Response(stream, { + status: 200, + headers: { 'content-type': 'text/event-stream' }, + }) +}) + upstream.get('/api/gzip-lie', (c) => { // Upstream actually gzip-compresses the body and declares content-encoding // for the COMPRESSED bytes. undici transparently decompresses on the @@ -174,4 +192,68 @@ describe('arenaProxy', () => { const body = await res.json() as { ok: boolean } expect(body.ok).toBe(true) }) + + it('streams SSE incrementally through the proxy', async () => { + process.env.ARENA_UPSTREAM_URL = upstreamUrl + const app = buildProxyingApp() + + const start = Date.now() + const res = await app.request('/api/sse') + expect(res.status).toBe(200) + expect(res.body).not.toBeNull() + + const reader = res.body!.getReader() + const decoder = new TextDecoder() + let firstFrameAt: number | null = null + let buffer = '' + let frameCount = 0 + + while (true) { + const { value, done } = await reader.read() + if (done) break + buffer += decoder.decode(value, { stream: true }) + const frames = buffer.split('\n\n').filter((f) => f.includes('event: frame')) + if (frames.length > 0 && firstFrameAt === null) { + firstFrameAt = Date.now() + } + frameCount = frames.length + } + + expect(frameCount).toBe(3) + // The first frame must have arrived well before the full ~120ms stream + // finished — proves the proxy piped the stream through instead of + // buffering the whole thing before responding. + expect(firstFrameAt).not.toBeNull() + expect(firstFrameAt! - start).toBeLessThan(100) + }) + + it('forwards the client address in x-forwarded-for', async () => { + process.env.ARENA_UPSTREAM_URL = upstreamUrl + const app = buildProxyingApp() + + const res = await app.request('/api/echo', { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({}), + }) + expect(res.status).toBe(200) + const body = await res.json() as { xff: string | null } + // app.request() drives the Hono app directly (no real Node socket), so the + // remoteAddress lookup this proxy relies on may legitimately be + // undetermined here — the behavioral contract is that the header is + // either forwarded (non-empty) or cleanly absent, never invented/garbage. + if (body.xff !== null) { + expect(body.xff.length).toBeGreaterThan(0) + } + }) + + it('answers 502 when the arena is unreachable', async () => { + process.env.ARENA_UPSTREAM_URL = 'http://127.0.0.1:1' + const app = buildProxyingApp() + + const res = await app.request('/api/bots') + expect(res.status).toBe(502) + const body = await res.json() as { error: string } + expect(body.error).toBeTruthy() + }) }) diff --git a/server/src/middleware/arena-proxy.ts b/server/src/middleware/arena-proxy.ts index a7492a9..94cbb71 100644 --- a/server/src/middleware/arena-proxy.ts +++ b/server/src/middleware/arena-proxy.ts @@ -1,10 +1,15 @@ import type { Context, Next } from 'hono' +import { logger } from '../lib/logger.js' // Requests answered locally even when ARENA_UPSTREAM_URL is set — the manifest // health check must never depend on the canonical arena being reachable, or a // perfectly healthy node container gets marked unhealthy and restart-looped. const LOCAL_BYPASS_PATHS = new Set(['/api/health']) +// SSE fight streams are long-lived by design — never time them out. +const SSE_STREAM_PATH = /^\/api\/fights\/[^/]+\/stream$/ +const NON_STREAM_TIMEOUT_MS = 30_000 + // Headers that must never be copied verbatim between hops (either because // they're connection-scoped, or because copying a stale value corrupts the // forwarded/returned message — e.g. content-length after undici recomputes @@ -37,6 +42,18 @@ function copyForwardHeaders(c: Context): Headers { // Ask the upstream for an uncompressed body — Node's fetch already handles // decoding for us, and forwarding compression bookkeeping is unnecessary. headers.set('accept-encoding', 'identity') + + // Forward the originating client IP so the canonical arena's per-IP rate + // limiting doesn't collapse an entire node's user base into one bucket. + // Skip both headers when the address can't be determined rather than + // inventing a value. + const remoteAddress = (c.env as Record | undefined)?.incoming?.socket?.remoteAddress + if (typeof remoteAddress === 'string' && remoteAddress.length > 0) { + const existingXff = headers.get('x-forwarded-for') + headers.set('x-forwarded-for', existingXff ? `${existingXff}, ${remoteAddress}` : remoteAddress) + if (!headers.has('x-real-ip')) headers.set('x-real-ip', remoteAddress) + } + return headers } @@ -64,8 +81,9 @@ export async function arenaProxy(c: Context, next: Next) { const target = buildTargetUrl(upstream, path, search) const headers = copyForwardHeaders(c) const method = c.req.method + const isStream = SSE_STREAM_PATH.test(path) - const upstreamRes = await fetch(target, { + const init: RequestInit = { method, headers, body: method === 'GET' || method === 'HEAD' ? undefined : c.req.raw.body, @@ -73,7 +91,19 @@ export async function arenaProxy(c: Context, next: Next) { // Node's undici fetch requires `duplex` whenever a streamed body is sent. // `@types/node` 22.13.14 already includes `duplex` on RequestInit. duplex: 'half', - }) + } + // SSE fight streams are long-lived by design — exempt from the timeout. + if (!isStream) { + init.signal = AbortSignal.timeout(NON_STREAM_TIMEOUT_MS) + } + + let upstreamRes: Response + try { + upstreamRes = await fetch(target, init) + } catch (err) { + logger.error('arena-proxy', `upstream unreachable: ${target.origin}`, err) + return c.json({ error: 'Arena unreachable.' }, 502) + } return new Response(upstreamRes.body, { status: upstreamRes.status, diff --git a/server/src/routes/fights.ts b/server/src/routes/fights.ts index 0ce9460..9235edc 100644 --- a/server/src/routes/fights.ts +++ b/server/src/routes/fights.ts @@ -420,6 +420,12 @@ fightsRouter.get('/:id/stream', (c) => { return c.json({ error: 'Too many SSE connections' }, 429) } + // nginx (e.g. nginx-proxy-manager fronting the canonical arena) buffers + // proxied responses by default, which would hold every SSE frame until the + // stream closes. This is the documented opt-out — harmless when no nginx + // sits in front of this instance. + c.header('X-Accel-Buffering', 'no') + return streamSSE(c, async (stream) => { // Track connections activeSSECount++