diff --git a/packages/ui/src/sync/event-pipeline.ts b/packages/ui/src/sync/event-pipeline.ts index 760eed0d..43e5e488 100644 --- a/packages/ui/src/sync/event-pipeline.ts +++ b/packages/ui/src/sync/event-pipeline.ts @@ -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 diff --git a/packages/web/server/lib/opencode/proxy.js b/packages/web/server/lib/opencode/proxy.js index 69a6c9ba..3b09fc50 100644 --- a/packages/web/server/lib/opencode/proxy.js +++ b/packages/web/server/lib/opencode/proxy.js @@ -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) { diff --git a/packages/web/server/opencode-proxy.test.js b/packages/web/server/opencode-proxy.test.js index 8b4e8564..78ff8189 100644 --- a/packages/web/server/opencode-proxy.test.js +++ b/packages/web/server/opencode-proxy.test.js @@ -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) => {