diff --git a/openapi/ga/individual/platform.openapi.yaml b/openapi/ga/individual/platform.openapi.yaml index 7384b337b7..ed2a4cf8cc 100644 --- a/openapi/ga/individual/platform.openapi.yaml +++ b/openapi/ga/individual/platform.openapi.yaml @@ -10628,6 +10628,14 @@ components: title: AtifImageSource AtifIngestRequest: properties: + experiment_context: + $ref: '#/components/schemas/ExperimentContext' + evaluation_context: + allOf: + - $ref: '#/components/schemas/EvaluationContext' + description: Deprecated. Use experiment_context; when both are sent, experiment_context + takes precedence. + deprecated: true schema_version: type: string enum: @@ -10643,8 +10651,6 @@ components: session_id: title: Session Id type: string - evaluation_context: - $ref: '#/components/schemas/EvaluationContext' agent: $ref: '#/components/schemas/AtifAgent' final_metrics: @@ -10675,7 +10681,7 @@ components: ATIF project scoping is intentionally not accepted here; use the workspace - route and ``evaluation_context`` for evaluation/run identity.' + route and ``experiment_context`` for experiment identity.' AtifMetrics: properties: prompt_tokens: @@ -12225,6 +12231,14 @@ components: description: User message parameter for chat completion. ChatCompletionsIngestRequest: properties: + experiment_context: + $ref: '#/components/schemas/ExperimentContext' + evaluation_context: + allOf: + - $ref: '#/components/schemas/EvaluationContext' + description: Deprecated. Use experiment_context; when both are sent, experiment_context + takes precedence. + deprecated: true request: $ref: '#/components/schemas/CapturedChatCompletionsRequest' response: @@ -12240,8 +12254,6 @@ components: This is not a grouping mechanism for chat-completions calls; use session_id to group related calls. type: string - evaluation_context: - $ref: '#/components/schemas/EvaluationContext' provider: title: Provider type: string @@ -14604,10 +14616,14 @@ components: p99: title: P99 type: number + count: + type: integer + title: Count + default: 0 type: object title: EvaluatorAggregate - description: Cross-run statistics for one evaluator. Populated by the rollup - path (later PR). + description: Aggregate statistics over evaluator scores or session-level metric + values. EvaluatorResult: properties: evaluator_result_id: @@ -15035,6 +15051,22 @@ components: - action_name title: ExecutedAction description: Information about an action that was executed. + ExperimentContext: + properties: + experiment_id: + type: string + title: Experiment Id + description: Name of an existing Experiment entity. + test_case_id: + title: Test Case Id + description: Optional producer-supplied test case id. + type: string + additionalProperties: false + type: object + required: + - experiment_id + title: ExperimentContext + description: Experiment context accepted by ingest endpoints. ExperimentFilter: additionalProperties: false description: Filter for listing Experiments. @@ -15251,7 +15283,10 @@ components: items: type: string type: array + uniqueItems: true title: Model Names + description: Distinct model names observed across ingested sessions for + this experiment. aggregate_scores: title: Aggregate Scores additionalProperties: @@ -15260,7 +15295,13 @@ components: run_count: type: integer title: Run Count + description: Number of distinct ingested experiment sessions; one session + is treated as one run. default: 0 + cost_usd: + $ref: '#/components/schemas/EvaluatorAggregate' + latency_ms: + $ref: '#/components/schemas/EvaluatorAggregate' type: object required: - id diff --git a/openapi/ga/openapi.yaml b/openapi/ga/openapi.yaml index 7384b337b7..ed2a4cf8cc 100644 --- a/openapi/ga/openapi.yaml +++ b/openapi/ga/openapi.yaml @@ -10628,6 +10628,14 @@ components: title: AtifImageSource AtifIngestRequest: properties: + experiment_context: + $ref: '#/components/schemas/ExperimentContext' + evaluation_context: + allOf: + - $ref: '#/components/schemas/EvaluationContext' + description: Deprecated. Use experiment_context; when both are sent, experiment_context + takes precedence. + deprecated: true schema_version: type: string enum: @@ -10643,8 +10651,6 @@ components: session_id: title: Session Id type: string - evaluation_context: - $ref: '#/components/schemas/EvaluationContext' agent: $ref: '#/components/schemas/AtifAgent' final_metrics: @@ -10675,7 +10681,7 @@ components: ATIF project scoping is intentionally not accepted here; use the workspace - route and ``evaluation_context`` for evaluation/run identity.' + route and ``experiment_context`` for experiment identity.' AtifMetrics: properties: prompt_tokens: @@ -12225,6 +12231,14 @@ components: description: User message parameter for chat completion. ChatCompletionsIngestRequest: properties: + experiment_context: + $ref: '#/components/schemas/ExperimentContext' + evaluation_context: + allOf: + - $ref: '#/components/schemas/EvaluationContext' + description: Deprecated. Use experiment_context; when both are sent, experiment_context + takes precedence. + deprecated: true request: $ref: '#/components/schemas/CapturedChatCompletionsRequest' response: @@ -12240,8 +12254,6 @@ components: This is not a grouping mechanism for chat-completions calls; use session_id to group related calls. type: string - evaluation_context: - $ref: '#/components/schemas/EvaluationContext' provider: title: Provider type: string @@ -14604,10 +14616,14 @@ components: p99: title: P99 type: number + count: + type: integer + title: Count + default: 0 type: object title: EvaluatorAggregate - description: Cross-run statistics for one evaluator. Populated by the rollup - path (later PR). + description: Aggregate statistics over evaluator scores or session-level metric + values. EvaluatorResult: properties: evaluator_result_id: @@ -15035,6 +15051,22 @@ components: - action_name title: ExecutedAction description: Information about an action that was executed. + ExperimentContext: + properties: + experiment_id: + type: string + title: Experiment Id + description: Name of an existing Experiment entity. + test_case_id: + title: Test Case Id + description: Optional producer-supplied test case id. + type: string + additionalProperties: false + type: object + required: + - experiment_id + title: ExperimentContext + description: Experiment context accepted by ingest endpoints. ExperimentFilter: additionalProperties: false description: Filter for listing Experiments. @@ -15251,7 +15283,10 @@ components: items: type: string type: array + uniqueItems: true title: Model Names + description: Distinct model names observed across ingested sessions for + this experiment. aggregate_scores: title: Aggregate Scores additionalProperties: @@ -15260,7 +15295,13 @@ components: run_count: type: integer title: Run Count + description: Number of distinct ingested experiment sessions; one session + is treated as one run. default: 0 + cost_usd: + $ref: '#/components/schemas/EvaluatorAggregate' + latency_ms: + $ref: '#/components/schemas/EvaluatorAggregate' type: object required: - id diff --git a/openapi/openapi.yaml b/openapi/openapi.yaml index 7384b337b7..ed2a4cf8cc 100644 --- a/openapi/openapi.yaml +++ b/openapi/openapi.yaml @@ -10628,6 +10628,14 @@ components: title: AtifImageSource AtifIngestRequest: properties: + experiment_context: + $ref: '#/components/schemas/ExperimentContext' + evaluation_context: + allOf: + - $ref: '#/components/schemas/EvaluationContext' + description: Deprecated. Use experiment_context; when both are sent, experiment_context + takes precedence. + deprecated: true schema_version: type: string enum: @@ -10643,8 +10651,6 @@ components: session_id: title: Session Id type: string - evaluation_context: - $ref: '#/components/schemas/EvaluationContext' agent: $ref: '#/components/schemas/AtifAgent' final_metrics: @@ -10675,7 +10681,7 @@ components: ATIF project scoping is intentionally not accepted here; use the workspace - route and ``evaluation_context`` for evaluation/run identity.' + route and ``experiment_context`` for experiment identity.' AtifMetrics: properties: prompt_tokens: @@ -12225,6 +12231,14 @@ components: description: User message parameter for chat completion. ChatCompletionsIngestRequest: properties: + experiment_context: + $ref: '#/components/schemas/ExperimentContext' + evaluation_context: + allOf: + - $ref: '#/components/schemas/EvaluationContext' + description: Deprecated. Use experiment_context; when both are sent, experiment_context + takes precedence. + deprecated: true request: $ref: '#/components/schemas/CapturedChatCompletionsRequest' response: @@ -12240,8 +12254,6 @@ components: This is not a grouping mechanism for chat-completions calls; use session_id to group related calls. type: string - evaluation_context: - $ref: '#/components/schemas/EvaluationContext' provider: title: Provider type: string @@ -14604,10 +14616,14 @@ components: p99: title: P99 type: number + count: + type: integer + title: Count + default: 0 type: object title: EvaluatorAggregate - description: Cross-run statistics for one evaluator. Populated by the rollup - path (later PR). + description: Aggregate statistics over evaluator scores or session-level metric + values. EvaluatorResult: properties: evaluator_result_id: @@ -15035,6 +15051,22 @@ components: - action_name title: ExecutedAction description: Information about an action that was executed. + ExperimentContext: + properties: + experiment_id: + type: string + title: Experiment Id + description: Name of an existing Experiment entity. + test_case_id: + title: Test Case Id + description: Optional producer-supplied test case id. + type: string + additionalProperties: false + type: object + required: + - experiment_id + title: ExperimentContext + description: Experiment context accepted by ingest endpoints. ExperimentFilter: additionalProperties: false description: Filter for listing Experiments. @@ -15251,7 +15283,10 @@ components: items: type: string type: array + uniqueItems: true title: Model Names + description: Distinct model names observed across ingested sessions for + this experiment. aggregate_scores: title: Aggregate Scores additionalProperties: @@ -15260,7 +15295,13 @@ components: run_count: type: integer title: Run Count + description: Number of distinct ingested experiment sessions; one session + is treated as one run. default: 0 + cost_usd: + $ref: '#/components/schemas/EvaluatorAggregate' + latency_ms: + $ref: '#/components/schemas/EvaluatorAggregate' type: object required: - id diff --git a/packages/nemo_platform_ext/src/nemo_platform_ext/cli/commands/api/intake/ingest/atif.py b/packages/nemo_platform_ext/src/nemo_platform_ext/cli/commands/api/intake/ingest/atif.py index 1c16f67b5c..d4102f5a2c 100644 --- a/packages/nemo_platform_ext/src/nemo_platform_ext/cli/commands/api/intake/ingest/atif.py +++ b/packages/nemo_platform_ext/src/nemo_platform_ext/cli/commands/api/intake/ingest/atif.py @@ -33,6 +33,10 @@ def create_atif( ] = None, continued_trajectory_ref: Annotated[str | None, typer.Option("--continued-trajectory-ref")] = None, evaluation_context: Annotated[str | None, typer.Option("--evaluation-context", help="JSON string")] = None, + experiment_context: Annotated[ + str | None, + typer.Option("--experiment-context", help="Experiment context accepted by ingest endpoints. (JSON string)"), + ] = None, extra: Annotated[str | None, typer.Option("--extra", help="JSON string")] = None, final_metrics: Annotated[str | None, typer.Option("--final-metrics", help="JSON string")] = None, notes: Annotated[str | None, typer.Option("--notes")] = None, @@ -75,6 +79,8 @@ def create_atif( input_payload["continued_trajectory_ref"] = continued_trajectory_ref if evaluation_context is not None: input_payload["evaluation_context"] = read_payload("evaluation_context", evaluation_context) + if experiment_context is not None: + input_payload["experiment_context"] = read_payload("experiment_context", experiment_context) if extra is not None: input_payload["extra"] = read_payload("extra", extra) if final_metrics is not None: diff --git a/packages/nemo_platform_ext/src/nemo_platform_ext/cli/commands/api/intake/ingest/chat_completions.py b/packages/nemo_platform_ext/src/nemo_platform_ext/cli/commands/api/intake/ingest/chat_completions.py index b02a9eee44..f57ade70b8 100644 --- a/packages/nemo_platform_ext/src/nemo_platform_ext/cli/commands/api/intake/ingest/chat_completions.py +++ b/packages/nemo_platform_ext/src/nemo_platform_ext/cli/commands/api/intake/ingest/chat_completions.py @@ -51,6 +51,10 @@ def create_chat_completions( ), ] = None, evaluation_context: Annotated[str | None, typer.Option("--evaluation-context", help="JSON string")] = None, + experiment_context: Annotated[ + str | None, + typer.Option("--experiment-context", help="Experiment context accepted by ingest endpoints. (JSON string)"), + ] = None, provider: Annotated[str | None, typer.Option("--provider")] = None, session_id: Annotated[ str | None, @@ -108,6 +112,8 @@ def create_chat_completions( input_payload["cost_usd"] = cost_usd if evaluation_context is not None: input_payload["evaluation_context"] = read_payload("evaluation_context", evaluation_context) + if experiment_context is not None: + input_payload["experiment_context"] = read_payload("experiment_context", experiment_context) if provider is not None: input_payload["provider"] = provider if session_id is not None: diff --git a/sdk/python/nemo-platform/.nmpcontext/openapi.yaml b/sdk/python/nemo-platform/.nmpcontext/openapi.yaml index 7384b337b7..ed2a4cf8cc 100644 --- a/sdk/python/nemo-platform/.nmpcontext/openapi.yaml +++ b/sdk/python/nemo-platform/.nmpcontext/openapi.yaml @@ -10628,6 +10628,14 @@ components: title: AtifImageSource AtifIngestRequest: properties: + experiment_context: + $ref: '#/components/schemas/ExperimentContext' + evaluation_context: + allOf: + - $ref: '#/components/schemas/EvaluationContext' + description: Deprecated. Use experiment_context; when both are sent, experiment_context + takes precedence. + deprecated: true schema_version: type: string enum: @@ -10643,8 +10651,6 @@ components: session_id: title: Session Id type: string - evaluation_context: - $ref: '#/components/schemas/EvaluationContext' agent: $ref: '#/components/schemas/AtifAgent' final_metrics: @@ -10675,7 +10681,7 @@ components: ATIF project scoping is intentionally not accepted here; use the workspace - route and ``evaluation_context`` for evaluation/run identity.' + route and ``experiment_context`` for experiment identity.' AtifMetrics: properties: prompt_tokens: @@ -12225,6 +12231,14 @@ components: description: User message parameter for chat completion. ChatCompletionsIngestRequest: properties: + experiment_context: + $ref: '#/components/schemas/ExperimentContext' + evaluation_context: + allOf: + - $ref: '#/components/schemas/EvaluationContext' + description: Deprecated. Use experiment_context; when both are sent, experiment_context + takes precedence. + deprecated: true request: $ref: '#/components/schemas/CapturedChatCompletionsRequest' response: @@ -12240,8 +12254,6 @@ components: This is not a grouping mechanism for chat-completions calls; use session_id to group related calls. type: string - evaluation_context: - $ref: '#/components/schemas/EvaluationContext' provider: title: Provider type: string @@ -14604,10 +14616,14 @@ components: p99: title: P99 type: number + count: + type: integer + title: Count + default: 0 type: object title: EvaluatorAggregate - description: Cross-run statistics for one evaluator. Populated by the rollup - path (later PR). + description: Aggregate statistics over evaluator scores or session-level metric + values. EvaluatorResult: properties: evaluator_result_id: @@ -15035,6 +15051,22 @@ components: - action_name title: ExecutedAction description: Information about an action that was executed. + ExperimentContext: + properties: + experiment_id: + type: string + title: Experiment Id + description: Name of an existing Experiment entity. + test_case_id: + title: Test Case Id + description: Optional producer-supplied test case id. + type: string + additionalProperties: false + type: object + required: + - experiment_id + title: ExperimentContext + description: Experiment context accepted by ingest endpoints. ExperimentFilter: additionalProperties: false description: Filter for listing Experiments. @@ -15251,7 +15283,10 @@ components: items: type: string type: array + uniqueItems: true title: Model Names + description: Distinct model names observed across ingested sessions for + this experiment. aggregate_scores: title: Aggregate Scores additionalProperties: @@ -15260,7 +15295,13 @@ components: run_count: type: integer title: Run Count + description: Number of distinct ingested experiment sessions; one session + is treated as one run. default: 0 + cost_usd: + $ref: '#/components/schemas/EvaluatorAggregate' + latency_ms: + $ref: '#/components/schemas/EvaluatorAggregate' type: object required: - id diff --git a/sdk/python/nemo-platform/.nmpcontext/stainless.yaml b/sdk/python/nemo-platform/.nmpcontext/stainless.yaml index e761ae113e..71969fca67 100644 --- a/sdk/python/nemo-platform/.nmpcontext/stainless.yaml +++ b/sdk/python/nemo-platform/.nmpcontext/stainless.yaml @@ -1016,6 +1016,7 @@ resources: ingest: models: evaluation_context: EvaluationContext + experiment_context: ExperimentContext subresources: atif: models: diff --git a/sdk/python/nemo-platform/src/nemo_platform/cli/commands/api/intake/ingest/atif.py b/sdk/python/nemo-platform/src/nemo_platform/cli/commands/api/intake/ingest/atif.py index c596eb126d..62daa27916 100644 --- a/sdk/python/nemo-platform/src/nemo_platform/cli/commands/api/intake/ingest/atif.py +++ b/sdk/python/nemo-platform/src/nemo_platform/cli/commands/api/intake/ingest/atif.py @@ -33,6 +33,10 @@ def create_atif( ] = None, continued_trajectory_ref: Annotated[str | None, typer.Option("--continued-trajectory-ref")] = None, evaluation_context: Annotated[str | None, typer.Option("--evaluation-context", help="JSON string")] = None, + experiment_context: Annotated[ + str | None, + typer.Option("--experiment-context", help="Experiment context accepted by ingest endpoints. (JSON string)"), + ] = None, extra: Annotated[str | None, typer.Option("--extra", help="JSON string")] = None, final_metrics: Annotated[str | None, typer.Option("--final-metrics", help="JSON string")] = None, notes: Annotated[str | None, typer.Option("--notes")] = None, @@ -75,6 +79,8 @@ def create_atif( input_payload["continued_trajectory_ref"] = continued_trajectory_ref if evaluation_context is not None: input_payload["evaluation_context"] = read_payload("evaluation_context", evaluation_context) + if experiment_context is not None: + input_payload["experiment_context"] = read_payload("experiment_context", experiment_context) if extra is not None: input_payload["extra"] = read_payload("extra", extra) if final_metrics is not None: diff --git a/sdk/python/nemo-platform/src/nemo_platform/cli/commands/api/intake/ingest/chat_completions.py b/sdk/python/nemo-platform/src/nemo_platform/cli/commands/api/intake/ingest/chat_completions.py index 0970a4f09f..1e91a12650 100644 --- a/sdk/python/nemo-platform/src/nemo_platform/cli/commands/api/intake/ingest/chat_completions.py +++ b/sdk/python/nemo-platform/src/nemo_platform/cli/commands/api/intake/ingest/chat_completions.py @@ -51,6 +51,10 @@ def create_chat_completions( ), ] = None, evaluation_context: Annotated[str | None, typer.Option("--evaluation-context", help="JSON string")] = None, + experiment_context: Annotated[ + str | None, + typer.Option("--experiment-context", help="Experiment context accepted by ingest endpoints. (JSON string)"), + ] = None, provider: Annotated[str | None, typer.Option("--provider")] = None, session_id: Annotated[ str | None, @@ -108,6 +112,8 @@ def create_chat_completions( input_payload["cost_usd"] = cost_usd if evaluation_context is not None: input_payload["evaluation_context"] = read_payload("evaluation_context", evaluation_context) + if experiment_context is not None: + input_payload["experiment_context"] = read_payload("experiment_context", experiment_context) if provider is not None: input_payload["provider"] = provider if session_id is not None: diff --git a/sdk/python/nemo-platform/src/nemo_platform/resources/experiments/experiments.py b/sdk/python/nemo-platform/src/nemo_platform/resources/experiments/experiments.py index 48bd162d47..fb84ef712d 100644 --- a/sdk/python/nemo-platform/src/nemo_platform/resources/experiments/experiments.py +++ b/sdk/python/nemo-platform/src/nemo_platform/resources/experiments/experiments.py @@ -33,7 +33,6 @@ async_to_streamed_response_wrapper, ) from ...pagination import SyncDefaultPagination, AsyncDefaultPagination -from ..._exceptions import ConflictError from ..._base_client import AsyncPaginator, make_request_options from ...types.experiments import ( experiment_list_params, @@ -42,6 +41,7 @@ ) from ...types.experiments.experiment_response import ExperimentResponse from ...types.experiments.experiment_filter_param import ExperimentFilterParam +from ..._exceptions import ConflictError __all__ = ["ExperimentsResource", "AsyncExperimentsResource"] @@ -155,7 +155,7 @@ def create( except ConflictError: if not exist_ok: raise - return self.retrieve(name=name, workspace=workspace) + return self.retrieve(name = name, workspace = workspace) def retrieve( self, @@ -492,7 +492,7 @@ async def create( except ConflictError: if not exist_ok: raise - return await self.retrieve(name=name, workspace=workspace) + return await self.retrieve(name = name, workspace = workspace) async def retrieve( self, diff --git a/sdk/python/nemo-platform/src/nemo_platform/resources/intake/api.md b/sdk/python/nemo-platform/src/nemo_platform/resources/intake/api.md index 65f63df3b1..8384b9222c 100644 --- a/sdk/python/nemo-platform/src/nemo_platform/resources/intake/api.md +++ b/sdk/python/nemo-platform/src/nemo_platform/resources/intake/api.md @@ -58,7 +58,7 @@ Methods: Types: ```python -from nemo_platform.types.intake import EvaluationContext +from nemo_platform.types.intake import EvaluationContext, ExperimentContext ``` ### Atif diff --git a/sdk/python/nemo-platform/src/nemo_platform/resources/intake/ingest/atif.py b/sdk/python/nemo-platform/src/nemo_platform/resources/intake/ingest/atif.py index e1301481bd..af199756af 100644 --- a/sdk/python/nemo-platform/src/nemo_platform/resources/intake/ingest/atif.py +++ b/sdk/python/nemo-platform/src/nemo_platform/resources/intake/ingest/atif.py @@ -37,6 +37,7 @@ from ....types.intake.ingest.atif_step_param import AtifStepParam from ....types.intake.ingest.atif_agent_param import AtifAgentParam from ....types.intake.evaluation_context_param import EvaluationContextParam +from ....types.intake.experiment_context_param import ExperimentContextParam from ....types.intake.ingest.atif_final_metrics_param import AtifFinalMetricsParam __all__ = ["AtifResource", "AsyncAtifResource"] @@ -72,6 +73,7 @@ def create( ], continued_trajectory_ref: str | Omit = omit, evaluation_context: EvaluationContextParam | Omit = omit, + experiment_context: ExperimentContextParam | Omit = omit, extra: Dict[str, object] | Omit = omit, final_metrics: AtifFinalMetricsParam | Omit = omit, notes: str | Omit = omit, @@ -88,6 +90,8 @@ def create( Ingest Atif Args: + experiment_context: Experiment context accepted by ingest endpoints. + extra_headers: Send extra headers extra_query: Add additional query parameters to the request @@ -109,6 +113,7 @@ def create( "schema_version": schema_version, "continued_trajectory_ref": continued_trajectory_ref, "evaluation_context": evaluation_context, + "experiment_context": experiment_context, "extra": extra, "final_metrics": final_metrics, "notes": notes, @@ -154,6 +159,7 @@ async def create( ], continued_trajectory_ref: str | Omit = omit, evaluation_context: EvaluationContextParam | Omit = omit, + experiment_context: ExperimentContextParam | Omit = omit, extra: Dict[str, object] | Omit = omit, final_metrics: AtifFinalMetricsParam | Omit = omit, notes: str | Omit = omit, @@ -170,6 +176,8 @@ async def create( Ingest Atif Args: + experiment_context: Experiment context accepted by ingest endpoints. + extra_headers: Send extra headers extra_query: Add additional query parameters to the request @@ -191,6 +199,7 @@ async def create( "schema_version": schema_version, "continued_trajectory_ref": continued_trajectory_ref, "evaluation_context": evaluation_context, + "experiment_context": experiment_context, "extra": extra, "final_metrics": final_metrics, "notes": notes, diff --git a/sdk/python/nemo-platform/src/nemo_platform/resources/intake/ingest/chat_completions.py b/sdk/python/nemo-platform/src/nemo_platform/resources/intake/ingest/chat_completions.py index 6e9c03dbc6..44216e1f6a 100644 --- a/sdk/python/nemo-platform/src/nemo_platform/resources/intake/ingest/chat_completions.py +++ b/sdk/python/nemo-platform/src/nemo_platform/resources/intake/ingest/chat_completions.py @@ -36,6 +36,7 @@ chat_completion_create_params, ) from ....types.intake.evaluation_context_param import EvaluationContextParam +from ....types.intake.experiment_context_param import ExperimentContextParam from ....types.intake.ingest.chat_completions_ingest_response import ChatCompletionsIngestResponse from ....types.intake.ingest.captured_chat_completions_request_param import CapturedChatCompletionsRequestParam from ....types.intake.ingest.captured_chat_completions_response_param import CapturedChatCompletionsResponseParam @@ -74,6 +75,7 @@ def create( cost_output_usd: float | Omit = omit, cost_usd: float | Omit = omit, evaluation_context: EvaluationContextParam | Omit = omit, + experiment_context: ExperimentContextParam | Omit = omit, provider: str | Omit = omit, session_id: str | Omit = omit, trace_id: str | Omit = omit, @@ -101,6 +103,8 @@ def create( cost_usd: Total estimated cost of this model call in USD. This matches ATIF step metrics; Intake stores it as semantic cost_total_usd on spans. + experiment_context: Experiment context accepted by ingest endpoints. + session_id: Groups related chat-completions calls without forcing them into the same trace. trace_id: Opt into joining an existing trace built via OTel or ATIF. This is not a @@ -130,6 +134,7 @@ def create( "cost_output_usd": cost_output_usd, "cost_usd": cost_usd, "evaluation_context": evaluation_context, + "experiment_context": experiment_context, "provider": provider, "session_id": session_id, "trace_id": trace_id, @@ -174,6 +179,7 @@ async def create( cost_output_usd: float | Omit = omit, cost_usd: float | Omit = omit, evaluation_context: EvaluationContextParam | Omit = omit, + experiment_context: ExperimentContextParam | Omit = omit, provider: str | Omit = omit, session_id: str | Omit = omit, trace_id: str | Omit = omit, @@ -201,6 +207,8 @@ async def create( cost_usd: Total estimated cost of this model call in USD. This matches ATIF step metrics; Intake stores it as semantic cost_total_usd on spans. + experiment_context: Experiment context accepted by ingest endpoints. + session_id: Groups related chat-completions calls without forcing them into the same trace. trace_id: Opt into joining an existing trace built via OTel or ATIF. This is not a @@ -230,6 +238,7 @@ async def create( "cost_output_usd": cost_output_usd, "cost_usd": cost_usd, "evaluation_context": evaluation_context, + "experiment_context": experiment_context, "provider": provider, "session_id": session_id, "trace_id": trace_id, diff --git a/sdk/python/nemo-platform/src/nemo_platform/types/experiments/evaluator_aggregate.py b/sdk/python/nemo-platform/src/nemo_platform/types/experiments/evaluator_aggregate.py index 8551b4ff43..1b1c95a196 100644 --- a/sdk/python/nemo-platform/src/nemo_platform/types/experiments/evaluator_aggregate.py +++ b/sdk/python/nemo-platform/src/nemo_platform/types/experiments/evaluator_aggregate.py @@ -23,12 +23,9 @@ class EvaluatorAggregate(BaseModel): - """Cross-run statistics for one evaluator. + """Aggregate statistics over evaluator scores or session-level metric values.""" - Populated by the rollup path (later PR). - """ - - sum: Optional[float] = None + count: Optional[int] = None mean: Optional[float] = None @@ -39,3 +36,5 @@ class EvaluatorAggregate(BaseModel): p95: Optional[float] = None p99: Optional[float] = None + + sum: Optional[float] = None diff --git a/sdk/python/nemo-platform/src/nemo_platform/types/experiments/experiment_response.py b/sdk/python/nemo-platform/src/nemo_platform/types/experiments/experiment_response.py index 33c7ec3500..6ac77ec12d 100644 --- a/sdk/python/nemo-platform/src/nemo_platform/types/experiments/experiment_response.py +++ b/sdk/python/nemo-platform/src/nemo_platform/types/experiments/experiment_response.py @@ -18,6 +18,7 @@ from typing import Dict, List, Optional from datetime import datetime +from ..._compat import PYDANTIC_V1, ConfigDict from ..._models import BaseModel from .evaluator_aggregate import EvaluatorAggregate @@ -41,6 +42,9 @@ class ExperimentResponse(BaseModel): aggregate_scores: Optional[Dict[str, EvaluatorAggregate]] = None + cost_usd: Optional[EvaluatorAggregate] = None + """Aggregate statistics over evaluator scores or session-level metric values.""" + created_at: Optional[datetime] = None dataset_version: Optional[str] = None @@ -55,14 +59,26 @@ class ExperimentResponse(BaseModel): Soft reference, not validated. """ + latency_ms: Optional[EvaluatorAggregate] = None + """Aggregate statistics over evaluator scores or session-level metric values.""" + metadata: Optional[Dict[str, object]] = None model_names: Optional[List[str]] = None + """Distinct model names observed across ingested sessions for this experiment.""" run_count: Optional[int] = None + """ + Number of distinct ingested experiment sessions; one session is treated as one + run. + """ source_link: Optional[str] = None summary: Optional[str] = None updated_at: Optional[datetime] = None + + if not PYDANTIC_V1: + # allow fields with a `model_` prefix + model_config = ConfigDict(protected_namespaces=tuple()) diff --git a/sdk/python/nemo-platform/src/nemo_platform/types/intake/__init__.py b/sdk/python/nemo-platform/src/nemo_platform/types/intake/__init__.py index 0dfcdafc78..f9405580c2 100644 --- a/sdk/python/nemo-platform/src/nemo_platform/types/intake/__init__.py +++ b/sdk/python/nemo-platform/src/nemo_platform/types/intake/__init__.py @@ -50,6 +50,7 @@ from .span_evaluation_context import SpanEvaluationContext as SpanEvaluationContext from .annotation_create_params import AnnotationCreateParams as AnnotationCreateParams from .evaluation_context_param import EvaluationContextParam as EvaluationContextParam +from .experiment_context_param import ExperimentContextParam as ExperimentContextParam from .feedback_annotation_param import FeedbackAnnotationParam as FeedbackAnnotationParam from .metadata_annotation_param import MetadataAnnotationParam as MetadataAnnotationParam from .evaluator_result_data_type import EvaluatorResultDataType as EvaluatorResultDataType diff --git a/sdk/python/nemo-platform/src/nemo_platform/types/intake/experiment_context_param.py b/sdk/python/nemo-platform/src/nemo_platform/types/intake/experiment_context_param.py new file mode 100644 index 0000000000..ca6a7b3850 --- /dev/null +++ b/sdk/python/nemo-platform/src/nemo_platform/types/intake/experiment_context_param.py @@ -0,0 +1,32 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +# File generated from our OpenAPI spec by Stainless. See CONTRIBUTING.md for details. + +from __future__ import annotations + +from typing_extensions import Required, TypedDict + +__all__ = ["ExperimentContextParam"] + + +class ExperimentContextParam(TypedDict, total=False): + """Experiment context accepted by ingest endpoints.""" + + experiment_id: Required[str] + """Name of an existing Experiment entity.""" + + test_case_id: str + """Optional producer-supplied test case id.""" diff --git a/sdk/python/nemo-platform/src/nemo_platform/types/intake/ingest/atif_create_params.py b/sdk/python/nemo-platform/src/nemo_platform/types/intake/ingest/atif_create_params.py index b3296a06c8..d25fab2a17 100644 --- a/sdk/python/nemo-platform/src/nemo_platform/types/intake/ingest/atif_create_params.py +++ b/sdk/python/nemo-platform/src/nemo_platform/types/intake/ingest/atif_create_params.py @@ -24,6 +24,7 @@ from .atif_agent_param import AtifAgentParam from .atif_final_metrics_param import AtifFinalMetricsParam from ..evaluation_context_param import EvaluationContextParam +from ..experiment_context_param import ExperimentContextParam __all__ = ["AtifCreateParams"] @@ -41,6 +42,9 @@ class AtifCreateParams(TypedDict, total=False): evaluation_context: EvaluationContextParam + experiment_context: ExperimentContextParam + """Experiment context accepted by ingest endpoints.""" + extra: Dict[str, object] final_metrics: AtifFinalMetricsParam diff --git a/sdk/python/nemo-platform/src/nemo_platform/types/intake/ingest/chat_completion_create_params.py b/sdk/python/nemo-platform/src/nemo_platform/types/intake/ingest/chat_completion_create_params.py index 49bb4d6168..8071f2dba0 100644 --- a/sdk/python/nemo-platform/src/nemo_platform/types/intake/ingest/chat_completion_create_params.py +++ b/sdk/python/nemo-platform/src/nemo_platform/types/intake/ingest/chat_completion_create_params.py @@ -21,6 +21,7 @@ from typing_extensions import Required, TypedDict from ..evaluation_context_param import EvaluationContextParam +from ..experiment_context_param import ExperimentContextParam from .captured_chat_completions_request_param import CapturedChatCompletionsRequestParam from .captured_chat_completions_response_param import CapturedChatCompletionsResponseParam @@ -54,6 +55,9 @@ class ChatCompletionCreateParams(TypedDict, total=False): evaluation_context: EvaluationContextParam + experiment_context: ExperimentContextParam + """Experiment context accepted by ingest endpoints.""" + provider: str session_id: str diff --git a/sdk/python/nemo-platform/tests/api_resources/intake/ingest/test_atif.py b/sdk/python/nemo-platform/tests/api_resources/intake/ingest/test_atif.py index 745f896bbf..797c3d460a 100644 --- a/sdk/python/nemo-platform/tests/api_resources/intake/ingest/test_atif.py +++ b/sdk/python/nemo-platform/tests/api_resources/intake/ingest/test_atif.py @@ -68,6 +68,10 @@ def test_method_create_with_all_params(self, client: NeMoPlatform) -> None: "metadata": {"foo": "bar"}, "test_case_id": "test_case_id", }, + experiment_context={ + "experiment_id": "experiment_id", + "test_case_id": "test_case_id", + }, extra={"foo": "bar"}, final_metrics={ "extra": {"foo": "bar"}, @@ -185,6 +189,10 @@ async def test_method_create_with_all_params(self, async_client: AsyncNeMoPlatfo "metadata": {"foo": "bar"}, "test_case_id": "test_case_id", }, + experiment_context={ + "experiment_id": "experiment_id", + "test_case_id": "test_case_id", + }, extra={"foo": "bar"}, final_metrics={ "extra": {"foo": "bar"}, diff --git a/sdk/python/nemo-platform/tests/api_resources/intake/ingest/test_chat_completions.py b/sdk/python/nemo-platform/tests/api_resources/intake/ingest/test_chat_completions.py index 5c9a0a039c..2b5b08ef11 100644 --- a/sdk/python/nemo-platform/tests/api_resources/intake/ingest/test_chat_completions.py +++ b/sdk/python/nemo-platform/tests/api_resources/intake/ingest/test_chat_completions.py @@ -74,6 +74,10 @@ def test_method_create_with_all_params(self, client: NeMoPlatform) -> None: "metadata": {"foo": "bar"}, "test_case_id": "test_case_id", }, + experiment_context={ + "experiment_id": "experiment_id", + "test_case_id": "test_case_id", + }, provider="provider", session_id="session_id", trace_id="trace_id", @@ -175,6 +179,10 @@ async def test_method_create_with_all_params(self, async_client: AsyncNeMoPlatfo "metadata": {"foo": "bar"}, "test_case_id": "test_case_id", }, + experiment_context={ + "experiment_id": "experiment_id", + "test_case_id": "test_case_id", + }, provider="provider", session_id="session_id", trace_id="trace_id", diff --git a/sdk/python/nemo-platform/tests/api_resources/test_experiments.py b/sdk/python/nemo-platform/tests/api_resources/test_experiments.py index b1dd5efc73..b001e0d939 100644 --- a/sdk/python/nemo-platform/tests/api_resources/test_experiments.py +++ b/sdk/python/nemo-platform/tests/api_resources/test_experiments.py @@ -60,7 +60,7 @@ def test_method_create_with_all_params(self, client: NeMoPlatform) -> None: description="description", experiment_group_id="experiment_group_id", metadata={"foo": "bar"}, - source_link="https://example.com/experiments/source", + source_link="https://example.com", summary="summary", ) assert_matches_type(ExperimentResponse, experiment, path=["response"]) @@ -190,7 +190,7 @@ def test_method_update_with_all_params(self, client: NeMoPlatform) -> None: description="description", experiment_group_id="experiment_group_id", metadata={"foo": "bar"}, - source_link="https://example.com/experiments/source", + source_link="https://example.com", summary="summary", ) assert_matches_type(ExperimentResponse, experiment, path=["response"]) @@ -396,7 +396,7 @@ async def test_method_create_with_all_params(self, async_client: AsyncNeMoPlatfo description="description", experiment_group_id="experiment_group_id", metadata={"foo": "bar"}, - source_link="https://example.com/experiments/source", + source_link="https://example.com", summary="summary", ) assert_matches_type(ExperimentResponse, experiment, path=["response"]) @@ -526,7 +526,7 @@ async def test_method_update_with_all_params(self, async_client: AsyncNeMoPlatfo description="description", experiment_group_id="experiment_group_id", metadata={"foo": "bar"}, - source_link="https://example.com/experiments/source", + source_link="https://example.com", summary="summary", ) assert_matches_type(ExperimentResponse, experiment, path=["response"]) diff --git a/sdk/stainless.yaml b/sdk/stainless.yaml index e761ae113e..71969fca67 100644 --- a/sdk/stainless.yaml +++ b/sdk/stainless.yaml @@ -1016,6 +1016,7 @@ resources: ingest: models: evaluation_context: EvaluationContext + experiment_context: ExperimentContext subresources: atif: models: diff --git a/services/intake/README.md b/services/intake/README.md index 46366c3c46..bda01b8139 100644 --- a/services/intake/README.md +++ b/services/intake/README.md @@ -53,6 +53,19 @@ Read it back: curl -i "http://127.0.0.1:8080/apis/intake/v2/workspaces/default/spans?filter[session_id]=sample-session" ``` +Seed an Experiment rollup and read it back: + +```bash +uv run services/intake/scripts/spans/seed_experiment_rollup_data.py +curl -s "http://127.0.0.1:8080/apis/intake/v2/workspaces/default/experiments/rollup-smoke-exp" | jq + +# Optional larger local workload. +uv run services/intake/scripts/spans/seed_experiment_rollup_data.py \ + --experiment rollup-perf-exp \ + --runs 100 \ + --cases-per-run 10 +``` + ## Testing Focused route-surface test: diff --git a/services/intake/scripts/spans/seed_experiment_rollup_data.py b/services/intake/scripts/spans/seed_experiment_rollup_data.py new file mode 100644 index 0000000000..44324e2989 --- /dev/null +++ b/services/intake/scripts/spans/seed_experiment_rollup_data.py @@ -0,0 +1,227 @@ +#!/usr/bin/env python +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +"""Seed local Intake with valid experiment telemetry and verify API rollups.""" + +from __future__ import annotations + +import argparse +import json +import time +from datetime import datetime, timedelta, timezone +from typing import Any +from urllib.parse import urlsplit, urlunsplit + +import httpx + +DEFAULT_BASE_URL = "http://127.0.0.1:8080" +DEFAULT_WORKSPACE = "default" +DEFAULT_EXPERIMENT = "rollup-smoke-exp" +DATASET_NAME = "rollup-smoke-dataset" +AGENT_NAME = "sample-agent" +AGENT_VERSION = "1.0.0" + +SAMPLE_ROWS = [ + ("run-1", "case-1a", 0.4, 0.05, 500), + ("run-1", "case-1b", 0.8, 0.10, 1500), + ("run-2", "case-2", 0.8, 0.20, 2000), + ("run-3", "case-3", 1.0, 0.30, 3000), +] + + +def main() -> None: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--base-url", default=DEFAULT_BASE_URL) + parser.add_argument("--workspace", default=DEFAULT_WORKSPACE) + parser.add_argument("--experiment", default=DEFAULT_EXPERIMENT) + parser.add_argument( + "--runs", type=int, help="Generate this many synthetic experiment runs instead of the smoke set." + ) + parser.add_argument("--cases-per-run", type=int, default=1, help="Synthetic test cases per run when --runs is set.") + args = parser.parse_args() + if args.runs is not None and args.runs < 1: + raise SystemExit("--runs must be >= 1") + if args.cases_per_run < 1: + raise SystemExit("--cases-per-run must be >= 1") + + sample_rows = _sample_rows(runs=args.runs, cases_per_run=args.cases_per_run) + + base_url = args.base_url.rstrip("/") + _preflight(base_url) + + with httpx.Client(timeout=10.0) as client: + _upsert_experiment(client, base_url, args.workspace, args.experiment) + started_at = datetime.now(timezone.utc).replace(microsecond=0) + for index, (run_id, test_case_id, score, cost_usd, latency_ms) in enumerate(sample_rows): + response = client.post( + _intake_url(base_url, args.workspace, "/ingest/atif"), + json=_atif_body( + started_at=started_at, + experiment_id=args.experiment, + run_id=run_id, + test_case_id=test_case_id, + score=score, + cost_usd=cost_usd, + latency_ms=latency_ms, + offset_seconds=index * 10, + ), + ) + response.raise_for_status() + if (index + 1) % 100 == 0: + print(f"seeded {index + 1} sessions") + + expected_session_count = len(sample_rows) + experiment = _wait_for_rollup( + client, + base_url, + args.workspace, + args.experiment, + expected_run_count=expected_session_count, + expected_session_count=expected_session_count, + ) + print(json.dumps(experiment, indent=2, sort_keys=True)) + + +def _preflight(base_url: str) -> None: + try: + response = httpx.get(_replace_path(base_url, "/openapi.json"), timeout=2.0) + response.raise_for_status() + except Exception as exc: + raise SystemExit(f"Cannot reach NeMo Platform at {base_url}: {exc}") from exc + + +def _upsert_experiment(client: httpx.Client, base_url: str, workspace: str, experiment: str) -> None: + body = { + "name": experiment, + "agent_name": AGENT_NAME, + "agent_version": AGENT_VERSION, + "dataset_name": DATASET_NAME, + "dataset_version": "v1", + "metadata": {"seeded_by": "services/intake/scripts/spans/seed_experiment_rollup_data.py"}, + } + response = client.post(_intake_url(base_url, workspace, "/experiments"), json=body) + if response.status_code == 409: + response = client.put(_intake_url(base_url, workspace, f"/experiments/{experiment}"), json=body) + response.raise_for_status() + + +def _wait_for_rollup( + client: httpx.Client, + base_url: str, + workspace: str, + experiment: str, + *, + expected_run_count: int, + expected_session_count: int, +) -> dict[str, Any]: + url = _intake_url(base_url, workspace, f"/experiments/{experiment}") + last_response: httpx.Response | None = None + for _ in range(20): + response = client.get(url) + last_response = response + response.raise_for_status() + payload = response.json() + score = (payload.get("aggregate_scores") or {}).get("harbor.verifier") or {} + if payload.get("run_count") == expected_run_count and score.get("count") == expected_session_count: + return payload + time.sleep(0.25) + detail = last_response.text if last_response is not None else "" + raise SystemExit(f"Experiment rollup did not become visible at {url}: {detail}") + + +def _atif_body( + *, + started_at: datetime, + experiment_id: str, + run_id: str, + test_case_id: str, + score: float, + cost_usd: float, + latency_ms: int, + offset_seconds: int, +) -> dict[str, Any]: + session_started_at = started_at + timedelta(seconds=offset_seconds) + finished_at = session_started_at + timedelta(milliseconds=latency_ms) + session_id = f"{experiment_id}-{run_id}-{test_case_id}" + return { + "schema_version": "ATIF-v1.7", + "session_id": session_id, + "experiment_context": { + "experiment_id": experiment_id, + "test_case_id": test_case_id, + }, + "extra": { + "task_id": test_case_id, + "task_name": test_case_id, + "verifier": { + "started_at": _iso(session_started_at), + "finished_at": _iso(finished_at), + }, + "verifier_result": {"rewards": {"reward": score}}, + }, + "agent": { + "name": AGENT_NAME, + "version": AGENT_VERSION, + "model_name": "provider/sample-model", + }, + "steps": [ + { + "step_id": 1, + "timestamp": _iso(session_started_at), + "source": "agent", + "model_name": "provider/sample-model", + "message": f"solved {test_case_id}", + "metrics": { + "prompt_tokens": 100, + "completion_tokens": 10, + "cost_usd": cost_usd, + }, + } + ], + } + + +def _sample_rows(runs: int | None, cases_per_run: int) -> list[tuple[str, str, float, float, int]]: + if runs is None: + return SAMPLE_ROWS + return [ + ( + f"run-{run_index}", + f"case-{run_index}-{case_index}", + _synthetic_score(run_index=run_index, case_index=case_index), + _synthetic_cost(case_index=case_index), + _synthetic_latency_ms(run_index=run_index, case_index=case_index), + ) + for run_index in range(1, runs + 1) + for case_index in range(1, cases_per_run + 1) + ] + + +def _synthetic_score(*, run_index: int, case_index: int) -> float: + return round(0.4 + (((run_index * 17 + case_index * 11) % 60) / 100), 3) + + +def _synthetic_cost(*, case_index: int) -> float: + return round(0.001 * case_index, 6) + + +def _synthetic_latency_ms(*, run_index: int, case_index: int) -> int: + return 250 + run_index + case_index * 10 + + +def _intake_url(base_url: str, workspace: str, suffix: str) -> str: + return f"{base_url}/apis/intake/v2/workspaces/{workspace}{suffix}" + + +def _replace_path(base_url: str, path: str) -> str: + parts = urlsplit(base_url) + return urlunsplit((parts.scheme, parts.netloc, path, "", "")) + + +def _iso(value: datetime) -> str: + return value.isoformat().replace("+00:00", "Z") + + +if __name__ == "__main__": + main() diff --git a/services/intake/src/nmp/intake/api/v2/experiments/endpoints.py b/services/intake/src/nmp/intake/api/v2/experiments/endpoints.py index 052220b523..913ebbe987 100644 --- a/services/intake/src/nmp/intake/api/v2/experiments/endpoints.py +++ b/services/intake/src/nmp/intake/api/v2/experiments/endpoints.py @@ -3,16 +3,16 @@ """Create, list, get, and delete endpoints for Experiments and ExperimentGroups. -Entity-store (Postgres) operations wired directly onto ``EntityClient``, following -the inline pattern used by the core services. PUT updates only the mutable fields -(group membership, summary, description, metadata); an Experiment's identity and the -dataset/agent it ran against are fixed and changing them is rejected. Rollup fields on -the read models are hydrated from ClickHouse in a later PR; for now they return defaults. +Entity-store (Postgres) operations are wired directly onto ``EntityClient``, +following the inline pattern used by the core services. PUT updates only the +mutable fields; an Experiment's identity and the dataset/agent it ran against +are fixed. Rollup fields on read models are hydrated from ClickHouse. """ from __future__ import annotations -from typing import Annotated, Literal +import logging +from typing import Annotated, Literal, TypeVar from fastapi import APIRouter, Depends, HTTPException, Query, Request, status from nmp.common.api.common import Page, PaginationData @@ -21,6 +21,7 @@ from nmp.common.entities.client import EntityClient, EntityConflictError, EntityNotFoundError from nmp.common.service.dependencies import get_entity_client from nmp.intake.api.v2.experiments.schemas import ( + EvaluatorAggregate, ExperimentFilter, ExperimentGroupFilter, ExperimentGroupRequest, @@ -30,22 +31,42 @@ ) from nmp.intake.entities.experiments import Experiment, ExperimentGroup from nmp.intake.spans.api.dependencies import require_workspace_access, validate_list_query_params +from nmp.intake.spans.experiment_rollup_repository import ( + ExperimentRollup, + ExperimentRollupRepository, + ScoreRollup, +) +logger = logging.getLogger(__name__) router = APIRouter(dependencies=[Depends(require_workspace_access)]) GROUPS_TAG = "Experiment Groups" EXPERIMENTS_TAG = "Experiments" SortField = Literal["-created_at", "created_at", "-updated_at", "updated_at", "-name", "name"] +EntityT = TypeVar("EntityT", Experiment, ExperimentGroup) EntityClientDep = Annotated[EntityClient, Depends(get_entity_client)] ExperimentGroupFilterDep = Annotated[ParsedFilter, Depends(make_filter_dep(ExperimentGroupFilter))] ExperimentFilterDep = Annotated[ParsedFilter, Depends(make_filter_dep(ExperimentFilter))] -# ============================================================================= -# Experiment Groups -# ============================================================================= +def get_experiment_rollup_repository(request: Request) -> ExperimentRollupRepository | None: + # Rollups are enrichment only. Experiment entity reads should continue when + # ClickHouse is disabled or temporarily unavailable. + service = getattr(request.app.state, "intake_service", None) + if service is None: + service = getattr(request.app.state, "service", None) + if service is None: + return None + + service_client = getattr(service, "clickhouse_client", None) + if service_client is None: + return None + return ExperimentRollupRepository(service_client) + + +ExperimentRollupRepositoryDep = Annotated[ExperimentRollupRepository | None, Depends(get_experiment_rollup_repository)] @router.post( @@ -117,13 +138,13 @@ async def get_experiment_group( name: str, entity_client: EntityClientDep, ) -> ExperimentGroupResponse: - try: - entity = await entity_client.get(ExperimentGroup, name=name, workspace=workspace) - except EntityNotFoundError as e: - raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND, - detail=f"Experiment group '{workspace}/{name}' not found.", - ) from e + entity = await _get_or_404( + entity_client, + ExperimentGroup, + workspace=workspace, + name=name, + label="Experiment group", + ) return ExperimentGroupResponse.from_entity(entity) @@ -142,13 +163,13 @@ async def update_experiment_group( body: ExperimentGroupRequest, entity_client: EntityClientDep, ) -> ExperimentGroupResponse: - try: - existing = await entity_client.get(ExperimentGroup, name=name, workspace=workspace) - except EntityNotFoundError as e: - raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND, - detail=f"Experiment group '{workspace}/{name}' not found.", - ) from e + existing = await _get_or_404( + entity_client, + ExperimentGroup, + workspace=workspace, + name=name, + label="Experiment group", + ) if body.name != name: raise HTTPException( status_code=status.HTTP_409_CONFLICT, @@ -170,18 +191,13 @@ async def delete_experiment_group( name: str, entity_client: EntityClientDep, ) -> None: - try: - await entity_client.delete(ExperimentGroup, name=name, workspace=workspace) - except EntityNotFoundError as e: - raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND, - detail=f"Experiment group '{workspace}/{name}' not found.", - ) from e - - -# ============================================================================= -# Experiments -# ============================================================================= + await _delete_or_404( + entity_client, + ExperimentGroup, + workspace=workspace, + name=name, + label="Experiment group", + ) @router.post( @@ -232,6 +248,7 @@ async def list_experiments( workspace: str, request: Request, entity_client: EntityClientDep, + rollup_repository: ExperimentRollupRepositoryDep, parsed: ExperimentFilterDep, page: int = Query(default=1, ge=1, description="Page number."), page_size: int = Query(default=100, ge=1, le=1000, description="Page size."), @@ -246,8 +263,10 @@ async def list_experiments( page=page, page_size=page_size, ) + responses = [ExperimentResponse.from_entity(e) for e in result.data] + await _hydrate_rollups(workspace=workspace, responses=responses, rollup_repository=rollup_repository) return Page( - data=[ExperimentResponse.from_entity(e) for e in result.data], + data=responses, pagination=PaginationData(**result.pagination.model_dump()), sort=sort, filter=parsed.to_response(), @@ -264,15 +283,18 @@ async def get_experiment( workspace: str, name: str, entity_client: EntityClientDep, + rollup_repository: ExperimentRollupRepositoryDep, ) -> ExperimentResponse: - try: - entity = await entity_client.get(Experiment, name=name, workspace=workspace) - except EntityNotFoundError as e: - raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND, - detail=f"Experiment '{workspace}/{name}' not found.", - ) from e - return ExperimentResponse.from_entity(entity) + entity = await _get_or_404( + entity_client, + Experiment, + workspace=workspace, + name=name, + label="Experiment", + ) + response = ExperimentResponse.from_entity(entity) + await _hydrate_rollups(workspace=workspace, responses=[response], rollup_repository=rollup_repository) + return response # Identity and the dataset/agent it was run against are fixed for the life of an @@ -295,14 +317,15 @@ async def update_experiment( name: str, body: ExperimentRequest, entity_client: EntityClientDep, + rollup_repository: ExperimentRollupRepositoryDep, ) -> ExperimentResponse: - try: - existing = await entity_client.get(Experiment, name=name, workspace=workspace) - except EntityNotFoundError as e: - raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND, - detail=f"Experiment '{workspace}/{name}' not found.", - ) from e + existing = await _get_or_404( + entity_client, + Experiment, + workspace=workspace, + name=name, + label="Experiment", + ) changed = [f for f in _IMMUTABLE_EXPERIMENT_FIELDS if getattr(body, f) != getattr(existing, f)] if changed: @@ -320,7 +343,9 @@ async def update_experiment( existing.description = body.description existing.summary = body.summary updated = await entity_client.update(existing) - return ExperimentResponse.from_entity(updated) + response = ExperimentResponse.from_entity(updated) + await _hydrate_rollups(workspace=workspace, responses=[response], rollup_repository=rollup_repository) + return response @router.delete( @@ -334,10 +359,86 @@ async def delete_experiment( name: str, entity_client: EntityClientDep, ) -> None: + await _delete_or_404( + entity_client, + Experiment, + workspace=workspace, + name=name, + label="Experiment", + ) + + +async def _get_or_404( + entity_client: EntityClient, + entity_type: type[EntityT], + *, + workspace: str, + name: str, + label: str, +) -> EntityT: try: - await entity_client.delete(Experiment, name=name, workspace=workspace) + return await entity_client.get(entity_type, name=name, workspace=workspace) except EntityNotFoundError as e: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, - detail=f"Experiment '{workspace}/{name}' not found.", + detail=f"{label} '{workspace}/{name}' not found.", ) from e + + +async def _delete_or_404( + entity_client: EntityClient, + entity_type: type[EntityT], + *, + workspace: str, + name: str, + label: str, +) -> None: + try: + await entity_client.delete(entity_type, name=name, workspace=workspace) + except EntityNotFoundError as e: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, + detail=f"{label} '{workspace}/{name}' not found.", + ) from e + + +async def _hydrate_rollups( + *, + workspace: str, + responses: list[ExperimentResponse], + rollup_repository: ExperimentRollupRepository | None, +) -> None: + if rollup_repository is None or not responses: + return + try: + rollups = await rollup_repository.get_rollups( + workspace=workspace, experiment_ids=[response.name for response in responses] + ) + except Exception: + logger.exception("Skipping experiment rollup hydration because ClickHouse is unavailable") + return + for response in responses: + rollup = rollups.get(response.name) + if rollup is not None: + _apply_rollup(response, rollup) + + +def _apply_rollup(response: ExperimentResponse, rollup: ExperimentRollup) -> None: + response.evaluator_names = rollup.evaluator_names + response.model_names = rollup.model_names + response.aggregate_scores = {name: _aggregate(score) for name, score in rollup.evaluator_scores.items()} or None + response.run_count = rollup.run_count + response.cost_usd = _aggregate(rollup.cost_usd) if rollup.cost_usd is not None else None + response.latency_ms = _aggregate(rollup.latency_ms) if rollup.latency_ms is not None else None + + +def _aggregate(rollup: ScoreRollup) -> EvaluatorAggregate: + return EvaluatorAggregate( + sum=rollup.sum, + mean=rollup.mean, + median=rollup.median, + p90=rollup.p90, + p95=rollup.p95, + p99=rollup.p99, + count=rollup.count, + ) diff --git a/services/intake/src/nmp/intake/api/v2/experiments/schemas.py b/services/intake/src/nmp/intake/api/v2/experiments/schemas.py index d342594ceb..4ae6470369 100644 --- a/services/intake/src/nmp/intake/api/v2/experiments/schemas.py +++ b/services/intake/src/nmp/intake/api/v2/experiments/schemas.py @@ -3,10 +3,8 @@ """Request, response, and filter schemas for the Experiments API. -Response models are standalone (not entity subclasses): they translate from the -stored entity via ``from_entity`` and carry rollup fields that are hydrated from -ClickHouse at read time. In this PR the rollups are always defaults; the -hydration path lands in a later PR. +Response models are standalone: they translate from the stored entity via +``from_entity`` and carry rollup fields hydrated from ClickHouse at read time. """ from __future__ import annotations @@ -18,10 +16,6 @@ from nmp.intake.entities.experiments import Experiment, ExperimentGroup from pydantic import AnyUrl, BaseModel, ConfigDict, Field -# ============================================================================= -# Requests (workspace comes from the route parameter) -# ============================================================================= - class ExperimentGroupRequest(BaseModel): """Request body for creating an ExperimentGroup.""" @@ -52,11 +46,6 @@ class ExperimentRequest(BaseModel): summary: str | None = Field(default=None, description="Human-authored summary of results.") -# ============================================================================= -# Responses -# ============================================================================= - - class ExperimentGroupResponse(BaseModel): """ExperimentGroup as served by the API.""" @@ -80,7 +69,7 @@ def from_entity(cls, entity: ExperimentGroup) -> ExperimentGroupResponse: class EvaluatorAggregate(BaseModel): - """Cross-run statistics for one evaluator. Populated by the rollup path (later PR).""" + """Aggregate statistics over evaluator scores or session-level metric values.""" sum: float | None = None mean: float | None = None @@ -88,6 +77,7 @@ class EvaluatorAggregate(BaseModel): p90: float | None = None p95: float | None = None p99: float | None = None + count: int = 0 class ExperimentResponse(BaseModel): @@ -111,11 +101,19 @@ class ExperimentResponse(BaseModel): created_at: datetime | None = None updated_at: datetime | None = None - # Hydrated from ClickHouse at read time in a later PR; defaults until then. evaluator_names: list[str] = Field(default_factory=list) - model_names: list[str] = Field(default_factory=list) + model_names: list[str] = Field( + default_factory=list, + description="Distinct model names observed across ingested sessions for this experiment.", + json_schema_extra={"uniqueItems": True}, + ) aggregate_scores: dict[str, EvaluatorAggregate] | None = None - run_count: int = 0 + run_count: int = Field( + default=0, + description="Number of distinct ingested experiment sessions; one session is treated as one run.", + ) + cost_usd: EvaluatorAggregate | None = None + latency_ms: EvaluatorAggregate | None = None @classmethod def from_entity(cls, entity: Experiment) -> ExperimentResponse: @@ -137,11 +135,6 @@ def from_entity(cls, entity: Experiment) -> ExperimentResponse: ) -# ============================================================================= -# List filters (declarative; the entity store applies them) -# ============================================================================= - - class ExperimentGroupFilter(Filter): """Filter for listing ExperimentGroups.""" diff --git a/services/intake/src/nmp/intake/entities/experiments.py b/services/intake/src/nmp/intake/entities/experiments.py index 3725030c4d..a51815a4eb 100644 --- a/services/intake/src/nmp/intake/entities/experiments.py +++ b/services/intake/src/nmp/intake/entities/experiments.py @@ -3,14 +3,9 @@ """Experiment and ExperimentGroup entity definitions for the Intake service. -These are entity-store (Postgres) entities, distinct from the ClickHouse-backed -telemetry (spans, evaluator_results). They hold the durable, producer-supplied -metadata that organizes telemetry into leaderboard-shaped views. - -Cross-run rollups (per-evaluator aggregate scores, run count, and the unions of -evaluator/model names) are intentionally *not* stored here. They are derived from -ClickHouse and hydrated onto the read model at query time; see -``nmp.intake.api.v2.experiments.schemas.ExperimentResponse``. +These are entity-store rows, distinct from ClickHouse telemetry. They hold the +durable, producer-supplied metadata that organizes telemetry into leaderboard +views. Rollups are derived from ClickHouse at read time. """ from __future__ import annotations @@ -35,8 +30,7 @@ class ExperimentGroup(EntityBase): class Experiment(EntityBase): """A single agent/config run against a dataset: one row on a leaderboard. - ``name`` is the producer-supplied, workspace-unique experiment id (e.g. - ``"terminal-bench-2_claude-code_opus_baseline"``); create is keyed on it. + ``name`` is the producer-supplied, workspace-unique experiment id. """ __entity_type__: ClassVar[str] = "experiment" diff --git a/services/intake/src/nmp/intake/service.py b/services/intake/src/nmp/intake/service.py index 26b5ee4c67..b529177eb4 100644 --- a/services/intake/src/nmp/intake/service.py +++ b/services/intake/src/nmp/intake/service.py @@ -51,6 +51,11 @@ def get_routers(self) -> List[RouterConfig]: tag="Annotations", description="Post-hoc annotation endpoints (feedback, labels, notes, metadata)", ), + RouterConfig( + experiments.router, + tag="Experiments", + description="Create, list, get, and delete Experiments and Experiment Groups", + ), RouterConfig(otlp.router, tag="Ingest", description="OTLP/HTTP trace ingest endpoints"), RouterConfig(atif.router, tag="Ingest", description="ATIF trajectory ingest endpoints"), RouterConfig( @@ -58,11 +63,6 @@ def get_routers(self) -> List[RouterConfig]: tag="Ingest", description="OpenAI-compatible chat-completion ingest endpoint", ), - RouterConfig( - experiments.router, - tag="Experiments", - description="Create, list, get, and delete Experiments and Experiment Groups", - ), ] async def on_startup(self) -> None: diff --git a/services/intake/src/nmp/intake/spans/clickhouse_migrations.py b/services/intake/src/nmp/intake/spans/clickhouse_migrations.py index 0ed5fe0e15..3a0a33c4d2 100644 --- a/services/intake/src/nmp/intake/spans/clickhouse_migrations.py +++ b/services/intake/src/nmp/intake/spans/clickhouse_migrations.py @@ -239,11 +239,90 @@ def _add_evaluator_results_skip_indexes(client, settings: ClickHouseMigrationSet client.command(f"ALTER TABLE {table} MATERIALIZE INDEX idx_created_at") +def _create_experiment_sessions_schema(client, settings: ClickHouseMigrationSettings) -> None: + """Create the root-span experiment membership table and insert-time projection.""" + + table = _table(settings, "experiment_sessions") + view = _table(settings, "experiment_sessions_mv") + client.command(f"DROP TABLE IF EXISTS {view}") + client.command(f"DROP TABLE IF EXISTS {table}") + + # Note this is logically a single table. CH requires creating an underlying table and then a view that writes to that table. + client.command( + f""" + CREATE TABLE {table} + ( + workspace LowCardinality(String), + experiment_id String, + session_id String, + test_case_id String DEFAULT '', + + source_format LowCardinality(String), + trace_id String, + root_span_id String, + + start_time DateTime64(6) CODEC(Delta(8), ZSTD(1)), + end_time Nullable(DateTime64(6)) CODEC(Delta(8), ZSTD(1)), + latency_ms Nullable(Float64), + + event_ts DateTime64(6), + is_deleted UInt8 DEFAULT 0 + ) + ENGINE = ReplacingMergeTree(event_ts, is_deleted) + PARTITION BY toYYYYMM(start_time) + PRIMARY KEY (workspace, experiment_id, session_id) + ORDER BY (workspace, experiment_id, session_id, root_span_id) + TTL toDate(start_time) + INTERVAL 90 DAY + SETTINGS + index_granularity = 256, + ttl_only_drop_parts = 1 + """ + ) + experiment_sessions_select_sql = f""" + SELECT + workspace, + attributes_string['experiment.id'] AS experiment_id, + session_id, + if( + has(mapKeys(attributes_string), 'test_case.id'), + attributes_string['test_case.id'], + '' + ) AS test_case_id, + source_format, + trace_id, + external_span_id AS root_span_id, + start_time, + nullIf(end_time, toDateTime64(0, 6)) AS end_time, + if(end_time = toDateTime64(0, 6), NULL, dateDiff('millisecond', start_time, end_time)) AS latency_ms, + event_ts, + is_deleted + FROM {_table(settings, "spans")} + WHERE external_parent_span_id = '' + AND has(mapKeys(attributes_string), 'experiment.id') + AND attributes_string['experiment.id'] != '' + """ + client.command( + f""" + CREATE MATERIALIZED VIEW {view} + TO {table} + AS + {experiment_sessions_select_sql} + """ + ) + client.command( + f""" + INSERT INTO {table} + {experiment_sessions_select_sql} + """ + ) + + _MIGRATIONS: list[tuple[str, Callable[..., None]]] = [ ("ch_spans_0002", _create_spans_schema), ("ch_evaluator_results_0001", _create_evaluator_results_schema), ("ch_annotations_0001", _create_annotations_schema), ("ch_evaluator_results_0002", _add_evaluator_results_skip_indexes), + ("ch_experiment_sessions_0002", _create_experiment_sessions_schema), ] CURRENT_SCHEMA_VERSION = _MIGRATIONS[-1][0] diff --git a/services/intake/src/nmp/intake/spans/experiment_rollup_repository.py b/services/intake/src/nmp/intake/spans/experiment_rollup_repository.py new file mode 100644 index 0000000000..61fae0320a --- /dev/null +++ b/services/intake/src/nmp/intake/spans/experiment_rollup_repository.py @@ -0,0 +1,293 @@ +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +"""ClickHouse rollups for Experiment read models.""" + +from __future__ import annotations + +from dataclasses import dataclass, field +from typing import Any + +from nmp.intake.spans.clickhouse_client import ClickHouseSpanClient +from nmp.intake.spans.span_attribute_catalog import COST_SCALE, SpanAttributeField, spec_for_field +from nmp.intake.spans.storage import result_rows +from nmp.intake.spans.trace_repository import current_spans_sql + + +@dataclass(frozen=True) +class ScoreRollup: + sum: float | None + mean: float | None + median: float | None + p90: float | None + p95: float | None + p99: float | None + count: int + + +@dataclass +class ExperimentRollup: + experiment_id: str + run_count: int = 0 + model_names: list[str] = field(default_factory=list) + evaluator_scores: dict[str, ScoreRollup] = field(default_factory=dict) + cost_usd: ScoreRollup | None = None + latency_ms: ScoreRollup | None = None + + @property + def evaluator_names(self) -> list[str]: + return sorted(self.evaluator_scores) + + +class ExperimentRollupRepository: + def __init__(self, client: ClickHouseSpanClient) -> None: + self._client = client + + async def get_rollups(self, *, workspace: str, experiment_ids: list[str]) -> dict[str, ExperimentRollup]: + experiment_ids = list(dict.fromkeys(experiment_ids)) + rollups = {experiment_id: ExperimentRollup(experiment_id=experiment_id) for experiment_id in experiment_ids} + if not experiment_ids: + return rollups + + experiment_names_sql, experiment_parameters = _experiment_id_parameters(experiment_ids) + parameters = {"workspace": workspace, **experiment_parameters} + sessions_table = self._client.table("experiment_sessions") + + for row in result_rows( + await self._client.query( + _run_counts_sql(sessions_table, experiment_names_sql), + parameters=parameters, + ) + ): + rollups[row["experiment_id"]].run_count = int(row["run_count"]) + + for row in result_rows( + await self._client.query( + _score_rollups_sql( + sessions_table=sessions_table, + evaluator_results_table=self._client.table("evaluator_results"), + experiment_names_sql=experiment_names_sql, + ), + parameters=parameters, + ) + ): + rollups[row["experiment_id"]].evaluator_scores[row["evaluator_name"]] = ScoreRollup( + sum=_float_or_none(row["sum"]), + mean=_float_or_none(row["mean"]), + median=_float_or_none(row["median"]), + p90=_float_or_none(row["p90"]), + p95=_float_or_none(row["p95"]), + p99=_float_or_none(row["p99"]), + count=int(row["count"]), + ) + + for row in result_rows( + await self._client.query( + _metric_rollups_sql( + sessions_table=sessions_table, + spans_table=self._client.table("spans"), + experiment_names_sql=experiment_names_sql, + ), + parameters={ + **parameters, + "cost_key": spec_for_field(SpanAttributeField.COST_TOTAL_USD).bag_key, + "model_key": spec_for_field(SpanAttributeField.MODEL).bag_key, + }, + ) + ): + rollup = rollups[row["experiment_id"]] + rollup.model_names = _string_list(row["model_names"]) + rollup.cost_usd = _score_rollup(row, "cost") + rollup.latency_ms = _score_rollup(row, "latency") + + return rollups + + +def _experiment_id_parameters(experiment_ids: list[str]) -> tuple[str, dict[str, str]]: + parameters = {f"experiment_id_{index}": experiment_id for index, experiment_id in enumerate(experiment_ids)} + return ", ".join(f"%({name})s" for name in parameters), parameters + + +def _scoped_sessions_sql(sessions_table: str, experiment_names_sql: str) -> str: + return f""" + SELECT workspace, experiment_id, session_id, latency_ms + FROM {sessions_table} FINAL + WHERE workspace = %(workspace)s + AND is_deleted = 0 + AND experiment_id IN ({experiment_names_sql}) + ORDER BY start_time ASC, root_span_id ASC + LIMIT 1 BY workspace, session_id, experiment_id + """ + + +def _run_counts_sql(sessions_table: str, experiment_names_sql: str) -> str: + return f""" + WITH scoped_sessions AS ( + {_scoped_sessions_sql(sessions_table, experiment_names_sql)} + ) + SELECT + experiment_id, + count() AS run_count + FROM scoped_sessions + GROUP BY experiment_id + ORDER BY experiment_id ASC + """ + + +# Quantile name -> ClickHouse probability argument, shared by every distribution rollup. +_STAT_QUANTILES = {"median": "0.5", "p90": "0.9", "p95": "0.95", "p99": "0.99"} + + +def _stat_columns(value_expr: str, *, prefix: str = "", guarded: bool = False) -> str: + """Build the sum/mean/median/p90/p95/p99/count column list for one value expression. + + ``guarded`` wraps each aggregate so that an empty or all-NULL set yields NULL instead of + ``sumIf``'s 0 / ``avgIf``'s NaN; use it when ``value_expr`` can be NULL per input row. + """ + + label = f"{prefix}_" if prefix else "" + if guarded: + not_null = f"isNotNull({value_expr})" + guard = f"countIf({not_null}) = 0" + + def stat(expr: str) -> str: + return f"if({guard}, NULL, {expr})" + + columns = [ + f"{stat(f'sumIf({value_expr}, {not_null})')} AS {label}sum", + f"{stat(f'avgIf({value_expr}, {not_null})')} AS {label}mean", + *( + f"{stat(f'quantileExactIf({q})({value_expr}, {not_null})')} AS {label}{name}" + for name, q in _STAT_QUANTILES.items() + ), + f"countIf({not_null}) AS {label}count", + ] + else: + columns = [ + f"sum({value_expr}) AS {label}sum", + f"avg({value_expr}) AS {label}mean", + *(f"quantileExact({q})({value_expr}) AS {label}{name}" for name, q in _STAT_QUANTILES.items()), + f"count() AS {label}count", + ] + return ",\n ".join(columns) + + +def _score_rollups_sql(*, sessions_table: str, evaluator_results_table: str, experiment_names_sql: str) -> str: + # Each run (session) contributes one score per evaluator, so reduce the per-span + # evaluator_results rows to a single per-(experiment, session, evaluator) value before + # the distribution rollup. This keeps `count` aligned with run_count and the mean + # run-weighted rather than weighted by spans-per-session. + return f""" + WITH + scoped_sessions AS ( + {_scoped_sessions_sql(sessions_table, experiment_names_sql)} + ), + session_scores AS ( + SELECT + sessions.experiment_id AS experiment_id, + results.name AS evaluator_name, + avg(results.value) AS value + FROM scoped_sessions AS sessions + INNER JOIN ( + SELECT workspace, session_id, name, value + FROM {evaluator_results_table} FINAL + WHERE workspace = %(workspace)s + AND (workspace, session_id) IN ( + SELECT DISTINCT workspace, session_id + FROM scoped_sessions + ) + AND data_type IN ('NUMERIC', 'BOOLEAN') + AND value IS NOT NULL + ) AS results + ON sessions.workspace = results.workspace + AND sessions.session_id = results.session_id + GROUP BY sessions.experiment_id, sessions.session_id, results.name + ) + SELECT + experiment_id, + evaluator_name, + {_stat_columns("value")} + FROM session_scores + GROUP BY experiment_id, evaluator_name + ORDER BY experiment_id ASC, evaluator_name ASC + """ + + +def _metric_rollups_sql(*, sessions_table: str, spans_table: str, experiment_names_sql: str) -> str: + return f""" + WITH + scoped_sessions AS ( + {_scoped_sessions_sql(sessions_table, experiment_names_sql)} + ), + current_session_spans AS ( + { + current_spans_sql( + spans_table, + extra_where_sql=( + "(span_versions.workspace, span_versions.session_id) IN " + "(SELECT DISTINCT workspace, session_id FROM scoped_sessions)" + ), + ) + } + ), + session_costs AS ( + SELECT + sessions.experiment_id AS experiment_id, + sessions.session_id AS session_id, + sessions.latency_ms AS latency_ms, + if( + countIf(has(mapKeys(spans.attributes_number), %(cost_key)s)) = 0, + NULL, + sumIf( + spans.attributes_number[%(cost_key)s], + has(mapKeys(spans.attributes_number), %(cost_key)s) + ) / {COST_SCALE} + ) AS cost_usd, + groupUniqArrayIf( + spans.attributes_string[%(model_key)s], + has(mapKeys(spans.attributes_string), %(model_key)s) + AND spans.attributes_string[%(model_key)s] != '' + ) AS model_names + FROM scoped_sessions AS sessions + LEFT JOIN current_session_spans AS spans + ON sessions.workspace = spans.workspace + AND sessions.session_id = spans.session_id + AND spans.is_deleted = 0 + GROUP BY sessions.experiment_id, sessions.session_id, sessions.latency_ms + ) + SELECT + experiment_id, + arraySort(arrayDistinct(arrayFlatten(groupArray(model_names)))) AS model_names, + {_stat_columns("cost_usd", prefix="cost", guarded=True)}, + {_stat_columns("latency_ms", prefix="latency", guarded=True)} + FROM session_costs + GROUP BY experiment_id + ORDER BY experiment_id ASC + """ + + +def _score_rollup(row: dict[str, Any], prefix: str) -> ScoreRollup | None: + count = int(row[f"{prefix}_count"]) + if count == 0: + return None + return ScoreRollup( + sum=_float_or_none(row[f"{prefix}_sum"]), + mean=_float_or_none(row[f"{prefix}_mean"]), + median=_float_or_none(row[f"{prefix}_median"]), + p90=_float_or_none(row[f"{prefix}_p90"]), + p95=_float_or_none(row[f"{prefix}_p95"]), + p99=_float_or_none(row[f"{prefix}_p99"]), + count=count, + ) + + +def _string_list(value: Any) -> list[str]: + if value is None: + return [] + return sorted(str(item) for item in value if str(item)) + + +def _float_or_none(value: Any) -> float | None: + if value is None: + return None + return float(value) diff --git a/services/intake/src/nmp/intake/spans/ingest/atif.py b/services/intake/src/nmp/intake/spans/ingest/atif.py index b5bbdd0b15..81106d4b55 100644 --- a/services/intake/src/nmp/intake/spans/ingest/atif.py +++ b/services/intake/src/nmp/intake/spans/ingest/atif.py @@ -5,9 +5,11 @@ from __future__ import annotations -from typing import Any +from typing import Annotated, Any from fastapi import APIRouter, Depends, Response, status +from nmp.common.entities.client import EntityClient +from nmp.common.service.dependencies import get_entity_client from nmp.intake.spans.api.dependencies import SpansServiceDep, require_workspace_access from nmp.intake.spans.domain import TraceBatch from nmp.intake.spans.ingest.atif_domain import ( @@ -20,26 +22,27 @@ validate_atif_tool_call_references, ) from nmp.intake.spans.ingest.atif_mapping import trajectory_to_evaluator_results, trajectory_to_spans -from nmp.intake.spans.ingest.evaluation_context import EvaluationContext +from nmp.intake.spans.ingest.evaluation_context import ExperimentContextIngestModel +from nmp.intake.spans.ingest.experiment_context_validation import validate_experiment_context from nmp.intake.spans.storage import utc_now -from pydantic import BaseModel, ConfigDict, Field, model_validator +from pydantic import ConfigDict, Field, model_validator router = APIRouter(dependencies=[Depends(require_workspace_access)]) API_TAG = "Ingest" +EntityClientDep = Annotated[EntityClient, Depends(get_entity_client)] -class AtifIngestRequest(BaseModel): +class AtifIngestRequest(ExperimentContextIngestModel): """Span-based ATIF ingest request. ATIF project scoping is intentionally not accepted here; use the workspace - route and ``evaluation_context`` for evaluation/run identity. + route and ``experiment_context`` for experiment identity. """ model_config = ConfigDict(extra="forbid") schema_version: AtifSchemaVersion session_id: str | None = None - evaluation_context: EvaluationContext | None = None agent: AtifAgent final_metrics: AtifFinalMetrics | None = None continued_trajectory_ref: str | None = None @@ -63,7 +66,7 @@ def to_trajectory(self) -> AtifTrajectory: continued_trajectory_ref=self.continued_trajectory_ref, notes=self.notes, extra=self.extra, - evaluation_context=self.evaluation_context, + evaluation_context=self.resolved_evaluation_context(), **kwargs, ) @@ -78,7 +81,13 @@ async def ingest_atif( workspace: str, body: AtifIngestRequest, service: SpansServiceDep, + entity_client: EntityClientDep, ) -> Response: + await validate_experiment_context( + workspace=workspace, + context=body.resolved_evaluation_context(), + entity_client=entity_client, + ) ingested_at = utc_now() trajectory = body.to_trajectory() spans = trajectory_to_spans( diff --git a/services/intake/src/nmp/intake/spans/ingest/atif_mapping.py b/services/intake/src/nmp/intake/spans/ingest/atif_mapping.py index a638b17401..1d9a5e74bd 100644 --- a/services/intake/src/nmp/intake/spans/ingest/atif_mapping.py +++ b/services/intake/src/nmp/intake/spans/ingest/atif_mapping.py @@ -445,8 +445,8 @@ def _span_attributes( cost_total_usd=cost_total_usd, ) attribute_bags = semantic_attributes.to_bags() - if evaluation_context is not None: - attribute_bags.put_json("evaluation.metadata", evaluation_context.metadata) + if evaluation_context is not None and evaluation_context.metadata: + attribute_bags.put_json("experiment.metadata", evaluation_context.metadata) if raw_attributes is not None: attribute_bags.put_json("atif.raw", raw_attributes) return attribute_bags diff --git a/services/intake/src/nmp/intake/spans/ingest/chat_completions.py b/services/intake/src/nmp/intake/spans/ingest/chat_completions.py index a25bcd7b1b..967bf84181 100644 --- a/services/intake/src/nmp/intake/spans/ingest/chat_completions.py +++ b/services/intake/src/nmp/intake/spans/ingest/chat_completions.py @@ -17,9 +17,14 @@ from typing import Annotated, Any, Self from fastapi import APIRouter, Depends, status +from nmp.common.entities.client import EntityClient +from nmp.common.service.dependencies import get_entity_client from nmp.intake.spans.api.dependencies import SpansServiceDep, require_workspace_access from nmp.intake.spans.domain import IntakeSpan, SpanKind, SpanStatus, TraceBatch -from nmp.intake.spans.ingest.evaluation_context import EvaluationContext +from nmp.intake.spans.ingest.evaluation_context import ( + ExperimentContextIngestModel, +) +from nmp.intake.spans.ingest.experiment_context_validation import validate_experiment_context from nmp.intake.spans.span_attribute_bags import SpanAttributeBags from nmp.intake.spans.span_semantic_attributes import SpanSemanticAttributes from nmp.intake.spans.storage import json_dumps_preserve, stable_id, utc_now @@ -27,6 +32,7 @@ router = APIRouter(dependencies=[Depends(require_workspace_access)]) API_TAG = "Ingest" +EntityClientDep = Annotated[EntityClient, Depends(get_entity_client)] SOURCE_FORMAT = "chat_completions" NonNegativeFloat = Annotated[float, Field(ge=0)] @@ -91,7 +97,7 @@ def _require_choices_or_error(self) -> Self: return self -class ChatCompletionsIngestRequest(BaseModel): +class ChatCompletionsIngestRequest(ExperimentContextIngestModel): model_config = ConfigDict(extra="forbid") request: CapturedChatCompletionsRequest @@ -108,7 +114,6 @@ class ChatCompletionsIngestRequest(BaseModel): "chat-completions calls; use session_id to group related calls." ), ) - evaluation_context: EvaluationContext | None = None provider: str | None = None cost_usd: NonNegativeFloat | None = Field( default=None, @@ -146,7 +151,13 @@ async def ingest_chat_completion( workspace: str, body: ChatCompletionsIngestRequest, service: SpansServiceDep, + entity_client: EntityClientDep, ) -> ChatCompletionsIngestResponse: + await validate_experiment_context( + workspace=workspace, + context=body.resolved_evaluation_context(), + entity_client=entity_client, + ) ingested_at = utc_now() span = _chat_completion_to_span(workspace=workspace, body=body, ingested_at=ingested_at) await service.ingest_batch(TraceBatch(spans=[span])) @@ -217,7 +228,7 @@ def _build_attribute_bags( error_type = _as_str(error.get("type")) or _as_str(error.get("code")) error_message = _as_str(error.get("message")) - evaluation_context = body.evaluation_context + evaluation_context = body.resolved_evaluation_context() semantic = SpanSemanticAttributes( model=_as_str(response.get("model")) or _as_str(request.get("model")), provider=body.provider or _infer_provider(response), @@ -242,13 +253,13 @@ def _build_attribute_bags( cost_output_usd=_decimal_or_none(body.cost_output_usd), ) attribute_bags = semantic.to_bags() + if evaluation_context is not None and evaluation_context.metadata: + attribute_bags.put_json("experiment.metadata", evaluation_context.metadata) for key, value in body.cost_details.items(): bag_key = f"cost.{key}" if bag_key in attribute_bags.number: continue attribute_bags.put_unhandled_source_attribute(f"llm.cost.{key}", value) - if evaluation_context is not None: - attribute_bags.put_json("evaluation.metadata", evaluation_context.metadata) return attribute_bags diff --git a/services/intake/src/nmp/intake/spans/ingest/evaluation_context.py b/services/intake/src/nmp/intake/spans/ingest/evaluation_context.py index c16253807a..0b1b49510b 100644 --- a/services/intake/src/nmp/intake/spans/ingest/evaluation_context.py +++ b/services/intake/src/nmp/intake/spans/ingest/evaluation_context.py @@ -1,16 +1,32 @@ # SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. # SPDX-License-Identifier: Apache-2.0 -"""Shared evaluation context model for span ingest endpoints.""" +"""Shared experiment/evaluation context models for span ingest endpoints.""" from __future__ import annotations from typing import Any -from pydantic import BaseModel, ConfigDict, Field, model_validator +from pydantic import BaseModel, ConfigDict, Field + + +class ExperimentContext(BaseModel): + """Experiment context accepted by ingest endpoints.""" + + experiment_id: str = Field(description="Name of an existing Experiment entity.") + test_case_id: str | None = Field(default=None, description="Optional producer-supplied test case id.") + + model_config = ConfigDict(extra="forbid") + + def to_evaluation_context(self) -> EvaluationContext: + return EvaluationContext( + evaluation_id=self.experiment_id, + test_case_id=self.test_case_id, + ) class EvaluationContext(BaseModel): + # Deprecated pre-release ingest shape. Use ExperimentContext / experiment_context. evaluation_id: str | None = None evaluation_sha: str | None = None evaluation_run_id: str | None = None @@ -22,21 +38,20 @@ class EvaluationContext(BaseModel): model_config = ConfigDict(extra="forbid") - @model_validator(mode="after") - def require_run_id_when_context_is_set(self) -> EvaluationContext: - if self.evaluation_run_id is None and self._has_values(): - raise ValueError("evaluation_context.evaluation_run_id is required when evaluation_context fields are set") - return self - - def _has_values(self) -> bool: - return any( - ( - self.evaluation_id, - self.evaluation_sha, - self.dataset_id, - self.dataset_name, - self.dataset_version, - self.test_case_id, - self.metadata, - ) - ) + +class ExperimentContextIngestModel(BaseModel): + """Base model for ingest payloads that accept experiment context.""" + + experiment_context: ExperimentContext | None = None + # Deprecated pre-release field. Use experiment_context. + evaluation_context: EvaluationContext | None = Field( + default=None, + deprecated=True, + description="Deprecated. Use experiment_context; when both are sent, experiment_context takes precedence.", + ) + + def resolved_evaluation_context(self) -> EvaluationContext | None: + if self.experiment_context is not None: + return self.experiment_context.to_evaluation_context() + evaluation_context: EvaluationContext | None = self.__dict__.get("evaluation_context") + return evaluation_context diff --git a/services/intake/src/nmp/intake/spans/ingest/experiment_context_validation.py b/services/intake/src/nmp/intake/spans/ingest/experiment_context_validation.py new file mode 100644 index 0000000000..753ed85910 --- /dev/null +++ b/services/intake/src/nmp/intake/spans/ingest/experiment_context_validation.py @@ -0,0 +1,35 @@ +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +"""Validation helpers for experiment-scoped ingest.""" + +from fastapi import HTTPException, status +from nmp.common.entities.client import EntityClient, EntityNotFoundError +from nmp.intake.entities.experiments import Experiment +from nmp.intake.spans.ingest.evaluation_context import EvaluationContext, ExperimentContext + + +async def validate_experiment_context( + *, + workspace: str, + context: EvaluationContext | ExperimentContext | None, + entity_client: EntityClient, +) -> None: + if context is None: + return + experiment_id = _experiment_id(context) + if not experiment_id: + return + try: + await entity_client.get(Experiment, name=experiment_id, workspace=workspace) + except EntityNotFoundError as exc: + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail=f"Experiment '{experiment_id}' must be created before it can be logged.", + ) from exc + + +def _experiment_id(context: EvaluationContext | ExperimentContext) -> str | None: + if isinstance(context, ExperimentContext): + return context.experiment_id + return context.evaluation_id diff --git a/services/intake/src/nmp/intake/spans/span_attribute_bags.py b/services/intake/src/nmp/intake/spans/span_attribute_bags.py index 78fd150d65..016c17eac7 100644 --- a/services/intake/src/nmp/intake/spans/span_attribute_bags.py +++ b/services/intake/src/nmp/intake/spans/span_attribute_bags.py @@ -13,6 +13,7 @@ from nmp.intake.spans.span_attribute_catalog import ( COST_SCALE, KNOWN_BAG_KEYS, + LEGACY_BAG_KEY_ALIASES, SOURCE_ONLY_KEYS, SOURCE_ONLY_PREFIXES, AttributeBag, @@ -48,7 +49,14 @@ def from_domain_maps( def get_field(self, field: SpanAttributeField | str) -> str | int | float | bool | Decimal | None: spec = spec_for_field(field) bag = self._bag_for_spec(spec) - return from_bag(bag.get(spec.bag_key), spec) + value = from_bag(bag.get(spec.bag_key), spec) + if value is not None: + return value + for alias in LEGACY_BAG_KEY_ALIASES.get(spec.field, ()): + value = from_bag(bag.get(alias), spec) + if value is not None: + return value + return None def put_field(self, field: SpanAttributeField | str, value: Any) -> None: self.put_spec(spec_for_field(field), value) @@ -106,11 +114,11 @@ def raw_attributes_json(self) -> str | None: parsed_atif_raw = json.loads(atif_raw) if not isinstance(parsed_atif_raw, dict): raise TypeError("Expected atif.raw to contain a JSON object") - parsed_atif_raw.pop("evaluation.metadata", None) + parsed_atif_raw.pop("experiment.metadata", None) raw.update(parsed_atif_raw) for key, value in self.string.items(): - if key not in {"atif.raw", "evaluation.metadata"} and key not in KNOWN_BAG_KEYS: + if key not in {"atif.raw", "experiment.metadata"} and key not in KNOWN_BAG_KEYS: raw[key] = value for key, value in self.number.items(): if key not in KNOWN_BAG_KEYS: @@ -123,12 +131,12 @@ def raw_attributes_json(self) -> str | None: return json.dumps(raw, separators=(",", ":"), ensure_ascii=False) def evaluation_metadata(self) -> dict[str, Any] | None: - value = self.string.get("evaluation.metadata") + value = self.string.get("experiment.metadata") or self.string.get("evaluation.metadata") if value is None: return None parsed = json.loads(value) if not isinstance(parsed, dict): - raise TypeError("Expected evaluation.metadata to contain a JSON object") + raise TypeError("Expected experiment/evaluation metadata to contain a JSON object") return parsed def _bag_for_spec(self, spec: AttributeSpec) -> dict[str, str] | dict[str, float] | dict[str, bool]: diff --git a/services/intake/src/nmp/intake/spans/span_attribute_catalog.py b/services/intake/src/nmp/intake/spans/span_attribute_catalog.py index aac67eb964..e421ec612b 100644 --- a/services/intake/src/nmp/intake/spans/span_attribute_catalog.py +++ b/services/intake/src/nmp/intake/spans/span_attribute_catalog.py @@ -167,8 +167,11 @@ class AttributeSpec: AttributeSpec( field=SpanAttributeField.EVALUATION_ID, bag=AttributeBag.STRING, - bag_key="evaluation.id", + # Renamed from evaluation.id. + bag_key="experiment.id", source_keys=( + "experiment.id", + "experiment_id", "evaluation.id", "evaluation_id", ), @@ -176,8 +179,10 @@ class AttributeSpec: AttributeSpec( field=SpanAttributeField.EVALUATION_SHA, bag=AttributeBag.STRING, - bag_key="evaluation.sha", + bag_key="experiment.sha", source_keys=( + "experiment.sha", + "experiment_sha", "evaluation.sha", "evaluation_sha", ), @@ -185,8 +190,9 @@ class AttributeSpec: AttributeSpec( field=SpanAttributeField.EVALUATION_RUN_ID, bag=AttributeBag.STRING, - bag_key="evaluation.run_id", + bag_key="experiment.run_id", source_keys=( + "experiment.run_id", "evaluation.run_id", "evaluation_run_id", ), @@ -221,10 +227,11 @@ class AttributeSpec: AttributeSpec( field=SpanAttributeField.TEST_CASE_ID, bag=AttributeBag.STRING, - bag_key="dataset.test_case_id", + bag_key="test_case.id", source_keys=( - "dataset.test_case_id", + "test_case.id", "test_case_id", + "dataset.test_case_id", ), ), AttributeSpec( @@ -330,7 +337,16 @@ class AttributeSpec: SPECS_BY_FIELD = {spec.field: spec for spec in ATTRIBUTE_SPECS} SPECS_BY_FIELD_VALUE = {spec.field.value: spec for spec in ATTRIBUTE_SPECS} SPECS_BY_BAG_KEY = {spec.bag_key: spec for spec in ATTRIBUTE_SPECS} -KNOWN_BAG_KEYS = frozenset(SPECS_BY_BAG_KEY) +LEGACY_BAG_KEY_ALIASES = { + SpanAttributeField.EVALUATION_ID: ("evaluation.id",), + SpanAttributeField.EVALUATION_SHA: ("evaluation.sha",), + SpanAttributeField.EVALUATION_RUN_ID: ("evaluation.run_id",), + SpanAttributeField.TEST_CASE_ID: ("dataset.test_case_id",), +} +LEGACY_KNOWN_BAG_KEYS = frozenset(key for aliases in LEGACY_BAG_KEY_ALIASES.values() for key in aliases) | frozenset( + {"evaluation.metadata"} +) +KNOWN_BAG_KEYS = frozenset(SPECS_BY_BAG_KEY) | LEGACY_KNOWN_BAG_KEYS QUERYABLE_FIELDS = frozenset(SPECS_BY_FIELD_VALUE) # Source keys that populate top-level span fields or payloads rather than @@ -436,11 +452,24 @@ def where_clause( if bag_value is None: raise ValueError(f"Span attribute filter {field!r} does not support null values") - sql = ( - f"has(mapKeys({spec.bag.value}), %({key_param})s) " - f"AND {spec.bag.value}[%({key_param})s] {sql_operator} %({value_param})s" - ) - return sql, {key_param: spec.bag_key, value_param: bag_value} + alias_keys = LEGACY_BAG_KEY_ALIASES.get(spec.field, ()) + if not alias_keys: + sql = ( + f"has(mapKeys({spec.bag.value}), %({key_param})s) " + f"AND {spec.bag.value}[%({key_param})s] {sql_operator} %({value_param})s" + ) + return sql, {key_param: spec.bag_key, value_param: bag_value} + + clauses: list[str] = [] + parameters: dict[str, Any] = {value_param: bag_value} + for index, bag_key in enumerate((spec.bag_key, *alias_keys)): + aliased_key_param = key_param if index == 0 else f"{key_param}_{index}" + clauses.append( + f"(has(mapKeys({spec.bag.value}), %({aliased_key_param})s) " + f"AND {spec.bag.value}[%({aliased_key_param})s] {sql_operator} %({value_param})s)" + ) + parameters[aliased_key_param] = bag_key + return f"({' OR '.join(clauses)})", parameters def to_semantic_value(value: Any, spec: AttributeSpec) -> str | int | float | bool | Decimal | None: diff --git a/services/intake/src/nmp/intake/spans/trace_repository.py b/services/intake/src/nmp/intake/spans/trace_repository.py index 8356e03784..af226c9511 100644 --- a/services/intake/src/nmp/intake/spans/trace_repository.py +++ b/services/intake/src/nmp/intake/spans/trace_repository.py @@ -211,7 +211,7 @@ def _trace_summary_sql(table: str, filters: TraceListFilter) -> tuple[str, dict[ nullIf({root_alias}.end_time, {_ZERO_DATETIME_SQL}) AS ended_at, {root_alias}.attributes_string AS root_attributes_string, {root_alias}.event_ts AS ingested_at - FROM {_current_spans_sql(table)} AS {root_alias} + FROM {current_spans_sql(table)} AS {root_alias} WHERE {base_where_sql} ORDER BY {root_alias}.start_time ASC, {root_alias}.id ASC LIMIT 1 BY {root_alias}.workspace, {root_alias}.source_format, {root_alias}.trace_id @@ -228,7 +228,7 @@ def _trace_status_sql(table: str, filters: TraceListFilter) -> tuple[str, dict[s {source_alias}.source_format AS source_format, {source_alias}.trace_id AS trace_id, {_rolled_up_status_sql(source_alias)} AS status - FROM {_current_spans_sql(table)} AS {source_alias} + FROM {current_spans_sql(table)} AS {source_alias} WHERE {base_where_sql} GROUP BY {source_alias}.workspace, {source_alias}.source_format, {source_alias}.trace_id """ @@ -265,7 +265,7 @@ def _trace_aggregates_sql(table: str, filters: TraceListFilter) -> tuple[str, di )) AS providers, count() AS span_count, countIf({source_alias}.status = 'error') AS error_count - FROM {_current_spans_sql(table)} AS {source_alias} + FROM {current_spans_sql(table)} AS {source_alias} WHERE {base_where_sql} GROUP BY {source_alias}.workspace, {source_alias}.source_format, {source_alias}.trace_id """ @@ -381,14 +381,14 @@ def _candidate_subquery( return ( f""" SELECT workspace, source_format, trace_id - FROM {_current_spans_sql(table)} AS candidate_spans + FROM {current_spans_sql(table)} AS candidate_spans WHERE {" AND ".join(clauses)} """, parameters, ) -def _current_spans_sql(table: str) -> str: +def current_spans_sql(table: str, *, extra_where_sql: str | None = None) -> str: source_alias = "span_versions" columns = [ *[f"{source_alias}.{column} AS {column}" for column in _CURRENT_SPAN_IDENTITY_COLUMNS], @@ -399,12 +399,15 @@ def _current_spans_sql(table: str) -> str: ] columns_sql = ",\n ".join(columns) group_by_sql = ", ".join(f"{source_alias}.{column}" for column in _CURRENT_SPAN_IDENTITY_COLUMNS) + where_sql = f"{source_alias}.workspace = %(workspace)s" + if extra_where_sql is not None: + where_sql = f"{where_sql}\n AND {extra_where_sql}" return f""" ( SELECT {columns_sql} FROM {table} AS {source_alias} - WHERE {source_alias}.workspace = %(workspace)s + WHERE {where_sql} GROUP BY {group_by_sql} ) """ diff --git a/services/intake/tests/conftest.py b/services/intake/tests/conftest.py index e85d22ecc1..3844578f1c 100644 --- a/services/intake/tests/conftest.py +++ b/services/intake/tests/conftest.py @@ -5,6 +5,7 @@ import pytest from fastapi.testclient import TestClient +from nmp.intake.api.v2.experiments.endpoints import get_experiment_rollup_repository from nmp.intake.service import IntakeService from nmp.testing import create_test_client @@ -12,5 +13,9 @@ @pytest.fixture def client(): """Create test client with mocked entity client.""" - with create_test_client(IntakeService, client_type=TestClient) as tc: + with create_test_client( + IntakeService, + client_type=TestClient, + dependency_overrides={get_experiment_rollup_repository: lambda: None}, + ) as tc: yield tc diff --git a/services/intake/tests/integration/spans/conftest.py b/services/intake/tests/integration/spans/conftest.py index 70decb1bb4..602605d514 100644 --- a/services/intake/tests/integration/spans/conftest.py +++ b/services/intake/tests/integration/spans/conftest.py @@ -87,10 +87,10 @@ def clickhouse_client(clickhouse_settings: ClickHouseSettings): @pytest.fixture(autouse=True) def clean_clickhouse(clickhouse_client: ClickHouseSpanClient): - for table in ("spans", "evaluator_results"): + for table in ("spans", "evaluator_results", "experiment_sessions"): _run(clickhouse_client.command(f"TRUNCATE TABLE {clickhouse_client.table(table)}")) yield - for table in ("spans", "evaluator_results"): + for table in ("spans", "evaluator_results", "experiment_sessions"): _run(clickhouse_client.command(f"TRUNCATE TABLE {clickhouse_client.table(table)}")) diff --git a/services/intake/tests/integration/spans/test_atif_ingest.py b/services/intake/tests/integration/spans/test_atif_ingest.py index 3cde2bdd3a..58050ac85d 100644 --- a/services/intake/tests/integration/spans/test_atif_ingest.py +++ b/services/intake/tests/integration/spans/test_atif_ingest.py @@ -173,6 +173,7 @@ def test_atif_ingest_accepts_example_trajectory_and_reconstructs_read_side_data( "test_case_id": "sample-test-case-a", "metadata": {"trial": "sample-test-case-a__trial-a"}, } + _create_experiment(client, evaluation_context["evaluation_id"]) tool_call = { "tool_call_id": "tooluse_tuIapjh62ZTI1pildiC9sg", "function_name": "Bash", @@ -347,11 +348,21 @@ def test_atif_ingest_accepts_example_trajectory_and_reconstructs_read_side_data( assert "cached_tokens" not in trajectory assert "total_tokens" not in trajectory assert "cost_total_usd" not in trajectory - assert trajectory["evaluation_context"] == evaluation_context + assert trajectory["evaluation_context"] == { + "evaluation_id": evaluation_context["evaluation_id"], + "evaluation_sha": evaluation_context["evaluation_sha"], + "evaluation_run_id": evaluation_context["evaluation_run_id"], + "dataset_id": evaluation_context["dataset_id"], + "dataset_name": evaluation_context["dataset_name"], + "dataset_version": evaluation_context["dataset_version"], + "test_case_id": evaluation_context["test_case_id"], + "metadata": evaluation_context["metadata"], + } assert "attributes_string" not in trajectory trajectory_raw = json.loads(trajectory["raw_attributes"]) assert trajectory_raw["session_id"] == body["session_id"] assert "evaluation_context" not in trajectory_raw + assert "experiment.metadata" not in trajectory_raw assert "evaluation.metadata" not in trajectory_raw assert trajectory["started_at"] == "2026-05-04T18:57:59.943000" assert trajectory["ended_at"] == "2026-05-04T19:06:45.570079" @@ -363,7 +374,7 @@ def test_atif_ingest_accepts_example_trajectory_and_reconstructs_read_side_data( evaluation_response = client.get( "/apis/intake/v2/workspaces/default/spans", - params={"filter[evaluation_run_id]": evaluation_run_id, "page_size": 10}, + params={"filter[evaluation_id]": evaluation_context["evaluation_id"], "page_size": 10}, ) assert evaluation_response.status_code == 200, evaluation_response.text evaluation_spans = evaluation_response.json()["data"] @@ -372,11 +383,6 @@ def test_atif_ingest_accepts_example_trajectory_and_reconstructs_read_side_data( for field, value in { "evaluation_id": evaluation_context["evaluation_id"], - "evaluation_sha": evaluation_context["evaluation_sha"], - "evaluation_run_id": evaluation_context["evaluation_run_id"], - "dataset_id": evaluation_context["dataset_id"], - "dataset_name": evaluation_context["dataset_name"], - "dataset_version": evaluation_context["dataset_version"], "test_case_id": evaluation_context["test_case_id"], }.items(): filtered = client.get( @@ -484,13 +490,15 @@ def test_atif_ingest_accepts_example_trajectory_and_reconstructs_read_side_data( evaluation_roots_response = client.get( "/apis/intake/v2/workspaces/default/spans", - params={"filter[evaluation_run_id]": evaluation_run_id, "page_size": 20, "sort": "started_at"}, + params={"filter[evaluation_id]": evaluation_context["evaluation_id"], "page_size": 20, "sort": "started_at"}, ) assert evaluation_roots_response.status_code == 200, evaluation_roots_response.text evaluation_roots = evaluation_roots_response.json()["data"] assert len(evaluation_roots) == 2 assert {span["name"] for span in evaluation_roots} == {"sample-agent"} - assert {span["evaluation_context"]["evaluation_run_id"] for span in evaluation_roots} == {evaluation_run_id} + assert {span["evaluation_context"]["evaluation_id"] for span in evaluation_roots} == { + evaluation_context["evaluation_id"] + } assert {span["session_id"] for span in evaluation_roots} == { "d074dfb7-3691-443c-b137-720d75e40afa", "441e9149-e4e6-41c0-82b0-a36802f83d3a", @@ -513,10 +521,10 @@ def test_atif_ingest_accepts_example_trajectory_and_reconstructs_read_side_data( params={"filter[evaluation_run_id]": other_evaluation_run_id, "page_size": 10}, ) assert other_evaluation_response.status_code == 200, other_evaluation_response.text - other_spans = other_evaluation_response.json()["data"] - assert len(other_spans) == 1 - assert other_spans[0]["name"] == "sample-agent" - assert other_spans[0]["span_id"] == trajectory["span_id"] + other_evaluation_spans = other_evaluation_response.json()["data"] + assert len(other_evaluation_spans) == 1 + assert other_evaluation_spans[0]["name"] == "sample-agent" + assert other_evaluation_spans[0]["evaluation_context"]["evaluation_run_id"] == other_evaluation_run_id same_session_response = client.get( "/apis/intake/v2/workspaces/default/spans", @@ -640,3 +648,23 @@ def test_atif_trace_tokens_do_not_double_count_when_trajectory_and_steps_both_ca assert trace["output_tokens"] == 600 assert trace["total_tokens"] == 30600 assert trace["cost_usd"] == pytest.approx(0.45) + + +def _create_experiment(client: TestClient, name: str) -> str: + response = client.post( + "/apis/intake/v2/workspaces/default/experiments", + json={ + "name": name, + "agent_name": "sample-agent", + "agent_version": "1.0.0", + "dataset_name": "sample-dataset", + "dataset_version": "v1", + }, + ) + assert response.status_code in {201, 409}, response.text + if response.status_code == 201: + return response.json()["name"] + + existing = client.get(f"/apis/intake/v2/workspaces/default/experiments/{name}") + assert existing.status_code == 200, existing.text + return existing.json()["name"] diff --git a/services/intake/tests/integration/spans/test_chat_completions_ingest.py b/services/intake/tests/integration/spans/test_chat_completions_ingest.py index 552b58bf4e..5ca849d86a 100644 --- a/services/intake/tests/integration/spans/test_chat_completions_ingest.py +++ b/services/intake/tests/integration/spans/test_chat_completions_ingest.py @@ -14,7 +14,7 @@ INGEST_URL = "/apis/intake/v2/workspaces/default/ingest/chat-completions" SPANS_URL = "/apis/intake/v2/workspaces/default/spans" TRACES_URL = "/apis/intake/v2/workspaces/default/traces" -EVALUATION_CONTEXT = { +EVALUATION_CONTEXT: dict[str, Any] = { "evaluation_id": "chat-eval", "evaluation_sha": "chat-eval-sha", "evaluation_run_id": "evalrun-chat-001", @@ -24,6 +24,15 @@ "test_case_id": "chat-case-001", "metadata": {"source": "chat-completions-test"}, } +EXPERIMENT_CONTEXT: dict[str, Any] = { + "experiment_id": EVALUATION_CONTEXT["evaluation_id"], + "test_case_id": EVALUATION_CONTEXT["test_case_id"], +} +EXPECTED_CONTEXT_FROM_EXPERIMENT_CONTEXT: dict[str, Any] = { + "evaluation_id": EVALUATION_CONTEXT["evaluation_id"], + "test_case_id": EVALUATION_CONTEXT["test_case_id"], + "metadata": {}, +} def _openai_request(**overrides: Any) -> dict[str, Any]: @@ -72,11 +81,12 @@ def _openai_response(**overrides: Any) -> dict[str, Any]: def test_chat_completions_ingest_happy_path(client: TestClient): + experiment_id = _create_experiment(client, EXPERIMENT_CONTEXT["experiment_id"]) body = { "request": _openai_request(), "response": _openai_response(), "session_id": "session-happy", - "evaluation_context": EVALUATION_CONTEXT, + "experiment_context": {**EXPERIMENT_CONTEXT, "experiment_id": experiment_id}, "provider": "openai", } response = client.post(INGEST_URL, json=body) @@ -99,7 +109,7 @@ def test_chat_completions_ingest_happy_path(client: TestClient): assert span["name"] == "gpt-4o-mini-2024-08-06" assert span["model"] == "gpt-4o-mini-2024-08-06" assert span["provider"] == "openai" - assert span["evaluation_context"] == EVALUATION_CONTEXT + assert span["evaluation_context"] == EXPECTED_CONTEXT_FROM_EXPERIMENT_CONTEXT assert "evaluation_run_id" not in span assert "raw_attributes" not in span assert span["input_tokens"] == 24 @@ -112,7 +122,7 @@ def test_chat_completions_ingest_happy_path(client: TestClient): filtered = client.get( SPANS_URL, - params={"filter[evaluation_run_id]": EVALUATION_CONTEXT["evaluation_run_id"], "page_size": 10}, + params={"filter[evaluation_id]": EVALUATION_CONTEXT["evaluation_id"], "page_size": 10}, ) assert filtered.status_code == 200, filtered.text filtered_spans = filtered.json()["data"] @@ -120,6 +130,19 @@ def test_chat_completions_ingest_happy_path(client: TestClient): assert filtered_spans[0]["span_id"] == "chatcmpl-test-abc123" +def test_chat_completions_ingest_rejects_unknown_experiment_context(client: TestClient): + body = { + "request": _openai_request(), + "response": _openai_response(id="chatcmpl-missing-exp"), + "session_id": "session-missing-exp", + "experiment_context": {"experiment_id": "missing-exp", "test_case_id": "case-1"}, + } + response = client.post(INGEST_URL, json=body) + + assert response.status_code == 400, response.text + assert "must be created before it can be logged" in response.json()["detail"] + + def test_chat_completions_ingest_persists_cost_fields(client: TestClient): body = { "request": _openai_request(), @@ -220,12 +243,13 @@ def test_chat_completions_ingest_handles_missing_usage(client: TestClient): assert span["model"] == "gpt-4o-mini-2024-08-06" -def test_chat_completions_ingest_accepts_run_id_only_evaluation_context(client: TestClient): +def test_chat_completions_ingest_accepts_deprecated_evaluation_context(client: TestClient): + _create_experiment(client, EVALUATION_CONTEXT["evaluation_id"]) body = { "request": _openai_request(), "response": _openai_response(id="chatcmpl-run-id-only"), "session_id": "session-run-id-only", - "evaluation_context": {"evaluation_run_id": "evalrun-chat-run-only"}, + "evaluation_context": EVALUATION_CONTEXT, } response = client.post(INGEST_URL, json=body) assert response.status_code == 201, response.text @@ -233,7 +257,45 @@ def test_chat_completions_ingest_accepts_run_id_only_evaluation_context(client: listed = client.get(SPANS_URL, params={"filter[session_id]": "session-run-id-only"}) assert listed.status_code == 200, listed.text span = listed.json()["data"][0] - assert span["evaluation_context"] == {"evaluation_run_id": "evalrun-chat-run-only", "metadata": {}} + assert span["evaluation_context"] == { + "evaluation_id": EVALUATION_CONTEXT["evaluation_id"], + "evaluation_sha": EVALUATION_CONTEXT["evaluation_sha"], + "evaluation_run_id": EVALUATION_CONTEXT["evaluation_run_id"], + "dataset_id": EVALUATION_CONTEXT["dataset_id"], + "dataset_name": EVALUATION_CONTEXT["dataset_name"], + "dataset_version": EVALUATION_CONTEXT["dataset_version"], + "test_case_id": EVALUATION_CONTEXT["test_case_id"], + "metadata": EVALUATION_CONTEXT["metadata"], + } + + +def test_chat_completions_ingest_rejects_unknown_deprecated_evaluation_context(client: TestClient): + body = { + "request": _openai_request(), + "response": _openai_response(id="chatcmpl-missing-deprecated-exp"), + "session_id": "session-missing-deprecated-exp", + "evaluation_context": { + "evaluation_id": "missing-exp", + "evaluation_run_id": "missing-run", + "test_case_id": "case-1", + }, + } + response = client.post(INGEST_URL, json=body) + + assert response.status_code == 400, response.text + + +def test_chat_completions_ingest_accepts_partial_deprecated_evaluation_context(client: TestClient): + _create_experiment(client, "chat-eval") + body = { + "request": _openai_request(), + "response": _openai_response(id="chatcmpl-partial-deprecated-context"), + "session_id": "session-partial-deprecated-context", + "evaluation_context": {"evaluation_id": "chat-eval", "test_case_id": "case-1"}, + } + response = client.post(INGEST_URL, json=body) + + assert response.status_code == 201, response.text def test_chat_completions_ingest_falls_back_when_response_id_missing(client: TestClient): @@ -365,12 +427,40 @@ def test_chat_completions_ingest_rejects_legacy_context_fields(client: TestClien assert response.status_code == 422, response.text -def test_chat_completions_ingest_requires_run_id_when_evaluation_context_is_set(client: TestClient): +def test_chat_completions_ingest_accepts_both_context_shapes_with_experiment_context_precedence(client: TestClient): + _create_experiment(client, "chat-eval") body = { "request": _openai_request(), - "response": _openai_response(id="chatcmpl-invalid-eval-context"), - "session_id": "session-invalid-eval-context", - "evaluation_context": {"evaluation_id": "chat-eval"}, + "response": _openai_response(id="chatcmpl-both-contexts"), + "session_id": "session-both-contexts", + "experiment_context": {"experiment_id": "chat-eval", "test_case_id": "case-1"}, + "evaluation_context": {"evaluation_id": "legacy-eval", "test_case_id": "legacy-case"}, } response = client.post(INGEST_URL, json=body) - assert response.status_code == 422, response.text + assert response.status_code == 201, response.text + + listed = client.get(SPANS_URL, params={"filter[session_id]": "session-both-contexts"}) + assert listed.status_code == 200, listed.text + span = listed.json()["data"][0] + assert span["evaluation_context"]["evaluation_id"] == "chat-eval" + assert span["evaluation_context"]["test_case_id"] == "case-1" + + +def _create_experiment(client: TestClient, name: str) -> str: + response = client.post( + "/apis/intake/v2/workspaces/default/experiments", + json={ + "name": name, + "agent_name": "sample-agent", + "agent_version": "1.0.0", + "dataset_name": "chat-dataset", + "dataset_version": "v1", + }, + ) + assert response.status_code in {201, 409}, response.text + if response.status_code == 201: + return response.json()["name"] + + existing = client.get(f"/apis/intake/v2/workspaces/default/experiments/{name}") + assert existing.status_code == 200, existing.text + return existing.json()["name"] diff --git a/services/intake/tests/integration/spans/test_clickhouse_bootstrap.py b/services/intake/tests/integration/spans/test_clickhouse_bootstrap.py index ca34f33ea6..3d73dbc2c1 100644 --- a/services/intake/tests/integration/spans/test_clickhouse_bootstrap.py +++ b/services/intake/tests/integration/spans/test_clickhouse_bootstrap.py @@ -25,6 +25,7 @@ def test_clickhouse_bootstrap_is_idempotent(clickhouse_client: ClickHouseSpanCli ("ch_annotations_0001",), ("ch_evaluator_results_0001",), ("ch_evaluator_results_0002",), + ("ch_experiment_sessions_0002",), ("ch_spans_0002",), ] diff --git a/services/intake/tests/integration/spans/test_experiment_rollups.py b/services/intake/tests/integration/spans/test_experiment_rollups.py new file mode 100644 index 0000000000..58a65e99cd --- /dev/null +++ b/services/intake/tests/integration/spans/test_experiment_rollups.py @@ -0,0 +1,209 @@ +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +"""Experiment rollup integration tests.""" + +from __future__ import annotations + +from datetime import datetime, timedelta, timezone +from typing import Any + +import pytest +from fastapi.testclient import TestClient + +ATIF_INGEST = "/apis/intake/v2/workspaces/default/ingest/atif" +EXPERIMENTS = "/apis/intake/v2/workspaces/default/experiments" + + +def test_experiment_response_hydrates_clickhouse_rollups(client: TestClient) -> None: + experiment_id = "rollup-exp" + created = client.post( + EXPERIMENTS, + json={ + "name": experiment_id, + "agent_name": "sample-agent", + "agent_version": "1.0.0", + "dataset_name": "rollup-dataset", + "dataset_version": "v1", + }, + ) + assert created.status_code == 201, created.text + + rows = [ + ("run-1", "case-1a", 0.4, 0.05, 500), + ("run-1", "case-1b", 0.8, 0.10, 1500), + ("run-2", "case-2", 0.8, 0.20, 2000), + ("run-3", "case-3", 1.0, 0.30, 3000), + ] + started_at = datetime.now(timezone.utc).replace(microsecond=0) + for index, (run_id, test_case_id, score, cost_usd, latency_ms) in enumerate(rows): + response = client.post( + ATIF_INGEST, + json=_atif_body( + started_at=started_at, + experiment_id=experiment_id, + run_id=run_id, + test_case_id=test_case_id, + score=score, + cost_usd=cost_usd, + latency_ms=latency_ms, + offset_seconds=index * 10, + ), + ) + assert response.status_code == 201, response.text + + fetched = client.get(f"{EXPERIMENTS}/{experiment_id}") + assert fetched.status_code == 200, fetched.text + experiment = fetched.json() + + assert experiment["run_count"] == 4 + assert experiment["evaluator_names"] == ["harbor.verifier"] + assert experiment["model_names"] == ["provider/sample-model"] + + score = experiment["aggregate_scores"]["harbor.verifier"] + assert score["sum"] == pytest.approx(3.0) + assert score["mean"] == pytest.approx(0.75) + assert score["median"] == pytest.approx(0.8) + assert score["p90"] == pytest.approx(1.0) + assert score["p95"] == pytest.approx(1.0) + assert score["p99"] == pytest.approx(1.0) + assert score["count"] == 4 + + cost = experiment["cost_usd"] + assert cost["sum"] == pytest.approx(0.65) + assert cost["mean"] == pytest.approx(0.1625) + assert cost["median"] == pytest.approx(0.2) + assert cost["p90"] == pytest.approx(0.3) + assert cost["p95"] == pytest.approx(0.3) + assert cost["p99"] == pytest.approx(0.3) + assert cost["count"] == 4 + + latency = experiment["latency_ms"] + assert latency["sum"] == pytest.approx(7000.0) + assert latency["mean"] == pytest.approx(1750.0) + assert latency["median"] == pytest.approx(2000.0) + assert latency["p90"] == pytest.approx(3000.0) + assert latency["p95"] == pytest.approx(3000.0) + assert latency["p99"] == pytest.approx(3000.0) + assert latency["count"] == 4 + + listed = client.get(EXPERIMENTS) + assert listed.status_code == 200, listed.text + listed_experiment = next(item for item in listed.json()["data"] if item["name"] == experiment_id) + assert listed_experiment["aggregate_scores"]["harbor.verifier"]["mean"] == pytest.approx(0.75) + + +def test_atif_ingest_rejects_unknown_experiment_context(client: TestClient) -> None: + response = client.post( + ATIF_INGEST, + json=_atif_body( + started_at=datetime.now(timezone.utc).replace(microsecond=0), + experiment_id="missing-exp", + run_id="run-1", + test_case_id="case-1", + score=1.0, + cost_usd=0.01, + latency_ms=100, + offset_seconds=0, + ), + ) + + assert response.status_code == 400, response.text + assert "must be created before it can be logged" in response.json()["detail"] + + +def test_deprecated_evaluation_context_hydrates_experiment_rollups(client: TestClient) -> None: + experiment_id = "legacy-eval-context-exp" + created = client.post( + EXPERIMENTS, + json={ + "name": experiment_id, + "agent_name": "sample-agent", + "agent_version": "1.0.0", + "dataset_name": "rollup-dataset", + "dataset_version": "v1", + }, + ) + assert created.status_code == 201, created.text + + response = client.post( + ATIF_INGEST, + json={ + **_atif_body( + started_at=datetime.now(timezone.utc).replace(microsecond=0), + experiment_id=experiment_id, + run_id="run-1", + test_case_id="case-1", + score=1.0, + cost_usd=0.01, + latency_ms=100, + offset_seconds=0, + ), + "experiment_context": None, + "evaluation_context": {"evaluation_id": experiment_id, "test_case_id": "case-1"}, + }, + ) + + assert response.status_code == 201, response.text + + fetched = client.get(f"{EXPERIMENTS}/{experiment_id}") + assert fetched.status_code == 200, fetched.text + experiment = fetched.json() + assert experiment["run_count"] == 1 + assert experiment["aggregate_scores"]["harbor.verifier"]["mean"] == pytest.approx(1.0) + + +def _atif_body( + *, + started_at: datetime, + experiment_id: str, + run_id: str, + test_case_id: str, + score: float, + cost_usd: float, + latency_ms: int, + offset_seconds: int, +) -> dict[str, Any]: + session_started_at = started_at + timedelta(seconds=offset_seconds) + finished_at = session_started_at + timedelta(milliseconds=latency_ms) + session_id = f"{experiment_id}-{run_id}-{test_case_id}" + return { + "schema_version": "ATIF-v1.7", + "session_id": session_id, + "experiment_context": { + "experiment_id": experiment_id, + "test_case_id": test_case_id, + }, + "extra": { + "task_id": test_case_id, + "task_name": test_case_id, + "verifier": { + "started_at": _iso(session_started_at), + "finished_at": _iso(finished_at), + }, + "verifier_result": {"rewards": {"reward": score}}, + }, + "agent": { + "name": "sample-agent", + "version": "1.0.0", + "model_name": "provider/sample-model", + }, + "steps": [ + { + "step_id": 1, + "timestamp": _iso(session_started_at), + "source": "agent", + "model_name": "provider/sample-model", + "message": f"solved {test_case_id}", + "metrics": { + "prompt_tokens": 100, + "completion_tokens": 10, + "cost_usd": cost_usd, + }, + } + ], + } + + +def _iso(value: datetime) -> str: + return value.isoformat().replace("+00:00", "Z") diff --git a/services/intake/tests/integration/spans/test_traces_read.py b/services/intake/tests/integration/spans/test_traces_read.py index 305e26c9a4..6b2c6bb7d6 100644 --- a/services/intake/tests/integration/spans/test_traces_read.py +++ b/services/intake/tests/integration/spans/test_traces_read.py @@ -22,9 +22,9 @@ def test_traces_read_returns_core_trace_summary(client: TestClient, make_otlp_re "openinference.span.kind": "AGENT", "gen_ai.conversation.id": "trace-session", "project": "project-a", - "evaluation.run_id": "run-a", + "experiment.run_id": "run-a", "dataset.name": "dataset-a", - "evaluation.metadata": {"split": "dev"}, + "experiment.metadata": {"split": "dev"}, "deployment.environment.name": "prod", "tag.tags": ["trace-read"], "metadata": {"owner": "trace-test"}, diff --git a/services/intake/tests/integration/test_experiments_crud.py b/services/intake/tests/integration/test_experiments_crud.py index 494a78ebcb..dbcfffe59d 100644 --- a/services/intake/tests/integration/test_experiments_crud.py +++ b/services/intake/tests/integration/test_experiments_crud.py @@ -5,9 +5,11 @@ from __future__ import annotations -from typing import Any +from typing import Any, cast +from fastapi import FastAPI from fastapi.testclient import TestClient +from nmp.intake.api.v2.experiments.endpoints import get_experiment_rollup_repository GROUPS = "/apis/intake/v2/workspaces/default/experiment-groups" EXPERIMENTS = "/apis/intake/v2/workspaces/default/experiments" @@ -107,6 +109,26 @@ def test_experiment_crud_and_empty_rollups(client: TestClient) -> None: assert exp["run_count"] == 0 +def test_experiment_read_degrades_when_rollup_hydration_fails(client: TestClient) -> None: + class FailingRollupRepository: + async def get_rollups(self, *, workspace: str, experiment_ids: list[str]) -> dict: + raise RuntimeError("clickhouse unavailable") + + app = cast(FastAPI, client.app) + app.dependency_overrides[get_experiment_rollup_repository] = lambda: FailingRollupRepository() + try: + created = client.post(EXPERIMENTS, json=_experiment_body(name="exp-rollup-fails")) + assert created.status_code == 201, created.text + + fetched = client.get(f"{EXPERIMENTS}/exp-rollup-fails") + assert fetched.status_code == 200, fetched.text + assert fetched.json()["name"] == "exp-rollup-fails" + assert fetched.json()["run_count"] == 0 + assert fetched.json()["aggregate_scores"] is None + finally: + app.dependency_overrides.pop(get_experiment_rollup_repository, None) + + def test_experiment_group_ref_is_soft(client: TestClient) -> None: # experiment_group_id is a soft reference: a non-existent group id is accepted. created = client.post(EXPERIMENTS, json=_experiment_body(experiment_group_id="grp-does-not-exist")) diff --git a/services/intake/tests/test_atif_v17.py b/services/intake/tests/test_atif_v17.py index f76bf0fe92..13d9231eb0 100644 --- a/services/intake/tests/test_atif_v17.py +++ b/services/intake/tests/test_atif_v17.py @@ -19,7 +19,7 @@ AtifTrajectory, ) from nmp.intake.spans.ingest.atif_mapping import trajectory_to_spans -from nmp.intake.spans.ingest.evaluation_context import EvaluationContext +from nmp.intake.spans.ingest.evaluation_context import EvaluationContext, ExperimentContext from pydantic import ValidationError EVALUATION_CONTEXT: dict[str, Any] = { @@ -167,15 +167,19 @@ def test_atif_v17_subagent_ref_requires_resolution_key() -> None: assert AtifSubagentTrajectoryRef(trajectory_path="subagents/sub-trajectory.json").trajectory_path is not None -def test_evaluation_context_requires_run_id_when_any_context_field_is_set() -> None: - assert EvaluationContext() == EvaluationContext(metadata={}) - assert EvaluationContext(evaluation_run_id="evalrun-1").evaluation_run_id == "evalrun-1" +def test_experiment_context_maps_to_storage_context() -> None: + context = ExperimentContext(experiment_id="exp-1", test_case_id="case-1") + storage_context = context.to_evaluation_context() + + assert storage_context.evaluation_id == "exp-1" + assert storage_context.evaluation_run_id is None + assert storage_context.test_case_id == "case-1" - with pytest.raises(ValidationError, match="evaluation_run_id"): - EvaluationContext(evaluation_id="eval-1") - with pytest.raises(ValidationError, match="evaluation_run_id"): - EvaluationContext(metadata={"attempt": 1}) +def test_deprecated_evaluation_context_accepts_missing_run_id() -> None: + assert EvaluationContext() == EvaluationContext(metadata={}) + assert EvaluationContext(evaluation_run_id="evalrun-1").evaluation_run_id == "evalrun-1" + assert EvaluationContext(evaluation_id="eval-1").evaluation_id == "eval-1" def test_atif_ingest_request_rejects_legacy_top_level_project() -> None: @@ -215,14 +219,15 @@ def test_atif_mapping_writes_evaluation_context_only_on_root_span() -> None: root = next(span for span in spans if span.name == "sample-agent") child = next(span for span in spans if span.name == "user-1") - assert root.attributes_string["evaluation.id"] == EVALUATION_CONTEXT["evaluation_id"] - assert root.attributes_string["evaluation.sha"] == EVALUATION_CONTEXT["evaluation_sha"] - assert root.attributes_string["evaluation.run_id"] == EVALUATION_CONTEXT["evaluation_run_id"] + assert root.attributes_string["experiment.id"] == EVALUATION_CONTEXT["evaluation_id"] + assert root.attributes_string["experiment.sha"] == EVALUATION_CONTEXT["evaluation_sha"] + assert root.attributes_string["experiment.run_id"] == EVALUATION_CONTEXT["evaluation_run_id"] + assert "evaluation.id" not in root.attributes_string + assert root.attributes_string["test_case.id"] == EVALUATION_CONTEXT["test_case_id"] assert root.attributes_string["dataset.id"] == EVALUATION_CONTEXT["dataset_id"] assert root.attributes_string["dataset.name"] == EVALUATION_CONTEXT["dataset_name"] assert root.attributes_string["dataset.version"] == EVALUATION_CONTEXT["dataset_version"] - assert root.attributes_string["dataset.test_case_id"] == EVALUATION_CONTEXT["test_case_id"] - assert json.loads(root.attributes_string["evaluation.metadata"]) == EVALUATION_CONTEXT["metadata"] + assert json.loads(root.attributes_string["experiment.metadata"]) == EVALUATION_CONTEXT["metadata"] root_response = Span.from_domain(root) assert root_response.evaluation_context is not None @@ -238,11 +243,36 @@ def test_atif_mapping_writes_evaluation_context_only_on_root_span() -> None: root_raw = json.loads(root_response.raw_attributes) assert "evaluation_context" not in root_raw assert "evaluation.metadata" not in root_raw + assert "experiment.metadata" not in root_raw child_response = Span.from_domain(child) assert child_response.evaluation_context is None - assert "evaluation.run_id" not in child.attributes_string - assert "evaluation.metadata" not in child.attributes_string + assert "evaluation.id" not in child.attributes_string + assert "experiment.id" not in child.attributes_string + assert "test_case.id" not in child.attributes_string + + +def test_atif_mapping_writes_experiment_context_to_experiment_attributes() -> None: + body = AtifIngestRequest.model_validate( + { + "schema_version": "ATIF-v1.7", + "session_id": "trace-session-id", + "experiment_context": {"experiment_id": "exp-1", "test_case_id": "case-1"}, + "agent": {"name": "sample-agent", "version": "1.0.0"}, + "steps": [{"step_id": 1, "source": "user", "message": "solve"}], + } + ) + + spans = trajectory_to_spans( + workspace="default", + trajectory=body.to_trajectory(), + ingested_at=datetime(2026, 5, 18, tzinfo=timezone.utc), + ) + + root = next(span for span in spans if span.name == "sample-agent") + assert root.attributes_string["experiment.id"] == "exp-1" + assert root.attributes_string["test_case.id"] == "case-1" + assert "evaluation.id" not in root.attributes_string def test_atif_mapping_populates_root_content_and_rolls_child_errors() -> None: diff --git a/services/intake/tests/test_experiment_rollup_repository.py b/services/intake/tests/test_experiment_rollup_repository.py new file mode 100644 index 0000000000..d432d943fa --- /dev/null +++ b/services/intake/tests/test_experiment_rollup_repository.py @@ -0,0 +1,147 @@ +# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +"""Experiment rollup repository tests.""" + +from typing import cast + +import pytest +from nmp.intake.spans.clickhouse_client import ClickHouseSpanClient +from nmp.intake.spans.experiment_rollup_repository import ExperimentRollupRepository + + +class _QueryResult: + def __init__(self, rows: list[tuple[object, ...]], columns: list[str]) -> None: + self.result_rows = rows + self.column_names = columns + + +class _Client: + def __init__(self, query_results: list[_QueryResult]) -> None: + self.queries: list[str] = [] + self.parameters: list[dict[str, object]] = [] + self.query_results = query_results + + def table(self, name: str) -> str: + return name + + async def query(self, query: str, *, parameters: dict[str, object]) -> _QueryResult: + self.queries.append(query) + self.parameters.append(parameters) + return self.query_results.pop(0) + + +def _repository(client: _Client) -> ExperimentRollupRepository: + return ExperimentRollupRepository(cast(ClickHouseSpanClient, client)) + + +@pytest.mark.asyncio +async def test_experiment_rollups_anchor_on_root_session_membership(): + client = _Client( + [ + _QueryResult( + [("exp-a", 3)], + ["experiment_id", "run_count"], + ), + _QueryResult( + [("exp-a", "harbor.verifier", 3.0, 0.75, 0.8, 1.0, 1.0, 1.0, 4)], + ["experiment_id", "evaluator_name", "sum", "mean", "median", "p90", "p95", "p99", "count"], + ), + _QueryResult( + [ + ( + "exp-a", + ["model-b", "model-a"], + 0.65, + 0.1625, + 0.2, + 0.3, + 0.3, + 0.3, + 4, + 7000.0, + 1750.0, + 2000.0, + 3000.0, + 3000.0, + 3000.0, + 4, + ) + ], + [ + "experiment_id", + "model_names", + "cost_sum", + "cost_mean", + "cost_median", + "cost_p90", + "cost_p95", + "cost_p99", + "cost_count", + "latency_sum", + "latency_mean", + "latency_median", + "latency_p90", + "latency_p95", + "latency_p99", + "latency_count", + ], + ), + ] + ) + repository = _repository(client) + + rollups = await repository.get_rollups(workspace="default", experiment_ids=["exp-a"]) + + rollup = rollups["exp-a"] + assert rollup.run_count == 3 + assert rollup.evaluator_names == ["harbor.verifier"] + assert rollup.model_names == ["model-a", "model-b"] + assert rollup.evaluator_scores["harbor.verifier"].sum == 3.0 + assert rollup.evaluator_scores["harbor.verifier"].mean == 0.75 + assert rollup.evaluator_scores["harbor.verifier"].median == 0.8 + assert rollup.evaluator_scores["harbor.verifier"].p90 == 1.0 + assert rollup.evaluator_scores["harbor.verifier"].p95 == 1.0 + assert rollup.evaluator_scores["harbor.verifier"].p99 == 1.0 + assert rollup.evaluator_scores["harbor.verifier"].count == 4 + assert rollup.cost_usd is not None + assert rollup.cost_usd.sum == 0.65 + assert rollup.cost_usd.mean == 0.1625 + assert rollup.cost_usd.median == 0.2 + assert rollup.cost_usd.p90 == 0.3 + assert rollup.cost_usd.p95 == 0.3 + assert rollup.cost_usd.p99 == 0.3 + assert rollup.cost_usd.count == 4 + assert rollup.latency_ms is not None + assert rollup.latency_ms.sum == 7000 + assert rollup.latency_ms.mean == 1750 + assert rollup.latency_ms.median == 2000 + assert rollup.latency_ms.p90 == 3000 + assert rollup.latency_ms.p95 == 3000 + assert rollup.latency_ms.p99 == 3000 + assert rollup.latency_ms.count == 4 + + assert len(client.queries) == 3 + assert "FROM experiment_sessions FINAL" in client.queries[0] + assert "count() AS run_count" in client.queries[0] + assert "experiment_id IN (%(experiment_id_0)s)" in client.queries[0] + assert "FROM evaluator_results FINAL" in client.queries[1] + assert "quantileExact(0.5)(value) AS median" in client.queries[1] + assert "quantileExact(0.99)(value) AS p99" in client.queries[1] + assert "AND (workspace, session_id) IN (" in client.queries[1] + assert "sessions.session_id = results.session_id" in client.queries[1] + # Scores are reduced to one value per (experiment, session, evaluator) before the + # distribution rollup so count tracks runs and the mean is not span-weighted. + assert "GROUP BY sessions.experiment_id, sessions.session_id, results.name" in client.queries[1] + assert "current_session_spans AS" in client.queries[2] + assert "(span_versions.workspace, span_versions.session_id) IN" in client.queries[2] + assert "LEFT JOIN current_session_spans AS spans" in client.queries[2] + assert "sessions.session_id = spans.session_id" in client.queries[2] + assert "arraySort(arrayDistinct(arrayFlatten(groupArray(model_names)))) AS model_names" in client.queries[2] + assert "quantileExactIf(0.5)" in client.queries[2] + assert "cost_median" in client.queries[2] + assert "quantileExactIf(0.99)" in client.queries[2] + assert "latency_p99" in client.queries[2] + assert "sessions.trace_id = spans.trace_id" not in client.queries[2] + assert client.parameters[0]["experiment_id_0"] == "exp-a" + assert client.parameters[2]["model_key"] == "gen_ai.request.model" diff --git a/services/intake/tests/test_spans_clickhouse_migrations.py b/services/intake/tests/test_spans_clickhouse_migrations.py index 45207f4eee..1384786f8d 100644 --- a/services/intake/tests/test_spans_clickhouse_migrations.py +++ b/services/intake/tests/test_spans_clickhouse_migrations.py @@ -9,6 +9,7 @@ import nmp.intake.spans.clickhouse_migrations as clickhouse_migrations import pytest from nmp.intake.spans.clickhouse_migrations import parse_clickhouse_url +from nmp.intake.spans.span_attribute_catalog import SpanAttributeField, spec_for_field def test_parse_clickhouse_url_rejects_hostless_url(): @@ -32,3 +33,42 @@ def test_spans_schema_keeps_cityhash_identity_expression(): "trace_id", "external_span_id", ] + + +def test_experiment_sessions_schema_is_ordered_by_experiment(): + source = Path(clickhouse_migrations.__file__).read_text(encoding="utf-8") + function_match = re.search( + r"def _create_experiment_sessions_schema\(.*?^_MIGRATIONS", + source, + re.DOTALL | re.MULTILINE, + ) + assert function_match is not None + source = function_match.group(0) + + table_match = re.search( + r"CREATE TABLE \{table\}.*?ttl_only_drop_parts = 1", + source, + re.DOTALL, + ) + + assert table_match is not None + ddl = source + + assert '"experiment_sessions"' in source + assert '"experiment_sessions_mv"' in source + assert "CREATE TABLE {table}" in ddl + assert "CREATE MATERIALIZED VIEW {view}" in ddl + assert "TO {table}" in ddl + assert "INSERT INTO {table}" in source + assert "attributes_string['experiment.id'] AS experiment_id" in ddl + assert "attributes_string['test_case.id']" in ddl + assert "attributes_string['evaluation.id']" not in ddl + assert "experiment_run_id" not in ddl + assert "PRIMARY KEY (workspace, experiment_id, session_id)" in ddl + assert "ORDER BY (workspace, experiment_id, session_id, root_span_id)" in ddl + assert "index_granularity = 256" in ddl + + +def test_experiment_sessions_mv_keys_match_attribute_catalog(): + assert spec_for_field(SpanAttributeField.EVALUATION_ID).bag_key == "experiment.id" + assert spec_for_field(SpanAttributeField.TEST_CASE_ID).bag_key == "test_case.id" diff --git a/services/intake/tests/test_spans_schemas.py b/services/intake/tests/test_spans_schemas.py index 69c5231b28..3c73afbe32 100644 --- a/services/intake/tests/test_spans_schemas.py +++ b/services/intake/tests/test_spans_schemas.py @@ -41,10 +41,10 @@ def test_span_response_raw_attributes_merges_atif_raw_with_unknown_attributes(): event_ts=now, attributes_string={ "atif.raw": json_dumps_preserve( - {"source_session_id": "session-a", "evaluation.metadata": {"source": "atif.raw"}} + {"source_session_id": "session-a", "experiment.metadata": {"source": "atif.raw"}} ), "custom.string": "value-a", - "evaluation.metadata": json.dumps({"source": "attribute.bag"}), + "experiment.metadata": json.dumps({"source": "attribute.bag"}), "gen_ai.request.model": "model-a", }, attributes_number={"custom.number": 1.25, "llm.token_count.prompt": 42}, @@ -74,9 +74,6 @@ def test_trace_response_maps_core_trace_fields(): name="root", input="root input", output="root output", - environment="prod", - tags=["red", "blue"], - metadata={"owner": "intake"}, project="project-a", evaluation_context=TraceEvaluationContext( evaluation_id="eval-a", diff --git a/services/intake/tests/test_spans_span_attribute_catalog.py b/services/intake/tests/test_spans_span_attribute_catalog.py index 62bbfa8c3f..cb706597ee 100644 --- a/services/intake/tests/test_spans_span_attribute_catalog.py +++ b/services/intake/tests/test_spans_span_attribute_catalog.py @@ -166,6 +166,17 @@ def test_span_attribute_catalog_predicates_include_existence_guard(operator: str assert params["cost_total_usd_value"] == 1_000_000 +def test_span_attribute_catalog_groups_legacy_alias_predicates(): + sql, params = where_clause("evaluation_run_id", "$eq", "run-a", param_prefix="root_candidate_0") + + assert sql.startswith("(") + assert sql.endswith(")") + assert " OR " in sql + assert params["root_candidate_0_key"] == "experiment.run_id" + assert params["root_candidate_0_key_1"] == "evaluation.run_id" + assert params["root_candidate_0_value"] == "run-a" + + def test_span_attribute_catalog_rejects_ordering_on_string_fields(): with pytest.raises(ValueError, match="only supports equality"): where_clause("model", ">", "gpt-4") diff --git a/services/intake/tests/test_traces_clickhouse_repository.py b/services/intake/tests/test_traces_clickhouse_repository.py index 9bc9deaff5..ae45ff5056 100644 --- a/services/intake/tests/test_traces_clickhouse_repository.py +++ b/services/intake/tests/test_traces_clickhouse_repository.py @@ -208,7 +208,7 @@ async def test_root_and_any_span_filters_select_candidate_trace_ids(): assert "(trace_spans.workspace, trace_spans.source_format, trace_spans.trace_id) IN" in client.queries[0] assert "FINAL" not in client.queries[0] assert "external_parent_span_id = ''" in client.queries[0] - assert client.parameters[0]["root_candidate_0_key"] == "evaluation.run_id" + assert client.parameters[0]["root_candidate_0_key"] == "experiment.run_id" assert client.parameters[0]["root_candidate_0_value"] == "run-a" assert client.parameters[0]["span_candidate_0_key"] == "gen_ai.request.model" assert client.parameters[0]["span_candidate_0_value"] == "model-a" @@ -246,7 +246,7 @@ def _trace_row( "error_count": 1 if detailed else None, "root_attributes_string": { "project.name": "project-a", - "evaluation.run_id": "run-a", + "experiment.run_id": "run-a", }, "ingested_at": ingested_at, }