From 25502a52bb7c264230f4a05a46275ef386a9b5a8 Mon Sep 17 00:00:00 2001 From: Codex Date: Thu, 30 Jul 2026 22:11:48 +0800 Subject: [PATCH] fix: flush proxied event stream headers --- apps/api/src/canary-gateway.test.ts | 25 +++++++++++++++++++++++++ apps/api/src/canary-gateway.ts | 1 + 2 files changed, 26 insertions(+) diff --git a/apps/api/src/canary-gateway.test.ts b/apps/api/src/canary-gateway.test.ts index cd93bb5..e80632a 100644 --- a/apps/api/src/canary-gateway.test.ts +++ b/apps/api/src/canary-gateway.test.ts @@ -334,6 +334,31 @@ describe('same-origin canary gateway', () => { await reader?.cancel(); }); + it('flushes SSE response headers before the first upstream event', async () => { + const upstream = createServer((_request, response) => { + response.writeHead(200, { + 'content-type': 'text/event-stream', + 'cache-control': 'no-cache', + }); + response.flushHeaders(); + }); + servers.push(upstream); + const gateway = await startGateway(await fixtureDist(), await listen(upstream)); + + const controller = new AbortController(); + const timeout = setTimeout(() => controller.abort(), 500); + try { + const response = await fetch(`${gateway.origin}/api/v1/events`, { + signal: controller.signal, + }); + expect(response.status).toBe(200); + expect(response.headers.get('content-type')).toBe('text/event-stream'); + await response.body?.cancel(); + } finally { + clearTimeout(timeout); + } + }); + it('has reversible, idempotent start/stop lifecycle', async () => { const gateway = createCanaryGateway({ distDir: await fixtureDist(), port: 0, upstreamPort: 1 }); gateways.push(gateway); diff --git a/apps/api/src/canary-gateway.ts b/apps/api/src/canary-gateway.ts index 6a31df9..cc70f51 100644 --- a/apps/api/src/canary-gateway.ts +++ b/apps/api/src/canary-gateway.ts @@ -177,6 +177,7 @@ async function proxyRequest( upstreamResponse.statusMessage, withoutHopByHop(upstreamResponse.headers), ); + response.flushHeaders(); pipeline(upstreamResponse, response) .then(resolvePromise) .catch(() => {