diff --git a/scripts/e2e_eval/run_eval.py b/scripts/e2e_eval/run_eval.py index 8d3a59eb7..9f7c0f45b 100644 --- a/scripts/e2e_eval/run_eval.py +++ b/scripts/e2e_eval/run_eval.py @@ -204,9 +204,7 @@ def _is_eval_target_available(ep: str | None, device: str | None) -> bool: raise WinMLEPNotDiscovered(f"No locally installed EP found for {target.ep}.") registry.auto_device(target) except (DeviceNotFound, WinMLEPNotDiscovered, WinMLEPRegistrationFailed) as exc: - safe_print( - f"[SKIP] {target.ep}/{target.device} is not available on this machine: {exc}" - ) + safe_print(f"[SKIP] {target.ep}/{target.device} is not available on this machine: {exc}") return False return True @@ -901,6 +899,7 @@ def _run_subprocess(args: list[str], timeout: int) -> dict: (log_dir / "stderr.log").open("w+b") as stderr_file, (log_dir / "events.jsonl").open("a", encoding="utf-8") as events_file, ): + def report(event: str, message: str, **details) -> None: record = { "timestamp": _utc_now(), @@ -976,9 +975,8 @@ def report(event: str, message: str, **details) -> None: last_download_progress = snapshot["last_progress"] download_active = snapshot["active"] download_state_known = True - if ( - now - last_download_observation >= _HF_DOWNLOAD_MONITOR_TIMEOUT - or (monitor is not None and monitor.poll() is not None) + if now - last_download_observation >= _HF_DOWNLOAD_MONITOR_TIMEOUT or ( + monitor is not None and monitor.poll() is not None ): download_monitor_stalled = True download_active = False @@ -1022,9 +1020,12 @@ def report(event: str, message: str, **details) -> None: last_output_progress = now last_output_sizes = output_sizes monitor_state = ( - "disabled" if download_monitor_stalled - else "starting" if not download_state_known - else "downloading" if download_active + "disabled" + if download_monitor_stalled + else "starting" + if not download_state_known + else "downloading" + if download_active else "idle" ) execution_remaining = max(0.0, timeout - execution_elapsed) @@ -1053,8 +1054,10 @@ def report(event: str, message: str, **details) -> None: else f"execution timeout ({timeout:g}s)" ) report( - "timeout", reason, - timeout=timed_out, hf_download_stalled=hf_download_stalled, + "timeout", + reason, + timeout=timed_out, + hf_download_stalled=hf_download_stalled, ) exit_code = -1 break @@ -1452,32 +1455,32 @@ def _run_recipe_build( def _extract_onnx_path(build_proc: dict, hf_id: str, task: str | None) -> str | None: - """Extract ONNX path from build subprocess output.""" + """Extract the final or reused ONNX path from build subprocess output.""" # Rich may wrap a long artifact path across physical output lines. Rejoin # those fragments before falling back to cache discovery. - markers = ("Final artifact:", "Existing artifact found:", "Artifact:") + markers = ("Final artifact:", "Existing artifact found:") output = re.sub( r"\x1b\[[0-?]*[ -/]*[@-~]", "", - build_proc["stderr"] + build_proc["stdout"], + "\n".join((build_proc["stderr"], build_proc["stdout"])), ) lines = output.splitlines() - for index, line in enumerate(lines): - for marker in markers: + for marker in markers: + for index, line in enumerate(lines): if marker not in line: continue - fragments = [line.split(marker, 1)[1].strip()] - for continuation in lines[index + 1 : index + 11]: - candidate = "".join(fragments) - if candidate and Path(candidate).is_file(): - return candidate - if candidate.lower().endswith(".onnx"): - break - fragments.append(continuation.strip()) - - candidate = "".join(fragments) - if candidate and Path(candidate).is_file(): - return candidate + fragments = [line.split(marker, 1)[1], *lines[index + 1 : index + 11]] + candidate = "" + for fragment in fragments: + candidate += fragment.strip() + if not candidate.lower().endswith(".onnx"): + continue + try: + if Path(candidate).is_file(): + return candidate + except (OSError, ValueError) as exc: + logging.debug("Skipping invalid ONNX candidate path %r: %s", candidate, exc) + break return _find_cached_model(hf_id, build_proc, task) @@ -2223,11 +2226,7 @@ def _record( def _is_finite_number(value: object) -> bool: - return ( - isinstance(value, (int, float)) - and not isinstance(value, bool) - and math.isfinite(value) - ) + return isinstance(value, (int, float)) and not isinstance(value, bool) and math.isfinite(value) def _validate_single_perf_result(result: dict, context: str = "result") -> str | None: @@ -2378,9 +2377,7 @@ def _single_physical_dml_gpu_luid() -> str: if adapter.device_type == "GPU" } if len(native_gpus) != 1: - raise RuntimeError( - f"DML CI pin requires exactly one DXCore GPU; found {list(native_gpus)}" - ) + raise RuntimeError(f"DML CI pin requires exactly one DXCore GPU; found {list(native_gpus)}") luid, adapter = next(iter(native_gpus.items())) selected = WinMLEPRegistry.instance().auto_device(EPDeviceTarget(ep="dml", device="gpu")) advertised = { @@ -2856,9 +2853,7 @@ def _run_update_baseline(entries: list[ModelEntry], args: argparse.Namespace) -> ds_config = get_dataset_config(entry.hf_id, entry.task) or {} cached = _lookup_baseline_cache(entry.hf_id, entry.task, ds_config) if cached is not None and not args.retry_failed: - safe_print( - f"{_progress_prefix(i, len(entries))} {label} (cached {cached['metric']})" - ) + safe_print(f"{_progress_prefix(i, len(entries))} {label} (cached {cached['metric']})") continue safe_print(f"{_progress_prefix(i, len(entries))} {label} running baseline ...") @@ -3026,9 +3021,7 @@ def _matches_hf_fetch_retry(existing: dict) -> bool: accuracy = existing.get("accuracy") or {} perf_failed = bool(perf) and not perf.get("passed") accuracy_failed = ( - bool(accuracy) - and not accuracy.get("skipped") - and accuracy_status(accuracy) != "PASS" + bool(accuracy) and not accuracy.get("skipped") and accuracy_status(accuracy) != "PASS" ) if not perf_failed and not accuracy_failed: return False @@ -3236,8 +3229,7 @@ def _build_jobs( # explicit per-model precision (e.g. fp16) skips this and is honored # by the single-fallback branch below via _resolve_precision. jobs.extend( - EvalJob(entry, None, fallback_precision=prec) - for prec in _NPU_FALLBACK_PRECISIONS + EvalJob(entry, None, fallback_precision=prec) for prec in _NPU_FALLBACK_PRECISIONS ) else: jobs.append(EvalJob(entry, None)) @@ -3381,10 +3373,7 @@ def parse_args() -> argparse.Namespace: choices=["P0", "P1", "P2", "P3"], default=["P0", "P1", "P2", "P3"], metavar="{P0,P1,P2,P3}", - help=( - "Filter by priority. Pass one or more, e.g. --priority P0 P1. " - "Default: P0 P1 P2 P3." - ), + help=("Filter by priority. Pass one or more, e.g. --priority P0 P1. Default: P0 P1 P2 P3."), ) parser.add_argument( "--release", @@ -3639,10 +3628,9 @@ def main() -> None: clean_cache_targets = _resolve_clean_cache_targets(args.clean_cache) args.clean_cache_targets = clean_cache_targets - if ( - not (args.list or args.list_json or args.update_baseline or args.build_only) - and not _is_eval_target_available(args.ep, args.device) - ): + if not ( + args.list or args.list_json or args.update_baseline or args.build_only + ) and not _is_eval_target_available(args.ep, args.device): return # 1. Load registry @@ -3922,9 +3910,7 @@ def main() -> None: ) if timeout_rule is not None: reason = timeout_rule.get("reason") or "timeout" - safe_print( - f"\n{_progress_prefix(i, total_jobs)} {label} (SKIP - TIMEOUT: {reason})" - ) + safe_print(f"\n{_progress_prefix(i, total_jobs)} {label} (SKIP - TIMEOUT: {reason})") model_dir.mkdir(parents=True, exist_ok=True) timeout_result = build_eval_result( entry=entry, @@ -3982,16 +3968,14 @@ def main() -> None: else "?" ) safe_print( - f"\n{_progress_prefix(i, total_jobs)} {label} " - f"(RETRY - was {retry_label})" + f"\n{_progress_prefix(i, total_jobs)} {label} (RETRY - was {retry_label})" ) except (json.JSONDecodeError, KeyError): pass # Corrupted result file — re-run if backfill_existing is None: safe_print( - f"\n{_progress_prefix(i, total_jobs)} {label} " - f"({entry.priority}, {entry.group})" + f"\n{_progress_prefix(i, total_jobs)} {label} ({entry.priority}, {entry.group})" ) try: diff --git a/src/winml/modelkit/commands/perf.py b/src/winml/modelkit/commands/perf.py index bb3be9534..40ecad86b 100644 --- a/src/winml/modelkit/commands/perf.py +++ b/src/winml/modelkit/commands/perf.py @@ -2522,8 +2522,6 @@ def _autobuild_genai_bundle( existing bundle (in which case task/precision were not applied to it). """ from ..cache import get_cache_dir, get_model_dir - from ..loader import resolve_loader_config - from ..models.winml import build_genai_bundle, resolve_genai_bundle from ..session import EPDeviceTarget, ep_to_device, resolve_device, short_ep_name from ..utils.constants import normalize_ep_name @@ -2553,6 +2551,9 @@ def _autobuild_genai_bundle( console.print(f"[dim]Reusing cached genai bundle:[/dim] {bundle_dir}") return bundle_dir, False + from ..loader import resolve_loader_config + from ..models.winml import build_genai_bundle, resolve_genai_bundle + # Cache miss (or forced rebuild): resolve the model family so its # genai-bundle recipe can drive the build. try: diff --git a/tests/e2e/test_perf_e2e.py b/tests/e2e/test_perf_e2e.py index ffd1b7aea..4610cfd79 100644 --- a/tests/e2e/test_perf_e2e.py +++ b/tests/e2e/test_perf_e2e.py @@ -1549,6 +1549,8 @@ def genai_bundle(self, tmp_path_factory: pytest.TempPathFactory) -> Path: proc = subprocess.run( # noqa: S603 -- trusted args (sys.executable + constants) cmd, capture_output=True, + encoding="utf-8", + errors="replace", text=True, timeout=1500, check=False, diff --git a/tests/unit/commands/test_perf_genai.py b/tests/unit/commands/test_perf_genai.py index bee4593d6..ab5674ed2 100644 --- a/tests/unit/commands/test_perf_genai.py +++ b/tests/unit/commands/test_perf_genai.py @@ -13,6 +13,7 @@ from __future__ import annotations import json +import sys from io import StringIO from pathlib import Path from types import SimpleNamespace @@ -1664,26 +1665,29 @@ def test_autobuild_rejects_unsupported_recipe_target( builder.assert_not_called() assert "config" not in capture_run + @pytest.mark.parametrize( + ("arguments", "directory"), + [([], "genai-bundle"), (["--ep", "CPUExecutionProvider"], "genai-bundle-cpu-cpu")], + ) def test_autobuild_reuses_cached_bundle( - self, runner: CliRunner, tmp_path: Path, capture_run: dict, monkeypatch + self, runner: CliRunner, tmp_path: Path, capture_run: dict, monkeypatch, + arguments, directory, ) -> None: - import winml.modelkit.models.winml as winml_models from winml.modelkit.cache import get_model_dir monkeypatch.setenv("WINML_CACHE_DIR", str(tmp_path)) - cached = get_model_dir("Qwen/Qwen3-0.6B", cache_dir=tmp_path) / "genai-bundle" + cached = get_model_dir("Qwen/Qwen3-0.6B", cache_dir=tmp_path) / directory cached.mkdir(parents=True) (cached / "genai_config.json").write_text("{}", encoding="utf-8") - build_calls: dict = {} - monkeypatch.setattr( - winml_models, "build_genai_bundle", _fake_build_genai_bundle(build_calls) + monkeypatch.setitem(sys.modules, "winml.modelkit.loader", None) + monkeypatch.setitem(sys.modules, "winml.modelkit.models.winml", None) + result = runner.invoke( + perf, ["-m", "Qwen/Qwen3-0.6B", "--runtime", "ort-genai", *arguments] ) - result = runner.invoke(perf, ["-m", "Qwen/Qwen3-0.6B", "--runtime", "ort-genai"]) - assert result.exit_code == 0, result.output - assert "build" not in build_calls # cache hit: never rebuilt + assert "Reusing cached genai bundle" in result.output assert capture_run["config"].bundle_dir == cached def test_rebuild_forces_autobuild( diff --git a/tests/unit/eval/test_run_eval_script.py b/tests/unit/eval/test_run_eval_script.py index 96a61c90a..448e4f0dd 100644 --- a/tests/unit/eval/test_run_eval_script.py +++ b/tests/unit/eval/test_run_eval_script.py @@ -80,6 +80,7 @@ def _deterministic_ep_deduction(run_eval): deduction patch it themselves. """ run_eval._deduce_ep_for_device.cache_clear() + def resolve_target(ep, device): from winml.modelkit.session import default_device_for_ep, expand_ep_name @@ -204,8 +205,13 @@ def test_automatic_axes_keep_existing_policy(self, run_eval, ep, device): def test_main_skips_before_loading_models_or_creating_output(self, run_eval, tmp_path, release): output_dir = tmp_path / "results" argv = [ - "run_eval.py", "--ep", "openvino", "--device", "npu", - "--output-dir", str(output_dir), + "run_eval.py", + "--ep", + "openvino", + "--device", + "npu", + "--output-dir", + str(output_dir), ] if release: argv.append("--release") @@ -278,8 +284,7 @@ def classifier(self, run_eval): ) def test_hf_streaming_failures_are_retryable(self, classifier, output): assert ( - classifier.classify_failure(output, exit_code=1) - is classifier.FailureType.HF_FETCH_FAIL + classifier.classify_failure(output, exit_code=1) is classifier.FailureType.HF_FETCH_FAIL ) @pytest.mark.parametrize( @@ -311,9 +316,7 @@ def test_explicit_hf_fetch_retry_markers(self, classifier, output): "[WinError 10061] connection refused", ], ) - def test_other_hf_failure_patterns_are_not_explicit_retry_markers( - self, classifier, output - ): + def test_other_hf_failure_patterns_are_not_explicit_retry_markers(self, classifier, output): assert classifier.matches_hf_fetch_retry(output) is False @@ -696,7 +699,8 @@ def test_download_progress_timestamp_includes_scan_time(self, run_eval, tmp_path with ( patch.object( - run_eval, "_process_tree_open_paths", + run_eval, + "_process_tree_open_paths", return_value={run_eval._normalized_path(incomplete)}, ), patch.object(run_eval, "_snapshot_hf_downloads", return_value={incomplete: (1, 1)}), @@ -815,8 +819,7 @@ def replace_when_unlocked(source, target): def test_thread_start_cannot_block_subprocess_timeout(self, run_eval, tmp_path): original_start = threading.Thread.start script = ( - "import sys, time; sys.stderr.write('x' * 262144); " - "sys.stderr.flush(); time.sleep(5)" + "import sys, time; sys.stderr.write('x' * 262144); sys.stderr.flush(); time.sleep(5)" ) def delayed_start(thread): @@ -882,12 +885,7 @@ def test_inherited_output_handles_do_not_delay_collection(self, run_eval, tmp_pa def test_curated_target_models_preserve_existing_priorities(run_eval): - testsets_dir = ( - Path(__file__).resolve().parents[3] - / "scripts" - / "e2e_eval" - / "testsets" - ) + testsets_dir = Path(__file__).resolve().parents[3] / "scripts" / "e2e_eval" / "testsets" curated = json.loads((testsets_dir / "models_curated.json").read_text(encoding="utf-8")) target_entries = curated[-43:] target_keys = {(entry["hf_id"], entry["task"]) for entry in target_entries} @@ -907,9 +905,7 @@ def test_curated_target_models_preserve_existing_priorities(run_eval): assert len(generated) == 43 assert {entry["group"] for entry in target_entries} == {"Top200"} assert { - (entry["hf_id"], entry["task"]) - for entry in target_entries - if entry["priority"] == "P2" + (entry["hf_id"], entry["task"]) for entry in target_entries if entry["priority"] == "P2" } == expected_p2 assert sum(entry["priority"] == "P3" for entry in target_entries) == 37 assert { @@ -1225,20 +1221,62 @@ def test_run_build_uses_composite_onnx_without_subprocess(self, run_eval, tmp_pa class TestExtractOnnxPath: - def test_rejoins_rich_wrapped_artifact_path(self, run_eval, tmp_path): + @pytest.mark.parametrize("width", [80, 240]) + @pytest.mark.parametrize("with_final", [True, False]) + def test_ignores_stage_artifacts(self, run_eval, tmp_path, monkeypatch, width, with_final): + from rich.console import Console + + from winml.modelkit.utils.console import StageLive, print_final + + export = tmp_path / "model_export.onnx" + artifact = tmp_path / "model_model.onnx" + export.touch() + artifact.touch() + console = Console(width=width, force_terminal=False, color_system=None) + with console.capture() as capture: + with StageLive("export", console) as stage: + stage.artifact(str(export), export.stat().st_size) + stage.set_done(0.1) + console.print("WARNING Model producer not matched: Expected pytorch") + with StageLive("optimize", console) as stage: + stage.artifact(str(artifact), artifact.stat().st_size) + stage.set_done(0.1) + if with_final: + print_final(console, 0.2, str(artifact)) + + original_is_file = run_eval.Path.is_file + + def check_file(path): + if "WARNING" in str(path): + raise OSError( + 22, + "No mapping for the Unicode character exists " + "in the target multi-byte code page", + str(path), + 1113, + ) + return original_is_file(path) + + monkeypatch.setattr(run_eval.Path, "is_file", check_file) + + assert run_eval._extract_onnx_path( + {"stderr": capture.get(), "stdout": ""}, "test/model", None + ) == (str(artifact) if with_final else None) + + @pytest.mark.parametrize("marker", ["Final artifact:", "Existing artifact found:"]) + @pytest.mark.parametrize("stream", ["stderr", "stdout"]) + def test_rejoins_rich_wrapped_artifact_path(self, run_eval, tmp_path, marker, stream): artifact = tmp_path / "model-with-a-long-name_model.onnx" artifact.touch() path = str(artifact) split_at = len(path) - 12 - build_proc = { - "stderr": ( - "Existing artifact found:\n" - f"{path[:split_at]}\n" - f"{path[split_at:]}\n" - "Use --rebuild to force rebuild.\n" - ), - "stdout": "", - } + build_proc = {"stderr": "", "stdout": ""} + build_proc[stream] = ( + f"\x1b[36m{marker}\x1b[0m\n" + f"\x1b[1m{path[:split_at]}\n" + f"{path[split_at:]}\x1b[0m\n" + "Use --rebuild to force rebuild.\n" + ) assert ( run_eval._extract_onnx_path( @@ -1249,9 +1287,45 @@ def test_rejoins_rich_wrapped_artifact_path(self, run_eval, tmp_path): == path ) - def test_cache_fallback_rejects_multiple_task_candidates( - self, run_eval, tmp_path, monkeypatch + @pytest.mark.parametrize("cached_first", [True, False]) + def test_prefers_final_artifact_and_separates_output_streams( + self, run_eval, tmp_path, cached_first ): + cached = tmp_path / "cached_model.onnx" + artifact = tmp_path / "final_model.onnx" + cached.touch() + artifact.touch() + stderr = f"Existing artifact found: {cached}\n" if cached_first else "" + stderr += f"Final artifact: {artifact}" + + assert run_eval._extract_onnx_path( + {"stderr": stderr, "stdout": "Build completed successfully.\n"}, + "test/model", + None, + ) == str(artifact) + + @pytest.mark.parametrize("error_type", [OSError, ValueError]) + @pytest.mark.parametrize("cache_exists", [True, False]) + def test_invalid_artifact_path_uses_cache_fallback( + self, run_eval, tmp_path, monkeypatch, error_type, cache_exists + ): + cache_dir = tmp_path / ".cache" / "winml" / "artifacts" / "test_model" + cached = cache_dir / "imgcls_valid_model.onnx" + if cache_exists: + cache_dir.mkdir(parents=True) + cached.touch() + monkeypatch.setattr(run_eval.Path, "home", lambda: tmp_path) + build_proc = { + "stderr": f"Final artifact: {tmp_path / 'invalid.onnx'}\n", + "stdout": "", + } + + with patch.object(run_eval.Path, "is_file", side_effect=error_type("invalid path")): + result = run_eval._extract_onnx_path(build_proc, "test/model", "image-classification") + + assert result == (str(cached) if cache_exists else None) + + def test_cache_fallback_rejects_multiple_task_candidates(self, run_eval, tmp_path, monkeypatch): cache_dir = tmp_path / ".cache" / "winml" / "artifacts" / "microsoft_beit" cache_dir.mkdir(parents=True) (cache_dir / "imgcls_fp32_model.onnx").touch() @@ -1466,9 +1540,7 @@ def test_missing_quant_section_is_none(self, run_eval, tmp_path): # verbatim rather than inferred as fp32 (the graph's own dtype is not # something the config states). assert run_eval._precision_from_build_config(self._write(tmp_path, {})) is None - assert ( - run_eval._precision_from_build_config(self._write(tmp_path, {"quant": None})) is None - ) + assert run_eval._precision_from_build_config(self._write(tmp_path, {"quant": None})) is None @pytest.mark.parametrize( ("weight_type", "activation_type", "expected"), @@ -2229,9 +2301,7 @@ def fake_run(args, timeout): assert proc["result"] is None assert "Invalid structured winml perf output" in proc["stderr"] - def test_perf_result_with_embedded_op_trace_is_copied_before_cleanup( - self, run_eval, tmp_path - ): + def test_perf_result_with_embedded_op_trace_is_copied_before_cleanup(self, run_eval, tmp_path): trace_result = {"status": "ok", "operators": []} perf_result = { **_perf_result(), @@ -2456,9 +2526,7 @@ def test_non_npu_keeps_only_non_quantized_variants(self, run_eval, tmp_path): def test_non_npu_recipe_only_quantized_falls_back(self, run_eval, tmp_path): # A recipe with no non-quantized variant leaves nothing to run off-NPU, # so the model builds a single winml-config fallback. - self._make_single_recipe( - tmp_path, "microsoft_resnet-50", "image-classification", ["w8a16"] - ) + self._make_single_recipe(tmp_path, "microsoft_resnet-50", "image-classification", ["w8a16"]) entry = _entry() jobs = run_eval._build_jobs([entry], tmp_path, "cpu") assert len(jobs) == 1 @@ -2508,9 +2576,7 @@ def test_npu_skip_quant_ep_drops_quantized_recipe_variants(self, run_eval, tmp_p def test_npu_skip_quant_ep_recipe_only_quantized_falls_back(self, run_eval, tmp_path): # Dropping every quantized variant leaves nothing to build, so the model # goes through the single unquantized winml-config fallback. - self._make_single_recipe( - tmp_path, "microsoft_resnet-50", "image-classification", ["w8a16"] - ) + self._make_single_recipe(tmp_path, "microsoft_resnet-50", "image-classification", ["w8a16"]) entry = _entry() jobs = run_eval._build_jobs([entry], tmp_path, "npu", ep="vitisai") assert len(jobs) == 1 @@ -3297,7 +3363,8 @@ def test_checked_in_recipe_filenames_match_quant_config(run_eval): for config_path in config_paths: config = json.loads(config_path.read_text(encoding="utf-8")) expected = ( - "fp32" if config["quant"] is None + "fp32" + if config["quant"] is None else run_eval._precision_from_build_config(config_path) ) assert expected is not None, f"Unknown quant configuration: {config_path}" @@ -3316,8 +3383,11 @@ def test_release_keeps_renamed_recipe_configuration(run_eval, tmp_path): ) entry = next(entry for entry in entries if entry.hf_id == "microsoft/resnet-50") jobs = run_eval._build_jobs( - [entry], testsets.parents[2] / "examples" / "recipes", - "gpu", ep="openvino", release=True, + [entry], + testsets.parents[2] / "examples" / "recipes", + "gpu", + ep="openvino", + release=True, ) (job,) = jobs assert job.precision == "fp32" @@ -3330,7 +3400,8 @@ def test_release_keeps_renamed_recipe_configuration(run_eval, tmp_path): args = argparse.Namespace(ep="openvino", device="gpu", timeout=300) with ( patch.object( - run_eval, "_run_subprocess", + run_eval, + "_run_subprocess", return_value={"exit_code": 0, "stdout": "", "stderr": ""}, ) as subprocess_call, patch.object(run_eval, "_extract_onnx_path", return_value=str(tmp_path / "model.onnx")),