Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
44 changes: 43 additions & 1 deletion packages/nemo_evaluator_sdk/examples/plugin_examples.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
from tempfile import TemporaryDirectory
from typing import TYPE_CHECKING, Any, cast

from nemo_evaluator.jobs.evaluate import EvaluateSpec
from nemo_evaluator.sdk.resources import AsyncEvaluator
from nemo_evaluator.sdk.resources import Evaluator as SyncEvaluator
from nemo_evaluator.sdk.types import (
Expand All @@ -23,10 +24,12 @@
RunConfig,
RunConfigOnlineModel,
)
from nemo_evaluator.shared.metric_bundles.bundles import bundle_metric
from nemo_evaluator.shared.metric_bundles.cloudpickle import CloudpickleMetricBundlePackager
from nemo_evaluator_sdk.enums import MetricType
from nemo_evaluator_sdk.metrics.exact_match import ExactMatchMetric
from nemo_evaluator_sdk.metrics.llm_judge import LLMJudgeMetric
from nemo_evaluator_sdk.metrics.protocol import Metric
from nemo_evaluator_sdk.metrics.protocol import Metric, MetricInput, MetricOutput, MetricOutputSpec, MetricResult
from nemo_evaluator_sdk.values import (
InferenceParams,
JSONScoreParser,
Expand Down Expand Up @@ -74,6 +77,23 @@
)


class CustomResponseLengthMetric:
"""Tiny custom metric used to demonstrate code-generated metric bundles."""

type = "custom-response-length"
description = "Scores each row by response length."
labels = {"source": "plugin-example"}

def output_spec(self) -> list[MetricOutputSpec]:
"""Return the metric outputs recorded in the bundle metadata."""
return [MetricOutputSpec.continuous_score("response-length")]

async def compute_scores(self, input: MetricInput) -> MetricResult:
"""Score one row with a deterministic custom Python implementation."""
response = str(input.row.data.get("response", ""))
return MetricResult(outputs=[MetricOutput(name="response-length", value=float(len(response)))])


def configure_example_logging() -> None:
"""Enable SDK progress logs when this example file is executed directly."""
logging.basicConfig(level=logging.INFO, format="%(levelname)s:%(name)s:%(message)s")
Expand Down Expand Up @@ -306,6 +326,26 @@ def _online_exact_match_metric() -> ExactMatchMetric:
return ExactMatchMetric(type=MetricType.EXACT_MATCH, reference="{{item.response}}")


def build_custom_metric_submit_spec_example() -> dict[str, Any]:
"""Return the generated job spec for remote custom metric submission.

The cloudpickle payload contains base64-encoded Python bytes. It is not a
field users should hand-author; generate it from the metric object with a
metric payload packager, or pass the packager to ``submit``.
"""
metric = CustomResponseLengthMetric()
spec = EvaluateSpec.model_validate(
{
"metrics": [
bundle_metric(metric, CloudpickleMetricBundlePackager()).model_dump(mode="json"),
],
"dataset": [{"response": "Paris is the capital of France."}],
"params": RunConfig(limit_samples=1).model_dump(mode="json"),
}
)
return spec.model_dump(mode="json")


def _assert_exact_match_result(result: EvaluationResult, *, workflow: str, expected_rows: int) -> None:
"""Assert the deterministic offline exact-match examples scored every selected row."""
if len(result.row_scores) != expected_rows:
Expand Down Expand Up @@ -345,6 +385,7 @@ async def _evaluate_metric(
metric=metric,
dataset=dataset,
config=config,
metric_bundle_packager=CloudpickleMetricBundlePackager(),
**run_kwargs,
)
print(f"Submitted evaluator plugin job: {job.name}")
Expand Down Expand Up @@ -506,6 +547,7 @@ def run_nmp_online_metric_example_sync_client(
metric=metric,
dataset=dataset,
config=config,
metric_bundle_packager=CloudpickleMetricBundlePackager(),
**run_kwargs,
)
print(f"Submitted evaluator plugin job: {job.name}")
Expand Down
1 change: 1 addition & 0 deletions packages/nemo_platform/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -345,6 +345,7 @@ nemo-data-designer-plugin = [

# Generated from [tool.bundle-package]; do not edit by hand.
nemo-evaluator-plugin = [
"cloudpickle>=3.1.1",
"nemo-evaluator-sdk",
"nemo-platform-plugin",
"nmp-evaluator",
Expand Down
1 change: 1 addition & 0 deletions plugins/nemo-evaluator/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ description = "Evaluator plugin scaffold for NeMo Platform."
readme = "README.md"
requires-python = ">=3.11,<3.14"
dependencies = [
"cloudpickle>=3.1.1",
"nemo-evaluator-sdk",
"nemo-platform-plugin",
"nemo-platform",
Expand Down
94 changes: 94 additions & 0 deletions plugins/nemo-evaluator/src/nemo_evaluator/jobs/compiler.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,94 @@
# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0

"""Plugin-native evaluator job compiler."""

from __future__ import annotations

from nemo_evaluator.jobs.evaluate import EvaluateSpec
from nemo_evaluator_sdk.values import Agent, Model, RunConfig, RunConfigOnline, RunConfigOnlineModel
from nemo_platform_plugin.jobs.api_factory import (
ContainerSpec,
CPUExecutionProviderSpec,
EnvironmentVariable,
EnvironmentVariableFromSecret,
PlatformJobSpec,
PlatformJobStep,
)
from nemo_platform_plugin.jobs.constants import (
DEFAULT_JOB_STORAGE_PATH,
PERSISTENT_JOB_STORAGE_PATH_ENVVAR,
)
from nmp.common.jobs.image import get_qualified_image

EVALUATE_STEP_NAME = "evaluate"
_RESERVED_SECRET_ENV_NAMES = frozenset({PERSISTENT_JOB_STORAGE_PATH_ENVVAR})


def compile_evaluate_job(spec: EvaluateSpec, *, profile: str | None = None) -> PlatformJobSpec:
"""Compile a bundle-native evaluator plugin job."""
_validate_evaluate_spec(spec)
return PlatformJobSpec(steps=[_evaluate_step(spec, profile)])


def _validate_evaluate_spec(spec: EvaluateSpec) -> None:
if isinstance(spec.target, Model):
if spec.prompt_template is None:
raise ValueError("prompt_template is required when EvaluateSpec.target is a model")
if not isinstance(spec.params, RunConfigOnlineModel):
raise TypeError("model target requires RunConfigOnlineModel")
elif isinstance(spec.target, Agent):
if spec.prompt_template is None:
raise ValueError("prompt_template is required when EvaluateSpec.target is an agent")
if not isinstance(spec.params, RunConfigOnline):
raise TypeError("agent target requires RunConfigOnline")
elif not isinstance(spec.params, RunConfig):
raise TypeError("offline evaluation requires RunConfig")


def _add_secret_ref(secret_refs: dict[str, str], env_name: str, secret_name: str) -> None:
if env_name in _RESERVED_SECRET_ENV_NAMES:
raise ValueError(f"{env_name!r} is reserved and cannot be sourced from secret refs")
existing = secret_refs.get(env_name)
if existing is not None and existing != secret_name:
raise ValueError(f"conflicting secret references for environment variable {env_name!r}")
secret_refs[env_name] = secret_name
Comment thread
coderabbitai[bot] marked this conversation as resolved.


def _secret_environment(spec: EvaluateSpec) -> list[EnvironmentVariable]:
environment = [
EnvironmentVariable(
name=PERSISTENT_JOB_STORAGE_PATH_ENVVAR,
value=DEFAULT_JOB_STORAGE_PATH,
)
]
secret_refs: dict[str, str] = {}
for bundle in spec.metrics:
for env_name, secret_ref in bundle.secrets.items():
_add_secret_ref(secret_refs, env_name, secret_ref.root)

if isinstance(spec.target, Model | Agent) and spec.target.api_key_secret is not None and spec.target.api_key_env:
_add_secret_ref(secret_refs, spec.target.api_key_env, spec.target.api_key_secret.root)

environment.extend(
EnvironmentVariable(name=env_name, from_secret=EnvironmentVariableFromSecret(name=secret_name))
for env_name, secret_name in sorted(secret_refs.items())
)
return environment


def _evaluate_step(spec: EvaluateSpec, profile: str | None) -> PlatformJobStep:
return PlatformJobStep(
name=EVALUATE_STEP_NAME,
executor=CPUExecutionProviderSpec(
profile=profile or "default",
provider="cpu",
container=ContainerSpec(
image=get_qualified_image("nmp-cpu-tasks"),
entrypoint=["python", "-m"],
command=["nemo_evaluator.tasks.evaluate"],
Comment thread
SandyChapman marked this conversation as resolved.
Comment thread
SandyChapman marked this conversation as resolved.
),
),
config=spec.model_dump(mode="json"),
environment=_secret_environment(spec),
)
Loading
Loading