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
1 change: 1 addition & 0 deletions changelog.d/fixes/14793-recovery-trace-correlation.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **fix(sse):** mid-stream recovery lines carry the requesting call id and stitched outcomes surface at info ([#14793](https://github.com/diegosouzapw/OmniRoute/pull/14793)) — thanks @maxmad64bis
2 changes: 1 addition & 1 deletion open-sse/handlers/chatCore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3457,7 +3457,7 @@ export async function handleChatCore({
)
),
continueStream,
...buildContinuationLogHooks(log),
...buildContinuationLogHooks(log, correlationId),
throughputWatchdog,
onWatchdogAbort: () =>
log?.warn?.(
Expand Down
26 changes: 19 additions & 7 deletions open-sse/handlers/chatCore/recoveryTraceLogging.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,12 @@
*
* Levels: the continuation attempt line keeps its release wording at warn; a recovery that
* gives up (the continuation budget is spent, or the continuation request returned no
* stream) is warn; every other outcome — stitched suffix, overlap rejection, terminal or
* empty continuation, a cut refused because of a tool call — is debug, so a healthy stream
* never adds a warn line. Every line carries `attempt N/MAX` so it joins the attempt line.
* stream) is warn; a cut refused because a tool call is in flight is debug; every other
* outcome — stitched suffix, overlap rejection, terminal or empty continuation, a cut
* refused for another reason — is info, so a healthy stream never adds a warn line but a
* stitched recovery stays visible outside debug. Every line carries
* `correlationId=<id|none>` so it joins the requesting call, plus `attempt N/MAX` so it
* joins the attempt line.
*/
import { STREAM_RECOVERY } from "../../config/constants.ts";
import type {
Expand All @@ -18,10 +21,14 @@ type RecoveryLogger =
| {
warn?: (tag: string, message: string) => void;
debug?: (tag: string, message: string) => void;
info?: (tag: string, message: string) => void;
}
| null
| undefined;

/** Correlated log line suffix. Falls back to "none" like the empty-string fallback in formatBufferedVerdictLog. */
type CorrelationId = string | null | undefined;

const TAG = "STREAM_RECOVERY";
const MAX = STREAM_RECOVERY.EARLY_RETRY_MAX;

Expand All @@ -46,14 +53,19 @@ export function isContinuationGiveUp(event: ContinuationOutcome): boolean {
}

export function buildContinuationLogHooks(
log: RecoveryLogger
log: RecoveryLogger,
correlationId?: CorrelationId
): Pick<RecoverableStreamOptions, "onContinue" | "onContinueOutcome"> {
// Same empty-string-to-"none" fallback as formatBufferedVerdictLog, without trim.
const cid = correlationId && correlationId.length > 0 ? correlationId : "none";
return {
onContinue: (attempt) => log?.warn?.(TAG, `mid-stream continuation attempt ${attempt}/${MAX}`),
onContinue: (attempt) =>
log?.warn?.(TAG, `mid-stream continuation attempt ${attempt}/${MAX} correlationId=${cid}`),
onContinueOutcome: (event) => {
const line = formatContinuationOutcome(event);
const line = `${formatContinuationOutcome(event)} correlationId=${cid}`;
if (isContinuationGiveUp(event)) log?.warn?.(TAG, line);
else log?.debug?.(TAG, line);
else if (event.outcome === "refused" && event.reason === "tool-call") log?.debug?.(TAG, line);
else log?.info?.(TAG, line);
},
};
}
21 changes: 15 additions & 6 deletions tests/unit/chatcore-stream-recovery-log-wiring.test.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,8 @@
// handleChatCore wiring of the mid-stream continuation log hooks: a real streaming request
// with stream recovery + mid-stream continuation enabled, an upstream that commits the
// holdback window and then drops, and a continuation that finishes the answer. The injected
// log must receive the release attempt line at warn and the stitched outcome at debug.
// log must receive the release attempt line at warn and the stitched outcome at info,
// both carrying the requesting call's correlationId.
import { after, before, test } from "node:test";
import assert from "node:assert/strict";
import fs from "node:fs";
Expand Down Expand Up @@ -72,8 +73,9 @@ test("handleChatCore routes continuation logs through buildContinuationLogHooks"

const warn: string[] = [];
const debug: string[] = [];
const info: string[] = [];
const log = {
info() {},
info: (tag: string, msg: string) => info.push(`${tag} ${msg}`),
error() {},
warn: (tag: string, msg: string) => warn.push(`${tag} ${msg}`),
debug: (tag: string, msg: string) => debug.push(`${tag} ${msg}`),
Expand All @@ -95,6 +97,7 @@ test("handleChatCore routes continuation logs through buildContinuationLogHooks"
userAgent: "unit-test",
isCombo: false,
log,
correlationId: "req-wiring-1",
} as unknown as Parameters<typeof handleChatCore>[0]);

const response = (result as { response?: Response }).response;
Expand All @@ -104,11 +107,17 @@ test("handleChatCore routes continuation logs through buildContinuationLogHooks"
assert.equal(calls, 2, "one upstream request plus one continuation");
assert.match(text, /nice to meet you!/);
const recoveryWarns = warn.filter((l) => l.startsWith("STREAM_RECOVERY "));
assert.deepEqual(recoveryWarns, ["STREAM_RECOVERY mid-stream continuation attempt 1/4"]);
assert.deepEqual(recoveryWarns, [
"STREAM_RECOVERY mid-stream continuation attempt 1/4 correlationId=req-wiring-1",
]);
assert.deepEqual(
debug.filter((l) => l.startsWith("STREAM_RECOVERY ")),
[]
);
assert.ok(
debug.includes(
"STREAM_RECOVERY mid-stream continuation attempt 1/4 outcome=suffix suffixChars=19"
info.includes(
"STREAM_RECOVERY mid-stream continuation attempt 1/4 outcome=suffix suffixChars=19 correlationId=req-wiring-1"
),
debug.filter((l) => l.startsWith("STREAM_RECOVERY")).join(" | ")
info.filter((l) => l.startsWith("STREAM_RECOVERY")).join(" | ")
);
});
106 changes: 84 additions & 22 deletions tests/unit/stream-recovery-trace-logging.test.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
// Mid-stream continuation log wiring: buildContinuationLogHooks (the exact hooks chatCore
// spreads into createRecoverableStream) driven through the real recoverable stream. Warn is
// reserved for the attempt line (release wording) and for a recovery that gives up; every
// other outcome is debug, and a healthy or tool-call stream adds no line at all.
// reserved for the attempt line (release wording) and for a recovery that gives up; a cut
// refused because of a tool call stays debug; every other outcome is info, so a stitched
// recovery stays visible outside debug and a healthy stream adds no warn line at all.
// Every emitted line carries `correlationId=<id|none>` from the requesting call.
import { after, test } from "node:test";
import assert from "node:assert/strict";

Expand All @@ -22,6 +24,8 @@ after(() => {

const enc = new TextEncoder();

const CID = "cid-test";

function steppingClock() {
let t = 0;
return () => (t += 1000);
Expand Down Expand Up @@ -64,14 +68,16 @@ const TOOL_CALL =
'data: {"choices":[{"delta":{"tool_calls":[{"id":"c1","function":{"name":"f"}}]}}]}\n\n';
const FINISH_TOOL_CALLS = 'data: {"choices":[{"delta":{},"finish_reason":"tool_calls"}]}\n\n';

function capture() {
function capture(correlationId: string | null = CID) {
const warn: string[] = [];
const debug: string[] = [];
const info: string[] = [];
const log = {
warn: (tag: string, msg: string) => warn.push(`${tag} ${msg}`),
debug: (tag: string, msg: string) => debug.push(`${tag} ${msg}`),
info: (tag: string, msg: string) => info.push(`${tag} ${msg}`),
};
return { warn, debug, hooks: buildContinuationLogHooks(log) };
return { warn, debug, info, log, hooks: buildContinuationLogHooks(log, correlationId) };
}

async function run(
Expand All @@ -89,29 +95,33 @@ async function run(
);
}

test("a stitched continuation warns once with the release attempt wording, outcome at debug", async () => {
const { warn, debug, hooks } = capture();
test("a stitched continuation warns once with the release attempt wording, outcome at info", async () => {
const { warn, debug, info, hooks } = capture();
const out = await run(
streamFrom([ROLE, content("Hello there world")]),
async () => streamFrom([ROLE, content("there world, nice to meet you!"), DONE]),
hooks
);
assert.match(out, /nice to meet you!/);
assert.deepEqual(warn, ["STREAM_RECOVERY mid-stream continuation attempt 1/4"]);
assert.deepEqual(debug, [
"STREAM_RECOVERY mid-stream continuation attempt 1/4 outcome=suffix suffixChars=19",
assert.deepEqual(warn, [
"STREAM_RECOVERY mid-stream continuation attempt 1/4 correlationId=cid-test",
]);
assert.deepEqual(debug, []);
assert.deepEqual(info, [
"STREAM_RECOVERY mid-stream continuation attempt 1/4 outcome=suffix suffixChars=19 correlationId=cid-test",
]);
});

test("a streamed tool call that ends nominally logs nothing", async () => {
const { warn, debug, hooks } = capture();
const { warn, debug, info, hooks } = capture();
await run(streamFrom([ROLE, TOOL_CALL, FINISH_TOOL_CALLS, DONE]), async () => null, hooks);
assert.deepEqual(warn, []);
assert.deepEqual(debug, []);
assert.deepEqual(info, []);
});

test("a cut refused because a tool call is in flight is debug only", async () => {
const { warn, debug, hooks } = capture();
const { warn, debug, info, hooks } = capture();
let calls = 0;
await run(
streamFrom([ROLE, content("Let me check. "), TOOL_CALL], true),
Expand All @@ -123,13 +133,14 @@ test("a cut refused because a tool call is in flight is debug only", async () =>
);
assert.equal(calls, 0);
assert.deepEqual(warn, []);
assert.deepEqual(info, []);
assert.deepEqual(debug, [
"STREAM_RECOVERY mid-stream continuation attempt 0/4 outcome=refused reason=tool-call",
"STREAM_RECOVERY mid-stream continuation attempt 0/4 outcome=refused reason=tool-call correlationId=cid-test",
]);
});

test("a spent continuation budget warns that the recovery gave up", async () => {
const { warn, debug, hooks } = capture();
const { warn, debug, info, hooks } = capture();
let calls = 0;
await run(
streamFrom([ROLE, content("Hello there world")], true),
Expand All @@ -142,34 +153,85 @@ test("a spent continuation budget warns that the recovery gave up", async () =>
);
assert.equal(calls, 4);
assert.deepEqual(warn, [
"STREAM_RECOVERY mid-stream continuation attempt 1/4",
"STREAM_RECOVERY mid-stream continuation attempt 2/4",
"STREAM_RECOVERY mid-stream continuation attempt 3/4",
"STREAM_RECOVERY mid-stream continuation attempt 4/4",
"STREAM_RECOVERY mid-stream continuation attempt 4/4 outcome=refused reason=budget",
"STREAM_RECOVERY mid-stream continuation attempt 1/4 correlationId=cid-test",
"STREAM_RECOVERY mid-stream continuation attempt 2/4 correlationId=cid-test",
"STREAM_RECOVERY mid-stream continuation attempt 3/4 correlationId=cid-test",
"STREAM_RECOVERY mid-stream continuation attempt 4/4 correlationId=cid-test",
"STREAM_RECOVERY mid-stream continuation attempt 4/4 outcome=refused reason=budget correlationId=cid-test",
]);
assert.deepEqual(debug, []);
assert.deepEqual(info, []);
});

test("a continuation request that returns no stream warns that the recovery gave up", async () => {
const { warn, debug, hooks } = capture();
const { warn, debug, info, hooks } = capture();
await run(streamFrom([ROLE, content("Hello there world")], true), async () => null, hooks);
assert.deepEqual(warn, [
"STREAM_RECOVERY mid-stream continuation attempt 1/4",
"STREAM_RECOVERY mid-stream continuation attempt 1/4 outcome=no-stream",
"STREAM_RECOVERY mid-stream continuation attempt 1/4 correlationId=cid-test",
"STREAM_RECOVERY mid-stream continuation attempt 1/4 outcome=no-stream correlationId=cid-test",
]);
assert.deepEqual(debug, []);
assert.deepEqual(info, []);
});

test("a non-OpenAI body ending without an OpenAI terminal logs nothing", async () => {
const { warn, debug, hooks } = capture();
const { warn, debug, info, hooks } = capture();
await run(
streamFrom(['event: content_block_delta\ndata: {"type":"content_block_delta"}\n\n']),
async () => null,
hooks
);
assert.deepEqual(warn, []);
assert.deepEqual(debug, []);
assert.deepEqual(info, []);
});

test("an absent or empty correlationId falls back to none, never undefined/null", async () => {
const lines: string[] = [];
for (const cid of [undefined, null, ""] as const) {
const c = capture(CID);
const hooks =
cid === undefined ? buildContinuationLogHooks(c.log) : buildContinuationLogHooks(c.log, cid);
await run(
streamFrom([ROLE, content("Hello there world")]),
async () => streamFrom([ROLE, content("there world, nice to meet you!"), DONE]),
hooks
);
assert.deepEqual(c.warn, [
"STREAM_RECOVERY mid-stream continuation attempt 1/4 correlationId=none",
]);
assert.deepEqual(c.debug, []);
assert.deepEqual(c.info, [
"STREAM_RECOVERY mid-stream continuation attempt 1/4 outcome=suffix suffixChars=19 correlationId=none",
]);
lines.push(...c.warn, ...c.debug, ...c.info);
}
for (const line of lines) assert.doesNotMatch(line, /undefined|null/);
});

test("a cut refused for a non-tool-call reason is info, not debug", () => {
const { warn, debug, info, hooks } = capture();
hooks.onContinueOutcome?.({ attempt: 0, outcome: "refused", reason: "not-continuable" });
assert.deepEqual(warn, []);
assert.deepEqual(debug, []);
assert.deepEqual(info, [
"STREAM_RECOVERY mid-stream continuation attempt 0/4 outcome=refused reason=not-continuable correlationId=cid-test",
]);
});

test("a log without info still handles a stitched outcome without throwing", () => {
const warn: string[] = [];
const debug: string[] = [];
const legacy = {
warn: (tag: string, msg: string) => warn.push(`${tag} ${msg}`),
debug: (tag: string, msg: string) => debug.push(`${tag} ${msg}`),
};
const hooks = buildContinuationLogHooks(legacy, CID);
assert.doesNotThrow(() =>
hooks.onContinueOutcome?.({ attempt: 1, outcome: "suffix", suffixChars: 3 })
);
assert.deepEqual(warn, []);
assert.deepEqual(debug, []);
});

test("every outcome formats with the attempt token and no undefined fields", () => {
Expand Down
Loading