* fix: preserve lastEventId in SSE path and add proxy heartbeat - Extract event.id from SSE stream events in event-pipeline.ts so that reconnects carry the correct Last-Event-ID header for gapless replay. - Emit :heartbeat comment every 20s in the direct SSE proxy to keep the UI heartbeat watchdog from aborting idle connections. * fix: handle SSE metadata through SDK callback * Guard SSE proxy heartbeats --------- Co-authored-by: Bohdan Triapitsyn <artmore@protonmail.com>
152 lines
5.0 KiB
JavaScript
152 lines
5.0 KiB
JavaScript
import { afterEach, describe, expect, it } from 'vitest';
|
|
import { EventEmitter } from 'node:events';
|
|
import express from 'express';
|
|
import path from 'path';
|
|
|
|
import { createSseBoundaryTracker, registerOpenCodeProxy, writeSseChunkWithBackpressure } 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');
|
|
});
|
|
|
|
it('waits for drain when writing to a slow SSE response', async () => {
|
|
const writes = [];
|
|
const res = new EventEmitter();
|
|
res.writableEnded = false;
|
|
res.destroyed = false;
|
|
res.write = (value) => {
|
|
writes.push(value);
|
|
return false;
|
|
};
|
|
const controller = new AbortController();
|
|
|
|
const write = writeSseChunkWithBackpressure(res, Buffer.from('data: {"ok":true}\n\n'), controller.signal);
|
|
|
|
await new Promise((resolve) => setTimeout(resolve, 0));
|
|
expect(writes).toHaveLength(1);
|
|
|
|
res.emit('drain');
|
|
|
|
await expect(write).resolves.toBe(true);
|
|
});
|
|
|
|
it('tracks whether a raw SSE stream is between event blocks', () => {
|
|
const tracker = createSseBoundaryTracker();
|
|
|
|
expect(tracker.isAtBoundary()).toBe(true);
|
|
expect(tracker.observe(Buffer.from('id: evt-1\n'))).toBe(false);
|
|
expect(tracker.observe(Buffer.from('data: {"ok"'))).toBe(false);
|
|
expect(tracker.observe(Buffer.from(':true}\n'))).toBe(false);
|
|
expect(tracker.observe(Buffer.from('\n'))).toBe(true);
|
|
expect(tracker.observe(Buffer.from('data: next\r\n\r\n'))).toBe(true);
|
|
});
|
|
|
|
it('routes generic API requests through external OpenCode base URL', async () => {
|
|
const upstream = express();
|
|
upstream.get('/config/providers', (_req, res) => {
|
|
res.json({ ok: true, source: 'external-host' });
|
|
});
|
|
upstreamServer = await listen(upstream);
|
|
const upstreamPort = upstreamServer.address().port;
|
|
const externalBaseUrl = `http://127.0.0.1:${upstreamPort}`;
|
|
|
|
const app = express();
|
|
registerOpenCodeProxy(app, {
|
|
fs: {},
|
|
os: {},
|
|
path,
|
|
OPEN_CODE_READY_GRACE_MS: 0,
|
|
getRuntime: () => ({
|
|
openCodePort: 3902,
|
|
openCodeBaseUrl: externalBaseUrl,
|
|
isOpenCodeReady: true,
|
|
openCodeNotReadySince: 0,
|
|
isRestartingOpenCode: false,
|
|
}),
|
|
getOpenCodeAuthHeaders: () => ({}),
|
|
buildOpenCodeUrl: (requestPath) => `${externalBaseUrl}${requestPath}`,
|
|
ensureOpenCodeApiPrefix: () => {},
|
|
});
|
|
proxyServer = await listen(app);
|
|
const proxyPort = proxyServer.address().port;
|
|
|
|
const response = await fetch(`http://127.0.0.1:${proxyPort}/api/config/providers`);
|
|
|
|
expect(response.status).toBe(200);
|
|
expect(await response.json()).toEqual({ ok: true, source: 'external-host' });
|
|
});
|
|
});
|