refactor(open-sse): DRY SSE primitives into utils/sseConstants.js

Gom SSE_DONE + SSE_HEADERS + SSE_HEADERS_NO_BUFFER vào 1 file không
phụ thuộc (tránh kéo usageDb vào executor). Áp cho kiro, cursor, qoder,
github, commandcode, perplexity-web, grok-web. Byte-for-byte verified
(headers/DONE giữ nguyên; biến thể no-buffer dùng const riêng).

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
decolua
2026-06-13 21:42:59 +07:00
co-authored by Cursor
parent 0a8d92a6ba
commit 6ee4555821
9 changed files with 42 additions and 25 deletions
+2 -1
View File
@@ -2,6 +2,7 @@ import { randomUUID } from "crypto";
import { BaseExecutor } from "./base.js"; import { BaseExecutor } from "./base.js";
import { PROVIDERS } from "../config/providers.js"; import { PROVIDERS } from "../config/providers.js";
import { convertCommandCodeToOpenAI } from "../translator/response/commandcode-to-openai.js"; import { convertCommandCodeToOpenAI } from "../translator/response/commandcode-to-openai.js";
import { SSE_DONE } from "../utils/sseConstants.js";
/** /**
* CommandCodeExecutor — talks to https://api.commandcode.ai/alpha/generate * CommandCodeExecutor — talks to https://api.commandcode.ai/alpha/generate
@@ -78,7 +79,7 @@ function wrapNdjsonAsOpenAISse(originalResponse, model) {
if (trimmed) { if (trimmed) {
emitChunks(convertCommandCodeToOpenAI(trimmed, state), controller); emitChunks(convertCommandCodeToOpenAI(trimmed, state), controller);
} }
controller.enqueue(encoder.encode("data: [DONE]\n\n")); controller.enqueue(encoder.encode(SSE_DONE));
}, },
}); });
+3 -6
View File
@@ -8,6 +8,7 @@ import {
} from "../utils/cursorProtobuf.js"; } from "../utils/cursorProtobuf.js";
import { buildCursorHeaders } from "../utils/cursorChecksum.js"; import { buildCursorHeaders } from "../utils/cursorChecksum.js";
import { estimateUsage } from "../utils/usageTracking.js"; import { estimateUsage } from "../utils/usageTracking.js";
import { SSE_DONE, SSE_HEADERS } from "../utils/sseConstants.js";
import { FORMATS } from "../translator/formats.js"; import { FORMATS } from "../translator/formats.js";
import { proxyAwareFetch } from "../utils/proxyFetch.js"; import { proxyAwareFetch } from "../utils/proxyFetch.js";
import zlib from "zlib"; import zlib from "zlib";
@@ -775,15 +776,11 @@ export class CursorExecutor extends BaseExecutor {
usage usage
})}\n\n` })}\n\n`
); );
chunks.push("data: [DONE]\n\n"); chunks.push(SSE_DONE);
return new Response(chunks.join(""), { return new Response(chunks.join(""), {
status: 200, status: 200,
headers: { headers: { ...SSE_HEADERS }
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
"Connection": "keep-alive"
}
}); });
} }
+2 -1
View File
@@ -7,6 +7,7 @@ import { openaiResponsesToOpenAIResponse } from "../translator/response/openai-r
import { initState } from "../translator/index.js"; import { initState } from "../translator/index.js";
import { parseSSELine, formatSSE } from "../utils/streamHelpers.js"; import { parseSSELine, formatSSE } from "../utils/streamHelpers.js";
import { proxyAwareFetch } from "../utils/proxyFetch.js"; import { proxyAwareFetch } from "../utils/proxyFetch.js";
import { SSE_DONE } from "../utils/sseConstants.js";
import crypto from "crypto"; import crypto from "crypto";
export class GithubExecutor extends BaseExecutor { export class GithubExecutor extends BaseExecutor {
@@ -244,7 +245,7 @@ export class GithubExecutor extends BaseExecutor {
if (!parsed) continue; if (!parsed) continue;
if (parsed.done && stream === true) { if (parsed.done && stream === true) {
controller.enqueue(new TextEncoder().encode("data: [DONE]\n\n")); controller.enqueue(new TextEncoder().encode(SSE_DONE));
continue; continue;
} }
+4 -3
View File
@@ -1,5 +1,6 @@
import { BaseExecutor } from "./base.js"; import { BaseExecutor } from "./base.js";
import { PROVIDERS } from "../config/providers.js"; import { PROVIDERS } from "../config/providers.js";
import { SSE_DONE, SSE_HEADERS_NO_BUFFER } from "../utils/sseConstants.js";
const GROK_CHAT_API = PROVIDERS["grok-web"].baseUrl; const GROK_CHAT_API = PROVIDERS["grok-web"].baseUrl;
const GROK_USER_AGENT = "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/136.0.0.0 Safari/537.36"; const GROK_USER_AGENT = "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/136.0.0.0 Safari/537.36";
@@ -175,13 +176,13 @@ function buildStreamingResponse(eventStream, model, cid, created, isThinkingMode
id: cid, object: "chat.completion.chunk", created, model, system_fingerprint: fp || null, id: cid, object: "chat.completion.chunk", created, model, system_fingerprint: fp || null,
choices: [{ index: 0, delta: {}, finish_reason: "stop", logprobs: null }], choices: [{ index: 0, delta: {}, finish_reason: "stop", logprobs: null }],
}))); })));
controller.enqueue(encoder.encode("data: [DONE]\n\n")); controller.enqueue(encoder.encode(SSE_DONE));
} catch (err) { } catch (err) {
controller.enqueue(encoder.encode(sseChunk({ controller.enqueue(encoder.encode(sseChunk({
id: cid, object: "chat.completion.chunk", created, model, system_fingerprint: null, id: cid, object: "chat.completion.chunk", created, model, system_fingerprint: null,
choices: [{ index: 0, delta: { content: `[Stream error: ${err.message || String(err)}]` }, finish_reason: "stop", logprobs: null }], choices: [{ index: 0, delta: { content: `[Stream error: ${err.message || String(err)}]` }, finish_reason: "stop", logprobs: null }],
}))); })));
controller.enqueue(encoder.encode("data: [DONE]\n\n")); controller.enqueue(encoder.encode(SSE_DONE));
} finally { } finally {
controller.close(); controller.close();
} }
@@ -333,7 +334,7 @@ export class GrokWebExecutor extends BaseExecutor {
const sseStream = buildStreamingResponse(response.body, model, cid, created, isThinking, signal); const sseStream = buildStreamingResponse(response.body, model, cid, created, isThinking, signal);
finalResponse = new Response(sseStream, { finalResponse = new Response(sseStream, {
status: 200, status: 200,
headers: { "Content-Type": "text/event-stream", "Cache-Control": "no-cache", "X-Accel-Buffering": "no" }, headers: { ...SSE_HEADERS_NO_BUFFER },
}); });
} else { } else {
finalResponse = await buildNonStreamingResponse(response.body, model, cid, created, isThinking, signal); finalResponse = await buildNonStreamingResponse(response.body, model, cid, created, isThinking, signal);
+4 -7
View File
@@ -2,6 +2,7 @@ import { BaseExecutor } from "./base.js";
import { PROVIDERS } from "../config/providers.js"; import { PROVIDERS } from "../config/providers.js";
import { v4 as uuidv4 } from "uuid"; import { v4 as uuidv4 } from "uuid";
import { refreshKiroToken } from "../services/tokenRefresh.js"; import { refreshKiroToken } from "../services/tokenRefresh.js";
import { SSE_DONE, SSE_HEADERS } from "../utils/sseConstants.js";
/** /**
* KiroExecutor - Executor for Kiro AI (AWS CodeWhisperer) * KiroExecutor - Executor for Kiro AI (AWS CodeWhisperer)
@@ -372,24 +373,20 @@ export class KiroExecutor extends BaseExecutor {
} }
// Send final done message // Send final done message
controller.enqueue(new TextEncoder().encode("data: [DONE]\n\n")); controller.enqueue(new TextEncoder().encode(SSE_DONE));
} }
}); });
// Pipe response body through transform stream // Pipe response body through transform stream
if (!response.body) { if (!response.body) {
return new Response("data: [DONE]\n\n", { status: response.status, headers: { "Content-Type": "text/event-stream" } }); return new Response(SSE_DONE, { status: response.status, headers: { "Content-Type": "text/event-stream" } });
} }
const transformedStream = response.body.pipeThrough(transformStream); const transformedStream = response.body.pipeThrough(transformStream);
return new Response(transformedStream, { return new Response(transformedStream, {
status: response.status, status: response.status,
statusText: response.statusText, statusText: response.statusText,
headers: { headers: { ...SSE_HEADERS }
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
"Connection": "keep-alive"
}
}); });
} }
+4 -3
View File
@@ -1,5 +1,6 @@
import { BaseExecutor } from "./base.js"; import { BaseExecutor } from "./base.js";
import { PROVIDERS } from "../config/providers.js"; import { PROVIDERS } from "../config/providers.js";
import { SSE_DONE, SSE_HEADERS_NO_BUFFER } from "../utils/sseConstants.js";
const PPLX_SSE_ENDPOINT = PROVIDERS["perplexity-web"].baseUrl; const PPLX_SSE_ENDPOINT = PROVIDERS["perplexity-web"].baseUrl;
const PPLX_API_VERSION = "2.18"; const PPLX_API_VERSION = "2.18";
@@ -340,7 +341,7 @@ function buildStreamingResponse(eventStream, model, cid, created, history, curre
id: cid, object: "chat.completion.chunk", created, model, system_fingerprint: null, id: cid, object: "chat.completion.chunk", created, model, system_fingerprint: null,
choices: [{ index: 0, delta: {}, finish_reason: "stop", logprobs: null }], choices: [{ index: 0, delta: {}, finish_reason: "stop", logprobs: null }],
}))); })));
controller.enqueue(encoder.encode("data: [DONE]\n\n")); controller.enqueue(encoder.encode(SSE_DONE));
sessionStore(history, currentMsg, cleanResponse(fullAnswer), respBackendUuid); sessionStore(history, currentMsg, cleanResponse(fullAnswer), respBackendUuid);
} catch (err) { } catch (err) {
@@ -348,7 +349,7 @@ function buildStreamingResponse(eventStream, model, cid, created, history, curre
id: cid, object: "chat.completion.chunk", created, model, system_fingerprint: null, id: cid, object: "chat.completion.chunk", created, model, system_fingerprint: null,
choices: [{ index: 0, delta: { content: `[Stream error: ${err.message || String(err)}]` }, finish_reason: "stop", logprobs: null }], choices: [{ index: 0, delta: { content: `[Stream error: ${err.message || String(err)}]` }, finish_reason: "stop", logprobs: null }],
}))); })));
controller.enqueue(encoder.encode("data: [DONE]\n\n")); controller.enqueue(encoder.encode(SSE_DONE));
} finally { } finally {
controller.close(); controller.close();
} }
@@ -493,7 +494,7 @@ export class PerplexityWebExecutor extends BaseExecutor {
const sseStream = buildStreamingResponse(response.body, model, cid, created, parsed.history, parsed.currentMsg, signal); const sseStream = buildStreamingResponse(response.body, model, cid, created, parsed.history, parsed.currentMsg, signal);
finalResponse = new Response(sseStream, { finalResponse = new Response(sseStream, {
status: 200, status: 200,
headers: { "Content-Type": "text/event-stream", "Cache-Control": "no-cache", "X-Accel-Buffering": "no" }, headers: { ...SSE_HEADERS_NO_BUFFER },
}); });
} else { } else {
finalResponse = await buildNonStreamingResponse(response.body, model, cid, created, parsed.history, parsed.currentMsg, signal); finalResponse = await buildNonStreamingResponse(response.body, model, cid, created, parsed.history, parsed.currentMsg, signal);
+5 -4
View File
@@ -28,6 +28,7 @@ import { createHash } from "crypto";
import { BaseExecutor } from "./base.js"; import { BaseExecutor } from "./base.js";
import { PROVIDERS } from "../config/providers.js"; import { PROVIDERS } from "../config/providers.js";
import { proxyAwareFetch } from "../utils/proxyFetch.js"; import { proxyAwareFetch } from "../utils/proxyFetch.js";
import { SSE_DONE } from "../utils/sseConstants.js";
import { FETCH_CONNECT_TIMEOUT_MS } from "../config/runtimeConfig.js"; import { FETCH_CONNECT_TIMEOUT_MS } from "../config/runtimeConfig.js";
import { import {
QODER_CHAT_URL_ENCODED, QODER_CHAT_URL_ENCODED,
@@ -241,7 +242,7 @@ function wrapQoderSSE(response, model) {
const data = trimmed.slice(5).trimStart(); const data = trimmed.slice(5).trimStart();
if (data === "[DONE]") { if (data === "[DONE]") {
controller.enqueue(encoder.encode("data: [DONE]\n\n")); controller.enqueue(encoder.encode(SSE_DONE));
doneEmitted = true; doneEmitted = true;
return; return;
} }
@@ -260,13 +261,13 @@ function wrapQoderSSE(response, model) {
choices: [{ index: 0, delta: { content: `\n[qoder error ${statusVal}: ${truncate(msg, 200)}]` }, finish_reason: "stop" }], choices: [{ index: 0, delta: { content: `\n[qoder error ${statusVal}: ${truncate(msg, 200)}]` }, finish_reason: "stop" }],
}); });
controller.enqueue(encoder.encode(`data: ${errChunk}\n\n`)); controller.enqueue(encoder.encode(`data: ${errChunk}\n\n`));
controller.enqueue(encoder.encode("data: [DONE]\n\n")); controller.enqueue(encoder.encode(SSE_DONE));
doneEmitted = true; doneEmitted = true;
return; return;
} }
if (!inner) return; if (!inner) return;
if (inner === "[DONE]") { if (inner === "[DONE]") {
controller.enqueue(encoder.encode("data: [DONE]\n\n")); controller.enqueue(encoder.encode(SSE_DONE));
doneEmitted = true; doneEmitted = true;
return; return;
} }
@@ -301,7 +302,7 @@ function wrapQoderSSE(response, model) {
buffer = ""; buffer = "";
} }
if (!doneEmitted) { if (!doneEmitted) {
controller.enqueue(encoder.encode("data: [DONE]\n\n")); controller.enqueue(encoder.encode(SSE_DONE));
doneEmitted = true; doneEmitted = true;
} }
}, },
+15
View File
@@ -0,0 +1,15 @@
// Shared SSE primitives (no imports → safe for executors + stream.js)
export const SSE_DONE = "data: [DONE]\n\n";
export const SSE_HEADERS = {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
"Connection": "keep-alive"
};
// Variant for web-cookie executors behind nginx (disable proxy buffering)
export const SSE_HEADERS_NO_BUFFER = {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no"
};
+3
View File
@@ -6,7 +6,10 @@ import { parseSSELine, hasValuableContent, fixInvalidId, formatSSE } from "./str
import { getOpenAIResponsesEventName, isOpenAIResponsesTerminalEvent, formatIncompleteOpenAIResponsesStreamFailure } from "./responsesStreamHelpers.js"; import { getOpenAIResponsesEventName, isOpenAIResponsesTerminalEvent, formatIncompleteOpenAIResponsesStreamFailure } from "./responsesStreamHelpers.js";
import { dbg, isDebugEnabled } from "./debugLog.js"; import { dbg, isDebugEnabled } from "./debugLog.js";
import { SSE_DONE, SSE_HEADERS, SSE_HEADERS_NO_BUFFER } from "./sseConstants.js";
export { COLORS, formatSSE }; export { COLORS, formatSSE };
export { SSE_DONE, SSE_HEADERS, SSE_HEADERS_NO_BUFFER };
// sharedEncoder is stateless — safe to share across streams // sharedEncoder is stateless — safe to share across streams
const sharedEncoder = new TextEncoder(); const sharedEncoder = new TextEncoder();