fix: preserve lastEventId in SSE path and add proxy heartbeat (#1041)
* 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>
This commit is contained in:
committed by
GitHub
co-authored by
Bohdan Triapitsyn
parent
43c92babc3
commit
6470d9d205
@@ -340,6 +340,12 @@ export function createEventPipeline(input: EventPipelineInput) {
|
||||
const events = await sdk.global.event({
|
||||
signal,
|
||||
...(lastEventId && lastEventId.length > 0 ? { headers: { "Last-Event-ID": lastEventId } } : {}),
|
||||
onSseEvent: (event: { id?: unknown }) => {
|
||||
resetHeartbeat()
|
||||
if (typeof event.id === "string" && event.id.length > 0) {
|
||||
lastEventId = event.id
|
||||
}
|
||||
},
|
||||
onSseError: (error: unknown) => {
|
||||
if (isAbortError(error)) return
|
||||
if (streamErrorLogged) return
|
||||
@@ -356,6 +362,7 @@ export function createEventPipeline(input: EventPipelineInput) {
|
||||
for await (const event of events.stream) {
|
||||
resetHeartbeat()
|
||||
streamErrorLogged = false
|
||||
|
||||
const payload = resolveEventPayload((event as { payload?: Event }).payload ?? event)
|
||||
if (!payload) {
|
||||
continue
|
||||
|
||||
@@ -43,6 +43,31 @@ export const writeSseChunkWithBackpressure = async (res, value, 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');
|
||||
},
|
||||
};
|
||||
};
|
||||
|
||||
export const registerOpenCodeProxy = (app, deps) => {
|
||||
const {
|
||||
fs,
|
||||
@@ -113,6 +138,9 @@ export const registerOpenCodeProxy = (app, deps) => {
|
||||
const closeUpstream = () => abortController.abort();
|
||||
let upstream = null;
|
||||
let reader = null;
|
||||
let heartbeatTimer = null;
|
||||
let writeQueue = Promise.resolve(true);
|
||||
const sseBoundary = createSseBoundaryTracker();
|
||||
|
||||
req.on('close', closeUpstream);
|
||||
|
||||
@@ -161,6 +189,38 @@ export const registerOpenCodeProxy = (app, deps) => {
|
||||
res.socket.setNoDelay(true);
|
||||
}
|
||||
|
||||
const SSE_HEARTBEAT_INTERVAL_MS = 20_000;
|
||||
|
||||
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 enqueueSseWrite = (value) => {
|
||||
writeQueue = writeQueue
|
||||
.catch(() => false)
|
||||
.then((canContinue) => {
|
||||
if (!canContinue) {
|
||||
return false;
|
||||
}
|
||||
return writeSseChunkWithBackpressure(res, value, abortController.signal);
|
||||
});
|
||||
return writeQueue;
|
||||
};
|
||||
|
||||
scheduleHeartbeat();
|
||||
|
||||
reader = upstream.body.getReader();
|
||||
while (!abortController.signal.aborted) {
|
||||
const { done, value } = await reader.read();
|
||||
@@ -168,7 +228,8 @@ export const registerOpenCodeProxy = (app, deps) => {
|
||||
break;
|
||||
}
|
||||
if (value && value.length > 0) {
|
||||
const canContinue = await writeSseChunkWithBackpressure(res, value, abortController.signal);
|
||||
sseBoundary.observe(value);
|
||||
const canContinue = await enqueueSseWrite(value);
|
||||
if (!canContinue) {
|
||||
break;
|
||||
}
|
||||
@@ -187,6 +248,10 @@ export const registerOpenCodeProxy = (app, deps) => {
|
||||
res.end();
|
||||
}
|
||||
} finally {
|
||||
if (heartbeatTimer) {
|
||||
clearTimeout(heartbeatTimer);
|
||||
heartbeatTimer = null;
|
||||
}
|
||||
req.off('close', closeUpstream);
|
||||
try {
|
||||
if (reader) {
|
||||
|
||||
@@ -3,7 +3,7 @@ import { EventEmitter } from 'node:events';
|
||||
import express from 'express';
|
||||
import path from 'path';
|
||||
|
||||
import { registerOpenCodeProxy, writeSseChunkWithBackpressure } from './lib/opencode/proxy.js';
|
||||
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));
|
||||
@@ -103,6 +103,17 @@ describe('OpenCode proxy SSE forwarding', () => {
|
||||
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) => {
|
||||
|
||||
Reference in New Issue
Block a user