* fix(proxy): reuse upstream connections for OpenCode API requests `createProxyMiddleware` was constructed without an `agent`, so `http-proxy` fell back to `agent: false`. That disables connection pooling and forces `Connection: close` on every proxied request, consuming one ephemeral port per request. Measured against a real `opencode serve` instance, 200 sequential requests through the proxy created 201 TIME_WAIT entries (1.005 ports/request). With a keep-alive agent the same load creates 0. On macOS the ephemeral range is 16,384 ports and TIME_WAIT lasts 30s, so sustained traffic around 546 req/sec exhausts the pool — after which every process on the host fails to open outbound connections with EADDRNOTAVAIL. `maxSockets: Infinity` preserves the unbounded concurrency of `agent: false`, so this changes connection reuse only, not request throughput. Partially addresses #2915. * fix(proxy): derive proxy agent class from the target scheme Addresses review feedback on #2916. The first commit created an unconditional `http.Agent`, which regresses external OpenCode servers configured over https via `OPENCODE_HOST` (accepted by env-config.js). http-proxy dispatches through `https.request` when the target protocol is `https:` (http-proxy/lib/http-proxy/passes/web-incoming.js:126), and `http.Agent#createConnection` is plain `net.createConnection` — so an http.Agent would open a plaintext socket to a TLS port and fail every proxied request. `agent: false` previously worked for both schemes. `createOpenCodeProxyAgent(target)` now returns an `https.Agent` for https targets and an `http.Agent` otherwise, derived once from `resolveProxyTarget()` at registration so the single shared instance is preserved across `apiProxy` and `interactiveOAuthProxy`. Guarded in both test layers, verified to fail when the selection is reverted to an unconditional http.Agent. `https.Agent` extends `http.Agent`, so the http cases assert `not.toBeInstanceOf(https.Agent)`. * Round 2: fix: resolve the proxy agent lazily so cold starts honor https Addresses the round-2 blocker on #2916. Deriving the agent class at registration is too early: startup-pipeline-runtime.js calls setupProxy() (line 104) before bootstrapOpenCodeAtStartup() (line 141), so on a fresh process state.openCodePort is null, buildOpenCodeUrl() throws (network-runtime.js:86-88), and resolveProxyTarget() returns the http loopback fallback. An external server configured via OPENCODE_HOST=https:// only appears on state.openCodeBaseUrl after bootstrap, so it was still getting a plain http.Agent — the regression the previous commit intended to fix. `agent` is now a getter backed by a per-scheme memoizing resolver. http-proxy-middleware rebuilds per-request options with `Object.assign({}, this.proxyOptions)` in prepareProxyRequest, which invokes getters, so resolution happens at request time while still yielding one shared pool per scheme. Tests now model the production ordering — registration while the port is null and buildOpenCodeUrl throws, then an https base URL appearing after bootstrap — and fail against the eager implementation. A behavioral test pins the http-proxy-middleware option re-read the fix depends on, so a library change that froze options would fail loudly instead of silently regressing https targets. The resolver is module-private; `bun run dead-code` flagged it as an unused export when it was exported. * Round 3: docs(changelog): note upstream connection reuse under [Unreleased] Repo precedent adds [Unreleased] bullets for comparable proxy/stability fixes (1.18.4 Stability, 1.9.3 Reliability/Proxy). Non-blocker raised in review on #2916. * Round 3: docs(changelog): use repo-standard 'behavior' spelling * Round 4: docs(changelog): don't imply a restart is the only recovery The ephemeral port pool drains on its own once the exhausting traffic stops (TIME_WAIT expiry), so a restart is sufficient but not necessary. Optional nit raised in review on #2916. * Round 5: fix: construct the proxy agent through one factory; widen the pool Review found the https branch was mutation-uncovered: the resolver re-implemented agent construction inline instead of calling the exported `createOpenCodeProxyAgent(target)`, so replacing its https branch with `new https.Agent()` — dropping OPENCODE_AGENT_OPTIONS, and with it keep-alive — left the entire suite green. Since `createOpenCodeProxyAgent` also had no production callers, its four tests were pinning dead code. Delegating collapses both: the factory is now the single construction path, and the mutation fails 2 tests including the live resolver path. Also from review: - maxFreeSockets 32 -> 256 (Node's own default). The lower cap evicted pooled sockets under concurrency, reintroducing the churn this agent exists to prevent: at 64 concurrent requests it left 303 sockets in TIME_WAIT versus 0 at 256. - Added `timeout` to OPENCODE_AGENT_OPTIONS. Free-socket eviction is governed by agent.options.timeout, which was unset, so idle sockets persisted until the peer closed them. `keepAliveMsecs` is the TCP probe delay, not the idle lifetime. - resolveProxyTarget() now checks openCodePort before calling buildOpenCodeUrl instead of relying on it throwing. The port is nulled on several runtime paths (health-check failure, failed restart), so a degraded OpenCode made every proxied request pay for a thrown-and-caught exception — and the getter added a second call per request. - Test fixtures use :4096 rather than :443; WHATWG URL elides the default port, so parseInt('') is NaN and env-config rejects that host. The fixtures modeled a state that cannot reach production. - The getter-read assertion is now exact (0 at construction, 1, then 2) rather than >= 2, which would have passed if the getter were read twice at construction and never per-request. - listen() rejects on 'error' and servers start inside try/finally, so a bind failure fails the test instead of hanging to timeout.
958 lines
33 KiB
JavaScript
958 lines
33 KiB
JavaScript
import http from 'node:http';
|
|
import https from 'node:https';
|
|
|
|
import { createProxyMiddleware } from 'http-proxy-middleware';
|
|
|
|
import {
|
|
applyForwardProxyResponseHeaders,
|
|
collectForwardProxyHeaders,
|
|
shouldForwardProxyResponseHeader,
|
|
} from '../../proxy-headers.js';
|
|
import { createRealpathCache } from '../path-realpath-cache.js';
|
|
import { DEFAULT_UPSTREAM_STALL_TIMEOUT_MS } from '../event-stream/upstream-reader.js';
|
|
import { recordStartupPerformance } from './startup-performance.js';
|
|
|
|
const DEFAULT_SSE_HEARTBEAT_INTERVAL_MS = 20_000;
|
|
|
|
const OPENCODE_AGENT_KEEP_ALIVE_MS = 30_000;
|
|
// Node's own default. A lower cap evicts pooled sockets under concurrency,
|
|
// which reintroduces exactly the per-request connection churn this agent
|
|
// exists to prevent (measured: at 64 concurrent requests, a cap of 32 left
|
|
// 303 sockets in TIME_WAIT versus 0 at 256).
|
|
const OPENCODE_AGENT_MAX_FREE_SOCKETS = 256;
|
|
// Evicts idle free sockets from our side. Without it the only thing that
|
|
// retires an idle pooled socket is the upstream closing it. Note this is
|
|
// distinct from `keepAliveMsecs`, which is the TCP keep-alive probe delay.
|
|
const OPENCODE_AGENT_IDLE_TIMEOUT_MS = 60_000;
|
|
|
|
const OPENCODE_AGENT_OPTIONS = {
|
|
keepAlive: true,
|
|
keepAliveMsecs: OPENCODE_AGENT_KEEP_ALIVE_MS,
|
|
maxSockets: Infinity,
|
|
maxFreeSockets: OPENCODE_AGENT_MAX_FREE_SOCKETS,
|
|
timeout: OPENCODE_AGENT_IDLE_TIMEOUT_MS,
|
|
};
|
|
|
|
const isHttpsProxyTarget = (target) => {
|
|
if (typeof target !== 'string') {
|
|
return false;
|
|
}
|
|
try {
|
|
return new URL(target).protocol === 'https:';
|
|
} catch {
|
|
return /^https:/i.test(target.trim());
|
|
}
|
|
};
|
|
|
|
/**
|
|
* Agent for proxied OpenCode API requests.
|
|
*
|
|
* When no agent is supplied, `http-proxy` falls back to `agent: false`, which
|
|
* both disables connection pooling and forces `Connection: close` on every
|
|
* proxied request (http-proxy/lib/http-proxy/common.js). That consumes one
|
|
* ephemeral port per request, and sustained traffic can exhaust the host's
|
|
* ephemeral port range — after which every process on the machine fails to
|
|
* open outbound connections with EADDRNOTAVAIL.
|
|
*
|
|
* The agent must match the target scheme: http-proxy dispatches through
|
|
* `https.request` when `target.protocol === 'https:'`
|
|
* (http-proxy/lib/http-proxy/passes/web-incoming.js), and an `http.Agent`
|
|
* would open a plaintext socket to a TLS port. External servers may be
|
|
* configured over https via `OPENCODE_HOST` (see env-config.js), so derive the
|
|
* agent class from the resolved target.
|
|
*
|
|
* `maxSockets: Infinity` preserves the unbounded concurrency of `agent: false`,
|
|
* so this changes connection reuse only, not request throughput.
|
|
*/
|
|
export const createOpenCodeProxyAgent = (target) => (
|
|
isHttpsProxyTarget(target)
|
|
? new https.Agent(OPENCODE_AGENT_OPTIONS)
|
|
: new http.Agent(OPENCODE_AGENT_OPTIONS)
|
|
);
|
|
|
|
/**
|
|
* Lazily resolves the proxy agent, memoized per scheme.
|
|
*
|
|
* The scheme cannot be decided at registration time: `setupProxy()` runs before
|
|
* `bootstrapOpenCodeAtStartup()` (startup-pipeline-runtime.js), so on a cold
|
|
* start `state.openCodePort` is still null, `buildOpenCodeUrl()` throws
|
|
* (network-runtime.js) and `resolveProxyTarget()` falls back to the http
|
|
* loopback default. An external server configured over https via
|
|
* `OPENCODE_HOST` only becomes visible on `state.openCodeBaseUrl` after
|
|
* bootstrap completes.
|
|
*
|
|
* http-proxy-middleware rebuilds its per-request options with
|
|
* `Object.assign({}, this.proxyOptions)` inside `prepareProxyRequest`, which
|
|
* invokes getters, so exposing `agent` as a getter defers resolution to request
|
|
* time. Memoizing per scheme keeps a single shared pool per scheme rather than
|
|
* allocating an agent per request.
|
|
*/
|
|
const createOpenCodeProxyAgentResolver = (resolveTarget) => {
|
|
const agents = new Map();
|
|
|
|
return () => {
|
|
const target = resolveTarget();
|
|
const scheme = isHttpsProxyTarget(target) ? 'https:' : 'http:';
|
|
let agent = agents.get(scheme);
|
|
if (!agent) {
|
|
// Construct through the shared factory rather than inline, so both
|
|
// schemes are built from OPENCODE_AGENT_OPTIONS by the same code path.
|
|
agent = createOpenCodeProxyAgent(target);
|
|
agents.set(scheme, agent);
|
|
}
|
|
return agent;
|
|
};
|
|
};
|
|
|
|
export const createDirectoryQueryCanonicalizer = ({ realpath, ...cacheOptions } = {}) => {
|
|
const realpathCache = createRealpathCache({ fallbackOnError: true, realpath, ...cacheOptions });
|
|
|
|
return async (requestUrl) => {
|
|
if (typeof requestUrl !== 'string' || !requestUrl.includes('directory=')) {
|
|
return requestUrl;
|
|
}
|
|
|
|
const url = new URL(requestUrl, 'http://localhost');
|
|
const directory = url.searchParams.get('directory');
|
|
if (!directory) {
|
|
return requestUrl;
|
|
}
|
|
|
|
const canonicalDirectory = await realpathCache.resolve(directory);
|
|
if (!canonicalDirectory || canonicalDirectory === directory) {
|
|
return requestUrl;
|
|
}
|
|
|
|
url.searchParams.set('directory', canonicalDirectory);
|
|
return `${url.pathname}${url.search}`;
|
|
};
|
|
};
|
|
|
|
export const normalizeForwardedDirectoryHeaders = (headers) => {
|
|
const rawDirectory = headers?.['x-opencode-directory'];
|
|
if (typeof rawDirectory !== 'string') {
|
|
return headers;
|
|
}
|
|
|
|
if (headers['x-opencode-directory-encoding'] !== 'uri') {
|
|
return headers;
|
|
}
|
|
|
|
try {
|
|
headers['x-opencode-directory'] = decodeURIComponent(rawDirectory);
|
|
} catch {
|
|
// Leave malformed values untouched; upstream will reject invalid paths.
|
|
}
|
|
delete headers['x-opencode-directory-encoding'];
|
|
return headers;
|
|
};
|
|
|
|
const waitForSseDrain = (res, signal) => new Promise((resolve) => {
|
|
if (signal?.aborted || res.writableEnded || res.destroyed) {
|
|
resolve();
|
|
return;
|
|
}
|
|
|
|
const cleanup = () => {
|
|
res.off?.('drain', onDone);
|
|
res.off?.('close', onDone);
|
|
res.off?.('error', onDone);
|
|
signal?.removeEventListener?.('abort', onDone);
|
|
};
|
|
const onDone = () => {
|
|
cleanup();
|
|
resolve();
|
|
};
|
|
|
|
res.once?.('drain', onDone);
|
|
res.once?.('close', onDone);
|
|
res.once?.('error', onDone);
|
|
signal?.addEventListener?.('abort', onDone, { once: true });
|
|
});
|
|
|
|
export const writeSseChunkWithBackpressure = async (res, value, signal) => {
|
|
if (!value || value.length === 0 || signal?.aborted || res.writableEnded || res.destroyed) {
|
|
return false;
|
|
}
|
|
|
|
const flushed = res.write(value);
|
|
if (flushed !== false) {
|
|
return true;
|
|
}
|
|
|
|
await waitForSseDrain(res, signal);
|
|
return !signal?.aborted && !res.writableEnded && !res.destroyed;
|
|
};
|
|
|
|
export const createSseBoundaryTracker = () => {
|
|
const decoder = new TextDecoder();
|
|
let tail = '';
|
|
|
|
const normalize = (value) => value.replace(/\r\n/g, '\n').replace(/\r/g, '\n');
|
|
|
|
return {
|
|
observe(value) {
|
|
const text = typeof value === 'string'
|
|
? value
|
|
: decoder.decode(value, { stream: true });
|
|
if (text.length > 0) {
|
|
tail = `${tail}${normalize(text)}`;
|
|
if (tail.length > 4096) {
|
|
tail = tail.slice(-4096);
|
|
}
|
|
}
|
|
return this.isAtBoundary();
|
|
},
|
|
isAtBoundary() {
|
|
return tail.length === 0 || tail.endsWith('\n\n');
|
|
},
|
|
};
|
|
};
|
|
|
|
const SESSION_LIST_ALLOWED_FIELDS = [
|
|
'id',
|
|
'slug',
|
|
'projectID',
|
|
'workspaceID',
|
|
'directory',
|
|
'path',
|
|
'parentID',
|
|
'title',
|
|
'agent',
|
|
'model',
|
|
'version',
|
|
'time',
|
|
'cost',
|
|
'tokens',
|
|
'share',
|
|
'metadata',
|
|
'project',
|
|
];
|
|
|
|
const sanitizeSessionListItem = (session) => {
|
|
if (!session || typeof session !== 'object' || Array.isArray(session)) {
|
|
return session;
|
|
}
|
|
|
|
const sanitized = {};
|
|
for (const key of SESSION_LIST_ALLOWED_FIELDS) {
|
|
if (key in session) {
|
|
sanitized[key] = session[key];
|
|
}
|
|
}
|
|
|
|
const summary = session.summary;
|
|
if (summary && typeof summary === 'object' && !Array.isArray(summary)) {
|
|
const summaryWithoutDiffs = { ...summary };
|
|
delete summaryWithoutDiffs.diffs;
|
|
sanitized.summary = summaryWithoutDiffs;
|
|
}
|
|
|
|
const revert = session.revert;
|
|
if (revert && typeof revert === 'object' && !Array.isArray(revert)) {
|
|
const revertMarker = {};
|
|
if (typeof revert.messageID === 'string') {
|
|
revertMarker.messageID = revert.messageID;
|
|
}
|
|
if (typeof revert.partID === 'string') {
|
|
revertMarker.partID = revert.partID;
|
|
}
|
|
if (Object.keys(revertMarker).length > 0) {
|
|
sanitized.revert = revertMarker;
|
|
}
|
|
}
|
|
|
|
return sanitized;
|
|
};
|
|
|
|
const sanitizeSessionListPayload = (payload) => {
|
|
if (!Array.isArray(payload)) {
|
|
return payload;
|
|
}
|
|
return payload.map((session) => sanitizeSessionListItem(session));
|
|
};
|
|
|
|
export const registerOpenCodeProxy = (app, deps) => {
|
|
const {
|
|
fs,
|
|
os,
|
|
path,
|
|
OPEN_CODE_READY_GRACE_MS,
|
|
LONG_REQUEST_TIMEOUT_MS,
|
|
getRuntime,
|
|
getOpenCodeAuthHeaders,
|
|
buildOpenCodeUrl,
|
|
ensureOpenCodeApiPrefix,
|
|
SSE_HEARTBEAT_INTERVAL_MS = DEFAULT_SSE_HEARTBEAT_INTERVAL_MS,
|
|
SSE_UPSTREAM_STALL_TIMEOUT_MS = DEFAULT_UPSTREAM_STALL_TIMEOUT_MS,
|
|
getSseUpstreamStallTimeoutMs = () => SSE_UPSTREAM_STALL_TIMEOUT_MS,
|
|
} = deps;
|
|
|
|
if (app.get('opencodeProxyConfigured')) {
|
|
return;
|
|
}
|
|
|
|
const runtime = getRuntime();
|
|
if (runtime.openCodePort) {
|
|
console.log(`Setting up proxy to OpenCode on port ${runtime.openCodePort}`);
|
|
} else {
|
|
console.log('Setting up OpenCode API gate (OpenCode not started yet)');
|
|
}
|
|
app.set('opencodeProxyConfigured', true);
|
|
|
|
const isAbortError = (error) => error?.name === 'AbortError';
|
|
const FALLBACK_PROXY_TARGET = 'http://127.0.0.1:3902';
|
|
const canonicalizeDirectoryQuery = createDirectoryQueryCanonicalizer({
|
|
realpath: fs?.promises?.realpath?.bind(fs.promises),
|
|
});
|
|
|
|
const hasParsedBodyValue = (body) => {
|
|
if (body === undefined || body === null) return false;
|
|
if (Buffer.isBuffer(body)) return body.length > 0;
|
|
if (typeof body === 'string') return body.length > 0;
|
|
if (Array.isArray(body)) return body.length > 0;
|
|
if (typeof body === 'object') return Object.keys(body).length > 0;
|
|
return true;
|
|
};
|
|
|
|
const getContentType = (proxyReq, req) => {
|
|
const value = proxyReq.getHeader?.('content-type') ?? req.headers?.['content-type'] ?? '';
|
|
if (Array.isArray(value)) return value[0] || '';
|
|
return String(value || '');
|
|
};
|
|
|
|
const serializeUrlEncodedBody = (body) => {
|
|
if (!body || typeof body !== 'object' || Buffer.isBuffer(body)) {
|
|
return String(body ?? '');
|
|
}
|
|
|
|
const params = new URLSearchParams();
|
|
for (const [key, value] of Object.entries(body)) {
|
|
if (value === undefined || value === null) continue;
|
|
if (Array.isArray(value)) {
|
|
for (const entry of value) {
|
|
if (entry !== undefined && entry !== null) params.append(key, String(entry));
|
|
}
|
|
continue;
|
|
}
|
|
params.append(key, String(value));
|
|
}
|
|
return params.toString();
|
|
};
|
|
|
|
const serializeParsedBody = (req, proxyReq) => {
|
|
if (req.method === 'GET' || req.method === 'HEAD') return null;
|
|
if (req.body === undefined || req.body === null) return null;
|
|
const originalContentLength = Number.parseInt(req.headers?.['content-length'] || '0', 10) || 0;
|
|
if (!hasParsedBodyValue(req.body) && originalContentLength <= 0) return null;
|
|
|
|
const contentType = getContentType(proxyReq, req).toLowerCase();
|
|
if (Buffer.isBuffer(req.body)) return req.body;
|
|
if (contentType.includes('application/json')) return Buffer.from(JSON.stringify(req.body));
|
|
if (contentType.includes('application/x-www-form-urlencoded')) return Buffer.from(serializeUrlEncodedBody(req.body));
|
|
if (typeof req.body === 'string') return Buffer.from(req.body);
|
|
return null;
|
|
};
|
|
|
|
const replayParsedBody = (proxyReq, req) => {
|
|
const body = serializeParsedBody(req, proxyReq);
|
|
if (!body) return;
|
|
proxyReq.setHeader('content-length', String(body.length));
|
|
proxyReq.write(body);
|
|
};
|
|
|
|
const normalizeProxyTarget = (candidate) => {
|
|
if (typeof candidate !== 'string') {
|
|
return null;
|
|
}
|
|
|
|
const trimmed = candidate.trim();
|
|
if (!trimmed) {
|
|
return null;
|
|
}
|
|
|
|
return trimmed.replace(/\/+$/, '');
|
|
};
|
|
|
|
// Keep generic proxy requests on the same upstream base URL that health checks
|
|
// and direct fetch helpers use. This avoids split-brain state where /health
|
|
// succeeds against an external host but /api/* still proxies to 127.0.0.1.
|
|
const resolveProxyTarget = () => {
|
|
const runtimeState = getRuntime();
|
|
|
|
// `buildOpenCodeUrl` throws while the port is unknown, and the port is
|
|
// nulled on several runtime paths (health-check failure, failed restart),
|
|
// not just cold start. Checking first keeps a degraded OpenCode from
|
|
// making every proxied request pay for a thrown-and-caught exception.
|
|
if (runtimeState.openCodePort) {
|
|
try {
|
|
const resolved = normalizeProxyTarget(buildOpenCodeUrl('/', ''));
|
|
if (resolved) {
|
|
return resolved;
|
|
}
|
|
} catch {
|
|
}
|
|
}
|
|
|
|
const externalBase = normalizeProxyTarget(runtimeState.openCodeBaseUrl);
|
|
if (externalBase) {
|
|
return externalBase;
|
|
}
|
|
|
|
return FALLBACK_PROXY_TARGET;
|
|
};
|
|
|
|
const normalizeProxyTimeout = (value) => {
|
|
return Number.isFinite(value) && value > 0 ? value : 4 * 60 * 1000;
|
|
};
|
|
|
|
const PROXY_REQUEST_TIMEOUT_MS = normalizeProxyTimeout(LONG_REQUEST_TIMEOUT_MS);
|
|
const PROXY_TIMEOUT_MARKER = Symbol('openchamberProxyTimedOut');
|
|
|
|
// A provider OAuth callback blocks upstream for as long as the user takes to
|
|
// sign in in their browser (device-code polling, or a loopback redirect), so
|
|
// it cannot share the ordinary request deadline. Bounded by the shortest
|
|
// upstream expiry we know of — GitHub device codes last ~15 minutes.
|
|
const INTERACTIVE_OAUTH_TIMEOUT_MS = 15 * 60 * 1000;
|
|
const INTERACTIVE_OAUTH_PATH = /^\/provider\/[^/]+\/oauth\/callback\/?$/;
|
|
|
|
const isInteractiveOAuthCallback = (req) =>
|
|
req.method === 'POST' && INTERACTIVE_OAUTH_PATH.test(req.path);
|
|
|
|
const isProxyTimeoutError = (error) => {
|
|
const code = typeof error?.code === 'string' ? error.code : '';
|
|
const message = typeof error?.message === 'string' ? error.message.toLowerCase() : '';
|
|
return code === 'ETIMEDOUT'
|
|
|| code === 'ESOCKETTIMEDOUT'
|
|
|| message.includes('timeout')
|
|
|| message.includes('timed out');
|
|
};
|
|
|
|
const sendProxyErrorResponse = (res, statusCode) => {
|
|
if (!res || res.headersSent || res.writableEnded || typeof res.status !== 'function') {
|
|
return false;
|
|
}
|
|
res.status(statusCode).json({ error: statusCode === 504 ? 'OpenCode upstream timed out' : 'OpenCode service unavailable' });
|
|
return true;
|
|
};
|
|
|
|
const applyProxyResponseDeadline = (req, res, next) => {
|
|
if (isInteractiveOAuthCallback(req)) {
|
|
return next();
|
|
}
|
|
|
|
const timeout = setTimeout(() => {
|
|
req[PROXY_TIMEOUT_MARKER] = true;
|
|
if (sendProxyErrorResponse(res, 504)) {
|
|
res.once('finish', () => req.destroy?.());
|
|
}
|
|
}, PROXY_REQUEST_TIMEOUT_MS);
|
|
timeout.unref?.();
|
|
|
|
const clear = () => clearTimeout(timeout);
|
|
res.once('finish', clear);
|
|
res.once('close', clear);
|
|
next();
|
|
};
|
|
|
|
const forwardSseRequest = async (req, res) => {
|
|
const abortController = new AbortController();
|
|
const closeUpstream = () => abortController.abort();
|
|
let upstream = null;
|
|
let reader = null;
|
|
let heartbeatTimer = null;
|
|
let upstreamStallTimer = null;
|
|
let didUpstreamStall = false;
|
|
let writeQueue = Promise.resolve(true);
|
|
const sseBoundary = createSseBoundaryTracker();
|
|
|
|
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 = normalizeForwardedDirectoryHeaders(
|
|
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();
|
|
}
|
|
|
|
// Disable TCP Nagle's algorithm so small SSE chunks are sent immediately
|
|
// instead of being buffered up to ~200ms by the TCP stack.
|
|
if (res.socket && typeof res.socket.setNoDelay === 'function') {
|
|
res.socket.setNoDelay(true);
|
|
}
|
|
|
|
const scheduleHeartbeat = () => {
|
|
heartbeatTimer = setTimeout(async () => {
|
|
if (abortController.signal.aborted || res.writableEnded || res.destroyed) {
|
|
return;
|
|
}
|
|
if (!sseBoundary.isAtBoundary()) {
|
|
scheduleHeartbeat();
|
|
return;
|
|
}
|
|
const canContinue = await enqueueSseWrite(':heartbeat\n\n');
|
|
if (canContinue) {
|
|
scheduleHeartbeat();
|
|
}
|
|
}, SSE_HEARTBEAT_INTERVAL_MS);
|
|
};
|
|
|
|
const clearUpstreamStallTimer = () => {
|
|
clearTimeout(upstreamStallTimer);
|
|
upstreamStallTimer = null;
|
|
};
|
|
|
|
const resetUpstreamStallTimer = () => {
|
|
clearUpstreamStallTimer();
|
|
upstreamStallTimer = setTimeout(() => {
|
|
didUpstreamStall = true;
|
|
abortController.abort();
|
|
}, getSseUpstreamStallTimeoutMs());
|
|
upstreamStallTimer.unref?.();
|
|
};
|
|
|
|
const enqueueSseWrite = (value) => {
|
|
writeQueue = writeQueue
|
|
.catch(() => false)
|
|
.then((canContinue) => {
|
|
if (!canContinue) {
|
|
return false;
|
|
}
|
|
return writeSseChunkWithBackpressure(res, value, abortController.signal);
|
|
});
|
|
return writeQueue;
|
|
};
|
|
|
|
scheduleHeartbeat();
|
|
resetUpstreamStallTimer();
|
|
|
|
reader = upstream.body.getReader();
|
|
while (!abortController.signal.aborted) {
|
|
const { done, value } = await reader.read();
|
|
if (done) {
|
|
break;
|
|
}
|
|
if (value && value.length > 0) {
|
|
resetUpstreamStallTimer();
|
|
sseBoundary.observe(value);
|
|
const canContinue = await enqueueSseWrite(value);
|
|
if (!canContinue) {
|
|
break;
|
|
}
|
|
}
|
|
}
|
|
|
|
res.end();
|
|
} catch (error) {
|
|
if (isAbortError(error)) {
|
|
if (didUpstreamStall && !res.writableEnded && !res.destroyed) {
|
|
await writeQueue.catch(() => false);
|
|
res.end();
|
|
}
|
|
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 {
|
|
if (heartbeatTimer) {
|
|
clearTimeout(heartbeatTimer);
|
|
heartbeatTimer = null;
|
|
}
|
|
if (upstreamStallTimer) {
|
|
clearTimeout(upstreamStallTimer);
|
|
upstreamStallTimer = null;
|
|
}
|
|
req.off('close', closeUpstream);
|
|
try {
|
|
if (reader) {
|
|
await reader.cancel();
|
|
reader.releaseLock();
|
|
} else if (upstream?.body && !upstream.body.locked) {
|
|
await upstream.body.cancel();
|
|
}
|
|
} catch {
|
|
}
|
|
}
|
|
};
|
|
|
|
const fetchSessionListPayload = async (upstreamPath, { req = null, timeoutMs = null } = {}) => {
|
|
const headers = req
|
|
? {
|
|
...normalizeForwardedDirectoryHeaders(collectForwardProxyHeaders(req.headers, getOpenCodeAuthHeaders())),
|
|
accept: 'application/json',
|
|
'accept-encoding': 'identity',
|
|
}
|
|
: {
|
|
Accept: 'application/json',
|
|
...getOpenCodeAuthHeaders(),
|
|
'accept-encoding': 'identity',
|
|
};
|
|
const upstream = await fetch(buildOpenCodeUrl(upstreamPath, ''), {
|
|
method: 'GET',
|
|
headers,
|
|
...(typeof timeoutMs === 'number' ? { signal: AbortSignal.timeout(timeoutMs) } : {}),
|
|
});
|
|
const contentType = upstream.headers.get('content-type') || 'application/json; charset=utf-8';
|
|
const bodyText = await upstream.text();
|
|
const isJson = contentType.toLowerCase().includes('application/json');
|
|
|
|
if (!isJson) {
|
|
return { upstream, contentType, bodyText, payload: null, isJson: false };
|
|
}
|
|
|
|
try {
|
|
const payload = JSON.parse(bodyText);
|
|
return { upstream, contentType, bodyText, payload, isJson: true, parseError: null };
|
|
} catch (parseError) {
|
|
return { upstream, contentType, bodyText, payload: null, isJson: true, parseError };
|
|
}
|
|
};
|
|
|
|
const getRequestUpstreamPath = async (req) => {
|
|
const requestUrl = typeof req.originalUrl === 'string' && req.originalUrl.length > 0
|
|
? req.originalUrl
|
|
: (typeof req.url === 'string' ? req.url : '');
|
|
const upstreamPathRaw = requestUrl.startsWith('/api') ? requestUrl.slice(4) || '/' : requestUrl;
|
|
return canonicalizeDirectoryQuery(upstreamPathRaw);
|
|
};
|
|
|
|
const forwardSanitizedSessionListRequest = async (req, res, next, logLabel) => {
|
|
try {
|
|
const upstreamPath = await getRequestUpstreamPath(req);
|
|
const result = await fetchSessionListPayload(upstreamPath, { req });
|
|
|
|
res.status(result.upstream.status);
|
|
applyForwardProxyResponseHeaders(result.upstream.headers, res);
|
|
|
|
if (!result.isJson) {
|
|
res.setHeader('content-type', result.contentType);
|
|
res.end(result.bodyText);
|
|
return;
|
|
}
|
|
|
|
if (result.parseError || !Array.isArray(result.payload)) {
|
|
res.setHeader('content-type', result.contentType);
|
|
res.end(result.bodyText);
|
|
return;
|
|
}
|
|
|
|
res.setHeader('content-type', result.contentType);
|
|
res.json(sanitizeSessionListPayload(result.payload));
|
|
} catch (error) {
|
|
if (isAbortError(error)) {
|
|
return;
|
|
}
|
|
console.error(`[proxy] OpenCode ${logLabel} proxy error:`, error?.message ?? error);
|
|
if (!res.headersSent) {
|
|
res.status(503).json({ error: 'OpenCode service unavailable' });
|
|
return;
|
|
}
|
|
next(error);
|
|
}
|
|
};
|
|
|
|
// Ensure API prefix is detected before proxying
|
|
app.use('/api', (_req, _res, next) => {
|
|
ensureOpenCodeApiPrefix();
|
|
next();
|
|
});
|
|
|
|
// Readiness gate — while OpenCode is starting/restarting, HOLD the request and
|
|
// poll readiness instead of returning 503 immediately. A bare 503 pushes the
|
|
// client into an exponential-backoff retry loop (500ms → 1s → …) that wastes
|
|
// seconds of cold-start time and can fail bootstrap outright. Holding the
|
|
// request until OpenCode is ready (typically well under a second) lets the
|
|
// first call simply succeed. We still 503 if readiness doesn't arrive within a
|
|
// bounded window so genuinely-down servers fail fast.
|
|
const READINESS_HOLD_POLL_MS = 75;
|
|
const READINESS_HOLD_MAX_MS = 6000;
|
|
const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms));
|
|
const isStillWaiting = (runtimeState) => {
|
|
const waitElapsed = runtimeState.openCodeNotReadySince === 0 ? 0 : Date.now() - runtimeState.openCodeNotReadySince;
|
|
return (
|
|
(!runtimeState.isOpenCodeReady && (runtimeState.openCodeNotReadySince === 0 || waitElapsed < OPEN_CODE_READY_GRACE_MS)) ||
|
|
runtimeState.isRestartingOpenCode ||
|
|
!runtimeState.openCodePort
|
|
);
|
|
};
|
|
const classifyReadinessRoute = (requestPath) => {
|
|
if (/^\/session\/[^/]+\/message(?:\/|$)/.test(requestPath)) return 'session-messages';
|
|
if (requestPath === '/session' || requestPath.startsWith('/session/')) return 'session';
|
|
if (requestPath === '/event' || requestPath === '/global/event') return 'events';
|
|
return 'other';
|
|
};
|
|
|
|
app.use('/api', async (req, res, next) => {
|
|
if (
|
|
req.path.startsWith('/themes/custom') ||
|
|
req.path.startsWith('/push') ||
|
|
req.path.startsWith('/config/agents') ||
|
|
req.path.startsWith('/config/opencode-resolution') ||
|
|
req.path.startsWith('/config/settings') ||
|
|
req.path.startsWith('/config/skills') ||
|
|
req.path === '/config/reload' ||
|
|
req.path === '/health'
|
|
) {
|
|
return next();
|
|
}
|
|
|
|
if (!isStillWaiting(getRuntime())) {
|
|
return next();
|
|
}
|
|
|
|
const holdStartedAt = performance.now();
|
|
const routeClass = classifyReadinessRoute(req.path);
|
|
const deadline = Date.now() + Math.min(OPEN_CODE_READY_GRACE_MS, READINESS_HOLD_MAX_MS);
|
|
while (Date.now() < deadline) {
|
|
// Client gave up (closed/aborted) — stop holding.
|
|
if (res.writableEnded || req.aborted) {
|
|
recordStartupPerformance('proxy.readiness-hold', {
|
|
durationMs: performance.now() - holdStartedAt,
|
|
outcome: 'aborted',
|
|
routeClass,
|
|
});
|
|
return;
|
|
}
|
|
await sleep(READINESS_HOLD_POLL_MS);
|
|
if (!isStillWaiting(getRuntime())) {
|
|
recordStartupPerformance('proxy.readiness-hold', {
|
|
durationMs: performance.now() - holdStartedAt,
|
|
outcome: 'ready',
|
|
routeClass,
|
|
});
|
|
return next();
|
|
}
|
|
}
|
|
|
|
recordStartupPerformance('proxy.readiness-hold', {
|
|
durationMs: performance.now() - holdStartedAt,
|
|
outcome: 'timeout',
|
|
routeClass,
|
|
});
|
|
if (!res.headersSent) {
|
|
res.status(503).json({
|
|
error: 'OpenCode is restarting',
|
|
restarting: true,
|
|
});
|
|
}
|
|
});
|
|
|
|
// Windows: session merge for cross-directory session listing
|
|
if (process.platform === 'win32') {
|
|
app.get('/api/session', async (req, res, next) => {
|
|
const rawUrl = req.originalUrl || req.url || '';
|
|
if (rawUrl.includes('directory=')) return next();
|
|
|
|
const fetchWindowsSessionList = async (sessionPath) => {
|
|
const result = await fetchSessionListPayload(sessionPath, { req, timeoutMs: 10000 });
|
|
if (!result.upstream.ok || !Array.isArray(result.payload)) return null;
|
|
return sanitizeSessionListPayload(result.payload);
|
|
};
|
|
|
|
try {
|
|
const globalSessions = await fetchWindowsSessionList('/session').catch((error) => {
|
|
console.log(`[SessionMerge] Global session list failed: ${error.message}`);
|
|
return null;
|
|
});
|
|
|
|
const settingsPath = path.join(os.homedir(), '.config', 'openchamber', 'settings.json');
|
|
let projectDirs = [];
|
|
try {
|
|
const settingsRaw = fs.readFileSync(settingsPath, 'utf8');
|
|
const settings = JSON.parse(settingsRaw);
|
|
projectDirs = (settings.projects || [])
|
|
.map((project) => (typeof project?.path === 'string' ? project.path.trim() : ''))
|
|
.filter(Boolean);
|
|
} catch {
|
|
}
|
|
|
|
const seen = new Set(
|
|
(globalSessions || [])
|
|
.map((session) => (session && typeof session.id === 'string' ? session.id : null))
|
|
.filter((id) => typeof id === 'string')
|
|
);
|
|
const extraSessions = [];
|
|
let successfulProjectReads = 0;
|
|
for (const dir of projectDirs) {
|
|
const candidates = Array.from(new Set([
|
|
dir,
|
|
dir.replace(/\\/g, '/'),
|
|
dir.replace(/\//g, '\\'),
|
|
]));
|
|
for (const candidateDir of candidates) {
|
|
const encoded = encodeURIComponent(candidateDir);
|
|
try {
|
|
const dirSessions = await fetchWindowsSessionList(`/session?directory=${encoded}`);
|
|
if (dirSessions) {
|
|
successfulProjectReads += 1;
|
|
}
|
|
for (const session of dirSessions || []) {
|
|
const id = session && typeof session.id === 'string' ? session.id : null;
|
|
if (id && !seen.has(id)) {
|
|
seen.add(id);
|
|
extraSessions.push(session);
|
|
}
|
|
}
|
|
} catch {
|
|
}
|
|
}
|
|
}
|
|
|
|
if (!globalSessions && successfulProjectReads === 0) {
|
|
return res.status(504).json({ error: 'OpenCode session list timed out' });
|
|
}
|
|
|
|
const merged = [...(globalSessions || []), ...extraSessions];
|
|
merged.sort((a, b) => {
|
|
const aTime = a && typeof a.time_updated === 'number' ? a.time_updated : 0;
|
|
const bTime = b && typeof b.time_updated === 'number' ? b.time_updated : 0;
|
|
return bTime - aTime;
|
|
});
|
|
console.log(`[SessionMerge] ${globalSessions?.length || 0} global + ${extraSessions.length} extra = ${merged.length} total`);
|
|
return res.json(sanitizeSessionListPayload(merged));
|
|
} catch (error) {
|
|
console.log(`[SessionMerge] Error: ${error.message}`);
|
|
return res.status(500).json({ error: error.message || 'Failed to merge Windows sessions' });
|
|
}
|
|
});
|
|
}
|
|
|
|
app.get('/api/session', (req, res, next) => {
|
|
return forwardSanitizedSessionListRequest(req, res, next, 'session.list');
|
|
});
|
|
|
|
app.get('/api/global/event', forwardSseRequest);
|
|
app.get('/api/event', forwardSseRequest);
|
|
|
|
app.get('/api/experimental/session', (req, res, next) => {
|
|
return forwardSanitizedSessionListRequest(req, res, next, 'experimental.session');
|
|
});
|
|
|
|
// Generic proxy for non-SSE OpenCode API routes.
|
|
// The agent is exposed as a getter so its class is resolved per request, not
|
|
// at registration: the proxy is registered before OpenCode bootstraps, so an
|
|
// https target configured via OPENCODE_HOST is not yet visible here. Agents
|
|
// are memoized per scheme, so this is still one shared pool per scheme across
|
|
// `apiProxy` and `interactiveOAuthProxy`.
|
|
const resolveOpenCodeProxyAgent = createOpenCodeProxyAgentResolver(resolveProxyTarget);
|
|
|
|
const createApiProxy = (timeoutMs) => createProxyMiddleware({
|
|
target: resolveProxyTarget(),
|
|
get agent() {
|
|
return resolveOpenCodeProxyAgent();
|
|
},
|
|
changeOrigin: true,
|
|
pathRewrite: { '^/api': '' },
|
|
timeout: timeoutMs,
|
|
proxyTimeout: timeoutMs,
|
|
// Dynamic target — port can change after restart
|
|
router: () => resolveProxyTarget(),
|
|
on: {
|
|
proxyReq: (proxyReq, req) => {
|
|
// Inject OpenCode auth headers
|
|
const authHeaders = getOpenCodeAuthHeaders();
|
|
if (authHeaders.Authorization) {
|
|
proxyReq.setHeader('Authorization', authHeaders.Authorization);
|
|
}
|
|
|
|
if (req.headers?.['x-opencode-directory-encoding'] === 'uri') {
|
|
const rawDirectory = req.headers['x-opencode-directory'];
|
|
if (typeof rawDirectory === 'string') {
|
|
try {
|
|
proxyReq.setHeader('x-opencode-directory', decodeURIComponent(rawDirectory));
|
|
} catch {
|
|
proxyReq.setHeader('x-opencode-directory', rawDirectory);
|
|
}
|
|
}
|
|
proxyReq.removeHeader?.('x-opencode-directory-encoding');
|
|
}
|
|
|
|
// Defensive: request identity encoding from upstream OpenCode.
|
|
// This avoids compressed-body/header mismatches in multi-proxy setups.
|
|
proxyReq.setHeader('accept-encoding', 'identity');
|
|
|
|
replayParsedBody(proxyReq, req);
|
|
},
|
|
proxyRes: (proxyRes) => {
|
|
for (const key of Object.keys(proxyRes.headers || {})) {
|
|
if (!shouldForwardProxyResponseHeader(key)) {
|
|
delete proxyRes.headers[key];
|
|
}
|
|
}
|
|
},
|
|
error: (err, req, res) => {
|
|
console.error('[proxy] OpenCode proxy error:', err.message);
|
|
if (req?.[PROXY_TIMEOUT_MARKER]) {
|
|
return;
|
|
}
|
|
const statusCode = isProxyTimeoutError(err) ? 504 : 503;
|
|
sendProxyErrorResponse(res, statusCode);
|
|
},
|
|
},
|
|
});
|
|
|
|
const apiProxy = createApiProxy(PROXY_REQUEST_TIMEOUT_MS);
|
|
const interactiveOAuthProxy = createApiProxy(INTERACTIVE_OAUTH_TIMEOUT_MS);
|
|
|
|
// Best-effort fallback for stale clients still sending symlink paths.
|
|
// Settings and project selection normalize at source; this cached async path
|
|
// avoids blocking the proxy hot path on every directory-scoped request.
|
|
app.use('/api', async (req, _res, next) => {
|
|
try {
|
|
const rewrittenUrl = await canonicalizeDirectoryQuery(req.url);
|
|
if (rewrittenUrl !== req.url) {
|
|
req.url = rewrittenUrl;
|
|
}
|
|
} catch {
|
|
// Pass through as-is if URL parsing or realpath resolution fails.
|
|
}
|
|
next();
|
|
});
|
|
|
|
app.use('/api', applyProxyResponseDeadline);
|
|
app.post('/api/provider/:providerID/oauth/callback', interactiveOAuthProxy);
|
|
// OpenCode's native MCP OAuth flow: the request blocks until the user
|
|
// finishes authorization in the browser (up to OpenCode's 5-minute callback
|
|
// timeout), so it needs the interactive-OAuth deadline, not the default one.
|
|
app.post('/api/mcp/:name/auth/authenticate', interactiveOAuthProxy);
|
|
app.use('/api', apiProxy);
|
|
};
|