Skip to content
Draft
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
244 changes: 0 additions & 244 deletions plugins/nemo-evaluator/src/nemo_evaluator/sdk/_executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@

from __future__ import annotations

import asyncio
from collections.abc import Sequence
from typing import Any, TypeAlias, cast

Expand All @@ -15,21 +14,18 @@
from nemo_evaluator.jobs.evaluate import EvaluateInputSpec, EvaluateJob, EvaluateSpec, TargetSpec
from nemo_evaluator.resolvers import PlatformModelResolver
from nemo_evaluator.sdk import http_utils
from nemo_evaluator.sdk.fs_utils import EvaluatorLocalRunResult, local_result_path
from nemo_evaluator.sdk.job_resources import (
AsyncEvaluatorJobResource,
EvaluatorJob,
EvaluatorJobResource,
)
from nemo_evaluator.sdk.types import PluginDatasetInput
from nemo_evaluator.sdk.utils import filter_benchmark_result, filter_evaluation_result
from nemo_evaluator.shared.metric_bundles.bundles import (
MetricBundle,
MetricBundlePackager,
MetricBundlePackagerPolicyError,
bundle_metric,
)
from nemo_evaluator.shared.metric_bundles.defaults import resolve_default_metric_bundle_packager
from nemo_evaluator_sdk.datasets.loader import prepare_dataset_rows
from nemo_evaluator_sdk.execution.config import resolve_params
from nemo_evaluator_sdk.execution.metric_execution import run_sync
Expand All @@ -44,10 +40,7 @@
RunConfigOnline,
RunConfigOnlineModel,
)
from nemo_evaluator_sdk.values.multi_metric_results import BenchmarkEvaluationResult
from nemo_evaluator_sdk.values.results import AggregateFieldName, EvaluationResult
from nemo_platform import AsyncNeMoPlatform, NeMoPlatform
from nemo_platform_plugin.scheduler import NemoJobScheduler

_DEFAULT_POLL_INTERVAL_SECONDS = 10.0
_DEFAULT_JOB_TIMEOUT_SECONDS = 3600.0
Expand Down Expand Up @@ -247,90 +240,6 @@ def create(
)
return job_resource

def run_local(self, *, spec: EvaluateRequestSpec, workspace: str | None = None) -> EvaluatorLocalRunResult:
"""Run an evaluator plugin job locally with a sync platform client."""
resolved_workspace = http_utils.resolve_workspace(self._platform, workspace)
canonical_spec = _resolve_sync_local_spec(
spec,
platform=self._platform,
workspace=resolved_workspace,
)
payload = NemoJobScheduler().run_local(
EvaluateJob,
canonical_spec.model_dump(mode="json"),
workspace=resolved_workspace,
sdk=self._platform,
)

return EvaluatorLocalRunResult.model_validate(payload)

def evaluate_remote(
self,
*,
metric: Metric,
dataset: PluginDatasetInput,
params: RunConfig | RunConfigOnline | RunConfigOnlineModel,
target: Model | Agent | None = None,
field_mapping: FieldMapping | None = None,
prompt_template: str | dict[str, Any] | None = None,
aggregate_fields: tuple[AggregateFieldName, ...] | None = None,
metric_bundle_packager: MetricBundlePackager | None = None,
) -> EvaluationResult:
"""Submit, poll, and download a remote evaluator plugin metric job."""
normalized_params = resolve_params(params, target)
spec = _build_evaluate_spec(
metrics=metric,
dataset=dataset,
params=normalized_params,
target=target,
field_mapping=field_mapping,
prompt_template=prompt_template,
metric_bundle_packager=metric_bundle_packager,
)

job = self.create(
spec=spec, workspace=http_utils.resolve_workspace(self._platform, self._workspace, strict=True)
)
job.wait_until_done(
poll_interval_seconds=self._poll_interval_seconds,
job_timeout_seconds=self._job_timeout_seconds,
pending_timeout_seconds=self._pending_timeout_seconds,
)

return job.get_result(aggregate_fields=aggregate_fields)

def evaluate(
self,
*,
metric: Metric,
dataset: PluginDatasetInput,
params: RunConfig | RunConfigOnline | RunConfigOnlineModel | None = None,
target: Model | Agent | None = None,
field_mapping: FieldMapping | None = None,
prompt_template: str | dict[str, Any] | None = None,
aggregate_fields: tuple[AggregateFieldName, ...] | None = None,
) -> EvaluationResult:
"""Evaluate one metric through local plugin job execution."""
normalized_params = resolve_params(params, target)
spec = _build_evaluate_spec(
metrics=metric,
dataset=dataset,
params=normalized_params,
target=target,
field_mapping=field_mapping,
prompt_template=prompt_template,
metric_bundle_packager=resolve_default_metric_bundle_packager(
metric, None, allow_cloudpickle_fallback=True, action="Running"
),
)
payload = self.run_local(
spec=spec,
workspace=http_utils.resolve_workspace(self._platform, self._workspace, strict=True),
)
result_path = local_result_path(payload)
result = EvaluationResult.model_validate_json(result_path.read_text(encoding="utf-8"))
return filter_evaluation_result(result, aggregate_fields)

def submit(
self,
*,
Expand Down Expand Up @@ -361,38 +270,6 @@ def submit(

return job

def evaluate_benchmark(
self,
*,
metrics: Sequence[Metric],
dataset: PluginDatasetInput,
params: RunConfig | RunConfigOnline | RunConfigOnlineModel,
target: Model | Agent | None = None,
field_mapping: FieldMapping | None = None,
prompt_template: str | dict[str, Any] | None = None,
aggregate_fields: tuple[AggregateFieldName, ...] | None = None,
) -> BenchmarkEvaluationResult:
"""Evaluate multiple metrics through local plugin job execution."""
normalized_params = resolve_params(params, target)
spec = _build_evaluate_spec(
metrics=metrics,
dataset=dataset,
params=normalized_params,
target=target,
field_mapping=field_mapping,
prompt_template=prompt_template,
metric_bundle_packager=resolve_default_metric_bundle_packager(
metrics, None, allow_cloudpickle_fallback=True, action="Running"
),
)
payload = self.run_local(
spec=spec,
workspace=http_utils.resolve_workspace(self._platform, self._workspace, strict=True),
)
result_path = local_result_path(payload)
result = BenchmarkEvaluationResult.model_validate_json(result_path.read_text(encoding="utf-8"))
return filter_benchmark_result(result, aggregate_fields)


class _AsyncEvaluatorPluginExecutor:
"""Async evaluator plugin executor used by the async SDK resource."""
Expand Down Expand Up @@ -449,26 +326,6 @@ async def create(
)
return job_resource

async def run_local(self, *, spec: EvaluateRequestSpec, workspace: str | None = None) -> EvaluatorLocalRunResult:
"""Run an evaluator plugin job locally without blocking the event loop."""
resolved_workspace = http_utils.resolve_workspace(self._platform, workspace)
canonical_spec = await _resolve_async_local_spec(
spec,
platform=self._platform,
workspace=resolved_workspace,
)
scheduler = NemoJobScheduler()
# Leverages programmatic dispatch as described in
# packages/nemo_platform_plugin/src/nemo_platform_plugin/docs/ARCHITECTURE.md#job-entry-point-keys
payload = await asyncio.to_thread(
scheduler.run_local,
EvaluateJob,
canonical_spec.model_dump(mode="json"),
workspace=resolved_workspace,
async_sdk=self._platform,
)
return EvaluatorLocalRunResult.model_validate(payload)

async def submit(
self,
*,
Expand Down Expand Up @@ -499,107 +356,6 @@ async def submit(

return job

async def evaluate_remote(
self,
*,
metric: Metric,
dataset: PluginDatasetInput,
params: RunConfig | RunConfigOnline | RunConfigOnlineModel,
target: Model | Agent | None = None,
field_mapping: FieldMapping | None = None,
prompt_template: str | dict[str, Any] | None = None,
aggregate_fields: tuple[AggregateFieldName, ...] | None = None,
metric_bundle_packager: MetricBundlePackager | None = None,
) -> EvaluationResult:
"""Submit, poll, and download a remote evaluator plugin metric job."""
normalized_params = resolve_params(params, target)
spec = _build_evaluate_spec(
metrics=metric,
dataset=dataset,
params=normalized_params,
target=target,
field_mapping=field_mapping,
prompt_template=prompt_template,
metric_bundle_packager=metric_bundle_packager,
)

job = await self.create(
spec=spec, workspace=http_utils.resolve_workspace(self._platform, self._workspace, strict=True)
)
await job.wait_until_done(
poll_interval_seconds=self._poll_interval_seconds,
job_timeout_seconds=self._job_timeout_seconds,
pending_timeout_seconds=self._pending_timeout_seconds,
)

return await job.get_result(aggregate_fields=aggregate_fields)

async def evaluate(
self,
*,
metric: Metric,
dataset: PluginDatasetInput,
params: RunConfig | RunConfigOnline | RunConfigOnlineModel | None = None,
target: Model | Agent | None = None,
field_mapping: FieldMapping | None = None,
prompt_template: str | dict[str, Any] | None = None,
aggregate_fields: tuple[AggregateFieldName, ...] | None = None,
) -> EvaluationResult:
"""Evaluate one metric through local plugin job execution."""
normalized_params = resolve_params(params, target)
spec = _build_evaluate_spec(
metrics=metric,
dataset=dataset,
params=normalized_params,
target=target,
field_mapping=field_mapping,
prompt_template=prompt_template,
metric_bundle_packager=resolve_default_metric_bundle_packager(
metric, None, allow_cloudpickle_fallback=True, action="Running"
),
)
payload = await self.run_local(
spec=spec,
workspace=http_utils.resolve_workspace(self._platform, self._workspace, strict=True),
)
result_path = local_result_path(payload)
result_text = await asyncio.to_thread(result_path.read_text, encoding="utf-8")
result = EvaluationResult.model_validate_json(result_text)
return filter_evaluation_result(result, aggregate_fields)

async def evaluate_benchmark(
self,
*,
metrics: Sequence[Metric],
dataset: PluginDatasetInput,
params: RunConfig | RunConfigOnline | RunConfigOnlineModel,
target: Model | Agent | None = None,
field_mapping: FieldMapping | None = None,
prompt_template: str | dict[str, Any] | None = None,
aggregate_fields: tuple[AggregateFieldName, ...] | None = None,
) -> BenchmarkEvaluationResult:
"""Evaluate multiple metrics through local plugin job execution."""
normalized_params = resolve_params(params, target)
spec = _build_evaluate_spec(
metrics=metrics,
dataset=dataset,
params=normalized_params,
target=target,
field_mapping=field_mapping,
prompt_template=prompt_template,
metric_bundle_packager=resolve_default_metric_bundle_packager(
metrics, None, allow_cloudpickle_fallback=True, action="Running"
),
)
payload = await self.run_local(
spec=spec,
workspace=http_utils.resolve_workspace(self._platform, self._workspace, strict=True),
)
result_path = local_result_path(payload)
result_text = await asyncio.to_thread(result_path.read_text, encoding="utf-8")
result = BenchmarkEvaluationResult.model_validate_json(result_text)
return filter_benchmark_result(result, aggregate_fields)


def bundle_metrics_for_spec(
metrics: Metric | Sequence[Metric], *, metric_bundle_packager: MetricBundlePackager
Expand Down
54 changes: 0 additions & 54 deletions plugins/nemo-evaluator/src/nemo_evaluator/sdk/fs_utils.py

This file was deleted.

Loading
Loading