mirror of
https://github.com/teamchong/pxpipe.git
synced 2026-07-22 02:02:51 +02:00
feat(telemetry): log stop_reason + safety-flag per request
This commit is contained in:
+71
-11
@@ -48,6 +48,12 @@ export interface ProxyEvent {
|
||||
info?: TransformInfo;
|
||||
/** Usage block from Anthropic's response — input/output/cache tokens. */
|
||||
usage?: Usage;
|
||||
/** Model stop reason from the response ("end_turn", "tool_use", "max_tokens",
|
||||
* "refusal", …). "refusal" = safety classifier fired on the transformed request —
|
||||
* scorers must fail cost comparisons on refusal rows (refusals emit almost no
|
||||
* output and would otherwise look "cheaper"). OpenAI finish_reason ("stop",
|
||||
* "length", "content_filter", …) is normalized into the same field. */
|
||||
stopReason?: string;
|
||||
error?: string;
|
||||
/** First ~2 KiB of the upstream 4xx body (not captured on 2xx or 5xx). */
|
||||
errorBody?: string;
|
||||
@@ -112,7 +118,7 @@ export interface OutputMeasurement {
|
||||
function processSseEvent(
|
||||
block: string,
|
||||
m: OutputMeasurement,
|
||||
state: { usage: Usage | undefined },
|
||||
state: { usage: Usage | undefined; stopReason: string | undefined },
|
||||
): void {
|
||||
// Parse `event:` + `data:` lines; continuation data: lines concatenate per SSE spec.
|
||||
let event = '';
|
||||
@@ -136,11 +142,25 @@ function processSseEvent(
|
||||
// OpenAI Responses API streams usage nested under `response` on the terminal
|
||||
// `response.completed` (or `.incomplete`) event — not at the top level.
|
||||
if (event === 'response.completed' || event === 'response.incomplete') {
|
||||
const resp = obj.response as { usage?: unknown } | undefined;
|
||||
const resp = obj.response as
|
||||
| { usage?: unknown; incomplete_details?: { reason?: unknown } }
|
||||
| undefined;
|
||||
const respUsage = normalizeUsage(resp?.usage);
|
||||
if (respUsage) state.usage = respUsage;
|
||||
// Responses API has no stop_reason; normalize the terminal status/reason instead.
|
||||
const reason = resp?.incomplete_details?.reason;
|
||||
state.stopReason = typeof reason === 'string' ? reason
|
||||
: event === 'response.incomplete' ? 'incomplete' : 'stop';
|
||||
}
|
||||
measureOpenAIChoices(obj, m);
|
||||
// OpenAI chat chunks: the final chunk carries choices[].finish_reason (earlier chunks ship null).
|
||||
const choices = obj.choices;
|
||||
if (Array.isArray(choices)) {
|
||||
for (const c of choices) {
|
||||
const fr = (c as { finish_reason?: unknown } | undefined)?.finish_reason;
|
||||
if (typeof fr === 'string') state.stopReason = fr;
|
||||
}
|
||||
}
|
||||
|
||||
if (event === 'message_start') {
|
||||
const msg = obj.message as { usage?: Usage } | undefined;
|
||||
@@ -161,6 +181,9 @@ function processSseEvent(
|
||||
m.toolUseChars += d.partial_json.length;
|
||||
}
|
||||
} else if (event === 'message_delta') {
|
||||
// Anthropic ships the final stop_reason here ("end_turn", "refusal", …).
|
||||
const d = obj.delta as { stop_reason?: unknown } | undefined;
|
||||
if (typeof d?.stop_reason === 'string') state.stopReason = d.stop_reason;
|
||||
// Authoritative final output_tokens; merge over message_start (which ships output_tokens: 1).
|
||||
const u = obj.usage as Partial<Usage> | undefined;
|
||||
if (u) {
|
||||
@@ -260,6 +283,30 @@ function measureFromMessageJson(j: unknown): OutputMeasurement {
|
||||
return m;
|
||||
}
|
||||
|
||||
/** Stop reason from a non-streaming response JSON: Anthropic `stop_reason`,
|
||||
* OpenAI chat `choices[].finish_reason`, Responses `incomplete_details.reason`. */
|
||||
function readStopReasonFromJson(j: unknown): string | undefined {
|
||||
if (!j || typeof j !== 'object') return undefined;
|
||||
const obj = j as {
|
||||
stop_reason?: unknown;
|
||||
choices?: unknown;
|
||||
status?: unknown;
|
||||
incomplete_details?: { reason?: unknown };
|
||||
};
|
||||
if (typeof obj.stop_reason === 'string') return obj.stop_reason;
|
||||
if (Array.isArray(obj.choices)) {
|
||||
for (const c of obj.choices) {
|
||||
const fr = (c as { finish_reason?: unknown } | undefined)?.finish_reason;
|
||||
if (typeof fr === 'string') return fr;
|
||||
}
|
||||
}
|
||||
if (obj.status === 'incomplete') {
|
||||
const reason = obj.incomplete_details?.reason;
|
||||
return typeof reason === 'string' ? reason : 'incomplete';
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
/**
|
||||
* Tee the response body to extract usage + output measurement without blocking the client.
|
||||
* Streams are scanned to EOF (final output_tokens is in message_delta; redacted_thinking
|
||||
@@ -270,6 +317,7 @@ function teeForUsage(res: Response): {
|
||||
usagePromise: Promise<Usage | undefined>;
|
||||
errorBodyPromise: Promise<string | undefined>;
|
||||
measurementPromise: Promise<OutputMeasurement | undefined>;
|
||||
stopReasonPromise: Promise<string | undefined>;
|
||||
} {
|
||||
// No body at all: nothing to extract on either path.
|
||||
if (!res.body) {
|
||||
@@ -278,6 +326,7 @@ function teeForUsage(res: Response): {
|
||||
usagePromise: Promise.resolve(undefined),
|
||||
errorBodyPromise: Promise.resolve(undefined),
|
||||
measurementPromise: Promise.resolve(undefined),
|
||||
stopReasonPromise: Promise.resolve(undefined),
|
||||
};
|
||||
}
|
||||
// 4xx: tee for the error body but skip usage scanning entirely.
|
||||
@@ -313,6 +362,7 @@ function teeForUsage(res: Response): {
|
||||
usagePromise: Promise.resolve(undefined),
|
||||
errorBodyPromise,
|
||||
measurementPromise: Promise.resolve(undefined),
|
||||
stopReasonPromise: Promise.resolve(undefined),
|
||||
};
|
||||
}
|
||||
// 5xx: skip both (the host already synthesizes an error message).
|
||||
@@ -322,15 +372,17 @@ function teeForUsage(res: Response): {
|
||||
usagePromise: Promise.resolve(undefined),
|
||||
errorBodyPromise: Promise.resolve(undefined),
|
||||
measurementPromise: Promise.resolve(undefined),
|
||||
stopReasonPromise: Promise.resolve(undefined),
|
||||
};
|
||||
}
|
||||
const ct = (res.headers.get('content-type') ?? '').toLowerCase();
|
||||
const [forClient, forUs] = res.body.tee();
|
||||
|
||||
// Single read loop resolves both; exposed as separate promises for call-site readability.
|
||||
// Single read loop resolves all three; exposed as separate promises for call-site readability.
|
||||
const scanResult = (async (): Promise<{
|
||||
usage: Usage | undefined;
|
||||
measurement: OutputMeasurement | undefined;
|
||||
stopReason: string | undefined;
|
||||
}> => {
|
||||
const reader = forUs.getReader();
|
||||
const decoder = new TextDecoder();
|
||||
@@ -345,7 +397,10 @@ function teeForUsage(res: Response): {
|
||||
toolUseChars: 0,
|
||||
redactedBlockCount: 0,
|
||||
};
|
||||
const state: { usage: Usage | undefined } = { usage: undefined };
|
||||
const state: { usage: Usage | undefined; stopReason: string | undefined } = {
|
||||
usage: undefined,
|
||||
stopReason: undefined,
|
||||
};
|
||||
while (true) {
|
||||
const { done, value } = await reader.read();
|
||||
if (done) break;
|
||||
@@ -360,7 +415,7 @@ function teeForUsage(res: Response): {
|
||||
}
|
||||
buf += decoder.decode();
|
||||
if (buf.trim().length > 0) processSseEvent(buf, m, state); // trailing partial event
|
||||
return { usage: state.usage, measurement: m };
|
||||
return { usage: state.usage, measurement: m, stopReason: state.stopReason };
|
||||
}
|
||||
|
||||
if (ct.includes('application/json')) {
|
||||
@@ -376,9 +431,10 @@ function teeForUsage(res: Response): {
|
||||
return {
|
||||
usage: normalizeUsage(j?.usage),
|
||||
measurement: measureFromMessageJson(j),
|
||||
stopReason: readStopReasonFromJson(j),
|
||||
};
|
||||
} catch {
|
||||
return { usage: undefined, measurement: undefined };
|
||||
return { usage: undefined, measurement: undefined, stopReason: undefined };
|
||||
}
|
||||
}
|
||||
} catch {
|
||||
@@ -393,7 +449,7 @@ function teeForUsage(res: Response): {
|
||||
} catch {
|
||||
/* ignore */
|
||||
}
|
||||
return { usage: undefined, measurement: undefined };
|
||||
return { usage: undefined, measurement: undefined, stopReason: undefined };
|
||||
})();
|
||||
|
||||
return {
|
||||
@@ -405,6 +461,7 @@ function teeForUsage(res: Response): {
|
||||
usagePromise: scanResult.then((s) => s.usage),
|
||||
errorBodyPromise: Promise.resolve(undefined),
|
||||
measurementPromise: scanResult.then((s) => s.measurement),
|
||||
stopReasonPromise: scanResult.then((s) => s.stopReason),
|
||||
};
|
||||
}
|
||||
|
||||
@@ -573,6 +630,7 @@ export function createProxy(config: ProxyConfig = {}) {
|
||||
usage?: Usage,
|
||||
errorBody?: string,
|
||||
measurement?: OutputMeasurement,
|
||||
stopReason?: string,
|
||||
): void => {
|
||||
const is4xx = status >= 400 && status < 500;
|
||||
// Gzip body lazily (only on 4xx). Async IIFE keeps fire() synchronous.
|
||||
@@ -632,6 +690,7 @@ export function createProxy(config: ProxyConfig = {}) {
|
||||
reqBodySha8,
|
||||
reqBodyGz,
|
||||
measurement,
|
||||
stopReason,
|
||||
});
|
||||
};
|
||||
void finalize();
|
||||
@@ -761,16 +820,17 @@ export function createProxy(config: ProxyConfig = {}) {
|
||||
const firstByteMs = Date.now() - t0;
|
||||
|
||||
// Tee: client gets one side; scanner reads the other for usage/measurement/error body.
|
||||
const { response: teed, usagePromise, errorBodyPromise, measurementPromise } =
|
||||
const { response: teed, usagePromise, errorBodyPromise, measurementPromise, stopReasonPromise } =
|
||||
teeForUsage(upstreamRes);
|
||||
|
||||
// Fire event in background once all three resolve (measurementPromise shares the same stream read).
|
||||
// Fire event in background once all four resolve (all share the same stream read).
|
||||
void Promise.all([
|
||||
usagePromise.catch(() => undefined),
|
||||
errorBodyPromise.catch(() => undefined),
|
||||
measurementPromise.catch(() => undefined),
|
||||
]).then(([usage, errorBody, measurement]) =>
|
||||
fire(upstreamRes.status, info, undefined, firstByteMs, usage, errorBody, measurement),
|
||||
stopReasonPromise.catch(() => undefined),
|
||||
]).then(([usage, errorBody, measurement, stopReason]) =>
|
||||
fire(upstreamRes.status, info, undefined, firstByteMs, usage, errorBody, measurement, stopReason),
|
||||
);
|
||||
|
||||
return new Response(teed.body, {
|
||||
|
||||
@@ -112,6 +112,16 @@ export interface TrackEvent {
|
||||
/** Server-side web search calls billed per-request (not per-token). */
|
||||
web_search_requests?: number;
|
||||
|
||||
/** Model stop reason ("end_turn", "tool_use", "max_tokens", "refusal", …).
|
||||
* OpenAI finish_reason ("stop", "length", "content_filter", …) lands in the same field. */
|
||||
stop_reason?: string;
|
||||
/** True when the stop reason indicates a safety classifier fired ("refusal" /
|
||||
* "content_filter"). Refusal rows emit almost no output and would otherwise
|
||||
* read as "cheap" — scorers MUST fail cost comparisons on these rows, and a
|
||||
* cluster of them after a transform change means the imaged prompt itself is
|
||||
* tripping the classifier (see transform.ts reasoning_extraction notes). */
|
||||
safety_flagged?: boolean;
|
||||
|
||||
/** Ground-truth output chars measured by streaming the response body ourselves — independent of
|
||||
* usage.output_tokens. redacted_block_count_measured counts opaque server-encrypted blocks;
|
||||
* dashboard applies a low/mid/high estimate for those. Absent on non-scannable responses. */
|
||||
@@ -299,9 +309,16 @@ export function toTrackEvent(ev: ProxyEvent): TrackEvent {
|
||||
if (m.redactedBlockCount > 0)
|
||||
out.redacted_block_count_measured = m.redactedBlockCount;
|
||||
}
|
||||
if (ev.stopReason) {
|
||||
out.stop_reason = ev.stopReason;
|
||||
if (SAFETY_STOP_REASONS.has(ev.stopReason)) out.safety_flagged = true;
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
/** Stop reasons that mean a safety classifier fired (Anthropic / OpenAI spellings). */
|
||||
const SAFETY_STOP_REASONS = new Set(['refusal', 'content_filter']);
|
||||
|
||||
/** Writes one JSON line per event. Worker host uses console.log; Node host uses a file-backed variant. */
|
||||
export class JsonLogTracker implements Tracker {
|
||||
constructor(private readonly sink: (line: string) => void = (s) => console.log(s)) {}
|
||||
|
||||
@@ -1230,4 +1230,160 @@ describe('proxy usage extraction', () => {
|
||||
expect(captured).toBeDefined();
|
||||
expect(captured!.measurement?.textChars).toBe(8);
|
||||
});
|
||||
|
||||
it('extracts stop_reason from the SSE message_delta event', async () => {
|
||||
const sseBody =
|
||||
'event: message_start\n' +
|
||||
'data: {"type":"message_start","message":{"id":"x","type":"message","role":"assistant","content":[],"usage":{"input_tokens":10,"output_tokens":1}}}\n\n' +
|
||||
'event: content_block_delta\n' +
|
||||
'data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"ok"}}\n\n' +
|
||||
'event: message_delta\n' +
|
||||
'data: {"type":"message_delta","delta":{"stop_reason":"end_turn","stop_sequence":null},"usage":{"output_tokens":9}}\n\n' +
|
||||
'event: message_stop\n' +
|
||||
'data: {"type":"message_stop"}\n\n';
|
||||
|
||||
const restore = mockUpstream(
|
||||
() =>
|
||||
new Response(sseBody, {
|
||||
status: 200,
|
||||
headers: { 'content-type': 'text/event-stream' },
|
||||
}),
|
||||
);
|
||||
|
||||
let captured: ProxyEvent | undefined;
|
||||
const proxy = createProxy({
|
||||
transform: {},
|
||||
onRequest: (e) => { captured = e; },
|
||||
});
|
||||
|
||||
const res = await proxy(
|
||||
new Request('http://localhost/v1/messages', {
|
||||
method: 'POST',
|
||||
headers: { 'content-type': 'application/json' },
|
||||
body: SAMPLE_REQ_BODY,
|
||||
}),
|
||||
);
|
||||
await res.text();
|
||||
await new Promise((r) => setTimeout(r, 20));
|
||||
restore();
|
||||
|
||||
expect(captured).toBeDefined();
|
||||
expect(captured!.stopReason).toBe('end_turn');
|
||||
// message_delta output_tokens must win over message_start's placeholder 1.
|
||||
expect(captured!.usage?.output_tokens).toBe(9);
|
||||
});
|
||||
|
||||
it('extracts stop_reason "refusal" from a non-stream JSON response', async () => {
|
||||
const restore = mockUpstream(
|
||||
() =>
|
||||
new Response(
|
||||
JSON.stringify({
|
||||
id: 'msg_r',
|
||||
type: 'message',
|
||||
role: 'assistant',
|
||||
content: [],
|
||||
stop_reason: 'refusal',
|
||||
usage: { input_tokens: 5, output_tokens: 2 },
|
||||
}),
|
||||
{ status: 200, headers: { 'content-type': 'application/json' } },
|
||||
),
|
||||
);
|
||||
|
||||
let captured: ProxyEvent | undefined;
|
||||
const proxy = createProxy({
|
||||
transform: {},
|
||||
onRequest: (e) => { captured = e; },
|
||||
});
|
||||
|
||||
const res = await proxy(
|
||||
new Request('http://localhost/v1/messages', {
|
||||
method: 'POST',
|
||||
headers: { 'content-type': 'application/json' },
|
||||
body: SAMPLE_REQ_BODY,
|
||||
}),
|
||||
);
|
||||
await res.text();
|
||||
await new Promise((r) => setTimeout(r, 20));
|
||||
restore();
|
||||
|
||||
expect(captured).toBeDefined();
|
||||
expect(captured!.stopReason).toBe('refusal');
|
||||
});
|
||||
|
||||
it('extracts OpenAI choices[].finish_reason from a JSON body', async () => {
|
||||
const restore = mockUpstream(
|
||||
() =>
|
||||
new Response(
|
||||
JSON.stringify({
|
||||
id: 'chatcmpl-1',
|
||||
object: 'chat.completion',
|
||||
choices: [
|
||||
{
|
||||
index: 0,
|
||||
message: { role: 'assistant', content: '' },
|
||||
finish_reason: 'content_filter',
|
||||
},
|
||||
],
|
||||
usage: { prompt_tokens: 11, completion_tokens: 3 },
|
||||
}),
|
||||
{ status: 200, headers: { 'content-type': 'application/json' } },
|
||||
),
|
||||
);
|
||||
|
||||
let captured: ProxyEvent | undefined;
|
||||
const proxy = createProxy({
|
||||
transform: {},
|
||||
onRequest: (e) => { captured = e; },
|
||||
});
|
||||
|
||||
const res = await proxy(
|
||||
new Request('http://localhost/v1/messages', {
|
||||
method: 'POST',
|
||||
headers: { 'content-type': 'application/json' },
|
||||
body: SAMPLE_REQ_BODY,
|
||||
}),
|
||||
);
|
||||
await res.text();
|
||||
await new Promise((r) => setTimeout(r, 20));
|
||||
restore();
|
||||
|
||||
expect(captured).toBeDefined();
|
||||
expect(captured!.stopReason).toBe('content_filter');
|
||||
});
|
||||
|
||||
it('leaves stopReason undefined when the stream never ships one', async () => {
|
||||
const sseBody =
|
||||
'event: message_start\n' +
|
||||
'data: {"type":"message_start","message":{"id":"x","type":"message","role":"assistant","content":[],"usage":{"input_tokens":10,"output_tokens":1}}}\n\n' +
|
||||
'event: message_stop\n' +
|
||||
'data: {"type":"message_stop"}\n\n';
|
||||
|
||||
const restore = mockUpstream(
|
||||
() =>
|
||||
new Response(sseBody, {
|
||||
status: 200,
|
||||
headers: { 'content-type': 'text/event-stream' },
|
||||
}),
|
||||
);
|
||||
|
||||
let captured: ProxyEvent | undefined;
|
||||
const proxy = createProxy({
|
||||
transform: {},
|
||||
onRequest: (e) => { captured = e; },
|
||||
});
|
||||
|
||||
const res = await proxy(
|
||||
new Request('http://localhost/v1/messages', {
|
||||
method: 'POST',
|
||||
headers: { 'content-type': 'application/json' },
|
||||
body: SAMPLE_REQ_BODY,
|
||||
}),
|
||||
);
|
||||
await res.text();
|
||||
await new Promise((r) => setTimeout(r, 20));
|
||||
restore();
|
||||
|
||||
expect(captured).toBeDefined();
|
||||
expect(captured!.stopReason).toBeUndefined();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -195,6 +195,44 @@ describe('toTrackEvent', () => {
|
||||
expect(out.cwd).toBeUndefined();
|
||||
expect(out.input_tokens).toBeUndefined();
|
||||
});
|
||||
|
||||
it('maps stopReason → stop_reason without flagging benign reasons', () => {
|
||||
const out = toTrackEvent({
|
||||
method: 'POST',
|
||||
path: '/v1/messages',
|
||||
status: 200,
|
||||
durationMs: 10,
|
||||
stopReason: 'end_turn',
|
||||
});
|
||||
expect(out.stop_reason).toBe('end_turn');
|
||||
expect(out.safety_flagged).toBeUndefined();
|
||||
});
|
||||
|
||||
it.each(['refusal', 'content_filter'])(
|
||||
'sets safety_flagged for stop reason %s',
|
||||
(reason) => {
|
||||
const out = toTrackEvent({
|
||||
method: 'POST',
|
||||
path: '/v1/messages',
|
||||
status: 200,
|
||||
durationMs: 10,
|
||||
stopReason: reason,
|
||||
});
|
||||
expect(out.stop_reason).toBe(reason);
|
||||
expect(out.safety_flagged).toBe(true);
|
||||
},
|
||||
);
|
||||
|
||||
it('omits stop_reason entirely when the event has none', () => {
|
||||
const out = toTrackEvent({
|
||||
method: 'POST',
|
||||
path: '/v1/messages',
|
||||
status: 200,
|
||||
durationMs: 10,
|
||||
});
|
||||
expect('stop_reason' in out).toBe(false);
|
||||
expect('safety_flagged' in out).toBe(false);
|
||||
});
|
||||
});
|
||||
|
||||
describe('JsonLogTracker', () => {
|
||||
|
||||
Reference in New Issue
Block a user