Skip to content
Merged
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
61 changes: 43 additions & 18 deletions src/bridge.ts
Original file line number Diff line number Diff line change
Expand Up @@ -133,8 +133,18 @@ export function bridgeToResponsesSSE(
* from this callback instead of re-parsing the bridged SSE.
*/
onUsage?: (usage: OcxUsage | undefined) => void;
/**
* Test seam for the wire/stall beat loop. Production omits this and uses the
* global timers; injecting here must not change scheduling semantics.
*/
timers?: {
setInterval: (handler: () => void, ms: number) => unknown;
clearInterval: (id: unknown) => void;
};
},
): ReadableStream<Uint8Array> {
const setBeatInterval = options?.timers?.setInterval ?? ((handler: () => void, ms: number) => setInterval(handler, ms));
const clearBeatInterval = options?.timers?.clearInterval ?? ((id: unknown) => clearInterval(id as ReturnType<typeof setInterval>));
// Freeform/custom tools (apply_patch) carry their body in `input`; the model is given a
// function with `{input:string}`, so unwrap it here when relaying back as a custom_tool_call.
const freeformInput = (args: string): string => {
Expand Down Expand Up @@ -194,7 +204,7 @@ export function bridgeToResponsesSSE(
// buffering + raw-byte progress). Upstream activity only resets the stall watchdog.
let upstreamActivity = false;
let wireActivity = false;
let beat: ReturnType<typeof setInterval> | undefined;
let beat: unknown;
let controller: ReadableStreamDefaultController<Uint8Array>;
let emittedFrames = 0;
let gated = false;
Expand Down Expand Up @@ -515,7 +525,11 @@ export function bridgeToResponsesSSE(
if (currentRawReasoning) closeCurrentRawReasoning();
flushHiddenRawReasoning();
if (currentToolCall) closeCurrentToolCall();
if (currentMsg && currentMsg.phase !== event.phase) closeCurrentMessage("commentary");
// Only flush on an explicit phase change. A later delta that omits `phase` must
// keep appending to the current message rather than wiping the earlier phase.
if (currentMsg && event.phase !== undefined && currentMsg.phase !== event.phase) {
closeCurrentMessage("commentary");
}
if (!currentMsg) {
const itemId = `msg_${uuid()}`;
const item = {
Expand Down Expand Up @@ -809,7 +823,7 @@ export function bridgeToResponsesSSE(
stepping = false;
return;
}
if (beat) { clearInterval(beat); beat = undefined; }
if (beat !== undefined) { clearBeatInterval(beat); beat = undefined; }

if (!terminated) {
// The adapter generator ended without an explicit done/error event. Mark as incomplete
Expand Down Expand Up @@ -848,7 +862,7 @@ export function bridgeToResponsesSSE(
// The default ReadableStream strategy has HWM=1. Once one event's frames fill that
// queue, pull stepping pauses; no custom FIFO or queuing strategy is layered on top.
gated = true;
beat = setInterval(() => {
beat = setBeatInterval(() => {
if (closed || gated) return;
if (upstreamActivity) {
upstreamActivity = false;
Expand All @@ -871,7 +885,7 @@ export function bridgeToResponsesSSE(
terminated = true;
returnIterator();
emitDone();
if (beat) clearInterval(beat);
if (beat !== undefined) clearBeatInterval(beat);
beat = undefined;
try { controller.close(); } catch { /* already closed */ }
closed = true;
Expand Down Expand Up @@ -902,14 +916,14 @@ export function bridgeToResponsesSSE(
cancel() {
// Client (Codex) disconnected. Stop emitting and let the caller abort the upstream fetch so a
// cancelled turn does not leak the upstream stream or keep draining tokens (RC2).
clientCancelled = true;
closed = true;
if (beat) clearInterval(beat);
onCancel?.();
returnIterator();
},
});
}
clientCancelled = true;
closed = true;
if (beat !== undefined) clearBeatInterval(beat);
onCancel?.();
returnIterator();
},
});
}

export function buildResponseJSON(
events: AdapterEvent[],
Expand Down Expand Up @@ -1048,15 +1062,17 @@ export function buildResponseJSON(
flushToolCall();
break;
case "text_delta":
if (currentText && currentTextPhase !== e.phase) flushText("commentary");
// Only flush on an explicit phase change. A later delta that omits `phase` must keep
// appending under the previously established phase.
if (currentText && e.phase !== undefined && currentTextPhase !== e.phase) flushText("commentary");
if (currentSummaryReasoning) flushSummaryReasoning();
if (currentRawReasoning) flushRawReasoning();
if (currentToolCallId) flushToolCall();
// Compaction turns keep the summary out of normal message output (replay dedup — see
// bridgeToResponsesSSE); it ships only inside the synthetic compaction item below.
if (options?.compaction) compactionText += e.text;
else {
currentTextPhase = e.phase;
if (e.phase !== undefined) currentTextPhase = e.phase;
currentText += e.text;
}
break;
Expand Down Expand Up @@ -1131,7 +1147,8 @@ export function buildResponseJSON(
endTurn = e.endTurn;
cleanDone = e.stopReason === undefined;
if (e.providerState) options?.onProviderState?.(e.providerState);
if (e.stopReason === "max_tokens") stopReason = "max_tokens";
// Match streaming: max_tokens and content_filter both terminate as incomplete.
if (e.stopReason === "max_tokens" || e.stopReason === "content_filter") stopReason = e.stopReason;
break;
}
}
Expand All @@ -1141,14 +1158,20 @@ export function buildResponseJSON(
flushToolCall();
// A truncated turn must never be installed as replacement history: emit the
// compaction item only when the turn actually completed (#422).
if (options?.compaction && !errorEvent && !incompleteEvent && stopReason !== "max_tokens") {
if (
options?.compaction
&& !errorEvent
&& !incompleteEvent
&& stopReason !== "max_tokens"
&& stopReason !== "content_filter"
) {
output.push({ type: "compaction", id: `cmp_${uuid()}`, encrypted_content: encodeCompactionSummary(compactionText) });
}

const failure = errorEvent ? adapterFailureFromEvent(errorEvent) : undefined;
const status = errorEvent
? "failed"
: incompleteEvent || stopReason === "max_tokens"
: incompleteEvent || stopReason === "max_tokens" || stopReason === "content_filter"
? "incomplete"
: "completed";
options?.onUsage?.(incompleteEvent?.usage ?? usage);
Expand All @@ -1168,6 +1191,8 @@ export function buildResponseJSON(
},
} : stopReason === "max_tokens" ? {
incomplete_details: { reason: "max_output_tokens" },
} : stopReason === "content_filter" ? {
incomplete_details: { reason: "content_filter" },
} : {}),
usage: responsesUsage(incompleteEvent?.usage ?? usage),
};
Expand Down
171 changes: 156 additions & 15 deletions tests/bridge.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -529,6 +529,34 @@ describe("Responses bridge reasoning and usage parity", () => {
expect((explicit.output as Record<string, unknown>[])[0]).toMatchObject({ phase: "commentary" });
});

test("later text_delta omitting phase keeps the prior explicit phase", async () => {
const frames = await collectSse(bridgeToResponsesSSE(replay([
{ type: "text_delta", text: "Hello ", phase: "final_answer" },
{ type: "text_delta", text: "world." },
{ type: "done" },
]), "routed/chat-model"));
const completed = frames.find(frame => frame.event === "response.completed")?.data.response as Record<string, unknown>;
const output = completed.output as Record<string, unknown>[];
expect(output).toHaveLength(1);
expect(output[0]).toMatchObject({
type: "message",
phase: "final_answer",
content: [{ type: "output_text", text: "Hello world." }],
});

const batch = buildResponseJSON([
{ type: "text_delta", text: "Hello ", phase: "final_answer" },
{ type: "text_delta", text: "world." },
{ type: "done" },
], "routed/chat-model");
expect((batch.output as Record<string, unknown>[])).toHaveLength(1);
expect((batch.output as Record<string, unknown>[])[0]).toMatchObject({
type: "message",
phase: "final_answer",
content: [{ type: "output_text", text: "Hello world." }],
});
});

test("structured adapter errors override message heuristics", () => {
const json = buildResponseJSON([
{
Expand Down Expand Up @@ -684,44 +712,146 @@ describe("Responses bridge reasoning and usage parity", () => {
// calls, the adapter emits `heartbeat` events. They must keep the stall watchdog alive (no
// upstream_stall_timeout). Adapter heartbeats themselves are not translated into Responses
// protocol items; wire keepalives use a separate `response.heartbeat` frame (see next test).
//
// resolveStallTimeoutSec ceils to a minimum of 1s, so sub-second stallTimeoutSec values cannot
// prove the reset. Drive the beat loop through a test clock seam and run adapter-only progress
// past the effective deadline.
const heartbeatMs = 50;
const stallTimeoutSec = 1; // effective after resolveStallTimeoutSec
const maxStallTicks = Math.ceil((stallTimeoutSec * 1000) / heartbeatMs);
const cycles = maxStallTicks + 5; // wall-clock equivalent >> stall deadline

let beatTick: (() => void) | undefined;
const timers = {
setInterval(handler: () => void, _ms: number) {
beatTick = handler;
return 1;
},
clearInterval(_id: unknown) {
beatTick = undefined;
},
};

let waitResolve: (() => void) | undefined;
const waitDelay = () => new Promise<void>(resolve => { waitResolve = resolve; });
const releaseDelay = () => {
const resolve = waitResolve;
waitResolve = undefined;
resolve?.();
};
const flush = async () => {
for (let i = 0; i < 20; i++) await Promise.resolve();
};

async function* heartbeatsThenDone(): AsyncGenerator<AdapterEvent> {
// More heartbeats than maxStallTicks would allow if they did NOT reset the counter.
for (let i = 0; i < 6; i++) {
for (let i = 0; i < cycles; i++) {
yield { type: "heartbeat" };
await new Promise(r => setTimeout(r, 12));
await waitDelay();
}
yield { type: "text_delta", text: "ok" };
yield { type: "done" };
}
// heartbeatMs=10ms, stallTimeoutSec=0.03s -> maxStallTicks=3. 6 spaced heartbeats only survive
// if each one resets stallTicks.
const frames = await collectSse(bridgeToResponsesSSE(
heartbeatsThenDone(), "model", undefined, undefined, undefined, undefined, 10, { stallTimeoutSec: 0.03 },

const framesPromise = collectSse(bridgeToResponsesSSE(
heartbeatsThenDone(),
"model",
undefined,
undefined,
undefined,
undefined,
heartbeatMs,
{ stallTimeoutSec, timers },
));
expect(frames.some(f => (f.data.response as Record<string, unknown> | undefined)?.incomplete_details)).toBe(false);

// First heartbeat is pulled; step blocks on the delay gate with gated=false.
await flush();
for (let i = 0; i < cycles; i++) {
// maxStallTicks silent ticks would trip the watchdog if the preceding heartbeat did not
// count as upstream activity (first tick clears the flag; the rest must not reach the limit).
// Heartbeats between cycles reset the counter, so the stream must survive the full run.
for (let t = 0; t < maxStallTicks; t++) beatTick?.();
releaseDelay();
await flush();
}

const frames = await framesPromise;
expect(frames.some(f => {
const response = f.data.response as Record<string, unknown> | undefined;
const details = response?.incomplete_details as Record<string, unknown> | undefined;
return details?.reason === "upstream_stall_timeout";
})).toBe(false);
expect(frames.some(f => f.event === "response.completed")).toBe(true);
// Adapter heartbeats must not be mis-translated into a rich protocol event of their own.
expect(frames.some(f => f.event === "response.heartbeat" && f.data.type === "heartbeat" && Object.keys(f.data).length > 2)).toBe(false);
// Adapter heartbeats must not be mis-translated into a protocol event of their own.
expect(frames.some(f => f.data.type === "heartbeat")).toBe(false);
});

test("wire response.heartbeat keeps firing while only adapter heartbeats flow", async () => {
// Issue #521: web-search buffers semantic events and yields invisible adapter heartbeats from
// raw-byte progress. Those must not suppress wire keepalives, or Codex Desktop idle-timeouts
// (~5 min) while OCX still considers the upstream alive.
const heartbeatMs = 50;
const stallTimeoutSec = 1;
const cycles = 4;

let beatTick: (() => void) | undefined;
const timers = {
setInterval(handler: () => void, _ms: number) {
beatTick = handler;
return 1;
},
clearInterval(_id: unknown) {
beatTick = undefined;
},
};

let waitResolve: (() => void) | undefined;
const waitDelay = () => new Promise<void>(resolve => { waitResolve = resolve; });
const releaseDelay = () => {
const resolve = waitResolve;
waitResolve = undefined;
resolve?.();
};
const flush = async () => {
for (let i = 0; i < 20; i++) await Promise.resolve();
};

async function* adapterHeartbeatsOnly(): AsyncGenerator<AdapterEvent> {
for (let i = 0; i < 8; i++) {
for (let i = 0; i < cycles; i++) {
yield { type: "heartbeat" };
await new Promise(r => setTimeout(r, 15));
await waitDelay();
}
yield { type: "text_delta", text: "ok" };
yield { type: "done" };
}
const frames = await collectSse(bridgeToResponsesSSE(
adapterHeartbeatsOnly(), "model", undefined, undefined, undefined, undefined, 10, { stallTimeoutSec: 1 },

const framesPromise = collectSse(bridgeToResponsesSSE(
adapterHeartbeatsOnly(),
"model",
undefined,
undefined,
undefined,
undefined,
heartbeatMs,
{ stallTimeoutSec, timers },
));
expect(frames.some(f => f.event === "response.heartbeat" && f.data.type === "response.heartbeat")).toBe(true);

await flush();
for (let i = 0; i < cycles; i++) {
// Several silent beat ticks per adapter-only gap → multiple wire keepalives.
for (let t = 0; t < 3; t++) beatTick?.();
releaseDelay();
await flush();
}

const frames = await framesPromise;
const wireHeartbeats = frames.filter(f =>
f.event === "response.heartbeat" && f.data.type === "response.heartbeat"
);
expect(wireHeartbeats.length).toBeGreaterThan(1);
expect(frames.some(f => f.event === "response.completed")).toBe(true);
expect(frames.some(f => (f.data.response as Record<string, unknown> | undefined)?.incomplete_details)).toBe(false);
// Reject every adapter-shaped heartbeat payload, regardless of event name or field count.
expect(frames.some(f => f.data.type === "heartbeat")).toBe(false);
});
});

Expand Down Expand Up @@ -872,6 +1002,17 @@ describe("Responses bridge stopReason threading (issue #246)", () => {
expect(json.incomplete_details).toEqual({ reason: "max_output_tokens" });
});

test("batch buildResponseJSON with stopReason content_filter returns incomplete status", () => {
const json = buildResponseJSON([
{ type: "text_delta", text: "partial" },
{ type: "done", stopReason: "content_filter" },
], "routed/model", { compaction: true });
expect(json.status).toBe("incomplete");
expect(json.incomplete_details).toEqual({ reason: "content_filter" });
// Truncated turns must not install a compaction replacement (#422).

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🎯 Functional Correctness | 🟠 Major | ⚡ Quick win

Cover and suppress streaming compaction for content-filter stops.

This test only exercises buildResponseJSON. bridgeToResponsesSSE still emits its synthetic compaction item before handling event.stopReason === "content_filter", so streamed responses can contain the base64 compaction payload while batch responses correctly omit it. Guard the streaming compaction branch for max_tokens/content_filter and add a streamed { compaction: true } regression assertion.

Proposed fix
- if (options?.compaction) {
+ if (
+   options?.compaction &&
+   event.stopReason !== "max_tokens" &&
+   event.stopReason !== "content_filter"
+ ) {

As per path instructions, “Tests are flat Bun tests under tests/. A behavior change in src/ should come with a focused regression test near the existing tests for that subsystem.”

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@tests/bridge.test.ts` at line 1012, Update bridgeToResponsesSSE so its
synthetic compaction branch does not emit a compaction item when the stop reason
is max_tokens or content_filter, matching buildResponseJSON behavior. Extend the
nearby streaming tests with a regression assertion confirming the streamed
result contains no { compaction: true } item for a content-filter stop.

Source: Path instructions

expect((json.output as Record<string, unknown>[]).some(item => item.type === "compaction")).toBe(false);
});
Comment thread
coderabbitai[bot] marked this conversation as resolved.

test("batch buildResponseJSON without stopReason returns completed status", () => {
const json = buildResponseJSON([
{ type: "text_delta", text: "hello" },
Expand Down
Loading