Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 8 additions & 1 deletion src/chat/outbound.ts
Original file line number Diff line number Diff line change
Expand Up @@ -138,6 +138,13 @@ export function dataFrame(payload: Rec | "[DONE]"): string {
return `data: ${JSON.stringify(payload)}\n\n`;
}

/**
* The Chat wire's keepalive: an SSE comment, so idle-sensitive clients receive bytes without a
* semantic chunk parsers or usage counters could mistake for model output. Shared with the
* direct encoder (src/protocols/encoders/chat.ts) so both paths emit byte-identical keepalives.
*/
export const CHAT_COMPLETIONS_HEARTBEAT_COMMENT = ": opencodex heartbeat\n\n";

export function chunkBase(id: string, model: string, created: number): Rec {
return {
id,
Expand Down Expand Up @@ -618,7 +625,7 @@ export function responsesSseToChatCompletionsSse(
// that signal as an SSE comment so idle-sensitive Chat clients receive bytes without
// inventing a semantic chunk that parsers, usage counters, or progress watchdogs could
// mistake for model output.
enqueueLiveFrame(encoder.encode(": opencodex heartbeat\n\n"));
enqueueLiveFrame(encoder.encode(CHAT_COMPLETIONS_HEARTBEAT_COMMENT));
emittedFrames++;
break;
case "response.output_text.delta": {
Expand Down
8 changes: 5 additions & 3 deletions src/protocols/encoders/adapter-events.ts
Original file line number Diff line number Diff line change
Expand Up @@ -186,7 +186,7 @@ export function encodeAdapterEventStream(
const relayed = (observation?: RelayedEventObservation) => {
try { hooks.onRelayed?.(observation ?? {}); } catch { /* counters never break the stream */ }
};
const enqueueCharged = (text: string, observation?: RelayedEventObservation, activity = true): void => {
const enqueueCharged = (text: string, observation?: RelayedEventObservation, activity = true, countAsRelayed = true): void => {
if (closed) return;
const frame = textEncoder.encode(text);
const reservation = budget.reserveTransient(frame.byteLength, { kind: "live_transient" });
Expand All @@ -201,7 +201,9 @@ export function encodeAdapterEventStream(
queuedFrameBytes.push(frame.byteLength);
emittedFrames++;
if (activity) wireActivity = true;
relayed(observation);
// Keepalives are wire bytes, not delivered model events — the legacy bridge's wire-silence
// heartbeat likewise bypasses the relayed-event counter.
if (countAsRelayed) relayed(observation);
};
const sink: ClientFrameSink = {
budget,
Expand Down Expand Up @@ -244,7 +246,7 @@ export function encodeAdapterEventStream(
closed = true;
}
},
emitKeepalive: text => enqueueCharged(text, undefined, false),
emitKeepalive: text => enqueueCharged(text, undefined, false, false),
desiredSize: () => controller.desiredSize ?? 0,
};
const writer = createWriter(sink);
Expand Down
7 changes: 6 additions & 1 deletion src/protocols/encoders/chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import type { AdapterEvent } from "../../types";
import type { TranslatorBudget } from "../../lib/translator-budget";
import { responsesUsage } from "../../bridge/internal";
import {
CHAT_COMPLETIONS_HEARTBEAT_COMMENT,
chatCompletionsErrorResponse,
chatCompletionsFailedResponse,
chatCompletionsIncompleteOutcome,
Expand Down Expand Up @@ -122,7 +123,11 @@ function createChatCompletionWriter(sink: ClientFrameSink, model: string): Clien

return {
start: ensureRole,
heartbeat: ensureRole,
heartbeat() {
// The converter emits the role chunk, then the wire keepalive as an SSE comment.
ensureRole();
sink.emitKeepalive(CHAT_COMPLETIONS_HEARTBEAT_COMMENT);
},
Comment thread
devin-ai-integration[bot] marked this conversation as resolved.
text(delta) {
if (!delta) return;
ensureRole();
Expand Down
46 changes: 46 additions & 0 deletions tests/responses/protocol-direct-encoders-chat.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,8 @@ function normalizeFrames(text: string): unknown[] {
return text.split("\n\n").filter(block => block.trim().length > 0).map(block => {
const data = block.split("\n").filter(line => line.startsWith("data:")).map(line => line.slice(5).trim()).join("");
if (data === "[DONE]") return "[DONE]";
// Comment-only blocks (the wire heartbeat) stay comparable as raw text.
if (data === "") return block;
const parsed = JSON.parse(data) as Record<string, unknown>;
if ("id" in parsed) parsed.id = "ID";
if ("created" in parsed) parsed.created = 0;
Expand Down Expand Up @@ -384,4 +386,48 @@ describe("direct Chat encoder stream lifecycle", () => {
expect(directFrames).toEqual(legacyFrames);
expect(JSON.stringify(directFrames.at(-1))).toContain("upstream_stall_timeout");
});

test("an idle Chat heartbeat stays off the relayed-event counter", async () => {
const ticks: (() => void)[] = [];
const timers = {
setInterval: (handler: () => void) => { ticks.push(handler); return ticks.length - 1; },
clearInterval: () => {},
};
async function* hang(): AsyncGenerator<AdapterEvent> {
yield { type: "text_delta", text: "waiting" };
await new Promise(() => {});
}
let relayedEvents = 0;
const stream = encodeChatCompletionSse(hang(), {
...directOptions(),
heartbeatMs: 1_000,
stallTimeoutSec: 60,
timers,
hooks: { onRelayed: () => { relayedEvents++; } },
});
const reader = stream.getReader();
const decoder = new TextDecoder();
let body = "";
// Carry one pending read across drains; a timed-out read is not abandoned but resumed next drain.
let pending: Promise<ReadableStreamReadResult<Uint8Array>> | undefined;
const drain = async () => {
while (true) {
const read = pending ?? reader.read();
pending = undefined;
const settled = await Promise.race([read, Bun.sleep(10).then(() => undefined)]);
if (settled === undefined) { pending = read; break; }
if (settled.done) return;
body += decoder.decode(settled.value);
}
};
await drain();
const baseline = relayedEvents;
// The first tick only clears the just-consumed upstream/wire activity flags; the heartbeats
// the watchdog emits afterwards are the keepalive comments under test.
for (let i = 0; i < 3; i++) for (const tick of ticks) tick();
await drain();
expect(body).toContain(": opencodex heartbeat");
expect(relayedEvents).toBe(baseline);
await reader.cancel();
});
});
Loading