From c701222299a89d23d900a6c171ff3b7c0474d16a Mon Sep 17 00:00:00 2001 From: Rodolphe Lefebvre Date: Mon, 27 Jul 2026 21:19:58 +0200 Subject: [PATCH] feat: add canonical observation-aware aggregation --- CHANGELOG.md | 1 + .../benchmark-observation-aggregation.ts | 532 ++++++++++++++++ backend/src/services/benchmark-runner.ts | 110 +++- .../integration/benchmark-runner.test.ts | 11 +- .../benchmark-observation-aggregation.test.ts | 594 ++++++++++++++++++ 5 files changed, 1238 insertions(+), 10 deletions(-) create mode 100644 backend/src/services/benchmark-observation-aggregation.ts create mode 100644 backend/tests/unit/benchmark-observation-aggregation.test.ts diff --git a/CHANGELOG.md b/CHANGELOG.md index 15fc626..d282cc9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -32,6 +32,7 @@ The format is based on Keep a Changelog and this project follows Semantic Versio - **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`. +- **Canonical observation-aware aggregation** — benchmark execution now transiently aggregates canonical request observations with explicit coverage, applicability, execution-error, provenance, and paired-member accounting while persisted results remain on `metrics-v1`. ### Fixed diff --git a/backend/src/services/benchmark-observation-aggregation.ts b/backend/src/services/benchmark-observation-aggregation.ts new file mode 100644 index 0000000..c112dff --- /dev/null +++ b/backend/src/services/benchmark-observation-aggregation.ts @@ -0,0 +1,532 @@ +import type { + MetricObservation, + MetricObservationSource, + ProviderProtocol +} from './benchmark-metric-observations.js'; + +export type CanonicalMetricValueType = 'number' | 'boolean'; + +export type CanonicalNumericAggregation = + | 'mean' + | 'median' + | 'min' + | 'max' + | 'sum' + | 'stddev' + | 'variance' + | 'p50' + | 'p90' + | 'p95' + | 'p99'; + +export interface MetricObservationSample { + stage_id: string; + item_index: number; + iteration: number; + pair_member_id: string | null; + streaming: boolean; + expected: true; + attempted: boolean; + completed: boolean; + attempt_observations: MetricObservation[][]; + terminal_observations: MetricObservation[] | null; +} + +export interface MetricObservationStagePlan { + stage_id: string; + stage_type: 'dataset_loop' | 'single_request' | 'paired_request_loop'; + item_count: number; + iterations_per_item: number; + pair_member_ids: string[]; + record_metrics: boolean; +} + +export interface CanonicalMetricIntent { + stage_id: string; + pair_member_id: string | null; + metric_id: string; + value_type: CanonicalMetricValueType; + unit: string; +} + +export interface CanonicalMetricIntentContext { + stage_id: string; + pair_member_id: string | null; + requested_metrics: string[]; + derived_metric_references?: string[]; + streaming: boolean; +} + +export interface CanonicalProvenanceSignature { + source: MetricObservationSource; + provider_id: string | null; + provider_protocol: ProviderProtocol | null; + provider_version: string | null; + native_field: string | null; + accounting_scope: Record | null; +} + +export interface CanonicalMetricAggregate { + stage_id: string; + pair_member_id: string | null; + metric_id: string; + metric_version: 'metrics-v2'; + value_type: CanonicalMetricValueType; + unit: string; + expected_sample_count: number; + attempted_sample_count: number; + completed_sample_count: number; + valid_sample_count: number; + passed_sample_count: number; + unavailable_sample_count: number; + not_applicable_sample_count: number; + execution_error_sample_count: number; + observed_pass_rate: number | null; + coverage_rate: number; + end_to_end_pass_rate: number | null; + statistics: Partial>; + aggregation_eligible: boolean; + provenance_signatures: CanonicalProvenanceSignature[]; + warnings: string[]; +} + +interface MetricDefinition { + value_type: CanonicalMetricValueType; + unit: string; +} + +const METRIC_DEFINITIONS: Record = { + request_success: { value_type: 'boolean', unit: 'boolean' }, + stream_completed: { value_type: 'boolean', unit: 'boolean' }, + response_normalization_success: { value_type: 'boolean', unit: 'boolean' }, + attempt_count: { value_type: 'number', unit: 'attempts' }, + timeout_occurred: { value_type: 'boolean', unit: 'boolean' }, + retry_overhead_ms: { value_type: 'number', unit: 'milliseconds' }, + operation_elapsed_ms: { value_type: 'number', unit: 'milliseconds' }, + successful_attempt_latency_ms: { value_type: 'number', unit: 'milliseconds' }, + time_to_first_chunk_ms: { value_type: 'number', unit: 'milliseconds' }, + time_to_first_output_ms: { value_type: 'number', unit: 'milliseconds' }, + time_to_first_tool_call_ms: { value_type: 'number', unit: 'milliseconds' }, + time_to_tool_calls_ready_ms: { value_type: 'number', unit: 'milliseconds' }, + generation_window_ms: { value_type: 'number', unit: 'milliseconds' }, + model_load_time_ms: { value_type: 'number', unit: 'milliseconds' }, + server_total_time_ms: { value_type: 'number', unit: 'milliseconds' }, + server_prefill_time_ms: { value_type: 'number', unit: 'milliseconds' }, + server_decode_time_ms: { value_type: 'number', unit: 'milliseconds' }, + input_tokens: { value_type: 'number', unit: 'tokens' }, + output_tokens: { value_type: 'number', unit: 'tokens' }, + total_tokens: { value_type: 'number', unit: 'tokens' }, + per_request_output_tokens_per_second: { value_type: 'number', unit: 'tokens_per_second' }, + time_per_output_token_ms: { value_type: 'number', unit: 'milliseconds_per_token' }, + decode_output_tokens_per_second: { value_type: 'number', unit: 'tokens_per_second' }, + server_prefill_tokens_per_second: { value_type: 'number', unit: 'tokens_per_second' }, + server_decode_tokens_per_second: { value_type: 'number', unit: 'tokens_per_second' }, + output_input_token_ratio: { value_type: 'number', unit: 'ratio' } +}; + +const LEGACY_PERFORMANCE_ALIASES: Record = { + input_tokens: 'input_tokens', + output_tokens: 'output_tokens', + total_tokens: 'total_tokens', + elapsed_ms: 'operation_elapsed_ms', + first_token_ms: 'time_to_first_output_ms', + tokens_per_second: 'per_request_output_tokens_per_second', + decode_tokens_per_second: 'decode_output_tokens_per_second', + prefill_tokens_per_second: 'server_prefill_tokens_per_second', + output_input_token_ratio: 'output_input_token_ratio', + load_duration_ms: 'model_load_time_ms', + server_total_time_ms: 'server_total_time_ms', + server_prompt_eval_ms: 'server_prefill_time_ms', + server_eval_ms: 'server_decode_time_ms' +}; + +const EXECUTION_DEFAULTS = [ + 'request_success', + 'response_normalization_success', + 'attempt_count', + 'timeout_occurred', + 'retry_overhead_ms', + 'operation_elapsed_ms', + 'successful_attempt_latency_ms' +]; + +const STREAMING_DEFAULTS = [ + 'stream_completed', + 'time_to_first_chunk_ms', + 'time_to_first_output_ms', + 'time_per_output_token_ms', + 'decode_output_tokens_per_second' +]; + +const LEGACY_TOOL_METRICS = new Set([ + 'tool_call_count', + 'tool_selected_correctly', + 'tool_arguments_valid', + 'tool_call_assertion_pass', + 'missing_tool_call', + 'hallucinated_tool_call' +]); + +export function planMetricObservationSamples(input: { + stages: MetricObservationStagePlan[]; + streaming: boolean; +}): MetricObservationSample[] { + return input.stages.flatMap((stage) => { + if (!stage.record_metrics) return []; + const itemCount = stage.stage_type === 'single_request' + ? Math.min(stage.item_count, 1) + : stage.item_count; + const pairMembers = stage.stage_type === 'paired_request_loop' + ? stage.pair_member_ids + : [null]; + const samples: MetricObservationSample[] = []; + for (let itemIndex = 0; itemIndex < itemCount; itemIndex += 1) { + for (let iteration = 0; iteration < stage.iterations_per_item; iteration += 1) { + for (const pairMemberId of pairMembers) { + samples.push({ + stage_id: stage.stage_id, + item_index: itemIndex, + iteration, + pair_member_id: pairMemberId, + streaming: input.streaming, + expected: true, + attempted: false, + completed: false, + attempt_observations: [], + terminal_observations: null + }); + } + } + } + return samples; + }); +} + +function requestedMetricForContext( + reference: string, + pairMemberId: string | null +): string | null { + const parts = reference.split('.'); + if (pairMemberId === null) { + return parts.length === 1 ? parts[0] : null; + } + return parts.length === 3 && parts[0] === 'pair' && parts[1] === pairMemberId + ? parts[2] + : null; +} + +function intent( + context: CanonicalMetricIntentContext, + metricId: string +): CanonicalMetricIntent { + const definition = METRIC_DEFINITIONS[metricId]; + if (!definition) { + throw new Error(`Canonical metric definition is missing: ${metricId}`); + } + return { + stage_id: context.stage_id, + pair_member_id: context.pair_member_id, + metric_id: metricId, + ...definition + }; +} + +export function resolveCanonicalMetricIntents( + context: CanonicalMetricIntentContext +): CanonicalMetricIntent[] { + const metricIds = new Set(EXECUTION_DEFAULTS); + if (context.streaming) { + for (const metricId of STREAMING_DEFAULTS) metricIds.add(metricId); + } + + const references = [ + ...context.requested_metrics, + ...(context.derived_metric_references ?? []) + ]; + let requestsToolMetrics = false; + for (const reference of references) { + const requestedMetric = requestedMetricForContext(reference, context.pair_member_id); + if (!requestedMetric) continue; + if (LEGACY_TOOL_METRICS.has(requestedMetric)) { + requestsToolMetrics = true; + continue; + } + const canonicalMetric = LEGACY_PERFORMANCE_ALIASES[requestedMetric]; + if (canonicalMetric) metricIds.add(canonicalMetric); + } + + if (context.streaming && requestsToolMetrics) { + metricIds.add('time_to_first_tool_call_ms'); + metricIds.add('time_to_tool_calls_ready_ms'); + } + + return [...metricIds].map((metricId) => intent(context, metricId)); +} + +function stableValue(value: unknown): unknown { + if (Array.isArray(value)) return value.map(stableValue); + if (!value || typeof value !== 'object') return value; + return Object.fromEntries( + Object.entries(value as Record) + .sort(([left], [right]) => left.localeCompare(right)) + .map(([key, entry]) => [key, stableValue(entry)]) + ); +} + +function provenanceSignature(observation: MetricObservation): CanonicalProvenanceSignature { + return { + source: observation.source, + provider_id: observation.provider_id, + provider_protocol: observation.provider_protocol, + provider_version: observation.provider_version, + native_field: observation.native_field, + accounting_scope: observation.accounting_scope + ? stableValue(observation.accounting_scope) as Record + : null + }; +} + +function uniqueProvenanceSignatures( + observations: MetricObservation[] +): CanonicalProvenanceSignature[] { + const signatures = new Map(); + for (const observation of observations) { + const signature = provenanceSignature(observation); + signatures.set(JSON.stringify(signature), signature); + } + return [...signatures.values()]; +} + +function computeStat( + aggregation: CanonicalNumericAggregation, + values: number[] +): number | null { + if (values.length === 0) return null; + const sorted = [...values].sort((left, right) => left - right); + const mean = values.reduce((sum, value) => sum + value, 0) / values.length; + switch (aggregation) { + case 'mean': + return mean; + case 'median': + case 'p50': + return interpolatedPercentile(sorted, 50); + case 'min': + return sorted[0]; + case 'max': + return sorted[sorted.length - 1]; + case 'sum': + return values.reduce((sum, value) => sum + value, 0); + case 'stddev': + return Math.sqrt( + values.reduce((sum, value) => sum + (value - mean) ** 2, 0) / values.length + ); + case 'variance': + return values.reduce((sum, value) => sum + (value - mean) ** 2, 0) / values.length; + case 'p90': + return interpolatedPercentile(sorted, 90); + case 'p95': + return interpolatedPercentile(sorted, 95); + case 'p99': + return interpolatedPercentile(sorted, 99); + } +} + +function interpolatedPercentile(sorted: number[], percentile: number): number { + if (sorted.length === 1) return sorted[0]; + const index = (percentile / 100) * (sorted.length - 1); + const lower = Math.floor(index); + const upper = Math.ceil(index); + if (lower === upper) return sorted[lower]; + return sorted[lower] + (sorted[upper] - sorted[lower]) * (index - lower); +} + +function percentileWarnings( + aggregations: CanonicalNumericAggregation[], + validSampleCount: number +): string[] { + const warnings: string[] = []; + const thresholds: Partial> = { + p90: 10, + p95: 20, + p99: 100 + }; + for (const aggregation of aggregations) { + const threshold = thresholds[aggregation]; + if (threshold !== undefined && validSampleCount < threshold) { + warnings.push(`insufficient_valid_samples_for_${aggregation}`); + } + } + return warnings; +} + +function observationForIntent( + sample: MetricObservationSample, + metricId: string +): MetricObservation | null { + const matches = (sample.terminal_observations ?? []).filter( + (observation) => observation.metric_id === metricId + ); + if (matches.length > 1) { + throw new Error( + `Duplicate metric observation for ${sample.stage_id}/${sample.item_index}/${sample.iteration}/${String(sample.pair_member_id)}: ${metricId}` + ); + } + return matches[0] ?? null; +} + +function validateObservation( + observation: MetricObservation, + intentValueType: CanonicalMetricValueType, + unit: string +): void { + if (observation.metric_version !== 'metrics-v2') { + throw new Error(`Unexpected metric version for ${observation.metric_id}`); + } + if (observation.unit !== unit) { + throw new Error(`Conflicting unit for ${observation.metric_id}`); + } + if (observation.status !== 'measured') return; + if (intentValueType === 'boolean' && typeof observation.value !== 'boolean') { + throw new Error(`Conflicting value type for ${observation.metric_id}`); + } + if ( + intentValueType === 'number' + && (typeof observation.value !== 'number' || !Number.isFinite(observation.value)) + ) { + throw new Error(`Conflicting value type for ${observation.metric_id}`); + } +} + +export function aggregateMetricObservations(input: { + samples: MetricObservationSample[]; + intents: CanonicalMetricIntent[]; + requestedAggregations: string[]; +}): CanonicalMetricAggregate[] { + const requestedAggregations = [...new Set(input.requestedAggregations)] + .filter((aggregation): aggregation is CanonicalNumericAggregation => ( + aggregation !== 'count' + && [ + 'mean', + 'median', + 'min', + 'max', + 'sum', + 'stddev', + 'variance', + 'p50', + 'p90', + 'p95', + 'p99' + ].includes(aggregation) + )); + const uniqueIntents = new Map(); + for (const candidate of input.intents) { + const key = [ + candidate.stage_id, + candidate.pair_member_id ?? '', + candidate.metric_id + ].join('\u0000'); + const existing = uniqueIntents.get(key); + if ( + existing + && (existing.unit !== candidate.unit || existing.value_type !== candidate.value_type) + ) { + throw new Error(`Conflicting intent definition for ${candidate.metric_id}`); + } + uniqueIntents.set(key, candidate); + } + + return [...uniqueIntents.values()].map((metricIntent) => { + const samples = input.samples.filter((sample) => ( + sample.stage_id === metricIntent.stage_id + && sample.pair_member_id === metricIntent.pair_member_id + )); + const measuredObservations: MetricObservation[] = []; + const numericValues: number[] = []; + let attemptedSampleCount = 0; + let completedSampleCount = 0; + let validSampleCount = 0; + let passedSampleCount = 0; + let unavailableSampleCount = 0; + let notApplicableSampleCount = 0; + let executionErrorSampleCount = 0; + + for (const sample of samples) { + if (!sample.attempted) continue; + attemptedSampleCount += 1; + if (sample.completed) completedSampleCount += 1; + const observation = observationForIntent(sample, metricIntent.metric_id); + if (!observation) { + if (sample.completed) { + unavailableSampleCount += 1; + } else { + executionErrorSampleCount += 1; + } + continue; + } + validateObservation(observation, metricIntent.value_type, metricIntent.unit); + switch (observation.status) { + case 'measured': + validSampleCount += 1; + measuredObservations.push(observation); + if (metricIntent.value_type === 'boolean') { + if (observation.value === true) passedSampleCount += 1; + } else { + numericValues.push(observation.value as number); + } + break; + case 'unavailable': + unavailableSampleCount += 1; + break; + case 'not_applicable': + notApplicableSampleCount += 1; + break; + case 'execution_error': + executionErrorSampleCount += 1; + break; + } + } + + const provenanceSignatures = uniqueProvenanceSignatures(measuredObservations); + const aggregationEligible = provenanceSignatures.length <= 1; + const warnings = percentileWarnings(requestedAggregations, validSampleCount); + if (!aggregationEligible) { + warnings.push('incompatible_provenance_or_accounting_scope'); + } + const statistics: Partial> = {}; + if (metricIntent.value_type === 'number' && aggregationEligible) { + for (const aggregation of requestedAggregations) { + const value = computeStat(aggregation, numericValues); + if (value !== null) statistics[aggregation] = value; + } + } + const expectedSampleCount = samples.length; + return { + stage_id: metricIntent.stage_id, + pair_member_id: metricIntent.pair_member_id, + metric_id: metricIntent.metric_id, + metric_version: 'metrics-v2', + value_type: metricIntent.value_type, + unit: metricIntent.unit, + expected_sample_count: expectedSampleCount, + attempted_sample_count: attemptedSampleCount, + completed_sample_count: completedSampleCount, + valid_sample_count: validSampleCount, + passed_sample_count: passedSampleCount, + unavailable_sample_count: unavailableSampleCount, + not_applicable_sample_count: notApplicableSampleCount, + execution_error_sample_count: executionErrorSampleCount, + observed_pass_rate: metricIntent.value_type === 'boolean' + ? (validSampleCount > 0 ? passedSampleCount / validSampleCount : null) + : null, + coverage_rate: expectedSampleCount > 0 ? validSampleCount / expectedSampleCount : 0, + end_to_end_pass_rate: metricIntent.value_type === 'boolean' + ? (expectedSampleCount > 0 ? passedSampleCount / expectedSampleCount : 0) + : null, + statistics, + aggregation_eligible: aggregationEligible, + provenance_signatures: provenanceSignatures, + warnings + }; + }); +} diff --git a/backend/src/services/benchmark-runner.ts b/backend/src/services/benchmark-runner.ts index 2f4031a..064ff29 100644 --- a/backend/src/services/benchmark-runner.ts +++ b/backend/src/services/benchmark-runner.ts @@ -25,6 +25,13 @@ import { type ClientAttemptTelemetry } from './benchmark-client-metrics.js'; import { measuredMetricValue } from './benchmark-metric-observations.js'; +import { + aggregateMetricObservations, + planMetricObservationSamples, + resolveCanonicalMetricIntents, + type CanonicalMetricIntent, + type MetricObservationSample +} from './benchmark-observation-aggregation.js'; import { composeRequestMetricObservations } from './benchmark-request-metrics.js'; import { classifyStreamSemanticEvent, @@ -1365,7 +1372,7 @@ async function executeItem( normalizedResponse: Record | null; metrics: Record | null; error: Record | null; - observations: MetricObservation[][]; + observationSample: MetricObservationSample; }> { const operationSpec = objectAt(instantiation, 'operation_spec'); const url = textFromValue(operationSpec?.url); @@ -1407,6 +1414,18 @@ async function executeItem( measuredMetricValue(observations, 'operation_elapsed_ms') ?? (operationEndedAtMs - operationStartedAtMs) ); + const observationSample = (completed: boolean): MetricObservationSample => ({ + stage_id: stage.id, + item_index: executable.itemIndex, + iteration: executable.iteration, + pair_member_id: executable.pairMemberId ?? null, + streaming, + expected: true, + attempted: true, + completed, + attempt_observations: observationSets, + terminal_observations: observationSets[observationSets.length - 1] ?? null + }); for (let attempt = 1; attempt <= maxAttempts; attempt += 1) { const attemptStartedAtMs = performance.now(); @@ -1574,7 +1593,7 @@ async function executeItem( server_eval_ms: null }, error: issue, - observations: observationSets + observationSample: observationSample(false) }; } const normalized = normalizedResult.response; @@ -1661,7 +1680,7 @@ async function executeItem( }, metrics, error: null, - observations: observationSets + observationSample: observationSample(true) }; } @@ -1703,7 +1722,7 @@ async function executeItem( }, metrics, error: issue, - observations: observationSets + observationSample: observationSample(false) }; } } catch (error) { @@ -1754,7 +1773,7 @@ async function executeItem( normalizedResponse: null, metrics: null, error: issue, - observations: observationSets + observationSample: observationSample(false) }; } } finally { @@ -1778,6 +1797,42 @@ function executableItemsForStage(stage: BenchmarkStage, items: Record { + if (stage.record_metrics === false) return []; + const pairMembers = stage.type === 'paired_request_loop' + ? (stage.pair ?? []).map((member) => member.id) + : [null]; + const derivedMetricReferences = (stage.derived_metrics ?? []) + .flatMap((metric) => [metric.left, metric.right]); + return pairMembers.flatMap((pairMemberId) => resolveCanonicalMetricIntents({ + stage_id: stage.id, + pair_member_id: pairMemberId, + requested_metrics: requestedMetrics, + derived_metric_references: derivedMetricReferences, + streaming + })); + }); +} + function skippedResult(executable: ExecutableItem, reason: string): Record { return { item_index: executable.itemIndex, @@ -1870,7 +1925,7 @@ async function executePair( normalizedResponses: Record[]; metrics: Record | null; errors: Record[]; - observations: MetricObservation[][]; + observationSamples: MetricObservationSample[]; }> { const pair = stage.pair ?? []; const memberRequestedMetrics = memberMetricRequest(stage, requestedMetrics); @@ -1879,7 +1934,7 @@ async function executePair( const rawResponses: Record[] = []; const normalizedResponses: Record[] = []; const errors: Record[] = []; - const observations: MetricObservation[][] = []; + const observationSamples: MetricObservationSample[] = []; const startedAt = nowIso(); await sleep(stage.pre_iteration_delay_ms ?? 0); @@ -1898,7 +1953,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); + observationSamples.push(execution.observationSample); members[member.id] = { role: member.role ?? null, result: execution.result, @@ -1940,7 +1995,7 @@ async function executePair( normalizedResponses, metrics: pairMetricRow, errors, - observations + observationSamples }; } @@ -1985,6 +2040,26 @@ export async function runBenchmarkInstantiation(instantiationId: string) { const policy = parseExecutionPolicy(record.document); const requestedMetrics = templateMetricsFromInstantiation(record.document); const requestedAggregations = templateAggregationsFromInstantiation(record.document); + const streaming = runtimeParameters(record.document).stream === true; + const plannedObservationSamples = planMetricObservationSamples({ + stages: stages.map((stage) => ({ + stage_id: stage.id, + stage_type: stage.type, + item_count: items.length, + iterations_per_item: stage.iterations_per_item ?? 1, + pair_member_ids: (stage.pair ?? []).map((member) => member.id), + record_metrics: stage.record_metrics !== false + })), + streaming + }); + const observationSamplesByKey = new Map( + plannedObservationSamples.map((sample) => [metricObservationSampleKey(sample), sample]) + ); + const canonicalMetricIntents = canonicalMetricIntentsForRun( + stages, + requestedMetrics, + streaming + ); const server = getInferenceServerById(record.server_id); if (!server) { throw new BenchmarkNotFoundError(`Inference server not found: ${record.server_id}`); @@ -2026,6 +2101,15 @@ export async function runBenchmarkInstantiation(instantiationId: string) { } else if (execution.normalizedResponse) { normalizedResponses.push(execution.normalizedResponse); } + const executionObservationSamples = 'observationSamples' in execution + ? execution.observationSamples + : [execution.observationSample]; + for (const sample of executionObservationSamples) { + const key = metricObservationSampleKey(sample); + if (observationSamplesByKey.has(key)) { + observationSamplesByKey.set(key, sample); + } + } if (execution.metrics && (stage.record_metrics !== false)) metricResults.push(execution.metrics); const executionErrors = 'errors' in execution ? execution.errors : execution.error ? [execution.error] : []; if (executionErrors.length > 0) { @@ -2070,6 +2154,14 @@ export async function runBenchmarkInstantiation(instantiationId: string) { } } + const canonicalAggregates = aggregateMetricObservations({ + samples: [...observationSamplesByKey.values()], + intents: canonicalMetricIntents, + requestedAggregations + }); + // Canonical aggregates remain transient until the atomic metrics-v2 persistence cutover. + void canonicalAggregates; + const resultDocument = { kind: 'test_run_result', schema_version: 'benchmark_test_run_result_v1', diff --git a/backend/tests/integration/benchmark-runner.test.ts b/backend/tests/integration/benchmark-runner.test.ts index aa85f87..91c6632 100644 --- a/backend/tests/integration/benchmark-runner.test.ts +++ b/backend/tests/integration/benchmark-runner.test.ts @@ -28,6 +28,12 @@ function expectNoCanonicalDerivedMetrics(value: Record): void { } } +function expectNoCanonicalObservationState(value: Record): void { + expect(value).not.toHaveProperty('metric_observations'); + expect(value).not.toHaveProperty('canonical_aggregates'); + expect(value).not.toHaveProperty('metric_observation_aggregates'); +} + function streamFixture(name: string): string { return fs.readFileSync(new URL(`../fixtures/streams/${name}`, import.meta.url), 'utf8'); } @@ -636,6 +642,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'); + expectNoCanonicalObservationState(result.document); expectNoCanonicalDerivedMetrics(result.document.metric_results[0]); expect(result.document.aggregated_metrics.elapsed_ms.valid_sample_count).toBe(1); expect(mockServer.requests).toHaveLength(1); @@ -699,6 +706,8 @@ describe('benchmark runner API', () => { expect(result.document.metric_results[0].cold_token_delta).toBe(6); expect(result.document.aggregated_metrics.cold_token_delta.mean).toBe(5.5); expect(result.document.aggregated_metrics.cold_token_delta.valid_sample_count).toBe(2); + expect(result.document.metric_version).toBe('metrics-v1'); + expectNoCanonicalObservationState(result.document); await app.close(); }); @@ -967,7 +976,7 @@ describe('benchmark runner API', () => { output_tokens: 9 }); expect(result.document.metric_version).toBe('metrics-v1'); - expect(result.document).not.toHaveProperty('metric_observations'); + expectNoCanonicalObservationState(result.document); 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'); diff --git a/backend/tests/unit/benchmark-observation-aggregation.test.ts b/backend/tests/unit/benchmark-observation-aggregation.test.ts new file mode 100644 index 0000000..fd6a011 --- /dev/null +++ b/backend/tests/unit/benchmark-observation-aggregation.test.ts @@ -0,0 +1,594 @@ +import { describe, expect, it } from 'vitest'; + +import { + aggregateMetricObservations, + planMetricObservationSamples, + resolveCanonicalMetricIntents, + type CanonicalMetricIntent, + type MetricObservationSample +} from '../../src/services/benchmark-observation-aggregation.js'; +import type { + MetricObservation, + MetricObservationStatus +} from '../../src/services/benchmark-metric-observations.js'; + +function observation(input: { + metricId: string; + value: number | boolean | null; + unit: string; + status?: MetricObservationStatus; + providerId?: string | null; + accountingScope?: Record | null; +}): MetricObservation { + return { + metric_id: input.metricId, + value: input.value, + unit: input.unit, + status: input.status ?? 'measured', + reason: input.status && input.status !== 'measured' ? 'Test status.' : null, + source: input.providerId === undefined ? 'client_observed' : 'provider_reported', + metric_version: 'metrics-v2', + provider_id: input.providerId ?? null, + provider_protocol: input.providerId === undefined ? null : 'openai_chat', + provider_version: input.providerId === undefined ? null : '2026-07-01', + native_field: input.providerId === undefined ? null : `usage.${input.metricId}`, + native_value: input.value, + native_unit: input.unit, + normalization: 'identity', + accounting_scope: input.accountingScope ?? null + }; +} + +function sample(input: { + stageId?: string; + pairMemberId?: string | null; + itemIndex?: number; + iteration?: number; + attempted?: boolean; + completed?: boolean; + attempts?: MetricObservation[][]; + terminal?: MetricObservation[] | null; +} = {}): MetricObservationSample { + const terminal = input.terminal === undefined ? [] : input.terminal; + return { + stage_id: input.stageId ?? 'chat', + item_index: input.itemIndex ?? 0, + iteration: input.iteration ?? 0, + pair_member_id: input.pairMemberId ?? null, + streaming: false, + expected: true, + attempted: input.attempted ?? true, + completed: input.completed ?? true, + attempt_observations: input.attempts ?? (terminal ? [terminal] : []), + terminal_observations: terminal + }; +} + +function intent(input: { + metricId: string; + valueType?: 'number' | 'boolean'; + unit?: string; + stageId?: string; + pairMemberId?: string | null; +}): CanonicalMetricIntent { + return { + stage_id: input.stageId ?? 'chat', + pair_member_id: input.pairMemberId ?? null, + metric_id: input.metricId, + value_type: input.valueType ?? 'number', + unit: input.unit ?? 'milliseconds' + }; +} + +function metricIds(intents: CanonicalMetricIntent[]): string[] { + return intents.map((candidate) => candidate.metric_id); +} + +describe('resolveCanonicalMetricIntents', () => { + it('adds execution defaults and maps legacy performance metrics', () => { + const intents = resolveCanonicalMetricIntents({ + stage_id: 'chat', + pair_member_id: null, + requested_metrics: [ + 'input_tokens', + 'output_tokens', + 'elapsed_ms', + 'tokens_per_second', + 'json_valid' + ], + streaming: false + }); + + expect(metricIds(intents)).toEqual(expect.arrayContaining([ + 'request_success', + 'response_normalization_success', + 'attempt_count', + 'timeout_occurred', + 'retry_overhead_ms', + 'operation_elapsed_ms', + 'successful_attempt_latency_ms', + 'input_tokens', + 'output_tokens', + 'per_request_output_tokens_per_second' + ])); + expect(metricIds(intents)).not.toContain('json_syntax_valid'); + expect(metricIds(intents)).not.toContain('stream_completed'); + }); + + it('adds streaming and tool timing defaults without canonicalizing correctness aliases', () => { + const intents = resolveCanonicalMetricIntents({ + stage_id: 'tools', + pair_member_id: null, + requested_metrics: ['first_token_ms', 'tool_call_assertion_pass'], + streaming: true + }); + + expect(metricIds(intents)).toEqual(expect.arrayContaining([ + 'stream_completed', + 'time_to_first_chunk_ms', + 'time_to_first_output_ms', + 'time_per_output_token_ms', + 'decode_output_tokens_per_second', + 'time_to_first_tool_call_ms', + 'time_to_tool_calls_ready_ms' + ])); + expect(metricIds(intents)).not.toContain('tool_call_assertion_pass'); + }); + + it('resolves only the selected pair member and includes derived-metric operands', () => { + const baseline = resolveCanonicalMetricIntents({ + stage_id: 'pair', + pair_member_id: 'baseline', + requested_metrics: [ + 'pair.baseline.output_tokens', + 'pair.candidate.input_tokens', + 'exact_match' + ], + derived_metric_references: [ + 'pair.baseline.elapsed_ms', + 'pair.candidate.elapsed_ms' + ], + streaming: false + }); + const candidate = resolveCanonicalMetricIntents({ + stage_id: 'pair', + pair_member_id: 'candidate', + requested_metrics: [ + 'pair.baseline.output_tokens', + 'pair.candidate.input_tokens' + ], + derived_metric_references: [ + 'pair.baseline.elapsed_ms', + 'pair.candidate.elapsed_ms' + ], + streaming: false + }); + + expect(metricIds(baseline)).toContain('output_tokens'); + expect(metricIds(baseline)).not.toContain('input_tokens'); + expect(metricIds(candidate)).toContain('input_tokens'); + expect(metricIds(candidate)).not.toContain('output_tokens'); + expect(metricIds(baseline)).toContain('operation_elapsed_ms'); + expect(metricIds(candidate)).toContain('operation_elapsed_ms'); + }); +}); + +describe('planMetricObservationSamples', () => { + it('plans request samples per stage and pair member while excluding non-recording stages', () => { + const samples = planMetricObservationSamples({ + stages: [ + { + stage_id: 'dataset', + stage_type: 'dataset_loop', + item_count: 2, + iterations_per_item: 2, + pair_member_ids: [], + record_metrics: true + }, + { + stage_id: 'single', + stage_type: 'single_request', + item_count: 2, + iterations_per_item: 1, + pair_member_ids: [], + record_metrics: true + }, + { + stage_id: 'pair', + stage_type: 'paired_request_loop', + item_count: 1, + iterations_per_item: 1, + pair_member_ids: ['cold', 'hot'], + record_metrics: true + }, + { + stage_id: 'ignored', + stage_type: 'dataset_loop', + item_count: 2, + iterations_per_item: 1, + pair_member_ids: [], + record_metrics: false + } + ], + streaming: true + }); + + expect(samples).toHaveLength(7); + expect(samples.filter(({ stage_id }) => stage_id === 'dataset')).toHaveLength(4); + expect(samples.filter(({ stage_id }) => stage_id === 'single')).toHaveLength(1); + expect(samples.filter(({ stage_id }) => stage_id === 'pair').map( + ({ pair_member_id }) => pair_member_id + )).toEqual(['cold', 'hot']); + expect(samples.every((candidate) => ( + candidate.streaming + && !candidate.attempted + && candidate.terminal_observations === null + ))).toBe(true); + }); + + it('keeps a pair member unattempted when an earlier member fails', () => { + const planned = planMetricObservationSamples({ + stages: [{ + stage_id: 'pair', + stage_type: 'paired_request_loop', + item_count: 1, + iterations_per_item: 1, + pair_member_ids: ['cold', 'hot'], + record_metrics: true + }], + streaming: false + }); + planned[0] = { + ...planned[0], + attempted: true, + completed: false, + terminal_observations: [] + }; + const aggregates = aggregateMetricObservations({ + samples: planned, + intents: [ + intent({ + stageId: 'pair', + pairMemberId: 'cold', + metricId: 'request_success', + valueType: 'boolean', + unit: 'boolean' + }), + intent({ + stageId: 'pair', + pairMemberId: 'hot', + metricId: 'request_success', + valueType: 'boolean', + unit: 'boolean' + }) + ], + requestedAggregations: [] + }); + + expect(aggregates[0]).toMatchObject({ + pair_member_id: 'cold', + expected_sample_count: 1, + attempted_sample_count: 1, + execution_error_sample_count: 1 + }); + expect(aggregates[1]).toMatchObject({ + pair_member_id: 'hot', + expected_sample_count: 1, + attempted_sample_count: 0, + execution_error_sample_count: 0, + coverage_rate: 0 + }); + }); +}); + +describe('aggregateMetricObservations', () => { + it('aggregates measured values and accounts for every sample status', () => { + const aggregate = aggregateMetricObservations({ + samples: [ + sample({ + itemIndex: 0, + terminal: [observation({ metricId: 'operation_elapsed_ms', value: 100, unit: 'milliseconds' })] + }), + sample({ + itemIndex: 1, + terminal: [observation({ metricId: 'operation_elapsed_ms', value: 200, unit: 'milliseconds' })] + }), + sample({ + itemIndex: 2, + terminal: [observation({ + metricId: 'operation_elapsed_ms', + value: null, + unit: 'milliseconds', + status: 'unavailable' + })] + }), + sample({ itemIndex: 3, attempted: false, completed: false, terminal: null }), + sample({ itemIndex: 4, completed: false, terminal: [] }) + ], + intents: [intent({ metricId: 'operation_elapsed_ms' })], + requestedAggregations: ['mean', 'sum', 'p50', 'p95', 'count'] + })[0]; + + expect(aggregate).toMatchObject({ + expected_sample_count: 5, + attempted_sample_count: 4, + completed_sample_count: 3, + valid_sample_count: 2, + passed_sample_count: 0, + unavailable_sample_count: 1, + not_applicable_sample_count: 0, + execution_error_sample_count: 1, + coverage_rate: 0.4, + statistics: { + mean: 150, + sum: 300, + p50: 150, + p95: 195 + } + }); + expect(aggregate.statistics).not.toHaveProperty('count'); + expect(aggregate.warnings).toContain('insufficient_valid_samples_for_p95'); + }); + + it('computes boolean rates without treating false as unavailable', () => { + const aggregate = aggregateMetricObservations({ + samples: [ + sample({ + itemIndex: 0, + terminal: [observation({ metricId: 'request_success', value: true, unit: 'boolean' })] + }), + sample({ + itemIndex: 1, + terminal: [observation({ metricId: 'request_success', value: false, unit: 'boolean' })] + }), + sample({ + itemIndex: 2, + terminal: [observation({ + metricId: 'request_success', + value: null, + unit: 'boolean', + status: 'not_applicable' + })] + }), + sample({ itemIndex: 3, attempted: false, completed: false, terminal: null }) + ], + intents: [intent({ + metricId: 'request_success', + valueType: 'boolean', + unit: 'boolean' + })], + requestedAggregations: ['mean'] + })[0]; + + expect(aggregate).toMatchObject({ + expected_sample_count: 4, + valid_sample_count: 2, + passed_sample_count: 1, + unavailable_sample_count: 0, + not_applicable_sample_count: 1, + observed_pass_rate: 0.5, + coverage_rate: 0.5, + end_to_end_pass_rate: 0.25, + statistics: {} + }); + }); + + it('reports zero coverage for a fully missing requested metric', () => { + const aggregate = aggregateMetricObservations({ + samples: [ + sample({ itemIndex: 0, terminal: [] }), + sample({ itemIndex: 1, attempted: false, completed: false, terminal: null }) + ], + intents: [intent({ metricId: 'server_decode_time_ms' })], + requestedAggregations: ['mean'] + })[0]; + + expect(aggregate).toMatchObject({ + expected_sample_count: 2, + attempted_sample_count: 1, + completed_sample_count: 1, + valid_sample_count: 0, + unavailable_sample_count: 1, + execution_error_sample_count: 0, + coverage_rate: 0, + statistics: {} + }); + }); + + it('uses only the terminal observation set for a retried request', () => { + const firstAttempt = [ + observation({ metricId: 'attempt_count', value: 1, unit: 'attempts' }) + ]; + const terminal = [ + observation({ metricId: 'attempt_count', value: 2, unit: 'attempts' }) + ]; + const aggregate = aggregateMetricObservations({ + samples: [sample({ + attempts: [firstAttempt, terminal], + terminal + })], + intents: [intent({ metricId: 'attempt_count', unit: 'attempts' })], + requestedAggregations: ['mean', 'sum'] + })[0]; + + expect(aggregate.valid_sample_count).toBe(1); + expect(aggregate.statistics).toEqual({ mean: 2, sum: 2 }); + }); + + it('keeps stages and pair members in separate aggregate groups', () => { + const aggregates = aggregateMetricObservations({ + samples: [ + sample({ + stageId: 'pair', + pairMemberId: 'baseline', + terminal: [observation({ metricId: 'output_tokens', value: 10, unit: 'tokens' })] + }), + sample({ + stageId: 'pair', + pairMemberId: 'candidate', + terminal: [observation({ metricId: 'output_tokens', value: 20, unit: 'tokens' })] + }), + sample({ + stageId: 'other', + terminal: [observation({ metricId: 'output_tokens', value: 30, unit: 'tokens' })] + }) + ], + intents: [ + intent({ + stageId: 'pair', + pairMemberId: 'baseline', + metricId: 'output_tokens', + unit: 'tokens' + }), + intent({ + stageId: 'pair', + pairMemberId: 'candidate', + metricId: 'output_tokens', + unit: 'tokens' + }), + intent({ stageId: 'other', metricId: 'output_tokens', unit: 'tokens' }) + ], + requestedAggregations: ['mean'] + }); + + expect(aggregates.map((aggregate) => [ + aggregate.stage_id, + aggregate.pair_member_id, + aggregate.statistics.mean + ])).toEqual([ + ['pair', 'baseline', 10], + ['pair', 'candidate', 20], + ['other', null, 30] + ]); + }); + + it('preserves zero numeric values as measured evidence', () => { + const aggregate = aggregateMetricObservations({ + samples: [sample({ + terminal: [observation({ metricId: 'retry_overhead_ms', value: 0, unit: 'milliseconds' })] + })], + intents: [intent({ metricId: 'retry_overhead_ms' })], + requestedAggregations: ['mean'] + })[0]; + + expect(aggregate.valid_sample_count).toBe(1); + expect(aggregate.statistics.mean).toBe(0); + }); + + it('suppresses statistics for incompatible provenance or accounting scope', () => { + const aggregate = aggregateMetricObservations({ + samples: [ + sample({ + itemIndex: 0, + terminal: [observation({ + metricId: 'output_tokens', + value: 10, + unit: 'tokens', + providerId: 'server-1', + accountingScope: { candidate_count: 1 } + })] + }), + sample({ + itemIndex: 1, + terminal: [observation({ + metricId: 'output_tokens', + value: 20, + unit: 'tokens', + providerId: 'server-1', + accountingScope: { candidate_count: 2 } + })] + }) + ], + intents: [intent({ metricId: 'output_tokens', unit: 'tokens' })], + requestedAggregations: ['mean', 'sum'] + })[0]; + + expect(aggregate).toMatchObject({ + valid_sample_count: 2, + aggregation_eligible: false, + statistics: {} + }); + expect(aggregate.provenance_signatures).toHaveLength(2); + expect(aggregate.warnings).toContain('incompatible_provenance_or_accounting_scope'); + }); + + it('computes R type 7 percentiles for compatible measured values', () => { + const samples = Array.from({ length: 10 }, (_, index) => sample({ + itemIndex: index, + terminal: [observation({ + metricId: 'operation_elapsed_ms', + value: index + 1, + unit: 'milliseconds' + })] + })); + const aggregate = aggregateMetricObservations({ + samples, + intents: [intent({ metricId: 'operation_elapsed_ms' })], + requestedAggregations: ['p90', 'p95', 'p99'] + })[0]; + + expect(aggregate.statistics.p90).toBeCloseTo(9.1); + expect(aggregate.statistics.p95).toBeCloseTo(9.55); + expect(aggregate.statistics.p99).toBeCloseTo(9.91); + expect(aggregate.warnings).not.toContain('insufficient_valid_samples_for_p90'); + expect(aggregate.warnings).toEqual(expect.arrayContaining([ + 'insufficient_valid_samples_for_p95', + 'insufficient_valid_samples_for_p99' + ])); + }); + + it('rejects duplicate observations and conflicting units, types, or versions', () => { + const duplicate = observation({ + metricId: 'operation_elapsed_ms', + value: 100, + unit: 'milliseconds' + }); + expect(() => aggregateMetricObservations({ + samples: [sample({ terminal: [duplicate, duplicate] })], + intents: [intent({ metricId: 'operation_elapsed_ms' })], + requestedAggregations: ['mean'] + })).toThrow('Duplicate metric observation'); + + expect(() => aggregateMetricObservations({ + samples: [sample({ + terminal: [observation({ + metricId: 'operation_elapsed_ms', + value: 100, + unit: 'seconds' + })] + })], + intents: [intent({ metricId: 'operation_elapsed_ms' })], + requestedAggregations: ['mean'] + })).toThrow('Conflicting unit'); + + expect(() => aggregateMetricObservations({ + samples: [sample({ + terminal: [observation({ + metricId: 'request_success', + value: 1, + unit: 'boolean' + })] + })], + intents: [intent({ + metricId: 'request_success', + valueType: 'boolean', + unit: 'boolean' + })], + requestedAggregations: ['mean'] + })).toThrow('Conflicting value type'); + + const wrongVersion = { + ...observation({ + metricId: 'operation_elapsed_ms', + value: 100, + unit: 'milliseconds' + }), + metric_version: 'metrics-v1' + } as unknown as MetricObservation; + expect(() => aggregateMetricObservations({ + samples: [sample({ terminal: [wrongVersion] })], + intents: [intent({ metricId: 'operation_elapsed_ms' })], + requestedAggregations: ['mean'] + })).toThrow('Unexpected metric version'); + }); +});