diff --git a/packages/web/server/lib/opencode/proxy.js b/packages/web/server/lib/opencode/proxy.js index 4547f115..eb4a7c5d 100644 --- a/packages/web/server/lib/opencode/proxy.js +++ b/packages/web/server/lib/opencode/proxy.js @@ -1,7 +1,10 @@ -import express from 'express'; import { createProxyMiddleware } from 'http-proxy-middleware'; -import { shouldForwardProxyResponseHeader } from '../../proxy-headers.js'; +import { + applyForwardProxyResponseHeaders, + collectForwardProxyHeaders, + shouldForwardProxyResponseHeader, +} from '../../proxy-headers.js'; export const registerOpenCodeProxy = (app, deps) => { const { @@ -27,6 +30,91 @@ export const registerOpenCodeProxy = (app, deps) => { } app.set('opencodeProxyConfigured', true); + const isAbortError = (error) => error?.name === 'AbortError'; + + const forwardSseRequest = async (req, res) => { + const abortController = new AbortController(); + const closeUpstream = () => abortController.abort(); + let upstream = null; + let reader = null; + + req.on('close', closeUpstream); + + try { + const requestUrl = typeof req.originalUrl === 'string' && req.originalUrl.length > 0 + ? req.originalUrl + : (typeof req.url === 'string' ? req.url : ''); + const upstreamPath = requestUrl.startsWith('/api') ? requestUrl.slice(4) || '/' : requestUrl; + const headers = collectForwardProxyHeaders(req.headers, getOpenCodeAuthHeaders()); + headers.accept ??= 'text/event-stream'; + headers['cache-control'] ??= 'no-cache'; + + upstream = await fetch(buildOpenCodeUrl(upstreamPath, ''), { + method: 'GET', + headers, + signal: abortController.signal, + }); + + res.status(upstream.status); + applyForwardProxyResponseHeaders(upstream.headers, res); + + const contentType = upstream.headers.get('content-type') || 'text/event-stream'; + const isEventStream = contentType.toLowerCase().includes('text/event-stream'); + + if (!upstream.body) { + res.end(await upstream.text().catch(() => '')); + return; + } + + if (!isEventStream) { + res.end(await upstream.text()); + return; + } + + res.setHeader('Content-Type', contentType); + res.setHeader('Cache-Control', 'no-cache'); + res.setHeader('Connection', 'keep-alive'); + res.setHeader('X-Accel-Buffering', 'no'); + if (typeof res.flushHeaders === 'function') { + res.flushHeaders(); + } + + reader = upstream.body.getReader(); + while (!abortController.signal.aborted) { + const { done, value } = await reader.read(); + if (done) { + break; + } + if (value && value.length > 0) { + res.write(value); + } + } + + res.end(); + } catch (error) { + if (isAbortError(error)) { + return; + } + console.error('[proxy] OpenCode SSE proxy error:', error?.message ?? error); + if (!res.headersSent) { + res.status(503).json({ error: 'OpenCode service unavailable' }); + } else { + res.end(); + } + } finally { + req.off('close', closeUpstream); + try { + if (reader) { + await reader.cancel(); + reader.releaseLock(); + } else if (upstream?.body && !upstream.body.locked) { + await upstream.body.cancel(); + } + } catch { + } + } + }; + // Ensure API prefix is detected before proxying app.use('/api', (_req, _res, next) => { ensureOpenCodeApiPrefix(); @@ -140,7 +228,10 @@ export const registerOpenCodeProxy = (app, deps) => { }); } - // http-proxy-middleware handles SSE, large bodies, timeouts correctly + app.get('/api/global/event', forwardSseRequest); + app.get('/api/event', forwardSseRequest); + + // Generic proxy for non-SSE OpenCode API routes. const apiProxy = createProxyMiddleware({ target: `http://127.0.0.1:${runtime.openCodePort || 3902}`, changeOrigin: true, diff --git a/packages/web/server/opencode-proxy.test.js b/packages/web/server/opencode-proxy.test.js new file mode 100644 index 00000000..b2c2490b --- /dev/null +++ b/packages/web/server/opencode-proxy.test.js @@ -0,0 +1,83 @@ +import { afterEach, describe, expect, it } from 'bun:test'; +import express from 'express'; +import path from 'path'; + +import { registerOpenCodeProxy } from './lib/opencode/proxy.js'; + +const listen = (app, host = '127.0.0.1') => new Promise((resolve, reject) => { + const server = app.listen(0, host, () => resolve(server)); + server.once('error', reject); +}); + +const closeServer = (server) => new Promise((resolve, reject) => { + if (!server) { + resolve(); + return; + } + server.close((error) => { + if (error) { + reject(error); + return; + } + resolve(); + }); +}); + +describe('OpenCode proxy SSE forwarding', () => { + let upstreamServer; + let proxyServer; + + afterEach(async () => { + await closeServer(proxyServer); + await closeServer(upstreamServer); + proxyServer = undefined; + upstreamServer = undefined; + }); + + it('forwards event streams with nginx-safe headers', async () => { + let seenAuthorization = null; + + const upstream = express(); + upstream.get('/global/event', (req, res) => { + seenAuthorization = req.headers.authorization ?? null; + res.setHeader('Content-Type', 'text/event-stream; charset=utf-8'); + res.setHeader('Cache-Control', 'private, max-age=0'); + res.setHeader('X-Upstream-Test', 'ok'); + res.write('data: {"ok":true}\n\n'); + res.end(); + }); + upstreamServer = await listen(upstream); + const upstreamPort = upstreamServer.address().port; + + const app = express(); + registerOpenCodeProxy(app, { + fs: {}, + os: {}, + path, + OPEN_CODE_READY_GRACE_MS: 0, + getRuntime: () => ({ + openCodePort: upstreamPort, + isOpenCodeReady: true, + openCodeNotReadySince: 0, + isRestartingOpenCode: false, + }), + getOpenCodeAuthHeaders: () => ({ Authorization: 'Bearer test-token' }), + buildOpenCodeUrl: (requestPath) => `http://127.0.0.1:${upstreamPort}${requestPath}`, + ensureOpenCodeApiPrefix: () => {}, + }); + proxyServer = await listen(app); + const proxyPort = proxyServer.address().port; + + const response = await fetch(`http://127.0.0.1:${proxyPort}/api/global/event`, { + headers: { Accept: 'text/event-stream' }, + }); + + expect(response.status).toBe(200); + expect(response.headers.get('content-type')).toContain('text/event-stream'); + expect(response.headers.get('cache-control')).toBe('no-cache'); + expect(response.headers.get('x-accel-buffering')).toBe('no'); + expect(response.headers.get('x-upstream-test')).toBe('ok'); + expect(await response.text()).toBe('data: {"ok":true}\n\n'); + expect(seenAuthorization).toBe('Bearer test-token'); + }); +});