|
5 | 5 |
|
6 | 6 | from __future__ import annotations |
7 | 7 |
|
8 | | -import asyncio |
9 | 8 | from collections.abc import Sequence |
10 | 9 | from typing import Any, TypeAlias, cast |
11 | 10 |
|
|
15 | 14 | from nemo_evaluator.jobs.evaluate import EvaluateInputSpec, EvaluateJob, EvaluateSpec, TargetSpec |
16 | 15 | from nemo_evaluator.resolvers import PlatformModelResolver |
17 | 16 | from nemo_evaluator.sdk import http_utils |
18 | | -from nemo_evaluator.sdk.fs_utils import EvaluatorLocalRunResult, local_result_path |
19 | 17 | from nemo_evaluator.sdk.job_resources import ( |
20 | 18 | AsyncEvaluatorJobResource, |
21 | 19 | EvaluatorJob, |
22 | 20 | EvaluatorJobResource, |
23 | 21 | ) |
24 | 22 | from nemo_evaluator.sdk.types import PluginDatasetInput |
25 | | -from nemo_evaluator.sdk.utils import filter_benchmark_result, filter_evaluation_result |
26 | 23 | from nemo_evaluator.shared.metric_bundles.bundles import ( |
27 | 24 | MetricBundle, |
28 | 25 | MetricBundlePackager, |
29 | 26 | MetricBundlePackagerPolicyError, |
30 | 27 | bundle_metric, |
31 | 28 | ) |
32 | | -from nemo_evaluator.shared.metric_bundles.defaults import resolve_default_metric_bundle_packager |
33 | 29 | from nemo_evaluator_sdk.datasets.loader import prepare_dataset_rows |
34 | 30 | from nemo_evaluator_sdk.execution.config import resolve_params |
35 | 31 | from nemo_evaluator_sdk.execution.metric_execution import run_sync |
|
44 | 40 | RunConfigOnline, |
45 | 41 | RunConfigOnlineModel, |
46 | 42 | ) |
47 | | -from nemo_evaluator_sdk.values.multi_metric_results import BenchmarkEvaluationResult |
48 | | -from nemo_evaluator_sdk.values.results import AggregateFieldName, EvaluationResult |
49 | 43 | from nemo_platform import AsyncNeMoPlatform, NeMoPlatform |
50 | | -from nemo_platform_plugin.scheduler import NemoJobScheduler |
51 | 44 |
|
52 | 45 | _DEFAULT_POLL_INTERVAL_SECONDS = 10.0 |
53 | 46 | _DEFAULT_JOB_TIMEOUT_SECONDS = 3600.0 |
@@ -247,90 +240,6 @@ def create( |
247 | 240 | ) |
248 | 241 | return job_resource |
249 | 242 |
|
250 | | - def run_local(self, *, spec: EvaluateRequestSpec, workspace: str | None = None) -> EvaluatorLocalRunResult: |
251 | | - """Run an evaluator plugin job locally with a sync platform client.""" |
252 | | - resolved_workspace = http_utils.resolve_workspace(self._platform, workspace) |
253 | | - canonical_spec = _resolve_sync_local_spec( |
254 | | - spec, |
255 | | - platform=self._platform, |
256 | | - workspace=resolved_workspace, |
257 | | - ) |
258 | | - payload = NemoJobScheduler().run_local( |
259 | | - EvaluateJob, |
260 | | - canonical_spec.model_dump(mode="json"), |
261 | | - workspace=resolved_workspace, |
262 | | - sdk=self._platform, |
263 | | - ) |
264 | | - |
265 | | - return EvaluatorLocalRunResult.model_validate(payload) |
266 | | - |
267 | | - def evaluate_remote( |
268 | | - self, |
269 | | - *, |
270 | | - metric: Metric, |
271 | | - dataset: PluginDatasetInput, |
272 | | - params: RunConfig | RunConfigOnline | RunConfigOnlineModel, |
273 | | - target: Model | Agent | None = None, |
274 | | - field_mapping: FieldMapping | None = None, |
275 | | - prompt_template: str | dict[str, Any] | None = None, |
276 | | - aggregate_fields: tuple[AggregateFieldName, ...] | None = None, |
277 | | - metric_bundle_packager: MetricBundlePackager | None = None, |
278 | | - ) -> EvaluationResult: |
279 | | - """Submit, poll, and download a remote evaluator plugin metric job.""" |
280 | | - normalized_params = resolve_params(params, target) |
281 | | - spec = _build_evaluate_spec( |
282 | | - metrics=metric, |
283 | | - dataset=dataset, |
284 | | - params=normalized_params, |
285 | | - target=target, |
286 | | - field_mapping=field_mapping, |
287 | | - prompt_template=prompt_template, |
288 | | - metric_bundle_packager=metric_bundle_packager, |
289 | | - ) |
290 | | - |
291 | | - job = self.create( |
292 | | - spec=spec, workspace=http_utils.resolve_workspace(self._platform, self._workspace, strict=True) |
293 | | - ) |
294 | | - job.wait_until_done( |
295 | | - poll_interval_seconds=self._poll_interval_seconds, |
296 | | - job_timeout_seconds=self._job_timeout_seconds, |
297 | | - pending_timeout_seconds=self._pending_timeout_seconds, |
298 | | - ) |
299 | | - |
300 | | - return job.get_result(aggregate_fields=aggregate_fields) |
301 | | - |
302 | | - def evaluate( |
303 | | - self, |
304 | | - *, |
305 | | - metric: Metric, |
306 | | - dataset: PluginDatasetInput, |
307 | | - params: RunConfig | RunConfigOnline | RunConfigOnlineModel | None = None, |
308 | | - target: Model | Agent | None = None, |
309 | | - field_mapping: FieldMapping | None = None, |
310 | | - prompt_template: str | dict[str, Any] | None = None, |
311 | | - aggregate_fields: tuple[AggregateFieldName, ...] | None = None, |
312 | | - ) -> EvaluationResult: |
313 | | - """Evaluate one metric through local plugin job execution.""" |
314 | | - normalized_params = resolve_params(params, target) |
315 | | - spec = _build_evaluate_spec( |
316 | | - metrics=metric, |
317 | | - dataset=dataset, |
318 | | - params=normalized_params, |
319 | | - target=target, |
320 | | - field_mapping=field_mapping, |
321 | | - prompt_template=prompt_template, |
322 | | - metric_bundle_packager=resolve_default_metric_bundle_packager( |
323 | | - metric, None, allow_cloudpickle_fallback=True, action="Running" |
324 | | - ), |
325 | | - ) |
326 | | - payload = self.run_local( |
327 | | - spec=spec, |
328 | | - workspace=http_utils.resolve_workspace(self._platform, self._workspace, strict=True), |
329 | | - ) |
330 | | - result_path = local_result_path(payload) |
331 | | - result = EvaluationResult.model_validate_json(result_path.read_text(encoding="utf-8")) |
332 | | - return filter_evaluation_result(result, aggregate_fields) |
333 | | - |
334 | 243 | def submit( |
335 | 244 | self, |
336 | 245 | *, |
@@ -361,38 +270,6 @@ def submit( |
361 | 270 |
|
362 | 271 | return job |
363 | 272 |
|
364 | | - def evaluate_benchmark( |
365 | | - self, |
366 | | - *, |
367 | | - metrics: Sequence[Metric], |
368 | | - dataset: PluginDatasetInput, |
369 | | - params: RunConfig | RunConfigOnline | RunConfigOnlineModel, |
370 | | - target: Model | Agent | None = None, |
371 | | - field_mapping: FieldMapping | None = None, |
372 | | - prompt_template: str | dict[str, Any] | None = None, |
373 | | - aggregate_fields: tuple[AggregateFieldName, ...] | None = None, |
374 | | - ) -> BenchmarkEvaluationResult: |
375 | | - """Evaluate multiple metrics through local plugin job execution.""" |
376 | | - normalized_params = resolve_params(params, target) |
377 | | - spec = _build_evaluate_spec( |
378 | | - metrics=metrics, |
379 | | - dataset=dataset, |
380 | | - params=normalized_params, |
381 | | - target=target, |
382 | | - field_mapping=field_mapping, |
383 | | - prompt_template=prompt_template, |
384 | | - metric_bundle_packager=resolve_default_metric_bundle_packager( |
385 | | - metrics, None, allow_cloudpickle_fallback=True, action="Running" |
386 | | - ), |
387 | | - ) |
388 | | - payload = self.run_local( |
389 | | - spec=spec, |
390 | | - workspace=http_utils.resolve_workspace(self._platform, self._workspace, strict=True), |
391 | | - ) |
392 | | - result_path = local_result_path(payload) |
393 | | - result = BenchmarkEvaluationResult.model_validate_json(result_path.read_text(encoding="utf-8")) |
394 | | - return filter_benchmark_result(result, aggregate_fields) |
395 | | - |
396 | 273 |
|
397 | 274 | class _AsyncEvaluatorPluginExecutor: |
398 | 275 | """Async evaluator plugin executor used by the async SDK resource.""" |
@@ -449,26 +326,6 @@ async def create( |
449 | 326 | ) |
450 | 327 | return job_resource |
451 | 328 |
|
452 | | - async def run_local(self, *, spec: EvaluateRequestSpec, workspace: str | None = None) -> EvaluatorLocalRunResult: |
453 | | - """Run an evaluator plugin job locally without blocking the event loop.""" |
454 | | - resolved_workspace = http_utils.resolve_workspace(self._platform, workspace) |
455 | | - canonical_spec = await _resolve_async_local_spec( |
456 | | - spec, |
457 | | - platform=self._platform, |
458 | | - workspace=resolved_workspace, |
459 | | - ) |
460 | | - scheduler = NemoJobScheduler() |
461 | | - # Leverages programmatic dispatch as described in |
462 | | - # packages/nemo_platform_plugin/src/nemo_platform_plugin/docs/ARCHITECTURE.md#job-entry-point-keys |
463 | | - payload = await asyncio.to_thread( |
464 | | - scheduler.run_local, |
465 | | - EvaluateJob, |
466 | | - canonical_spec.model_dump(mode="json"), |
467 | | - workspace=resolved_workspace, |
468 | | - async_sdk=self._platform, |
469 | | - ) |
470 | | - return EvaluatorLocalRunResult.model_validate(payload) |
471 | | - |
472 | 329 | async def submit( |
473 | 330 | self, |
474 | 331 | *, |
@@ -499,107 +356,6 @@ async def submit( |
499 | 356 |
|
500 | 357 | return job |
501 | 358 |
|
502 | | - async def evaluate_remote( |
503 | | - self, |
504 | | - *, |
505 | | - metric: Metric, |
506 | | - dataset: PluginDatasetInput, |
507 | | - params: RunConfig | RunConfigOnline | RunConfigOnlineModel, |
508 | | - target: Model | Agent | None = None, |
509 | | - field_mapping: FieldMapping | None = None, |
510 | | - prompt_template: str | dict[str, Any] | None = None, |
511 | | - aggregate_fields: tuple[AggregateFieldName, ...] | None = None, |
512 | | - metric_bundle_packager: MetricBundlePackager | None = None, |
513 | | - ) -> EvaluationResult: |
514 | | - """Submit, poll, and download a remote evaluator plugin metric job.""" |
515 | | - normalized_params = resolve_params(params, target) |
516 | | - spec = _build_evaluate_spec( |
517 | | - metrics=metric, |
518 | | - dataset=dataset, |
519 | | - params=normalized_params, |
520 | | - target=target, |
521 | | - field_mapping=field_mapping, |
522 | | - prompt_template=prompt_template, |
523 | | - metric_bundle_packager=metric_bundle_packager, |
524 | | - ) |
525 | | - |
526 | | - job = await self.create( |
527 | | - spec=spec, workspace=http_utils.resolve_workspace(self._platform, self._workspace, strict=True) |
528 | | - ) |
529 | | - await job.wait_until_done( |
530 | | - poll_interval_seconds=self._poll_interval_seconds, |
531 | | - job_timeout_seconds=self._job_timeout_seconds, |
532 | | - pending_timeout_seconds=self._pending_timeout_seconds, |
533 | | - ) |
534 | | - |
535 | | - return await job.get_result(aggregate_fields=aggregate_fields) |
536 | | - |
537 | | - async def evaluate( |
538 | | - self, |
539 | | - *, |
540 | | - metric: Metric, |
541 | | - dataset: PluginDatasetInput, |
542 | | - params: RunConfig | RunConfigOnline | RunConfigOnlineModel | None = None, |
543 | | - target: Model | Agent | None = None, |
544 | | - field_mapping: FieldMapping | None = None, |
545 | | - prompt_template: str | dict[str, Any] | None = None, |
546 | | - aggregate_fields: tuple[AggregateFieldName, ...] | None = None, |
547 | | - ) -> EvaluationResult: |
548 | | - """Evaluate one metric through local plugin job execution.""" |
549 | | - normalized_params = resolve_params(params, target) |
550 | | - spec = _build_evaluate_spec( |
551 | | - metrics=metric, |
552 | | - dataset=dataset, |
553 | | - params=normalized_params, |
554 | | - target=target, |
555 | | - field_mapping=field_mapping, |
556 | | - prompt_template=prompt_template, |
557 | | - metric_bundle_packager=resolve_default_metric_bundle_packager( |
558 | | - metric, None, allow_cloudpickle_fallback=True, action="Running" |
559 | | - ), |
560 | | - ) |
561 | | - payload = await self.run_local( |
562 | | - spec=spec, |
563 | | - workspace=http_utils.resolve_workspace(self._platform, self._workspace, strict=True), |
564 | | - ) |
565 | | - result_path = local_result_path(payload) |
566 | | - result_text = await asyncio.to_thread(result_path.read_text, encoding="utf-8") |
567 | | - result = EvaluationResult.model_validate_json(result_text) |
568 | | - return filter_evaluation_result(result, aggregate_fields) |
569 | | - |
570 | | - async def evaluate_benchmark( |
571 | | - self, |
572 | | - *, |
573 | | - metrics: Sequence[Metric], |
574 | | - dataset: PluginDatasetInput, |
575 | | - params: RunConfig | RunConfigOnline | RunConfigOnlineModel, |
576 | | - target: Model | Agent | None = None, |
577 | | - field_mapping: FieldMapping | None = None, |
578 | | - prompt_template: str | dict[str, Any] | None = None, |
579 | | - aggregate_fields: tuple[AggregateFieldName, ...] | None = None, |
580 | | - ) -> BenchmarkEvaluationResult: |
581 | | - """Evaluate multiple metrics through local plugin job execution.""" |
582 | | - normalized_params = resolve_params(params, target) |
583 | | - spec = _build_evaluate_spec( |
584 | | - metrics=metrics, |
585 | | - dataset=dataset, |
586 | | - params=normalized_params, |
587 | | - target=target, |
588 | | - field_mapping=field_mapping, |
589 | | - prompt_template=prompt_template, |
590 | | - metric_bundle_packager=resolve_default_metric_bundle_packager( |
591 | | - metrics, None, allow_cloudpickle_fallback=True, action="Running" |
592 | | - ), |
593 | | - ) |
594 | | - payload = await self.run_local( |
595 | | - spec=spec, |
596 | | - workspace=http_utils.resolve_workspace(self._platform, self._workspace, strict=True), |
597 | | - ) |
598 | | - result_path = local_result_path(payload) |
599 | | - result_text = await asyncio.to_thread(result_path.read_text, encoding="utf-8") |
600 | | - result = BenchmarkEvaluationResult.model_validate_json(result_text) |
601 | | - return filter_benchmark_result(result, aggregate_fields) |
602 | | - |
603 | 359 |
|
604 | 360 | def bundle_metrics_for_spec( |
605 | 361 | metrics: Metric | Sequence[Metric], *, metric_bundle_packager: MetricBundlePackager |
|
0 commit comments