diff --git a/CHANGELOG.md b/CHANGELOG.md index b43c64f..18032f1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -30,6 +30,7 @@ The format is based on Keep a Changelog and this project follows Semantic Versio - **Canonical provider metric normalization** — benchmark response normalization now uses an internal typed `metrics-v2` observation contract for registered Ollama, OpenAI Chat, Anthropic Messages, and Gemini GenerateContent usage and timing fields while persisted benchmark results remain on the existing `metrics-v1` contract. - **Canonical client metric normalization** — benchmark execution now records transient `metrics-v2` observations for operation and successful-attempt latency, retry overhead, attempt count, request and normalization health, terminal timeouts, stream completion, and first transport chunk timing while retaining the persisted `metrics-v1` result shape. - **Canonical semantic stream timing** — streamed benchmark execution now distinguishes first transport bytes from meaningful text or tool-call output and records transient first-output, first-tool-call, tool-readiness, and last-output timestamps without changing persisted `metrics-v1` results. +- **Canonical per-request derived observations** — benchmark execution now composes client and provider observations and transiently calculates registered generation-window, token-efficiency, server-rate, and output/input formulas while persisted results remain on `metrics-v1`. ### Fixed diff --git a/backend/src/services/benchmark-provider-metrics.ts b/backend/src/services/benchmark-provider-metrics.ts index a91f1fc..f258786 100644 --- a/backend/src/services/benchmark-provider-metrics.ts +++ b/backend/src/services/benchmark-provider-metrics.ts @@ -414,6 +414,14 @@ function normalizeGemini( accountingScope: { includes: ['prompt', 'thoughts', 'response_candidates'] } } ]); + const candidateCount = Array.isArray(metadata.candidates) ? metadata.candidates.length : null; + const outputTokens = observations.find((observation) => observation.metric_id === 'output_tokens'); + if (outputTokens?.accounting_scope) { + outputTokens.accounting_scope = { + ...outputTokens.accounting_scope, + candidate_count: candidateCount + }; + } const inputTokens = measuredValue(observations, 'input_tokens'); const cachedInputTokens = measuredValue(observations, 'cached_input_tokens'); if (inputTokens !== null && cachedInputTokens !== null && cachedInputTokens <= inputTokens) { diff --git a/backend/src/services/benchmark-request-metrics.ts b/backend/src/services/benchmark-request-metrics.ts new file mode 100644 index 0000000..ea90ad9 --- /dev/null +++ b/backend/src/services/benchmark-request-metrics.ts @@ -0,0 +1,341 @@ +import { + normalizeClientMetricObservations, + type ClientOperationTelemetry +} from './benchmark-client-metrics.js'; +import type { + MetricObservation, + MetricObservationStatus +} from './benchmark-metric-observations.js'; + +export interface ComposeRequestMetricObservationsInput { + clientTelemetry: ClientOperationTelemetry; + providerObservations: MetricObservation[]; +} + +function measuredObservation( + observations: MetricObservation[], + metricId: string +): MetricObservation | null { + return observations.find( + (candidate) => candidate.metric_id === metricId + && candidate.status === 'measured' + && typeof candidate.value === 'number' + ) ?? null; +} + +function numericValue(observation: MetricObservation | null): number | null { + return observation && typeof observation.value === 'number' + ? observation.value + : null; +} + +function derivedObservation(input: { + metricId: string; + unit: string; + value?: number; + status?: MetricObservationStatus; + reason?: string; + formula: string; + components: MetricObservation[]; + accountingScope?: Record; +}): MetricObservation { + const providerReference = input.components.find( + (component) => component.provider_protocol !== null || component.provider_id !== null + ); + const componentScopes = input.components + .filter((component) => component.accounting_scope !== null) + .map((component) => ({ + metric_id: component.metric_id, + accounting_scope: component.accounting_scope + })); + return { + metric_id: input.metricId, + value: input.value ?? null, + unit: input.unit, + status: input.status ?? 'measured', + reason: input.reason ?? null, + source: 'derived', + metric_version: 'metrics-v2', + provider_id: providerReference?.provider_id ?? null, + provider_protocol: providerReference?.provider_protocol ?? null, + provider_version: providerReference?.provider_version ?? null, + native_field: null, + native_value: null, + native_unit: null, + normalization: input.formula, + accounting_scope: { + component_metric_ids: input.components.map((component) => component.metric_id), + ...(componentScopes.length > 0 ? { component_accounting_scopes: componentScopes } : {}), + ...(input.accountingScope ?? {}) + } + }; +} + +function shareProviderScope(left: MetricObservation, right: MetricObservation): boolean { + if ( + left.provider_protocol !== null + && right.provider_protocol !== null + && left.provider_protocol !== right.provider_protocol + ) { + return false; + } + return !( + left.provider_id !== null + && right.provider_id !== null + && left.provider_id !== right.provider_id + ); +} + +function streamOutputScopeIssue( + outputTokens: MetricObservation, + observations: MetricObservation[] +): string | null { + const scope = outputTokens.accounting_scope; + if (scope?.candidate_scope === 'all_candidates' && scope.candidate_count !== 1) { + return 'Output token usage spans multiple or unknown response candidates.'; + } + const reasoningTokens = measuredObservation(observations, 'reasoning_tokens'); + if ( + numericValue(reasoningTokens) !== null + && (numericValue(reasoningTokens) as number) > 0 + && reasoningTokens?.accounting_scope?.relationship === 'subset_of' + && reasoningTokens.accounting_scope.parent_metric_id === 'output_tokens' + ) { + return 'Output token usage includes hidden reasoning tokens outside the measured stream.'; + } + return null; +} + +function composeGenerationObservations( + input: ComposeRequestMetricObservationsInput, + observations: MetricObservation[] +): MetricObservation[] { + const outputTokens = measuredObservation(observations, 'output_tokens'); + if (!outputTokens) return []; + + const outputCount = numericValue(outputTokens) as number; + const formula = 't_last_output - t_first_output'; + const generationComponents = [outputTokens]; + let generation: MetricObservation; + + if (!input.clientTelemetry.streaming) { + generation = derivedObservation({ + metricId: 'generation_window_ms', + unit: 'milliseconds', + status: 'not_applicable', + reason: 'The operation did not request a streaming response.', + formula, + components: generationComponents + }); + } else if (outputCount < 2) { + generation = derivedObservation({ + metricId: 'generation_window_ms', + unit: 'milliseconds', + status: 'not_applicable', + reason: 'Generation-window metrics require at least two output tokens.', + formula, + components: generationComponents + }); + } else { + const scopeIssue = streamOutputScopeIssue(outputTokens, observations); + const successfulAttempt = [...input.clientTelemetry.attempts] + .reverse() + .find((attempt) => attempt.request_succeeded); + if (scopeIssue) { + generation = derivedObservation({ + metricId: 'generation_window_ms', + unit: 'milliseconds', + status: 'unavailable', + reason: scopeIssue, + formula, + components: generationComponents + }); + } else if ( + successfulAttempt?.first_output_at_ms === null + || successfulAttempt?.first_output_at_ms === undefined + || successfulAttempt.last_output_at_ms === null + ) { + generation = derivedObservation({ + metricId: 'generation_window_ms', + unit: 'milliseconds', + status: 'unavailable', + reason: 'The stream did not provide complete first- and last-output timing.', + formula, + components: generationComponents + }); + } else { + const generationWindow = successfulAttempt.last_output_at_ms + - successfulAttempt.first_output_at_ms; + generation = generationWindow < 0 + ? derivedObservation({ + metricId: 'generation_window_ms', + unit: 'milliseconds', + status: 'unavailable', + reason: 'The measured last-output timestamp preceded first output.', + formula, + components: generationComponents + }) + : derivedObservation({ + metricId: 'generation_window_ms', + unit: 'milliseconds', + value: generationWindow, + formula, + components: generationComponents, + accountingScope: { timing_basis: ['t_first_output', 't_last_output'] } + }); + } + } + + const derived = [generation]; + if (generation.status !== 'measured') { + derived.push( + derivedObservation({ + metricId: 'time_per_output_token_ms', + unit: 'milliseconds_per_token', + status: generation.status, + reason: generation.reason ?? 'Generation window is unavailable.', + formula: 'generation_window_ms / (output_tokens - 1)', + components: [generation, outputTokens] + }), + derivedObservation({ + metricId: 'decode_output_tokens_per_second', + unit: 'tokens_per_second', + status: generation.status, + reason: generation.reason ?? 'Generation window is unavailable.', + formula: '(output_tokens - 1) / (generation_window_ms / 1000)', + components: [generation, outputTokens] + }) + ); + return derived; + } + + const generationWindow = numericValue(generation) as number; + derived.push(derivedObservation({ + metricId: 'time_per_output_token_ms', + unit: 'milliseconds_per_token', + value: generationWindow / (outputCount - 1), + formula: 'generation_window_ms / (output_tokens - 1)', + components: [generation, outputTokens] + })); + derived.push(generationWindow > 0 + ? derivedObservation({ + metricId: 'decode_output_tokens_per_second', + unit: 'tokens_per_second', + value: (outputCount - 1) / (generationWindow / 1000), + formula: '(output_tokens - 1) / (generation_window_ms / 1000)', + components: [generation, outputTokens] + }) + : derivedObservation({ + metricId: 'decode_output_tokens_per_second', + unit: 'tokens_per_second', + status: 'unavailable', + reason: 'Decode throughput requires a positive generation window.', + formula: '(output_tokens - 1) / (generation_window_ms / 1000)', + components: [generation, outputTokens] + })); + return derived; +} + +function composeRateObservation(input: { + metricId: string; + unit: string; + numerator: MetricObservation | null; + denominator: MetricObservation | null; + formula: string; + scale: number; +}): MetricObservation | null { + if (!input.numerator || !input.denominator) return null; + if (!shareProviderScope(input.numerator, input.denominator)) { + return derivedObservation({ + metricId: input.metricId, + unit: input.unit, + status: 'unavailable', + reason: 'Formula inputs do not share compatible provider provenance.', + formula: input.formula, + components: [input.numerator, input.denominator] + }); + } + const denominator = numericValue(input.denominator) as number; + if (denominator <= 0) { + return derivedObservation({ + metricId: input.metricId, + unit: input.unit, + status: 'unavailable', + reason: 'Formula requires a positive duration denominator.', + formula: input.formula, + components: [input.numerator, input.denominator] + }); + } + return derivedObservation({ + metricId: input.metricId, + unit: input.unit, + value: (numericValue(input.numerator) as number) / (denominator / input.scale), + formula: input.formula, + components: [input.numerator, input.denominator] + }); +} + +export function composeRequestMetricObservations( + input: ComposeRequestMetricObservationsInput +): MetricObservation[] { + const clientObservations = normalizeClientMetricObservations(input.clientTelemetry); + const observations = [...clientObservations, ...input.providerObservations]; + const outputTokens = measuredObservation(observations, 'output_tokens'); + const inputTokens = measuredObservation(observations, 'input_tokens'); + const successfulLatency = measuredObservation(observations, 'successful_attempt_latency_ms'); + const serverPrefill = measuredObservation(observations, 'server_prefill_time_ms'); + const serverDecode = measuredObservation(observations, 'server_decode_time_ms'); + const derived: MetricObservation[] = composeGenerationObservations(input, observations); + + const requestRate = composeRateObservation({ + metricId: 'per_request_output_tokens_per_second', + unit: 'tokens_per_second', + numerator: outputTokens, + denominator: successfulLatency, + formula: 'output_tokens / (successful_attempt_latency_ms / 1000)', + scale: 1000 + }); + if (requestRate) derived.push(requestRate); + + const prefillRate = composeRateObservation({ + metricId: 'server_prefill_tokens_per_second', + unit: 'tokens_per_second', + numerator: inputTokens, + denominator: serverPrefill, + formula: 'input_tokens / (server_prefill_time_ms / 1000)', + scale: 1000 + }); + if (prefillRate) derived.push(prefillRate); + + const decodeRate = composeRateObservation({ + metricId: 'server_decode_tokens_per_second', + unit: 'tokens_per_second', + numerator: outputTokens, + denominator: serverDecode, + formula: 'output_tokens / (server_decode_time_ms / 1000)', + scale: 1000 + }); + if (decodeRate) derived.push(decodeRate); + + if (inputTokens && outputTokens) { + const inputCount = numericValue(inputTokens) as number; + derived.push(inputCount > 0 + ? derivedObservation({ + metricId: 'output_input_token_ratio', + unit: 'ratio', + value: (numericValue(outputTokens) as number) / inputCount, + formula: 'output_tokens / input_tokens', + components: [outputTokens, inputTokens] + }) + : derivedObservation({ + metricId: 'output_input_token_ratio', + unit: 'ratio', + status: 'not_applicable', + reason: 'Output/input ratio requires a positive input token count.', + formula: 'output_tokens / input_tokens', + components: [outputTokens, inputTokens] + })); + } + + return [...observations, ...derived]; +} diff --git a/backend/src/services/benchmark-runner.ts b/backend/src/services/benchmark-runner.ts index 7903643..2f4031a 100644 --- a/backend/src/services/benchmark-runner.ts +++ b/backend/src/services/benchmark-runner.ts @@ -22,10 +22,10 @@ import { type ProviderProtocol } from './benchmark-provider-metrics.js'; import { - measuredClientMetric, - normalizeClientMetricObservations, type ClientAttemptTelemetry } from './benchmark-client-metrics.js'; +import { measuredMetricValue } from './benchmark-metric-observations.js'; +import { composeRequestMetricObservations } from './benchmark-request-metrics.js'; import { classifyStreamSemanticEvent, createStreamTimingTracker, @@ -85,6 +85,16 @@ interface DerivedMetric { unit?: string; } +interface NormalizedResponseResult { + response: NormalizedResponse; + provider_observations: MetricObservation[]; +} + +interface NormalizedProviderMetricResult { + metrics: NormalizedProviderMetrics; + observations: MetricObservation[]; +} + interface ExecutableItem { item: Record; itemIndex: number; @@ -751,33 +761,36 @@ function projectedProviderMetric( function normalizedProviderMetrics( context: ProviderMetricContext | null, metadata: Record | null -): NormalizedProviderMetrics { +): NormalizedProviderMetricResult { const observations = context ? normalizeProviderMetricObservations(context, metadata) : []; const legacyTokens = usageTokens(metadata); return { - input_tokens: projectedProviderMetric(observations, 'input_tokens', legacyTokens.input_tokens), - output_tokens: projectedProviderMetric(observations, 'output_tokens', legacyTokens.output_tokens), - total_tokens: projectedProviderMetric(observations, 'total_tokens', legacyTokens.total_tokens), - load_duration_ms: projectedProviderMetric( - observations, - 'model_load_time_ms', - extractLoadDurationMs(metadata) - ), - server_total_time_ms: projectedProviderMetric( - observations, - 'server_total_time_ms', - extractServerTotalTimeMs(metadata) - ), - server_prompt_eval_ms: projectedProviderMetric( - observations, - 'server_prefill_time_ms', - extractServerPromptEvalMs(metadata) - ), - server_eval_ms: projectedProviderMetric( - observations, - 'server_decode_time_ms', - extractServerEvalMs(metadata) - ) + metrics: { + input_tokens: projectedProviderMetric(observations, 'input_tokens', legacyTokens.input_tokens), + output_tokens: projectedProviderMetric(observations, 'output_tokens', legacyTokens.output_tokens), + total_tokens: projectedProviderMetric(observations, 'total_tokens', legacyTokens.total_tokens), + load_duration_ms: projectedProviderMetric( + observations, + 'model_load_time_ms', + extractLoadDurationMs(metadata) + ), + server_total_time_ms: projectedProviderMetric( + observations, + 'server_total_time_ms', + extractServerTotalTimeMs(metadata) + ), + server_prompt_eval_ms: projectedProviderMetric( + observations, + 'server_prefill_time_ms', + extractServerPromptEvalMs(metadata) + ), + server_eval_ms: projectedProviderMetric( + observations, + 'server_decode_time_ms', + extractServerEvalMs(metadata) + ) + }, + observations }; } @@ -1156,7 +1169,7 @@ function normalizeStreamResponse( contentType: string, responseText: string, metricContext: ProviderMetricContext | null -): NormalizedResponse { +): NormalizedResponseResult { const stream = protocol === 'ollama_chat' && !contentType.includes('text/event-stream') ? parseOllamaJsonlStream(responseText) : protocol === 'anthropic_messages' @@ -1164,25 +1177,29 @@ function normalizeStreamResponse( : protocol === 'gemini_generate_content' ? parseGeminiSseStream(responseText) : parseOpenAiSseStream(responseText); + const providerMetrics = normalizedProviderMetrics(metricContext, stream.final_metadata); return { - answer_text: stream.answer_text, - ...normalizedProviderMetrics(metricContext, stream.final_metadata), - tool_calls: stream.tool_calls, - body: stream.final_metadata, - text: null, - stream: { - format: stream.format, - events: stream.events, - done: stream.done, - final_metadata: stream.final_metadata - } + response: { + answer_text: stream.answer_text, + ...providerMetrics.metrics, + tool_calls: stream.tool_calls, + body: stream.final_metadata, + text: null, + stream: { + format: stream.format, + events: stream.events, + done: stream.done, + final_metadata: stream.final_metadata + } + }, + provider_observations: providerMetrics.observations }; } function normalizeAnthropicResponse( record: Record, metricContext: ProviderMetricContext | null -): NormalizedResponse { +): NormalizedResponseResult { const content = Array.isArray(record.content) ? record.content : []; const textParts: string[] = []; const toolCalls: unknown[] = []; @@ -1206,19 +1223,23 @@ function normalizeAnthropicResponse( } } } + const providerMetrics = normalizedProviderMetrics(metricContext, record); return { - answer_text: textParts.join(''), - ...normalizedProviderMetrics(metricContext, record), - tool_calls: toolCalls.length > 0 ? toolCalls : null, - body: record, - text: null + response: { + answer_text: textParts.join(''), + ...providerMetrics.metrics, + tool_calls: toolCalls.length > 0 ? toolCalls : null, + body: record, + text: null + }, + provider_observations: providerMetrics.observations }; } function normalizeGeminiResponse( record: Record, metricContext: ProviderMetricContext | null -): NormalizedResponse { +): NormalizedResponseResult { const candidates = Array.isArray(record.candidates) ? record.candidates : []; const firstCandidate = objectValue(candidates[0]); const content = objectValue(firstCandidate?.content); @@ -1243,12 +1264,16 @@ function normalizeGeminiResponse( }); } } + const providerMetrics = normalizedProviderMetrics(metricContext, record); return { - answer_text: textParts.join(''), - ...normalizedProviderMetrics(metricContext, record), - tool_calls: toolCalls.length > 0 ? toolCalls : null, - body: record, - text: null + response: { + answer_text: textParts.join(''), + ...providerMetrics.metrics, + tool_calls: toolCalls.length > 0 ? toolCalls : null, + body: record, + text: null + }, + provider_observations: providerMetrics.observations }; } @@ -1257,7 +1282,7 @@ function normalizeResponse( body: unknown, text: string | null, metricContext: ProviderMetricContext | null -): NormalizedResponse { +): NormalizedResponseResult { if (body && typeof body === 'object' && !Array.isArray(body)) { const record = body as Record; if (protocol === 'anthropic_messages') { @@ -1275,26 +1300,33 @@ function normalizeResponse( ? record.message as Record : null; const toolCalls = Array.isArray(message?.tool_calls) ? (message.tool_calls as unknown[]) : null; + const providerMetrics = normalizedProviderMetrics(metricContext, record); return { - answer_text: textFromValue(message?.content) ?? textFromValue(ollamaMessage?.content) ?? textFromValue(record.response) ?? '', - ...normalizedProviderMetrics(metricContext, record), - tool_calls: toolCalls, - body, - text: null + response: { + answer_text: textFromValue(message?.content) ?? textFromValue(ollamaMessage?.content) ?? textFromValue(record.response) ?? '', + ...providerMetrics.metrics, + tool_calls: toolCalls, + body, + text: null + }, + provider_observations: providerMetrics.observations }; } return { - answer_text: text ?? '', - input_tokens: null, - output_tokens: null, - total_tokens: null, - load_duration_ms: null, - server_total_time_ms: null, - server_prompt_eval_ms: null, - server_eval_ms: null, - tool_calls: null, - body, - text + response: { + answer_text: text ?? '', + input_tokens: null, + output_tokens: null, + total_tokens: null, + load_duration_ms: null, + server_total_time_ms: null, + server_prompt_eval_ms: null, + server_eval_ms: null, + tool_calls: null, + body, + text + }, + provider_observations: [] }; } @@ -1333,6 +1365,7 @@ async function executeItem( normalizedResponse: Record | null; metrics: Record | null; error: Record | null; + observations: MetricObservation[][]; }> { const operationSpec = objectAt(instantiation, 'operation_spec'); const url = textFromValue(operationSpec?.url); @@ -1348,20 +1381,32 @@ async function executeItem( let firstTokenMs: number | null = null; const attemptErrors: Record[] = []; const clientAttempts: ClientAttemptTelemetry[] = []; + const observationSets: MetricObservation[][] = []; const maxAttempts = policy.retry.max_retries + 1; const pairMeta = executable.pairMemberId ? { pair_member_id: executable.pairMemberId } : {}; - const operationElapsedMs = () => { - const observations = normalizeClientMetricObservations({ - operation_started_at_ms: operationStartedAtMs, - operation_ended_at_ms: performance.now(), - streaming, - attempts: clientAttempts + const composeAttemptObservations = ( + providerObservations: MetricObservation[], + operationEndedAtMs: number + ) => { + const observations = composeRequestMetricObservations({ + clientTelemetry: { + operation_started_at_ms: operationStartedAtMs, + operation_ended_at_ms: operationEndedAtMs, + streaming, + attempts: clientAttempts + }, + providerObservations }); - return roundMilliseconds( - measuredClientMetric(observations, 'operation_elapsed_ms') - ?? (performance.now() - operationStartedAtMs) - ); + observationSets.push(observations); + return observations; }; + const operationElapsedMs = ( + observations: MetricObservation[], + operationEndedAtMs: number + ) => roundMilliseconds( + measuredMetricValue(observations, 'operation_elapsed_ms') + ?? (operationEndedAtMs - operationStartedAtMs) + ); for (let attempt = 1; attempt <= maxAttempts; attempt += 1) { const attemptStartedAtMs = performance.now(); @@ -1419,9 +1464,9 @@ async function executeItem( responseBody = null; } } - let normalized: NormalizedResponse; + let normalizedResult: NormalizedResponseResult; try { - normalized = effectiveStreaming + normalizedResult = effectiveStreaming ? normalizeStreamResponse(operationSpec?.protocol, contentType, responseText, metricContext) : normalizeResponse( operationSpec?.protocol, @@ -1456,7 +1501,9 @@ async function executeItem( response_normalization_succeeded: false, stream_completed: effectiveStreaming ? false : null }); - const elapsedMs = operationElapsedMs(); + const operationEndedAtMs = performance.now(); + const observations = composeAttemptObservations([], operationEndedAtMs); + const elapsedMs = operationElapsedMs(observations, operationEndedAtMs); return { result: { item_index: executable.itemIndex, @@ -1526,9 +1573,11 @@ async function executeItem( server_prompt_eval_ms: null, server_eval_ms: null }, - error: issue + error: issue, + observations: observationSets }; } + const normalized = normalizedResult.response; const errorCode = response.ok ? null : `http_${response.status}`; const streamTiming = streamTimingTracker?.snapshot() ?? emptyStreamTimingTelemetry(); clientAttempts.push({ @@ -1541,7 +1590,12 @@ async function executeItem( response_normalization_succeeded: response.ok, stream_completed: streaming ? Boolean(normalized.stream?.done) : null }); - const elapsedMs = operationElapsedMs(); + const operationEndedAtMs = performance.now(); + const observations = composeAttemptObservations( + normalizedResult.provider_observations, + operationEndedAtMs + ); + const elapsedMs = operationElapsedMs(observations, operationEndedAtMs); const rawResponse = { stage_id: stage.id, ...pairMeta, @@ -1606,7 +1660,8 @@ async function executeItem( ...normalized }, metrics, - error: null + error: null, + observations: observationSets }; } @@ -1647,7 +1702,8 @@ async function executeItem( ...normalized }, metrics, - error: issue + error: issue, + observations: observationSets }; } } catch (error) { @@ -1675,9 +1731,11 @@ async function executeItem( response_normalization_succeeded: null, stream_completed: streaming ? false : null }); + const operationEndedAtMs = performance.now(); + const observations = composeAttemptObservations([], operationEndedAtMs); if (!issue.retryable || attempt >= maxAttempts) { const completedAt = nowIso(); - const elapsedMs = operationElapsedMs(); + const elapsedMs = operationElapsedMs(observations, operationEndedAtMs); return { result: { item_index: executable.itemIndex, @@ -1695,7 +1753,8 @@ async function executeItem( rawResponse: null, normalizedResponse: null, metrics: null, - error: issue + error: issue, + observations: observationSets }; } } finally { @@ -1811,6 +1870,7 @@ async function executePair( normalizedResponses: Record[]; metrics: Record | null; errors: Record[]; + observations: MetricObservation[][]; }> { const pair = stage.pair ?? []; const memberRequestedMetrics = memberMetricRequest(stage, requestedMetrics); @@ -1819,6 +1879,7 @@ async function executePair( const rawResponses: Record[] = []; const normalizedResponses: Record[] = []; const errors: Record[] = []; + const observations: MetricObservation[][] = []; const startedAt = nowIso(); await sleep(stage.pre_iteration_delay_ms ?? 0); @@ -1837,6 +1898,7 @@ async function executePair( if (execution.normalizedResponse) normalizedResponses.push(execution.normalizedResponse); if (execution.metrics) metricsByMember.set(member.id, execution.metrics); if (execution.error) errors.push(execution.error); + observations.push(...execution.observations); members[member.id] = { role: member.role ?? null, result: execution.result, @@ -1877,7 +1939,8 @@ async function executePair( rawResponses, normalizedResponses, metrics: pairMetricRow, - errors + errors, + observations }; } diff --git a/backend/tests/integration/benchmark-runner.test.ts b/backend/tests/integration/benchmark-runner.test.ts index 4124338..aa85f87 100644 --- a/backend/tests/integration/benchmark-runner.test.ts +++ b/backend/tests/integration/benchmark-runner.test.ts @@ -12,6 +12,21 @@ import { BENCHMARK_DATASET_ROOT_ENV } from '../../src/services/benchmark-dataset import { sha256Document } from '../../src/services/benchmark-schemas.js'; const AUTH_HEADERS = { 'x-api-token': 'test-token' }; +const CANONICAL_DERIVED_METRIC_IDS = [ + 'generation_window_ms', + 'per_request_output_tokens_per_second', + 'time_per_output_token_ms', + 'decode_output_tokens_per_second', + 'server_prefill_tokens_per_second', + 'server_decode_tokens_per_second', + 'output_input_token_ratio' +]; + +function expectNoCanonicalDerivedMetrics(value: Record): void { + for (const metricId of CANONICAL_DERIVED_METRIC_IDS) { + expect(value).not.toHaveProperty(metricId); + } +} function streamFixture(name: string): string { return fs.readFileSync(new URL(`../fixtures/streams/${name}`, import.meta.url), 'utf8'); @@ -621,6 +636,7 @@ describe('benchmark runner API', () => { expect(result.document.normalized_responses[0]).not.toHaveProperty('provider_version'); expect(result.document.metric_results[0].total_tokens).toBe(7); expect(result.document.metric_version).toBe('metrics-v1'); + expectNoCanonicalDerivedMetrics(result.document.metric_results[0]); expect(result.document.aggregated_metrics.elapsed_ms.valid_sample_count).toBe(1); expect(mockServer.requests).toHaveLength(1); expect((mockServer.requests[0] as Record).model).toBe('mock-chat'); @@ -955,6 +971,7 @@ describe('benchmark runner API', () => { expect(result.document.normalized_responses[0]).not.toHaveProperty('metric_observations'); expect(result.document.metric_results[0]).not.toHaveProperty('time_to_first_tool_call_ms'); expect(result.document.metric_results[0]).not.toHaveProperty('time_to_tool_calls_ready_ms'); + expectNoCanonicalDerivedMetrics(result.document.metric_results[0]); await app.close(); }); @@ -1100,6 +1117,7 @@ describe('benchmark runner API', () => { expect(result.document).not.toHaveProperty('metric_observations'); expect(result.document.normalized_responses[0]).not.toHaveProperty('metric_observations'); expect(result.document.metric_results[0]).not.toHaveProperty('time_to_first_output_ms'); + expectNoCanonicalDerivedMetrics(result.document.metric_results[0]); expect((mockServer.requests[0] as Record).stream).toBe(true); await app.close(); }); diff --git a/backend/tests/unit/benchmark-provider-metrics.test.ts b/backend/tests/unit/benchmark-provider-metrics.test.ts index e34e1fd..9b60ef7 100644 --- a/backend/tests/unit/benchmark-provider-metrics.test.ts +++ b/backend/tests/unit/benchmark-provider-metrics.test.ts @@ -166,6 +166,7 @@ describe('normalizeProviderMetricObservations', () => { provider_name: 'Google Gemini' }), { + candidates: [{ content: { parts: [{ text: 'Answer' }] } }], usageMetadata: { promptTokenCount: 30, cachedContentTokenCount: 10, @@ -187,6 +188,7 @@ describe('normalizeProviderMetricObservations', () => { accounting_scope: { mapping_classification: 'qualified', candidate_scope: 'all_candidates', + candidate_count: 1, comparison_requires_equivalent_candidate_scope: true } }); diff --git a/backend/tests/unit/benchmark-request-metrics.test.ts b/backend/tests/unit/benchmark-request-metrics.test.ts new file mode 100644 index 0000000..7b23a09 --- /dev/null +++ b/backend/tests/unit/benchmark-request-metrics.test.ts @@ -0,0 +1,333 @@ +import { describe, expect, it } from 'vitest'; + +import type { + ClientAttemptTelemetry, + ClientOperationTelemetry +} from '../../src/services/benchmark-client-metrics.js'; +import type { MetricObservation } from '../../src/services/benchmark-metric-observations.js'; +import { composeRequestMetricObservations } from '../../src/services/benchmark-request-metrics.js'; + +function attempt(overrides: Partial = {}): ClientAttemptTelemetry { + return { + started_at_ms: 0, + ended_at_ms: 500, + first_chunk_at_ms: 50, + first_output_at_ms: 100, + first_tool_call_at_ms: null, + tool_calls_ready_at_ms: null, + last_output_at_ms: 300, + tool_call_started: false, + tool_call_error: null, + request_succeeded: true, + timed_out: false, + response_normalization_succeeded: true, + stream_completed: true, + ...overrides + }; +} + +function telemetry(overrides: Partial = {}): ClientOperationTelemetry { + return { + operation_started_at_ms: 0, + operation_ended_at_ms: 500, + streaming: true, + attempts: [attempt()], + ...overrides + }; +} + +function providerObservation( + metricId: string, + value: number, + unit: string, + overrides: Partial = {} +): MetricObservation { + return { + metric_id: metricId, + value, + unit, + status: 'measured', + reason: null, + source: 'server_reported', + metric_version: 'metrics-v2', + provider_id: 'server-1', + provider_protocol: 'ollama_chat', + provider_version: '0.12.0', + native_field: metricId, + native_value: value, + native_unit: unit, + normalization: 'identity', + accounting_scope: null, + ...overrides + }; +} + +function metric(observations: MetricObservation[], metricId: string): MetricObservation { + const result = observations.find((candidate) => candidate.metric_id === metricId); + expect(result, `Missing observation: ${metricId}`).toBeDefined(); + return result as MetricObservation; +} + +function fullProviderObservations(): MetricObservation[] { + return [ + providerObservation('input_tokens', 100, 'tokens'), + providerObservation('output_tokens', 5, 'tokens'), + providerObservation('server_prefill_time_ms', 50, 'milliseconds'), + providerObservation('server_decode_time_ms', 200, 'milliseconds') + ]; +} + +describe('composeRequestMetricObservations', () => { + it('composes client and provider inputs and emits every registered request formula', () => { + const observations = composeRequestMetricObservations({ + clientTelemetry: telemetry(), + providerObservations: fullProviderObservations() + }); + + expect(metric(observations, 'generation_window_ms')).toMatchObject({ + value: 200, + unit: 'milliseconds', + source: 'derived', + normalization: 't_last_output - t_first_output', + native_field: null, + native_value: null, + metric_version: 'metrics-v2', + provider_id: 'server-1', + provider_protocol: 'ollama_chat' + }); + for (const [metricId, value, normalization] of [ + [ + 'per_request_output_tokens_per_second', + 10, + 'output_tokens / (successful_attempt_latency_ms / 1000)' + ], + ['time_per_output_token_ms', 50, 'generation_window_ms / (output_tokens - 1)'], + [ + 'decode_output_tokens_per_second', + 20, + '(output_tokens - 1) / (generation_window_ms / 1000)' + ], + [ + 'server_prefill_tokens_per_second', + 2000, + 'input_tokens / (server_prefill_time_ms / 1000)' + ], + [ + 'server_decode_tokens_per_second', + 25, + 'output_tokens / (server_decode_time_ms / 1000)' + ], + ['output_input_token_ratio', 0.05, 'output_tokens / input_tokens'] + ] as const) { + expect(metric(observations, metricId)).toMatchObject({ value, normalization }); + } + expect(metric(observations, 'decode_output_tokens_per_second').accounting_scope).toMatchObject({ + component_metric_ids: ['generation_window_ms', 'output_tokens'] + }); + }); + + it('uses the successful retry attempt for request latency and generation timing', () => { + const observations = composeRequestMetricObservations({ + clientTelemetry: telemetry({ + operation_ended_at_ms: 700, + attempts: [ + attempt({ + started_at_ms: 0, + ended_at_ms: 100, + first_chunk_at_ms: null, + first_output_at_ms: null, + last_output_at_ms: null, + request_succeeded: false, + response_normalization_succeeded: null, + stream_completed: false + }), + attempt({ + started_at_ms: 200, + ended_at_ms: 700, + first_output_at_ms: 300, + last_output_at_ms: 500 + }) + ] + }), + providerObservations: [providerObservation('output_tokens', 5, 'tokens')] + }); + + expect(metric(observations, 'attempt_count').value).toBe(2); + expect(metric(observations, 'per_request_output_tokens_per_second').value).toBe(10); + expect(metric(observations, 'generation_window_ms').value).toBe(200); + }); + + it('marks stream-window formulas not applicable for non-streaming requests', () => { + const observations = composeRequestMetricObservations({ + clientTelemetry: telemetry({ + streaming: false, + attempts: [attempt({ + first_chunk_at_ms: null, + first_output_at_ms: null, + last_output_at_ms: null, + stream_completed: null + })] + }), + providerObservations: [providerObservation('output_tokens', 5, 'tokens')] + }); + + for (const metricId of [ + 'generation_window_ms', + 'time_per_output_token_ms', + 'decode_output_tokens_per_second' + ]) { + expect(metric(observations, metricId).status).toBe('not_applicable'); + } + expect(metric(observations, 'per_request_output_tokens_per_second').value).toBe(10); + }); + + it('handles zero token counts and non-positive duration denominators', () => { + const observations = composeRequestMetricObservations({ + clientTelemetry: telemetry({ + attempts: [attempt({ started_at_ms: 500, ended_at_ms: 500 })] + }), + providerObservations: [ + providerObservation('input_tokens', 0, 'tokens'), + providerObservation('output_tokens', 0, 'tokens'), + providerObservation('server_prefill_time_ms', 0, 'milliseconds'), + providerObservation('server_decode_time_ms', 0, 'milliseconds') + ] + }); + + expect(metric(observations, 'generation_window_ms').status).toBe('not_applicable'); + expect(metric(observations, 'per_request_output_tokens_per_second').status).toBe('unavailable'); + expect(metric(observations, 'server_prefill_tokens_per_second').status).toBe('unavailable'); + expect(metric(observations, 'server_decode_tokens_per_second').status).toBe('unavailable'); + expect(metric(observations, 'output_input_token_ratio').status).toBe('not_applicable'); + }); + + it('omits formulas whose optional provider inputs are absent', () => { + const observations = composeRequestMetricObservations({ + clientTelemetry: telemetry(), + providerObservations: [] + }); + + expect(observations.some(({ metric_id }) => [ + 'generation_window_ms', + 'per_request_output_tokens_per_second', + 'server_prefill_tokens_per_second', + 'server_decode_tokens_per_second', + 'output_input_token_ratio' + ].includes(metric_id))).toBe(false); + }); + + it('does not derive from unavailable provider observations', () => { + const observations = composeRequestMetricObservations({ + clientTelemetry: telemetry(), + providerObservations: [providerObservation('output_tokens', 0, 'tokens', { + value: null, + status: 'unavailable', + reason: 'Native output token count was non-finite.', + native_value: null + })] + }); + + expect(observations.some(({ metric_id }) => metric_id === 'generation_window_ms')).toBe(false); + expect(observations.some( + ({ metric_id }) => metric_id === 'per_request_output_tokens_per_second' + )).toBe(false); + }); + + it('rejects multi-candidate Gemini output scope for stream-window metrics', () => { + const outputTokens = providerObservation('output_tokens', 5, 'tokens', { + provider_protocol: 'gemini_generate_content', + accounting_scope: { + candidate_scope: 'all_candidates', + candidate_count: 2 + } + }); + const observations = composeRequestMetricObservations({ + clientTelemetry: telemetry(), + providerObservations: [outputTokens] + }); + + expect(metric(observations, 'generation_window_ms')).toMatchObject({ + status: 'unavailable', + reason: 'Output token usage spans multiple or unknown response candidates.' + }); + }); + + it('accepts Gemini output tokens when usage covers exactly one candidate', () => { + const observations = composeRequestMetricObservations({ + clientTelemetry: telemetry(), + providerObservations: [providerObservation('output_tokens', 5, 'tokens', { + provider_protocol: 'gemini_generate_content', + accounting_scope: { + candidate_scope: 'all_candidates', + candidate_count: 1 + } + })] + }); + + expect(metric(observations, 'generation_window_ms').value).toBe(200); + }); + + it('rejects output scope containing hidden reasoning tokens', () => { + const observations = composeRequestMetricObservations({ + clientTelemetry: telemetry(), + providerObservations: [ + providerObservation('output_tokens', 5, 'tokens', { + provider_protocol: 'anthropic_messages' + }), + providerObservation('reasoning_tokens', 2, 'tokens', { + provider_protocol: 'anthropic_messages', + accounting_scope: { + relationship: 'subset_of', + parent_metric_id: 'output_tokens' + } + }) + ] + }); + + expect(metric(observations, 'generation_window_ms')).toMatchObject({ + status: 'unavailable', + reason: 'Output token usage includes hidden reasoning tokens outside the measured stream.' + }); + }); + + it('keeps a zero generation window but rejects its decode-rate denominator', () => { + const observations = composeRequestMetricObservations({ + clientTelemetry: telemetry({ + attempts: [attempt({ first_output_at_ms: 100, last_output_at_ms: 100 })] + }), + providerObservations: [providerObservation('output_tokens', 2, 'tokens')] + }); + + expect(metric(observations, 'generation_window_ms').value).toBe(0); + expect(metric(observations, 'time_per_output_token_ms').value).toBe(0); + expect(metric(observations, 'decode_output_tokens_per_second').status).toBe('unavailable'); + }); + + it('marks a reversed semantic timing window unavailable', () => { + const observations = composeRequestMetricObservations({ + clientTelemetry: telemetry({ + attempts: [attempt({ first_output_at_ms: 300, last_output_at_ms: 100 })] + }), + providerObservations: [providerObservation('output_tokens', 5, 'tokens')] + }); + + expect(metric(observations, 'generation_window_ms').status).toBe('unavailable'); + }); + + it('rejects server-rate inputs from incompatible provider scopes', () => { + const observations = composeRequestMetricObservations({ + clientTelemetry: telemetry(), + providerObservations: [ + providerObservation('input_tokens', 100, 'tokens'), + providerObservation('server_prefill_time_ms', 50, 'milliseconds', { + provider_id: 'server-2' + }) + ] + }); + + expect(metric(observations, 'server_prefill_tokens_per_second')).toMatchObject({ + status: 'unavailable', + reason: 'Formula inputs do not share compatible provider provenance.' + }); + }); +});