diff --git a/docs/commands/perf.md b/docs/commands/perf.md index 37c6efcea..7abe93dd0 100644 --- a/docs/commands/perf.md +++ b/docs/commands/perf.md @@ -53,6 +53,81 @@ Both runtime reports include `schema_version: 2` and a `benchmark_info.runtime` When `--memory` is enabled, both `winml-ort` and `ort-genai` reports use the same `memory` field names for shared concepts: RSS baseline, after-compile/load, after-inference, peak, model-load delta, inference/generation delta, and total delta; VRAM local/shared baseline, after-compile/load, after-inference, peak, model-load delta, inference/generation delta, and total delta. +### Classic memory lifecycle + +#### Model loading from a user's perspective + +`load_memory` is the load-only view (block `version: 1`). It answers how much +additional memory was observed while this artifact became ready in this +runtime, and how much remained at readiness. Runtime/EP setup precedes its +baseline; its endpoint is after session compilation but before allocating +benchmark inputs, warmup or inference. If the selected path requires online +compilation, that temporary cost is included because the user must get through +it to load the model. A precompiled artifact is a different configuration. + +Each RSS/local/shared block records baseline, absolute peak, absolute ready +value, `peak_extra_mb` and `ready_extra_mb`, plus the peak method and availability. +RSS uses the Windows high-water when a **new** process high-water occurs inside +the load window; otherwise only observed load samples/endpoints can be used. +GPU peaks are sampled. No inference peak or larger pre-existing process peak +is substituted into this view. Unknown baselines produce null increments. + +This is an **observed load footprint**, not a certified minimum device-memory +capacity: RSS is resident memory, allocator caches can remain, sampling may +miss short GPU allocations, and unified RAM/GPU counters may overlap. Report +the load peak and ready increment separately; do not reduce this to parameter +file size or sum local/shared GPU counters. Already-loaded composite children +cannot provide this window and return an explicit unavailable status. + +The enclosing perf document still uses schema version 2; the optional +`load_memory` block is explicitly versioned. Consumers must support nullable +memory values and use the recorded scope rather than comparing older deltas +as if they covered the same interval. + +For `winml-ort` and `winml-runtime`, RSS baseline is captured **before model +loading**, including eager session construction and Runtime pipeline creation. +GPU baseline is captured after resolving the bound EP/adapter but before +constructing the model. Adapter discovery must not read lazy model properties. +The added `*_after_load_mb` checkpoint separates loading from the explicit +compile step. Input generation follows compilation/readiness. `*_after_compile_mb` and `*_after_inference_mb` +retain their existing names. + +`memory_measurement` records PID, adapter LUID, lifecycle scope, units and +per-checkpoint missing reasons. Fields with the legacy `_mb` suffix use MiB +(bytes / 1,048,576). Deltas are signed process changes, not model-only allocation; +negative values can reflect released buffers or working-set changes. Library +imports, build caches and input preparation can contribute to the measured span. + +`process_memory` separately records sampled RSS/local/shared peaks, actual +sample counts, interval and duration over loading, compilation and inference. +These are sampled peaks, not an exact continuous maximum. The configured +polling delay is not a sampling frequency: `configured_poll_delay_sec` records +the wait after each observation. `observed_mean_interval_sec` and +`observed_max_interval_sec` record elapsed time between completed observations +(null with fewer than two observations). GPU discovery/retries also delay RSS +sampling; these intervals include that overhead and scheduling delays. The legacy +`*_checkpoint_peak_mb` remains the maximum of the three original checkpoints, +and is `null` if any required checkpoint is unavailable. + +On Windows, `process_memory.os_rss_peak_before_mb` and +`os_rss_peak_after_mb` also record the OS working-set high-water marks. These +catch short allocations missed by polling, including native code holding the +Python GIL. They cover the **process lifetime**, cannot be reset at baseline, +and do not replace the benchmark-scoped sampled peak. On unsupported systems +they are `null` with an explicit reason. + +Missing GPU counters produce `null`, not zero; deltas with missing endpoints +also remain `null`. A GPU process-memory instance may not exist before model +allocation, so valid later absolute readings do not establish a zero baseline. +Local/shared counters can overlap on unified-memory devices and must not be +summed. GenAI's shared GPU-memory consumer also preserves unavailable values. + +Composite components are already loaded when measured. Their explicit scope +is `already_loaded_component_through_inference`; their values include the +shared process and do not claim an independent component-load footprint. +`--no-memory` disables this collector. Background sampling adds overhead and +is separate from the native inference timer. + With `--runtime ort-genai`, `winml perf` benchmarks the onnxruntime-genai decoder pipeline rather than a single `session.run()`. The JSON report uses a phase-based schema: `load` contains startup spans, `requests` contains one warmup or timed generation sample per request, `aggregate` summarizes timed requests only, `memory` contains optional RAM/VRAM deltas, and `hw_monitor` contains optional monitor output. The optional `memory` and `hw_monitor` top-level names match the classic `winml-ort` perf report; GenAI keeps `load`/`requests`/`aggregate` instead of classic `latency_ms`/`throughput` because generation has distinct prompt, first-token, and decode phases. For model-ID auto-builds, the selected EP/device must be supported by the model's @@ -79,34 +154,19 @@ target validation. ### Memory measurement contract -With --memory, the legacy baseline stays **after the model factory, before -input generation and explicit session.compile()**, as on main. Existing -baseline/load/inference/total delta fields and the maximum-of-three checkpoint -peak keep that boundary. Eager model loading before this baseline is excluded. - -Single-model runs also take an earlier, separately named before_model_load -snapshot after device resolution and before the model factory. For each of -rss, vram_local and vram_shared, additive fields are: - -- *_before_model_load_mb: earlier absolute snapshot. -- *_model_factory_delta_mb: legacy baseline minus earlier snapshot. -- *_total_from_before_model_load_delta_mb: inference end minus earlier snapshot. +With `--memory`, classic and Runtime runs use the lifecycle boundaries described +above: process RSS starts before model construction, while the load-only window +starts after runtime/device setup and ends before benchmark input allocation. +Preloaded composite components cannot claim a load baseline; their `load_memory` +block is unavailable and their process scope explicitly starts after loading. -The new total includes model factory/build and input/session overhead. It is -not weights-only memory or a continuous peak. Preloaded composite components -have no observation before loading: all added fields are null with an explicit -reason. They must not inherit the parent's aggregate baseline. - -Legacy *_mb fields use MiB. memory_measurement schema_version 2 retains the -legacy baseline definition and adds the earlier boundary definition, PID, -process creation time, selected LUID and timestamped byte/status/source records. -Private commit is separate from RSS. The legacy checkpoint peak excludes the -new earlier snapshot, even if that snapshot is larger. Signed deltas can be negative. - -Unavailable readings and dependent deltas are null, never zero. If the earlier -GPU process instance is absent, only metrics needing that point are unavailable; -a valid legacy baseline and its deltas remain usable. CPU GPU memory is -not_applicable. Performance success does not certify memory. +The outer result remains `schema_version: 2`; `load_memory.version` is 1. +Classic `memory_measurement.schema_version` is now 3 (replacing the version 2 +byte-checkpoint structure). It describes PID, selected LUID, lifecycle scope and missing +counter reasons. The old after-factory baseline and its additive +`*_before_model_load_mb` fields are superseded by these explicit blocks. +Signed deltas can be negative. Missing readings and dependent deltas stay null. +Performance success does not certify memory requirements. GPU memory uses main's effective EP-device binding, including --device-luid and provider selectors resolved through the advertised device options. Unresolved @@ -121,6 +181,28 @@ not proof of zero; enumeration errors and invalid readings are separate states. The CLI cannot certify that a process has never used the GPU merely from an absent PDH instance, so it never fabricates a zero baseline from that condition. +For a Windows GPU target with a resolved LUID, if process memory instances are +absent before model loading, the CLI creates a model-free D3D12 device on that +exact adapter and retries counter discovery for up to two seconds. It creates +no model, buffers or command queue and submits no GPU work. The device is kept +alive through the final checkpoint and released on success or failure, so its +setup footprint is present in both endpoints rather than appearing as model +allocation. Load-only RAM is captured after this preparation too; process RSS +still starts before generic runtime setup. Already valid counters, +CPU/NPU targets, unresolved adapters and preloaded components do not create this +device. Initialization errors never prevent the model's own provider from running. + +This is a prepared-device measurement, not cold GPU initialization. The +process delta includes model/provider/input allocations and is +not a minimum VRAM requirement. The load-only delta excludes inputs and inference. +Legacy-named baseline and delta fields use the new boundaries documented above; +they must not be compared directly with version 2 memory measurements. The optional +memory_measurement.gpu_baseline_preparation block records method, timing, +initial/prepared observations, selected adapter and failures. If the subsequent +checkpoint remains unavailable, dependent deltas remain null. A later valid +checkpoint never replaces a missing pre-model baseline. Historical results +cannot be repaired without recollection. + The hardware monitor refreshes PID/LUID memory instances every 200 ms, including when monitoring starts before model load. Counter registration is deduplicated, failed registrations retry, and missing/disappeared instances stay unknown. diff --git a/src/winml/modelkit/commands/_perf_genai.py b/src/winml/modelkit/commands/_perf_genai.py index 5c64f06fd..785e19d8c 100644 --- a/src/winml/modelkit/commands/_perf_genai.py +++ b/src/winml/modelkit/commands/_perf_genai.py @@ -187,14 +187,14 @@ def _round_stats(values: dict[str, float]) -> dict[str, float]: def _get_rss_mb() -> float: """Return current process RSS in MB.""" - from ..session.monitor.memory_tracker import get_rss_mb + from ..session.monitor import get_rss_mb return get_rss_mb() def _get_vram_mb(adapter_luid: str | None) -> tuple[float | None, float | None]: """Return current process device-memory usage as local/shared MB.""" - from ..session.monitor.memory_tracker import get_vram_mb + from ..session.monitor import get_vram_mb return get_vram_mb(adapter_luid) diff --git a/src/winml/modelkit/commands/perf.py b/src/winml/modelkit/commands/perf.py index 40ecad86b..98061bb83 100644 --- a/src/winml/modelkit/commands/perf.py +++ b/src/winml/modelkit/commands/perf.py @@ -30,7 +30,6 @@ from rich.markup import escape from rich.table import Table -from ..session.monitor.memory_tracker import MemoryTracker from ..utils import cli as cli_utils from ..utils.console import SafeConsole from ..utils.constants import ( @@ -63,6 +62,7 @@ from ..models.winml.base import WinMLPreTrainedModel from ..models.winml.composite_model import WinMLCompositeModel from ..session import EPDeviceTarget, WinMLDevice, WinMLEPDevice + from ..session.monitor import ProcessMemoryTracker from ..session.monitor.ep_monitor import WinMLEPMonitor from ..session.monitor.op_metrics import TraceFallbackReason from ..session.stats import PerfStats @@ -661,6 +661,8 @@ class BenchmarkResult: # Memory profile dict (rss deltas from memory_tracker) memory_profile: dict[str, float | None] | None = None memory_measurement: dict[str, Any] | None = None + process_memory: dict[str, Any] | None = None + load_memory: dict[str, Any] | None = None def to_dict(self) -> dict[str, Any]: """Convert to dictionary for JSON serialization.""" @@ -725,6 +727,10 @@ def to_dict(self) -> dict[str, Any]: result["memory"] = self.memory_profile if self.memory_measurement: result["memory_measurement"] = self.memory_measurement + if self.process_memory: + result["process_memory"] = self.process_memory + if self.load_memory: + result["load_memory"] = self.load_memory return result @@ -958,7 +964,10 @@ def __init__(self, config: BenchmarkConfig) -> None: self._ep_device: WinMLEPDevice | None = None self._effective_batch: int = config.batch_size self._memory: dict[str, float | None] | None = None - self._memory_tracker: MemoryTracker | None = None + self._memory_measurement: dict[str, Any] | None = None + self._process_memory: dict[str, Any] | None = None + self._memory_tracker: ProcessMemoryTracker | None = None + self._load_memory: dict[str, Any] | None = None self._runtime_backend = resolve_runtime_api_backend( config.runtime, config.model_id, config.backend ) @@ -1081,6 +1090,14 @@ def _single(self) -> WinMLPreTrainedModel: return cast("WinMLPreTrainedModel", self._model) def run(self) -> BenchmarkResult | dict[str, BenchmarkResult]: + """Run with the model-free baseline device retained through final measurements.""" + try: + return self._run_with_memory() + finally: + if self._memory_tracker is not None: + self._memory_tracker.close() + + def _run_with_memory(self) -> BenchmarkResult | dict[str, BenchmarkResult]: """Execute full benchmark pipeline. Returns: @@ -1090,34 +1107,45 @@ def run(self) -> BenchmarkResult | dict[str, BenchmarkResult]: ORT session, so each sub-model is benchmarked individually rather than timing the aggregate ``forward()`` pass. """ - # Capture before an eager model factory or lazy property can build. + self._memory = self._memory_measurement = self._process_memory = None + self._load_memory = None if self.config.memory: - self._resolve_device_ep() - self._start_memory( - "after_model_factory_before_inputs_and_explicit_compile; legacy baseline", - phase="before_model_load", - ) - - # [1] Load model (build pipeline: optimize, cache, etc.) - logger.info("Loading model: %s", self.config.model_id) - self._load_model() - assert self._model is not None + from ..session.monitor import ProcessMemoryTracker - if self._is_composite: - # Composite-ness is only known after _load_model, so this guard - # can't live with the up-front --module / --runtime checks. Without - # it, each sub-model's child benchmark calls load_input_data with a - # single .npz that can't match two different encoders' input names, - # surfacing as a re-wrapped "Sub-model '…' failed" RuntimeError - # instead of a clean up-front error. - if self.config.input_data is not None: - raise click.UsageError( - "--input-data is not supported for composite (dual-encoder) " - "models; each sub-model has its own inputs that a single " - ".npz cannot address." - ) - return self._run_sub_models() - return self._run_single() + self._memory_tracker = ProcessMemoryTracker() + self._memory_tracker.start() + try: + # [1] Load model (build pipeline: optimize, cache, etc.) + logger.info("Loading model: %s", self.config.model_id) + self._load_model() + if self._memory_tracker is not None: + self._memory_tracker.checkpoint("after_load") + assert self._model is not None + + if self._is_composite: + # Composite-ness is only known after _load_model, so this guard + # can't live with the up-front --module / --runtime checks. Without + # it, each sub-model's child benchmark calls load_input_data with a + # single .npz that can't match two different encoders' input names, + # surfacing as a re-wrapped "Sub-model '…' failed" RuntimeError + # instead of a clean up-front error. + if self.config.input_data is not None: + raise click.UsageError( + "--input-data is not supported for composite (dual-encoder) " + "models; each sub-model has its own inputs that a single " + ".npz cannot address." + ) + # Components were constructed together; do not attribute the parent + # model-load baseline to each already-loaded child. + if self._memory_tracker is not None: + self._memory_tracker.stop() + self._memory_tracker = None + return self._run_sub_models() + return self._run_single() + finally: + if self._memory_tracker is not None: + self._memory_tracker.stop() + self._memory_tracker = None def _run_sub_models(self) -> dict[str, BenchmarkResult]: """Benchmark each sub-model of a composite individually. @@ -1149,36 +1177,39 @@ def _run_sub_models(self) -> dict[str, BenchmarkResult]: return results def _run_single(self) -> BenchmarkResult: - """Benchmark the loaded single-session model. - - Returns: - BenchmarkResult with timing statistics - """ - import gc - - assert self._model is not None - + """Benchmark a single session and clean up memory tracking on failure.""" if self.config.memory and self._memory_tracker is None: - # Preloaded components cannot claim a pre-model baseline. - self._start_memory( - "preloaded_component_before_inputs; excludes existing models; process scope" - ) - elif self._memory_tracker is not None: - # Preserve the original baseline and all deltas derived from it. - gc.collect() - self._memory_tracker.capture("baseline") + from ..session.monitor import ProcessMemoryTracker - # [2] Generate inputs - logger.info("Generating benchmark inputs") - self._generate_inputs() + self._memory_tracker = ProcessMemoryTracker( + adapter_luid=self._resolve_adapter_luid(), + scope="already_loaded_component_through_inference", + ) + self._memory_tracker.start() + self._memory_tracker.checkpoint("after_load") + try: + return self._run_single_impl() + finally: + if self._memory_tracker is not None: + self._memory_tracker.stop() + self._memory_tracker = None + + def _run_single_impl(self) -> BenchmarkResult: + """Run with the lifecycle baseline captured before model construction.""" + assert self._model is not None # Compile session early so model.device is resolved for display with suppress_native_warnings(enabled=True): self._single._session.compile() if self._memory_tracker is not None: - gc.collect() - self._memory_tracker.capture("after_compile") + self._memory_tracker.record_model_ready() + self._load_memory = self._memory_tracker.load_memory() + self._memory_tracker.checkpoint("after_compile") + + # Inputs must not inflate the model load footprint. + logger.info("Generating benchmark inputs") + self._generate_inputs() # Pre-benchmark identity block (model + device sub-blocks). # opset is not currently extracted on this path; pass None. @@ -1218,8 +1249,11 @@ def _run_single(self) -> BenchmarkResult: stats = self._run_benchmark() if self._memory_tracker is not None: - self._memory_tracker.capture("after_inference") - self._memory = self._memory_tracker.profile() + self._memory_tracker.checkpoint("after_inference") + self._memory_tracker.stop() + self._memory = self._memory_tracker.to_dict() + self._memory_measurement = self._memory_tracker.metadata() + self._process_memory = self._memory_tracker.sampled() # [4] Collect results logger.info("Collecting results") @@ -1240,6 +1274,12 @@ def _load_model(self) -> None: # Resolve the concrete device + EP first so a bad combo fails fast, # before from_pretrained/from_onnx kick off the build pipeline. self._resolve_device_ep() + if self._memory_tracker is not None: + _, bound_device = _get_ep_device_binding(self._ep_device, self.config.ep_options) + self._memory_tracker.bind_before_load( + self._resolve_adapter_luid(), + device=bound_device or self._resolved_device or self.config.device or "auto", + ) assert self._ep_device is not None model_id = self.config.model_id @@ -1392,20 +1432,6 @@ def _resolve_adapter_luid(self) -> str | None: ) return bound_luid - def _start_memory(self, baseline: str, *, phase: str = "baseline") -> None: - import gc - - luid, bound_device = _get_ep_device_binding(self._ep_device, self.config.ep_options) - reason = None if luid else "selected_device_has_no_monitorable_luid" - self._memory_tracker = MemoryTracker( - luid, - baseline=baseline, - device=bound_device or self._resolved_device or self.config.device or "auto", - adapter_reason=reason, - ) - gc.collect() - self._memory_tracker.capture(phase) - def _run_benchmark(self) -> PerfStats: """Execute benchmark iterations with timing. @@ -1594,7 +1620,9 @@ def _collect_results(self, stats: PerfStats) -> BenchmarkResult: hw_monitor=getattr(self, "_hw_metrics", None), # Memory profile (only present when --memory is used) memory_profile=self._memory, - memory_measurement=self._memory_tracker.evidence() if self._memory_tracker else None, + memory_measurement=self._memory_measurement, + process_memory=self._process_memory, + load_memory=self._load_memory, ) @@ -2087,6 +2115,24 @@ def display_console_report(result: BenchmarkResult, console: SafeConsole) -> Non console.print(f" CPU: {cpu.get('mean_pct', 0):.1f}% avg | RAM: {ram_text} MiB") # Memory section (only when --memory is enabled) + if result.load_memory and result.load_memory.get("status") == "measured": + console.print() + console.print("[bold]Model loading memory (observed):[/bold]") + for family, label in ( + ("rss", "RAM (RSS)"), + ("vram_local", "GPU local"), + ("vram_shared", "GPU shared"), + ): + observation = result.load_memory[family] + peak, ready = observation["peak_extra_mb"], observation["ready_extra_mb"] + peak_text = "N/A" if peak is None else f"{peak:.1f}" + ready_text = "N/A" if ready is None else f"{ready:+.1f}" + console.print( + f" {label}: peak extra {peak_text} MiB | ready net change {ready_text} MiB" + ) + console.print( + " Includes this path's required load/compile work; excludes inputs/inference." + ) if result.memory_profile: mem = result.memory_profile console.print() @@ -2094,34 +2140,25 @@ def display_console_report(result: BenchmarkResult, console: SafeConsole) -> Non def memory_value(key: str, *, signed: bool = False) -> str: value = mem.get(key) - if value is None: - return "unavailable" - return f"{value:+.1f}" if signed else f"{value:.1f}" + return "N/A" if value is None else format(value, "+.1f" if signed else ".1f") console.print( - f" RAM: {memory_value('rss_after_inference_mb')} MiB | " - f"total delta: {memory_value('rss_total_delta_mb', signed=True)} MiB" + f" RAM (RSS): {memory_value('rss_after_inference_mb')} MiB | " + f"load/build net change: {memory_value('rss_model_load_delta_mb', signed=True)} MiB | " + f"inference net change: {memory_value('rss_inference_delta_mb', signed=True)} MiB | " + f"total net change: {memory_value('rss_total_delta_mb', signed=True)} MiB" ) console.print( - f" GPU local/shared: {memory_value('vram_local_after_inference_mb')}/" - f"{memory_value('vram_shared_after_inference_mb')} MiB | total delta: " - f"{memory_value('vram_local_total_delta_mb', signed=True)}/" + f" GPU memory L/S: {memory_value('vram_local_after_inference_mb')}/" + f"{memory_value('vram_shared_after_inference_mb')} MiB | " + f"net change: {memory_value('vram_local_total_delta_mb', signed=True)}/" f"{memory_value('vram_shared_total_delta_mb', signed=True)} MiB" ) - if "rss_before_model_load_mb" in mem: - console.print( - " From before model load (additional): RAM delta " - f"{memory_value('rss_total_from_before_model_load_delta_mb', signed=True)} MiB; " - "GPU local/shared delta " - f"{memory_value('vram_local_total_from_before_model_load_delta_mb', signed=True)}/" - f"{memory_value('vram_shared_total_from_before_model_load_delta_mb', signed=True)}" - " MiB" - ) if result.memory_measurement: - console.print(f" Baseline: {escape(result.memory_measurement['baseline'])}") - console.print( - " Phase snapshots; unavailable counters are not zero. Not a continuous peak." - ) + console.print(f" Memory scope: {result.memory_measurement['scope']}") + missing = result.memory_measurement.get("missing_reasons", {}) + if any(missing.values()): + console.print(" N/A memory counters: see memory_measurement.missing_reasons") console.print() diff --git a/src/winml/modelkit/session/monitor/__init__.py b/src/winml/modelkit/session/monitor/__init__.py index 0674a39af..28e16ca35 100644 --- a/src/winml/modelkit/session/monitor/__init__.py +++ b/src/winml/modelkit/session/monitor/__init__.py @@ -11,6 +11,7 @@ if TYPE_CHECKING: from .ep_monitor import EPMonitor, NullEPMonitor, WinMLEPMonitor + from .memory_tracker import ProcessMemoryTracker, get_rss_mb, get_vram_mb from .op_metrics import OperatorMetrics, OpTraceResult from .openvino_monitor import OpenVinoMonitor from .report import display_op_trace_report, write_op_trace_json @@ -18,6 +19,9 @@ _LAZY_IMPORTS: dict[str, tuple[str, str]] = { + "ProcessMemoryTracker": (".memory_tracker", "ProcessMemoryTracker"), + "get_rss_mb": (".memory_tracker", "get_rss_mb"), + "get_vram_mb": (".memory_tracker", "get_vram_mb"), "EPMonitor": (".ep_monitor", "EPMonitor"), "NullEPMonitor": (".ep_monitor", "NullEPMonitor"), "WinMLEPMonitor": (".ep_monitor", "WinMLEPMonitor"), @@ -36,8 +40,11 @@ "OpTraceResult", "OpenVinoMonitor", "OperatorMetrics", + "ProcessMemoryTracker", "WinMLEPMonitor", "display_op_trace_report", + "get_rss_mb", + "get_vram_mb", "write_op_trace_json", ] diff --git a/src/winml/modelkit/session/monitor/_gpu_baseline.py b/src/winml/modelkit/session/monitor/_gpu_baseline.py new file mode 100644 index 000000000..35ec4be22 --- /dev/null +++ b/src/winml/modelkit/session/monitor/_gpu_baseline.py @@ -0,0 +1,100 @@ +# ------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. +# -------------------------------------------------------------------------- + +"""Keep a model-free D3D12 device alive on the explicitly selected GPU.""" + +from __future__ import annotations + +import ctypes +import re +import sys +import uuid +from typing import Any, ClassVar + + +class _Luid(ctypes.Structure): + _fields_: ClassVar = [("low", ctypes.c_uint32), ("high", ctypes.c_int32)] + + +def _guid(value: str) -> ctypes.Array[ctypes.c_ubyte]: + return (ctypes.c_ubyte * 16).from_buffer_copy(uuid.UUID(value).bytes_le) + + +def _method(pointer: ctypes.c_void_p, index: int, result: Any, *args: Any) -> Any: + table = ctypes.cast(pointer, ctypes.POINTER(ctypes.POINTER(ctypes.c_void_p))).contents + return ctypes.WINFUNCTYPE(result, ctypes.c_void_p, *args)(table[index]) + + +def _check(status: int) -> None: + if status < 0: + raise OSError(f"D3D12 baseline device failed: HRESULT 0x{status & 0xFFFFFFFF:08X}") + + +class GpuBaselineDevice: + """A retained device, no model, resources, command queue or submitted GPU work.""" + + def __init__(self, luid: str) -> None: + match = re.fullmatch(r"0x([0-9a-f]{8})_0x([0-9a-f]{8})", luid, re.IGNORECASE) + if not match: + raise ValueError("Invalid GPU adapter LUID") + if sys.platform != "win32": + raise OSError("D3D12 baseline preparation requires Windows") + high, low = (int(part, 16) for part in match.groups()) + self._device = ctypes.c_void_p() + factory, adapter = ctypes.c_void_p(), ctypes.c_void_p() + dxgi = ctypes.WinDLL("dxgi") + self._d3d12 = ctypes.WinDLL("d3d12") + dxgi.CreateDXGIFactory1.argtypes = [ctypes.c_void_p, ctypes.POINTER(ctypes.c_void_p)] + dxgi.CreateDXGIFactory1.restype = ctypes.c_int32 + self._d3d12.D3D12CreateDevice.argtypes = [ + ctypes.c_void_p, + ctypes.c_uint, + ctypes.c_void_p, + ctypes.POINTER(ctypes.c_void_p), + ] + self._d3d12.D3D12CreateDevice.restype = ctypes.c_int32 + try: + # IDXGIFactory4::EnumAdapterByLuid, never adapter ordinal zero/fallback. + _check( + dxgi.CreateDXGIFactory1( + _guid("1bc6ea02-ef36-464f-bf0c-21ca39e5168a"), ctypes.byref(factory) + ) + ) + _check( + _method( + factory, + 26, + ctypes.c_int32, + _Luid, + ctypes.c_void_p, + ctypes.POINTER(ctypes.c_void_p), + )( + factory, + _Luid(low, ctypes.c_int32(high).value), + _guid("2411e7e1-12ac-4ccf-bd14-9798e8534dc0"), + ctypes.byref(adapter), + ) + ) + _check( + self._d3d12.D3D12CreateDevice( + adapter, + 0xB000, + _guid("189819f1-1db6-4b57-be54-1821339b85f7"), + ctypes.byref(self._device), + ) + ) + except Exception: + self.close() + raise + finally: + for pointer in (adapter, factory): + if pointer.value: + _method(pointer, 2, ctypes.c_ulong)(pointer) + + def close(self) -> None: + """Release only our device reference after the last benchmark checkpoint.""" + if self._device.value: + _method(self._device, 2, ctypes.c_ulong)(self._device) + self._device = ctypes.c_void_p() diff --git a/src/winml/modelkit/session/monitor/memory_tracker.py b/src/winml/modelkit/session/monitor/memory_tracker.py index fa1cdea8d..f434715fd 100644 --- a/src/winml/modelkit/session/monitor/memory_tracker.py +++ b/src/winml/modelkit/session/monitor/memory_tracker.py @@ -10,12 +10,16 @@ import math import os import sys +import threading import time -from typing import Any +from typing import TYPE_CHECKING, Any import psutil +if TYPE_CHECKING: + from ._gpu_baseline import GpuBaselineDevice + logger = logging.getLogger(__name__) MIB = 1024 * 1024 @@ -133,9 +137,70 @@ def __init__( self.checkpoints: dict[str, dict[str, Any]] = {} self.pid = os.getpid() self.process_created_at = psutil.Process(self.pid).create_time() + self._gpu_baseline_device: GpuBaselineDevice | None = None + self.gpu_baseline_preparation: dict[str, Any] = {"status": "not_attempted"} + + def _prepare_gpu_baseline(self) -> None: + """Materialize missing process counters before any model work, never infer zero.""" + if ( + sys.platform != "win32" + or self.device != "gpu" + or not self.adapter_luid + or self.adapter_reason + ): + self.gpu_baseline_preparation = {"status": "not_applicable"} + return + start = time.monotonic_ns() + before = sample_vram(self.adapter_luid) + if not any(v.get("discovery_status") == "absent_unconfirmed" for v in before.values()): + self.gpu_baseline_preparation = { + "status": "existing_counters" + if all(v["value_bytes"] is not None for v in before.values()) + else "counter_unavailable", + "initial_observation": before, + } + return + try: + from ._gpu_baseline import GpuBaselineDevice + + self._gpu_baseline_device = GpuBaselineDevice(self.adapter_luid) + deadline = time.monotonic() + 2.0 + while True: + observed = sample_vram(self.adapter_luid) + valid = all(v["value_bytes"] is not None for v in observed.values()) + if valid or time.monotonic() >= deadline: + break + time.sleep(0.05) + self.gpu_baseline_preparation = { + "status": "initialized" if valid else "counter_unavailable", + "method": "D3D12 device on selected LUID; no model or GPU work", + "initial_observation": before, + "prepared_observation": observed, + "retained_until": "after final checkpoint or failure cleanup", + } + except Exception as exc: + # Measurement setup must not prevent a provider from running. + self.gpu_baseline_preparation = { + "status": "failed", + "reason": f"{type(exc).__name__}: {exc}", + "initial_observation": before, + } + self.gpu_baseline_preparation.update( + started_ns=start, completed_ns=time.monotonic_ns(), adapter_luid=self.adapter_luid + ) + + def close(self) -> None: + """Keep baseline overhead stable through inference, then release it.""" + device, self._gpu_baseline_device = self._gpu_baseline_device, None + if device is not None: + device.close() def capture(self, phase: str) -> None: """Capture one live phase boundary, including raw bytes and status.""" + if phase == "before_model_load": + if phase in self.checkpoints: + return # Never replace the original baseline with a later observation. + self._prepare_gpu_baseline() snapshot: dict[str, Any] = {"timestamp_ns": time.monotonic_ns()} try: info = psutil.Process(self.pid).memory_info() @@ -218,5 +283,333 @@ def evidence(self) -> dict[str, Any]: "adapter_luid": self.adapter_luid, "adapter_reason": self.adapter_reason, "peak_kind": "maximum_of_three_checkpoints", + "gpu_baseline_preparation": self.gpu_baseline_preparation, + "gpu_delta_scope": ( + "Measured process change from the recorded endpoint; includes model/provider " + "allocations, excludes retained model-free baseline device. Not minimum model VRAM." + ), "checkpoints": self.checkpoints, } + + +def _device_memory(adapter_luid: str | None) -> tuple[dict[str, float | None], dict[str, str]]: + """Use discovered process nodes while preserving missing-counter reasons.""" + sample = sample_vram(adapter_luid) + values: dict[str, float | None] = {} + errors: dict[str, str] = {} + for key in ("local", "shared"): + name = "vram_" + key + value = sample[key]["value_bytes"] + values[name] = value / MIB if value is not None else None + if value is None: + errors[name] = sample[key]["reason"] or "counter_unavailable_or_invalid" + return values, errors + + +def _rounded(value: float | None) -> float | None: + return round(value, 2) if value is not None else None + + +def _delta(after: float | None, before: float | None) -> float | None: + return _rounded(after - before) if after is not None and before is not None else None + + +def _rss_high_water() -> tuple[float | None, str | None]: + """Read Windows' process-lifetime peak, independently of Python polling.""" + try: + peak = getattr(psutil.Process(os.getpid()).memory_info(), "peak_wset", None) + except psutil.Error as exc: + return None, f"query_failed:{type(exc).__name__}" + if not isinstance(peak, (int, float)) or not math.isfinite(peak) or peak < 0: + return None, "process_high_water_unavailable" + return peak / (1024 * 1024), None + + +class ProcessMemoryTracker: + """Capture checkpoints and sampled peaks over one explicitly named lifecycle. + + The tracker does not resolve devices or read model properties. Binding must + happen before model construction; a late binding cannot backfill a baseline. + Sampling records aggregates rather than growing a per-iteration buffer. + """ + + def __init__( + self, + *, + adapter_luid: str | None = None, + scope: str = "before_model_load_through_inference", + interval: float = 0.05, + ) -> None: + self.adapter_luid = adapter_luid + self.scope = scope + self.interval = interval + self.checkpoints: dict[str, dict[str, float | None]] = {} + self.availability: dict[str, dict[str, str]] = {} + self._peaks: dict[str, float] = {} + self._counts: dict[str, int] = {} + self._last_sample_at: float | None = None + self._sample_interval_count = 0 + self._sample_interval_total = 0.0 + self._sample_interval_max: float | None = None + self._lock = threading.Lock() + self._stop = threading.Event() + self._thread: threading.Thread | None = None + self._started = time.monotonic() + self._ended: float | None = None + self._binding_before_load = False + self._os_peak_before: float | None = None + self._os_peak_after: float | None = None + self._os_peak_error: str | None = None + self._load_baseline: dict[str, float | None] | None = None + self._load_ready: dict[str, float | None] | None = None + self._load_peaks: dict[str, float] = {} + self._load_counts: dict[str, int] = {} + self._load_active = False + self._load_os_before: float | None = None + self._load_os_ready: float | None = None + self._baseline_owner: MemoryTracker | None = None + + def _observe(self) -> tuple[dict[str, float | None], dict[str, str]]: + values, errors = _device_memory(self.adapter_luid) + try: + values["rss"] = get_rss_mb() + except psutil.Error as exc: + values["rss"] = None + errors["rss"] = f"query_failed:{type(exc).__name__}" + return values, errors + + def checkpoint(self, phase: str) -> None: + """Record an explicit lifecycle boundary without triggering model work.""" + during_load = self._load_active + values, errors = self._observe() + self.checkpoints[phase] = values + self.availability[phase] = errors + # Explicit observations belong to the load peak, but are not polling samples. + with self._lock: + if during_load and self._load_active: + for key, value in values.items(): + if value is not None: + self._load_peaks[key] = max(value, self._load_peaks.get(key, value)) + + def bind_before_load(self, adapter_luid: str | None, *, device: str = "auto") -> None: + """Capture a GPU baseline after EP resolution but before model creation.""" + if self._binding_before_load: + return + self.adapter_luid = adapter_luid + self._binding_before_load = True + self._baseline_owner = MemoryTracker( + adapter_luid, baseline="before_model_load", device=device + ) + self._baseline_owner.capture("before_model_load") + values, errors = _device_memory(adapter_luid) + self.checkpoints["baseline"].update(values) + for key in values: + self.availability["baseline"].pop(key, None) + self.availability["baseline"].update(errors) + # The model-specific window excludes generic imports/EP discovery. + self._load_os_before, _ = _rss_high_water() + self._load_baseline, reasons = self._observe() + self.availability["load_baseline"] = reasons + with self._lock: + self._load_peaks = {k: v for k, v in self._load_baseline.items() if v is not None} + self._load_active = True + + def record_model_ready(self) -> None: + """Freeze load-only observations before inputs, warmup and inference.""" + if self._load_baseline is None: + return + self._load_ready, reasons = self._observe() + self.availability["model_ready"] = reasons + self._load_os_ready, _ = _rss_high_water() + with self._lock: + self._load_active = False + for key, value in self._load_ready.items(): + if value is not None: + self._load_peaks[key] = max(value, self._load_peaks.get(key, value)) + + def start(self) -> None: + """Start sampling before the model-loading call.""" + self._os_peak_before, self._os_peak_error = _rss_high_water() + self.checkpoint("baseline") + self._thread = threading.Thread(target=self._sample, name="winml-memory", daemon=True) + self._thread.start() + + def _sample(self) -> None: + while not self._stop.is_set(): + during_load = self._load_active + values, _ = self._observe() + sampled_at = time.monotonic() + with self._lock: + if self._last_sample_at is not None: + elapsed = sampled_at - self._last_sample_at + self._sample_interval_count += 1 + self._sample_interval_total += elapsed + self._sample_interval_max = max(self._sample_interval_max or 0.0, elapsed) + self._last_sample_at = sampled_at + for key, value in values.items(): + if value is not None: + self._counts[key] = self._counts.get(key, 0) + 1 + self._peaks[key] = max(value, self._peaks.get(key, value)) + if during_load and self._load_active: + self._load_counts[key] = self._load_counts.get(key, 0) + 1 + self._load_peaks[key] = max(value, self._load_peaks.get(key, value)) + self._stop.wait(self.interval) + + def stop(self) -> None: + """Stop sampling on both successful execution and exceptions.""" + self._stop.set() + if self._thread is not None: + self._thread.join() + if self._ended is None: + self._os_peak_after, self._os_peak_error = _rss_high_water() + self._ended = time.monotonic() + self.close() + + def close(self) -> None: + """Release retained baseline preparation after sampling has stopped.""" + if self._baseline_owner is not None: + self._baseline_owner.close() + + def to_dict(self) -> dict[str, float | None]: + """Return legacy checkpoint/delta fields without treating null as zero.""" + result = {} + for family in ("rss", "vram_local", "vram_shared"): + values = [ + self.checkpoints.get(p, {}).get(family) + for p in ("baseline", "after_compile", "after_inference") + ] + for phase in ("baseline", "after_load", "after_compile", "after_inference"): + result[f"{family}_{phase}_mb"] = _rounded( + self.checkpoints.get(phase, {}).get(family) + ) + observed = [value for value in values if value is not None] + result[f"{family}_checkpoint_peak_mb"] = ( + _rounded(max(observed)) if len(observed) == len(values) else None + ) + for name, end, start in (("model_load", 1, 0), ("inference", 2, 1), ("total", 2, 0)): + result[f"{family}_{name}_delta_mb"] = _delta(values[end], values[start]) + return result + + def metadata(self) -> dict[str, Any]: + """Describe the actual scope and missing checkpoint observations.""" + return { + "schema_version": 3, + "gpu_baseline_preparation": ( + self._baseline_owner.gpu_baseline_preparation + if self._baseline_owner is not None + else {"status": "not_applicable"} + ), + "unit": "MiB", + "pid": os.getpid(), + "adapter_luid": self.adapter_luid, + "scope": self.scope, + "rss_baseline": self.scope.split("_through_", 1)[0], + "gpu_baseline": ( + "after_ep_resolution_before_model_load" + if self._binding_before_load + else self.scope.split("_through_", 1)[0] + ), + "missing_reasons": self.availability, + "delta_semantics": "signed process change; not model-only allocation", + "checkpoint_peak_semantics": "maximum of baseline/after_compile/after_inference only", + "local_shared_semantics": ( + "separate counters; may overlap on unified memory; do not sum" + ), + } + + def load_memory(self) -> dict[str, Any]: + """Describe measured load cost, without claiming a universal minimum.""" + result: dict[str, Any] = { + "version": 1, + "unit": "MiB", + "pid": os.getpid(), + "adapter_luid": self.adapter_luid, + "scope": "after_runtime_setup_before_model_construction_to_ready_before_inputs", + "includes": ( + "loading, weight preparation and compilation required by this artifact/path" + ), + "excludes": "generic runtime setup, benchmark input allocation, warmup and inference", + "interpretation": ( + "observed load footprint for this configuration; " + "not a certified minimum hardware capacity" + ), + } + if self._load_baseline is None or self._load_ready is None: + result.update( + status="unavailable", + reason="model was already loaded or readiness was not observed", + ) + return result + result["status"] = "measured" + for family in ("rss", "vram_local", "vram_shared"): + baseline = self._load_baseline.get(family) + ready = self._load_ready.get(family) + peak = self._load_peaks.get(family) + method = "sampled_and_endpoint_maximum" + if ( + family == "rss" + and self._load_os_before is not None + and self._load_os_ready is not None + and self._load_os_ready > self._load_os_before + ): + peak = max(peak or 0, self._load_os_ready) + method = "new_os_process_high_water_during_load" + result[family] = { + "baseline_mb": _rounded(baseline), + "ready_mb": _rounded(ready), + "peak_mb": _rounded(peak), + "peak_extra_mb": _delta(peak, baseline), + "ready_extra_mb": _delta(ready, baseline), + "peak_method": method, + "sample_count": self._load_counts.get(family, 0), + "missing_reason": ( + self.availability["load_baseline"].get(family) + or self.availability["model_ready"].get(family) + ), + } + result["os_rss_high_water_before_mb"] = self._load_os_before + result["os_rss_high_water_at_ready_mb"] = self._load_os_ready + result["limitations"] = [ + "Sampling may miss short allocations; OS high-water is attributable only " + + "when it increases during the load window.", + "RSS measures resident pages, not private commit or exact required capacity.", + "GPU counters are sampled; local/shared counters may overlap on UMA " + + "and must not be summed.", + ] + return result + + def sampled(self) -> dict[str, Any]: + """Return sampling coverage and absolute peaks separately from deltas.""" + with self._lock: + return { + "scope": self.scope, + "pid": os.getpid(), + "adapter_luid": self.adapter_luid, + "unit": "MiB", + "configured_poll_delay_sec": self.interval, + "observed_mean_interval_sec": ( + self._sample_interval_total / self._sample_interval_count + if self._sample_interval_count + else None + ), + "observed_max_interval_sec": self._sample_interval_max, + "duration_sec": (self._ended or time.monotonic()) - self._started, + "sample_count": max(self._counts.values(), default=0), + "rss_samples": self._counts.get("rss", 0), + "local_samples": self._counts.get("vram_local", 0), + "shared_samples": self._counts.get("vram_shared", 0), + "rss_peak_mb": self._peaks.get("rss"), + "local_peak_mb": self._peaks.get("vram_local"), + "shared_peak_mb": self._peaks.get("vram_shared"), + "os_rss_peak_before_mb": self._os_peak_before, + "os_rss_peak_after_mb": self._os_peak_after, + "os_rss_peak_scope": ( + "process lifetime through measurement stop; not reset at baseline" + ), + "os_rss_peak_missing_reason": self._os_peak_error, + "limitations": ( + "sampled observations, not an exact continuous peak; " + "GPU samples begin after adapter binding; Python polling may miss " + "short allocations or pause while native code holds the GIL" + ), + } diff --git a/src/winml/modelkit/session/runtime_session.py b/src/winml/modelkit/session/runtime_session.py index 63dc9ed02..79fb92a2f 100644 --- a/src/winml/modelkit/session/runtime_session.py +++ b/src/winml/modelkit/session/runtime_session.py @@ -179,13 +179,9 @@ def _apply_io_metadata(io_config: dict[str, Any], model_path: Path) -> None: input_names = [item["name"] for item in metadata.get("inputs", [])] output_names = [item["name"] for item in metadata.get("outputs", [])] if len(input_names) != len(io_config["input_names"]): - raise click.ClickException( - f"Runtime input count does not match {metadata_path.name}." - ) + raise click.ClickException(f"Runtime input count does not match {metadata_path.name}.") if len(output_names) != len(io_config["output_names"]): - raise click.ClickException( - f"Runtime output count does not match {metadata_path.name}." - ) + raise click.ClickException(f"Runtime output count does not match {metadata_path.name}.") io_config["input_names"] = input_names io_config["output_names"] = output_names @@ -288,9 +284,7 @@ def _method( restype: Any, *argtypes: Any, ) -> Any: - vtable = ctypes.cast( - interface, ctypes.POINTER(ctypes.POINTER(ctypes.c_void_p)) - ).contents + vtable = ctypes.cast(interface, ctypes.POINTER(ctypes.POINTER(ctypes.c_void_p))).contents return ctypes.WINFUNCTYPE(restype, ctypes.c_void_p, *argtypes)(vtable[index]) @classmethod @@ -302,9 +296,7 @@ class LUID(ctypes.Structure): _fields_ = [("LowPart", ctypes.c_uint32), ("HighPart", ctypes.c_int32)] if not 0 <= luid_value <= 0xFFFFFFFFFFFFFFFF: - raise click.ClickException( - f"Adapter LUID is outside uint64 range: {luid_value!r}." - ) + raise click.ClickException(f"Adapter LUID is outside uint64 range: {luid_value!r}.") try: dxcore = ctypes.WinDLL("dxcore.dll") @@ -346,8 +338,7 @@ class LUID(ctypes.Structure): ) if hr < 0: raise click.ClickException( - f"DXCore could not resolve adapter LUID {luid_value} " - f"(0x{hr & 0xFFFFFFFF:08X})." + f"DXCore could not resolve adapter LUID {luid_value} (0x{hr & 0xFFFFFFFF:08X})." ) return cls(adapter, dxcore) finally: @@ -635,9 +626,7 @@ def __init__( self._ep_req = ep self._ep_source = ep_source if session_options is not None: - raise ValueError( - "session_options are not supported by the Windows ML Runtime backend." - ) + raise ValueError("session_options are not supported by the Windows ML Runtime backend.") if provider_options: raise click.ClickException( "--ep-options are not supported with --runtime winml-runtime because " @@ -889,8 +878,12 @@ def _do() -> dict[str, np.ndarray]: use_named_bindings=self._has_named_bindings, ) with _translate_native_errors("run"): - for index in range(output_count): - self._stage.request_output(index) + # Older projections (including 2.7.9) materialize outputs automatically. + # Newer projections require an explicit request on every execution. + request_output = getattr(self._stage, "request_output", None) + if callable(request_output): + for index in range(output_count): + request_output(index) self._pipeline.run() return self._read_outputs(named) diff --git a/tests/unit/commands/test_perf_cli.py b/tests/unit/commands/test_perf_cli.py index 9c206c821..8a4a994e2 100644 --- a/tests/unit/commands/test_perf_cli.py +++ b/tests/unit/commands/test_perf_cli.py @@ -2156,10 +2156,10 @@ def test_memory_profile_includes_additive_peak_and_compile_fields(self, monkeypa stats.samples_ms = [10.0] stats.all_samples_ms = [10.0] - rss_values = iter([100.0, 150.0, 180.0]) - vram_values = iter([(10.0, 20.0), (30.0, 50.0), (40.0, 70.0)]) + rss_values = iter([100.0, 100.0, 150.0, 180.0]) + vram_values = iter([(10.0, 20.0), (10.0, 20.0), (30.0, 50.0), (40.0, 70.0)]) - monkeypatch.setattr(perf_module, "_get_ep_device_binding", lambda *args: ("luid", "npu")) + monkeypatch.setattr(benchmark, "_resolve_adapter_luid", lambda: "luid") monkeypatch.setattr(benchmark, "_run_benchmark", lambda: stats) monkeypatch.setattr( benchmark, @@ -2170,28 +2170,28 @@ def test_memory_profile_includes_additive_peak_and_compile_fields(self, monkeypa {"pixel_values": MagicMock(shape=(1, 3, 224, 224))}, ), ) - from winml.modelkit.session.monitor import memory_tracker - - process = MagicMock() - process.create_time.return_value = 123.0 - process.memory_info.side_effect = lambda: SimpleNamespace( - rss=int(next(rss_values) * 1048576) + monkeypatch.setattr( + "winml.modelkit.session.monitor.memory_tracker.get_rss_mb", + lambda: next(rss_values), + ) + monkeypatch.setattr( + "winml.modelkit.session.monitor.memory_tracker._device_memory", + lambda _adapter_luid: ( + dict(zip(("vram_local", "vram_shared"), next(vram_values), strict=True)), + {}, + ), + ) + monkeypatch.setattr( + "winml.modelkit.session.monitor.ProcessMemoryTracker._sample", lambda _: None ) - monkeypatch.setattr(memory_tracker.psutil, "Process", lambda *_: process) - - def sample_vram(_luid): - local, shared = next(vram_values) - return { - key: memory_tracker.metric(int(value * 1048576), "test") - for key, value in (("local", local), ("shared", shared)) - } - - monkeypatch.setattr(memory_tracker, "sample_vram", sample_vram) monkeypatch.setattr("winml.modelkit.commands.perf._print_model_info", lambda *_, **__: None) result = benchmark._run_single() assert result.memory_profile == { + "rss_after_load_mb": 100.0, + "vram_local_after_load_mb": 10.0, + "vram_shared_after_load_mb": 20.0, "rss_baseline_mb": 100.0, "rss_after_compile_mb": 150.0, "rss_after_inference_mb": 180.0, @@ -2213,15 +2213,6 @@ def sample_vram(_luid): "vram_shared_inference_delta_mb": 20.0, "vram_local_total_delta_mb": 30.0, "vram_shared_total_delta_mb": 50.0, - **{ - f"{key}_{suffix}_mb": None - for key in ("rss", "vram_local", "vram_shared") - for suffix in ( - "before_model_load", - "model_factory_delta", - "total_from_before_model_load_delta", - ) - }, } @@ -2641,6 +2632,27 @@ def test_cli_closes_benchmark_after_success( class TestDisplayConsoleReport: + @pytest.mark.parametrize( + "ram, expected", + [ + ({}, "unavailable"), + ({"used_mb": None}, "unavailable"), + ({"used_mb": 0.0}, "0"), + ({"used_mb": 1024.0}, "1024"), + ], + ) + @pytest.mark.parametrize("device_kind", [None, "gpu"]) + def test_ram_availability(self, ram, expected, device_kind) -> None: + result = BenchmarkResult( + config=BenchmarkConfig(model_id="test"), + hw_monitor={"device_kind": device_kind, "ram": ram}, + ) + console = Console(file=StringIO(), width=200, record=True) + + display_console_report(result, console) + + assert f"RAM: {expected} MiB" in console.export_text() + class _FailingConsoleFile: encoding = "utf-8" diff --git a/tests/unit/commands/test_perf_memory_lifecycle.py b/tests/unit/commands/test_perf_memory_lifecycle.py new file mode 100644 index 000000000..2b658ea52 --- /dev/null +++ b/tests/unit/commands/test_perf_memory_lifecycle.py @@ -0,0 +1,422 @@ +# ------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. +# -------------------------------------------------------------------------- +"""Lifecycle and missing-counter regression tests without native EP dependencies.""" + +from __future__ import annotations + +import json +from types import SimpleNamespace +from unittest.mock import MagicMock, PropertyMock, patch + +import pytest + +from winml.modelkit.commands.perf import BenchmarkConfig, BenchmarkResult, PerfBenchmark +from winml.modelkit.session.monitor import ProcessMemoryTracker, get_vram_mb +from winml.modelkit.session.monitor import memory_tracker as memory + + +def test_missing_gpu_counter_is_null_and_query_always_closes(monkeypatch): + query = MagicMock() + query.collect.return_value = {"local_0": 0, "shared_0": None} + monkeypatch.setattr( + "winml.modelkit.session.monitor._pdh.memory_instances", lambda *args: ["node"] + ) + monkeypatch.setattr(memory.sys, "platform", "win32") + monkeypatch.setattr("winml.modelkit.session.monitor._pdh.PdhQuery", lambda: query) + values, errors = memory._device_memory("test-luid") + assert values == {"vram_local": 0.0, "vram_shared": None} + assert errors == {"vram_shared": "counter_data_unavailable"} + query.collect.assert_called_once_with(timeout=0.25, interval=0.05) + query.close.assert_called_once() + query.reset_mock() + query.collect.side_effect = OSError("unavailable") + assert get_vram_mb("test-luid") == (None, None) + query.close.assert_called_once() + + +@pytest.mark.parametrize("value", [-1, float("nan"), float("inf"), None]) +def test_invalid_gpu_readings_are_not_zero(monkeypatch, value): + query = MagicMock() + query.collect.return_value = {"local_0": value, "shared_0": 1048576} + monkeypatch.setattr( + "winml.modelkit.session.monitor._pdh.memory_instances", lambda *args: ["node"] + ) + monkeypatch.setattr(memory.sys, "platform", "win32") + monkeypatch.setattr("winml.modelkit.session.monitor._pdh.PdhQuery", lambda: query) + assert get_vram_mb("test-luid") == (None, 1.0) + + +def test_unknown_adapter_does_not_query_an_arbitrary_gpu(monkeypatch): + monkeypatch.setattr(memory.sys, "platform", "win32") + values, reasons = memory._device_memory(None) + assert set(values.values()) == {None} + assert set(reasons.values()) == {"adapter_unresolved"} + + +def test_snapshot_deltas_preserve_negative_and_missing_values(monkeypatch): + tracker = ProcessMemoryTracker() + observations = iter( + [ + ({"rss": 100.0, "vram_local": None, "vram_shared": 0.0}, {"vram_local": "missing"}), + ({"rss": 160.0, "vram_local": 12.0, "vram_shared": 8.0}, {}), + ({"rss": 80.0, "vram_local": 10.0, "vram_shared": 7.0}, {}), + ] + ) + monkeypatch.setattr(tracker, "_observe", lambda: next(observations)) + for phase in ("baseline", "after_compile", "after_inference"): + tracker.checkpoint(phase) + data = tracker.to_dict() + assert data["rss_total_delta_mb"] == -20 + assert data["rss_checkpoint_peak_mb"] == 160 + assert data["vram_local_total_delta_mb"] is None + assert data["vram_local_inference_delta_mb"] == -2 + assert data["vram_local_checkpoint_peak_mb"] is None + assert data["vram_shared_total_delta_mb"] == 7 + assert tracker.metadata()["missing_reasons"]["baseline"] == {"vram_local": "missing"} + json.dumps(data, allow_nan=False) + + +def test_sampler_tracks_temporary_peak_separately(monkeypatch): + tracker = ProcessMemoryTracker(interval=0.001) + observed = iter([12.0, 64.0, 15.0]) + + def sample(): + value = next(observed) + if value == 15: + tracker._stop.set() + return {"rss": value, "vram_local": None, "vram_shared": 0.0}, {} + + monkeypatch.setattr(tracker, "_observe", sample) + tracker._sample() + tracker.stop() + result = tracker.sampled() + assert result["rss_peak_mb"] == 64 + assert result["rss_samples"] == 3 + assert result["local_peak_mb"] is None and result["local_samples"] == 0 + assert result["shared_peak_mb"] == 0 and result["shared_samples"] == 3 + + +def test_sampling_reports_observed_cadence_not_configured_delay(monkeypatch): + tracker = ProcessMemoryTracker(interval=0.05) + clock = iter([10.0, 10.3, 10.8]) + monkeypatch.setattr(memory.time, "monotonic", lambda: next(clock)) + calls = [] + + def observe(): + calls.append(True) + if len(calls) == 3: + tracker._stop.set() + return {"rss": 100.0}, {} + + monkeypatch.setattr(tracker, "_observe", observe) + tracker._sample() + tracker._ended = tracker._started + 1 + result = tracker.sampled() + assert result["configured_poll_delay_sec"] == 0.05 + assert result["observed_mean_interval_sec"] == pytest.approx(0.4) + assert result["observed_max_interval_sec"] == pytest.approx(0.5) + assert "sampling_interval_sec" not in result + + +def test_sampling_without_two_observations_has_no_cadence(): + result = ProcessMemoryTracker().sampled() + assert result["observed_mean_interval_sec"] is None + assert result["observed_max_interval_sec"] is None + + +@pytest.mark.parametrize( + "runtime,artifact", + [ + ("winml-ort", "generated.onnx"), + ("winml-runtime", "generated.onnx"), + ("winml-runtime", "generated.mlir"), + ], +) +def test_baseline_precedes_eager_loading_and_device_properties(monkeypatch, runtime, artifact): + benchmark = PerfBenchmark(BenchmarkConfig(model_id=artifact, runtime=runtime, memory=True)) + state = {"rss": 100.0, "gpu": None} + events = [] + monkeypatch.setattr(memory, "get_rss_mb", lambda: state["rss"]) + monkeypatch.setattr(ProcessMemoryTracker, "_sample", lambda _: None) + monkeypatch.setattr( + memory, + "_device_memory", + lambda luid: ( + {"vram_local": state["gpu"], "vram_shared": state["gpu"]}, + {} + if state["gpu"] is not None + else {"vram_local": "unavailable", "vram_shared": "unavailable"}, + ), + ) + + def load(): + assert benchmark._memory_tracker.checkpoints["baseline"]["rss"] == 100 + benchmark._memory_tracker.bind_before_load("bound-luid") + state.update(rss=300.0, gpu=200.0) + events.append("load") + benchmark._model = MagicMock() + + def inputs(): + events.append("inputs") + state["rss"] = 310.0 + + def inference(): + events.append("inference") + state.update(rss=330.0, gpu=220.0) + return object() + + def collect(_): + return BenchmarkResult( + config=benchmark.config, + memory_profile=benchmark._memory, + memory_measurement=benchmark._memory_measurement, + process_memory=benchmark._process_memory, + load_memory=benchmark._load_memory, + ) + + monkeypatch.setattr(benchmark, "_load_model", load) + monkeypatch.setattr(benchmark, "_generate_inputs", inputs) + monkeypatch.setattr(benchmark, "_run_benchmark", inference) + monkeypatch.setattr(benchmark, "_collect_results", collect) + monkeypatch.setattr("winml.modelkit.commands.perf.print_pre_bench_block", lambda *a, **kw: None) + benchmark._ep_device = MagicMock() + with patch.object( + PerfBenchmark, "_is_composite", new_callable=PropertyMock, return_value=False + ): + result = benchmark.run() + assert events == ["load", "inputs", "inference"] + assert result.memory_profile["rss_baseline_mb"] == 100 + assert result.memory_profile["rss_after_load_mb"] == 300 + assert result.memory_profile["rss_model_load_delta_mb"] == 200 + assert result.load_memory["rss"]["ready_extra_mb"] == 200 + assert result.load_memory["rss"]["ready_mb"] == 300 + assert result.load_memory["status"] == "measured" + assert result.memory_profile["rss_total_delta_mb"] == 230 + assert result.memory_profile["vram_local_baseline_mb"] is None + assert result.memory_profile["vram_local_total_delta_mb"] is None + assert result.memory_measurement["scope"] == "before_model_load_through_inference" + assert benchmark._memory_tracker is None + + +def test_tracker_is_stopped_on_load_failure(monkeypatch): + benchmark = PerfBenchmark(BenchmarkConfig(model_id="generated.onnx")) + trackers = [] + original_start = ProcessMemoryTracker.start + + def start(tracker): + trackers.append(tracker) + original_start(tracker) + + def load(): + raise RuntimeError("load failed") + + monkeypatch.setattr(ProcessMemoryTracker, "start", start) + monkeypatch.setattr(benchmark, "_load_model", load) + with pytest.raises(RuntimeError, match="load failed"): + benchmark.run() + assert len(trackers) == 1 and not trackers[0]._thread.is_alive() + assert benchmark._memory_tracker is None + + +def test_disabled_memory_does_not_start_tracker(monkeypatch): + benchmark = PerfBenchmark(BenchmarkConfig(model_id="generated.onnx", memory=False)) + start = MagicMock() + monkeypatch.setattr(ProcessMemoryTracker, "start", start) + monkeypatch.setattr( + benchmark, "_load_model", lambda: setattr(benchmark, "_model", SimpleNamespace()) + ) + monkeypatch.setattr(benchmark, "_run_single", lambda: "result") + with patch.object( + PerfBenchmark, "_is_composite", new_callable=PropertyMock, return_value=False + ): + assert benchmark.run() == "result" + start.assert_not_called() + + +def test_load_binds_gpu_before_from_onnx(monkeypatch, tmp_path): + """Exercise the real loader's pre-construction binding hook.""" + from winml.modelkit.models import WinMLAutoModel + + source = tmp_path / "generated.onnx" + source.write_bytes(b"not read by the mocked model constructor") + benchmark = PerfBenchmark(BenchmarkConfig(model_id=str(source))) + benchmark._memory_tracker = ProcessMemoryTracker() + monkeypatch.setattr(ProcessMemoryTracker, "_sample", lambda _: None) + benchmark._memory_tracker.start() + benchmark._ep_device = MagicMock() + monkeypatch.setattr(benchmark, "_resolve_device_ep", lambda: None) + monkeypatch.setattr(benchmark, "_resolve_adapter_luid", lambda: "selected-luid") + monkeypatch.setattr( + memory, "_device_memory", lambda _: ({"vram_local": 0.0, "vram_shared": 0.0}, {}) + ) + + def construct(**_): + assert benchmark._memory_tracker.adapter_luid == "selected-luid" + assert benchmark._memory_tracker.checkpoints["baseline"]["vram_local"] == 0 + return MagicMock() + + monkeypatch.setattr(WinMLAutoModel, "from_onnx", construct) + try: + benchmark._load_model() + finally: + benchmark._memory_tracker.stop() + + +def test_nullable_memory_serialization_and_console(): + from io import StringIO + + from winml.modelkit.commands.perf import display_console_report + from winml.modelkit.utils.console import SafeConsole + + result = BenchmarkResult( + config=BenchmarkConfig(model_id="generated.onnx"), + memory_profile={ + "rss_after_inference_mb": 100.0, + "rss_total_delta_mb": -20.0, + "vram_local_after_inference_mb": None, + "vram_shared_after_inference_mb": 0.0, + }, + memory_measurement={ + "scope": "before_model_load_through_inference", + "missing_reasons": {"baseline": {"vram_local": "missing"}}, + }, + process_memory={"rss_peak_mb": 120.0, "sample_count": 3}, + ) + output = StringIO() + display_console_report(result, SafeConsole(file=output, width=200)) + assert "N/A" in output.getvalue() and "-20.0" in output.getvalue() + serialized = json.loads(json.dumps(result.to_dict(), allow_nan=False)) + assert serialized["memory"]["vram_local_after_inference_mb"] is None + assert serialized["process_memory"]["rss_peak_mb"] == 120.0 + + +def test_genai_nullable_counter_consumer(monkeypatch): + from winml.modelkit.commands._perf_genai import _GenaiMemoryTracker + + tracker = _GenaiMemoryTracker(adapter_luid="bound-luid") + monkeypatch.setattr("winml.modelkit.commands._perf_genai._get_rss_mb", lambda: 100.0) + monkeypatch.setattr("winml.modelkit.commands._perf_genai._get_vram_mb", lambda _: (None, 0.0)) + tracker.record_baseline() + tracker.record_after_load() + tracker.record_after_benchmark() + result = tracker.to_dict() + assert result["vram_local_total_delta_mb"] is None + assert result["vram_local_checkpoint_peak_mb"] is None + assert result["vram_shared_total_delta_mb"] == 0 + json.dumps(result, allow_nan=False) + + +def test_os_high_water_is_separate_from_sampled_peak(monkeypatch): + tracker = ProcessMemoryTracker() + peaks = iter([(120.0, None), (200.0, None)]) + monkeypatch.setattr(memory, "_rss_high_water", lambda: next(peaks)) + monkeypatch.setattr(tracker, "_sample", lambda: None) + monkeypatch.setattr(tracker, "_observe", lambda: ({"rss": 100.0}, {})) + tracker.start() + tracker.stop() + tracker.stop() + values = tracker.sampled() + assert values["rss_peak_mb"] is None + assert values["os_rss_peak_before_mb"] == 120 + assert values["os_rss_peak_after_mb"] == 200 + assert "process lifetime" in values["os_rss_peak_scope"] + + +def test_os_high_water_unavailable_is_not_zero(monkeypatch): + monkeypatch.setattr( + memory.psutil, + "Process", + lambda _: SimpleNamespace(memory_info=lambda: SimpleNamespace(rss=100)), + ) + assert memory._rss_high_water() == (None, "process_high_water_unavailable") + + +def test_load_only_peak_excludes_later_inference_and_old_high_water(monkeypatch): + tracker = ProcessMemoryTracker() + state = {"rss": 100.0, "vram_local": 20.0, "vram_shared": 20.0} + monkeypatch.setattr(tracker, "_observe", lambda: (dict(state), {})) + monkeypatch.setattr( + memory, + "_device_memory", + lambda _: ({"vram_local": state["vram_local"], "vram_shared": state["vram_shared"]}, {}), + ) + # A larger old peak must not be attributed to loading this model. + monkeypatch.setattr(memory, "_rss_high_water", lambda: (900.0, None)) + tracker.checkpoint("baseline") + tracker.bind_before_load("bound-luid") + state.update(rss=250.0, vram_local=60.0, vram_shared=60.0) + tracker.record_model_ready() + state.update(rss=1500.0, vram_local=500.0, vram_shared=500.0) + tracker.checkpoint("after_inference") + result = tracker.load_memory() + assert result["rss"]["peak_mb"] == 250 + assert result["rss"]["peak_extra_mb"] == 150 + assert result["rss"]["peak_method"] == "sampled_and_endpoint_maximum" + assert result["vram_local"]["ready_extra_mb"] == 40 + assert result["version"] == 1 + + +def test_new_os_peak_in_load_window_includes_temporary_compile_allocations(monkeypatch): + tracker = ProcessMemoryTracker() + state = {"rss": 100.0, "vram_local": None, "vram_shared": 0.0} + monkeypatch.setattr(tracker, "_observe", lambda: (dict(state), {"vram_local": "unavailable"})) + monkeypatch.setattr( + memory, + "_device_memory", + lambda _: ({"vram_local": None, "vram_shared": 0.0}, {"vram_local": "unavailable"}), + ) + peaks = iter([(110.0, None), (450.0, None)]) + monkeypatch.setattr(memory, "_rss_high_water", lambda: next(peaks)) + tracker.checkpoint("baseline") + tracker.bind_before_load("bound-luid") + state["rss"] = 200.0 + tracker.record_model_ready() + result = tracker.load_memory() + assert result["rss"]["peak_extra_mb"] == 350 + assert result["rss"]["ready_extra_mb"] == 100 + assert result["rss"]["peak_method"] == "new_os_process_high_water_during_load" + assert result["vram_local"]["peak_extra_mb"] is None + assert result["vram_local"]["missing_reason"] == "unavailable" + assert ProcessMemoryTracker().load_memory()["status"] == "unavailable" + + +@pytest.mark.parametrize("missing_local_baseline", [False, True]) +def test_load_peak_includes_checkpoints_without_counting_them_as_samples( + monkeypatch, missing_local_baseline +): + tracker = ProcessMemoryTracker() + state = { + "rss": 100.0, + "vram_local": None if missing_local_baseline else 10.0, + "vram_shared": 10.0, + } + monkeypatch.setattr(tracker, "_observe", lambda: (dict(state), {})) + monkeypatch.setattr( + memory, + "_device_memory", + lambda _: ({key: value for key, value in state.items() if key != "rss"}, {}), + ) + monkeypatch.setattr(memory, "_rss_high_water", lambda: (900.0, None)) + tracker.checkpoint("baseline") + tracker.bind_before_load("bound-luid") + state.update(rss=400.0, vram_local=200.0, vram_shared=100.0) + tracker.checkpoint("after_load") + state.update(rss=150.0, vram_local=30.0, vram_shared=20.0) + tracker.record_model_ready() + state.update(rss=1500.0, vram_local=500.0, vram_shared=300.0) + tracker.checkpoint("after_inference") + result = tracker.load_memory() + for family, peak, ready in ( + ("rss", 400.0, 150.0), + ("vram_local", 200.0, 30.0), + ("vram_shared", 100.0, 20.0), + ): + assert result[family]["peak_mb"] == peak + assert result[family]["ready_mb"] == ready + assert result[family]["sample_count"] == 0 + assert result["rss"]["peak_extra_mb"] == 300.0 + assert result["rss"]["peak_method"] == "sampled_and_endpoint_maximum" + assert result["vram_local"]["peak_extra_mb"] == (None if missing_local_baseline else 190.0) + assert tracker.sampled()["sample_count"] == 0 diff --git a/tests/unit/pattern/test_cgc_initializer_shape.py b/tests/unit/pattern/test_cgc_initializer_shape.py new file mode 100644 index 000000000..013571d22 --- /dev/null +++ b/tests/unit/pattern/test_cgc_initializer_shape.py @@ -0,0 +1,24 @@ +# ------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. +# -------------------------------------------------------------------------- + +import numpy as np +from onnx import TensorProto, helper, numpy_helper + +from winml.modelkit.pattern.cgc.cgc_constant_folding import _ConstantParameters + + +def test_initializer_shape_without_value_info(): + model = helper.make_model( + helper.make_graph( + [helper.make_node("Shape", ["weights"], ["shape"])], + "initializer_shape", + [], + [helper.make_tensor_value_info("shape", TensorProto.INT64, [2])], + [numpy_helper.from_array(np.zeros((2, 3), np.float32), "weights")], + ), + opset_imports=[helper.make_opsetid("", 17)], + ) + result = _ConstantParameters(model, static_shapes=True).evaluate("shape", set(), set()) + np.testing.assert_array_equal(result, np.asarray([2, 3], np.int64)) diff --git a/tests/unit/session/test_gpu_baseline.py b/tests/unit/session/test_gpu_baseline.py new file mode 100644 index 000000000..ac2c99f37 --- /dev/null +++ b/tests/unit/session/test_gpu_baseline.py @@ -0,0 +1,234 @@ +# ------------------------------------------------------------------------- +# Copyright (c) Microsoft Corporation. All rights reserved. +# Licensed under the MIT License. +# -------------------------------------------------------------------------- + +"""Baseline device lifetime and unknown counters do not contaminate deltas.""" + +import ctypes +from types import SimpleNamespace +from unittest.mock import Mock + +import pytest + +from winml.modelkit.session.monitor import _gpu_baseline as native +from winml.modelkit.session.monitor import memory_tracker as mt + + +def observation(value=None, absent=False): + return { + k: { + **mt.metric(value, "test", "absent" if absent else None), + **({"discovery_status": "absent_unconfirmed"} if absent else {}), + } + for k in ("local", "shared") + } + + +def test_device_created_before_snapshot_and_retained_until_close(monkeypatch): + active = [] + device = Mock() + + def create(luid): + assert luid == "selected" + active.append(True) + return device + + monkeypatch.setattr(native, "GpuBaselineDevice", create) + monkeypatch.setattr( + mt, + "sample_vram", + lambda _: observation(10 * mt.MIB) if active else observation(absent=True), + ) + tracker = mt.MemoryTracker("selected", baseline="legacy", device="gpu") + tracker.capture("before_model_load") + assert tracker.checkpoints["before_model_load"]["vram_local"]["value_bytes"] == 10 * mt.MIB + assert tracker.evidence()["gpu_baseline_preparation"]["status"] == "initialized" + device.close.assert_not_called() + monkeypatch.setattr(mt, "sample_vram", lambda _: observation(30 * mt.MIB)) + for phase in ("baseline", "after_compile", "after_inference"): + tracker.capture(phase) + tracker.capture("before_model_load") + assert tracker.profile()["vram_local_total_from_before_model_load_delta_mb"] == 20 + tracker.close() + tracker.close() + device.close.assert_called_once() + + +def test_existing_valid_zero_is_not_initialized(monkeypatch): + factory = Mock() + monkeypatch.setattr(native, "GpuBaselineDevice", factory) + monkeypatch.setattr(mt, "sample_vram", lambda _: observation(0)) + tracker = mt.MemoryTracker("selected", baseline="legacy", device="gpu") + tracker.capture("before_model_load") + factory.assert_not_called() + assert tracker.checkpoints["before_model_load"]["vram_local"]["value_origin"] == "measured_zero" + + +def test_failed_initialization_keeps_unknown_baseline_and_inference_can_proceed(monkeypatch): + monkeypatch.setattr(native, "GpuBaselineDevice", Mock(side_effect=OSError("unsupported"))) + monkeypatch.setattr(mt, "sample_vram", lambda _: observation(absent=True)) + tracker = mt.MemoryTracker("selected", baseline="legacy", device="gpu") + tracker.capture("before_model_load") + assert tracker.evidence()["gpu_baseline_preparation"]["status"] == "failed" + monkeypatch.setattr(mt, "sample_vram", lambda _: observation(100)) + for phase in ("baseline", "after_compile", "after_inference"): + tracker.capture(phase) + assert tracker.profile()["vram_local_total_from_before_model_load_delta_mb"] is None + tracker.close() + + +@pytest.mark.parametrize( + "device,phase", + [("cpu", "before_model_load"), ("npu", "before_model_load"), ("gpu", "baseline")], +) +def test_other_targets_and_preloaded_phase_do_not_create_device(monkeypatch, device, phase): + factory = Mock() + monkeypatch.setattr(native, "GpuBaselineDevice", factory) + monkeypatch.setattr(mt, "sample_vram", lambda _: observation(absent=True)) + tracker = mt.MemoryTracker("selected", baseline="legacy", device=device) + tracker.capture(phase) + factory.assert_not_called() + + +def test_preparation_timeout_retains_null_without_backfill(monkeypatch): + device = Mock() + monkeypatch.setattr(native, "GpuBaselineDevice", lambda _: device) + monkeypatch.setattr(mt, "sample_vram", lambda _: observation(absent=True)) + clock = iter([0, 3]) + monkeypatch.setattr(mt.time, "monotonic", lambda: next(clock)) + tracker = mt.MemoryTracker("selected", baseline="legacy", device="gpu") + tracker.capture("before_model_load") + assert tracker.evidence()["gpu_baseline_preparation"]["status"] == "counter_unavailable" + assert tracker.checkpoints["before_model_load"]["vram_local"]["value_bytes"] is None + tracker.close() + device.close.assert_called_once() + + +@pytest.mark.parametrize("failure", [False, True]) +def test_perf_releases_baseline_on_success_or_failure(monkeypatch, failure): + from winml.modelkit.commands.perf import BenchmarkConfig, PerfBenchmark + + bench = PerfBenchmark(BenchmarkConfig(model_id="unused")) + tracker = Mock() + bench._memory_tracker = tracker + + def run(): + tracker.close.assert_not_called() + if failure: + raise RuntimeError("model failed") + return "result" + + monkeypatch.setattr(bench, "_run_with_memory", run) + if failure: + with pytest.raises(RuntimeError, match="model failed"): + bench.run() + else: + assert bench.run() == "result" + tracker.close.assert_called_once() + + +def test_luid_rejected_before_native_calls(): + with pytest.raises(ValueError, match="LUID"): + native.GpuBaselineDevice("adapter-zero") + + +@pytest.mark.parametrize("failure", [None, "factory", "adapter", "device"]) +def test_native_device_releases_acquired_interfaces(monkeypatch, failure): + released = [] + + def acquire(stage, value, output): + if failure == stage: + return -2147467259 + ctypes.cast(output, ctypes.POINTER(ctypes.c_void_p))[0] = value + return 0 + + factory = Mock(side_effect=lambda iid, out: acquire("factory", 11, out)) + create = Mock(side_effect=lambda adapter, level, iid, out: acquire("device", 33, out)) + monkeypatch.setattr( + native.ctypes, + "WinDLL", + lambda name: ( + SimpleNamespace(CreateDXGIFactory1=factory) + if name == "dxgi" + else SimpleNamespace(D3D12CreateDevice=create) + ), + ) + + def method(pointer, index, result, *args): + if index == 2: + return lambda value: released.append(value.value) + assert pointer.value == 11 and index == 26 + + def adapter_by_luid(factory, luid, iid, output): + assert luid.low == 2 and luid.high == -1 + return acquire("adapter", 22, output) + + return adapter_by_luid + + monkeypatch.setattr(native, "_method", method) + if failure: + with pytest.raises(OSError, match="HRESULT"): + native.GpuBaselineDevice("0xffffffff_0x00000002") + assert released == {"factory": [], "adapter": [11], "device": [22, 11]}[failure] + else: + device = native.GpuBaselineDevice("0xffffffff_0x00000002") + assert released == [22, 11] + device.close() + device.close() + assert released == [22, 11, 33] + + +@pytest.mark.parametrize("failure", [False, True]) +def test_process_tracker_prepares_gpu_and_releases_after_model(monkeypatch, failure): + from winml.modelkit.commands.perf import BenchmarkConfig, PerfBenchmark + + device = Mock() + active = [] + + def create(luid): + active.append(luid) + return device + + monkeypatch.setattr(native, "GpuBaselineDevice", create) + monkeypatch.setattr( + mt, + "sample_vram", + lambda _: observation(10 * mt.MIB) if active else observation(absent=True), + ) + monkeypatch.setattr(mt.ProcessMemoryTracker, "_sample", lambda _: None) + monkeypatch.setattr( + mt, + "_device_memory", + lambda _: ( + {"vram_local": 10.0 if active else None, "vram_shared": 10.0 if active else None}, + {}, + ), + ) + bench = PerfBenchmark(BenchmarkConfig(model_id="unused", memory=True)) + + def load(): + tracker = bench._memory_tracker + tracker.bind_before_load("selected", device="gpu") + assert active == ["selected"] + assert tracker.metadata()["gpu_baseline_preparation"]["status"] == "initialized" + assert tracker.checkpoints["baseline"]["vram_local"] == 10 + device.close.assert_not_called() + if failure: + raise RuntimeError("model failed") + bench._model = Mock() + + monkeypatch.setattr(bench, "_load_model", load) + monkeypatch.setattr(bench, "_run_single", lambda: "result") + if failure: + with pytest.raises(RuntimeError, match="model failed"): + bench.run() + else: + assert bench.run() == "result" + device.close.assert_called_once() + + +@pytest.fixture(autouse=True) +def windows_measurement_gate(monkeypatch): + """Mock the platform gate without changing process-wide sys.platform.""" + monkeypatch.setattr(mt, "sys", SimpleNamespace(platform="win32")) diff --git a/tests/unit/session/test_memory_tracker.py b/tests/unit/session/test_memory_tracker.py index d3a314ec1..f272536bb 100644 --- a/tests/unit/session/test_memory_tracker.py +++ b/tests/unit/session/test_memory_tracker.py @@ -138,12 +138,8 @@ def test_perf_baseline_precedes_eager_factory(monkeypatch): events = [] bench = PerfBenchmark(BenchmarkConfig(model_id="test")) - monkeypatch.setattr(bench, "_resolve_device_ep", lambda: events.append("resolve")) - monkeypatch.setattr( - bench, - "_start_memory", - lambda baseline, **kwargs: events.append(kwargs.get("phase", "baseline")), - ) + monkeypatch.setattr(mt.ProcessMemoryTracker, "start", lambda self: events.append("baseline")) + monkeypatch.setattr(mt.ProcessMemoryTracker, "checkpoint", lambda *args: None) def load(): events.append("eager_load") @@ -153,7 +149,7 @@ def load(): monkeypatch.setattr(PerfBenchmark, "_is_composite", property(lambda self: False)) monkeypatch.setattr(bench, "_run_single", lambda: events.append("inference")) bench.run() - assert events == ["resolve", "before_model_load", "eager_load", "inference"] + assert events == ["baseline", "eager_load", "inference"] def test_adapter_resolution_never_reads_lazy_model(monkeypatch): @@ -271,8 +267,17 @@ def test_perf_records_both_boundaries_around_factory_and_inputs(monkeypatch): bench = perf_module.PerfBenchmark(perf_module.BenchmarkConfig(model_id="test")) monkeypatch.setattr(bench, "_resolve_device_ep", lambda: None) monkeypatch.setattr(perf_module, "_get_ep_device_binding", lambda *args: (None, "cpu")) - monkeypatch.setattr(mt.MemoryTracker, "capture", lambda self, phase: events.append(phase)) - monkeypatch.setattr(mt.MemoryTracker, "profile", lambda self: {}) + monkeypatch.setattr(mt.ProcessMemoryTracker, "_sample", lambda self: None) + original_checkpoint = mt.ProcessMemoryTracker.checkpoint + + def checkpoint(self, phase): + events.append(phase) + original_checkpoint(self, phase) + + monkeypatch.setattr(mt.ProcessMemoryTracker, "checkpoint", checkpoint) + monkeypatch.setattr( + mt.ProcessMemoryTracker, "record_model_ready", lambda self: events.append("model_ready") + ) single = MagicMock() single._session.compile.side_effect = lambda: events.append("compile") @@ -290,12 +295,13 @@ def factory(): monkeypatch.setattr(perf_module, "print_pre_bench_block", lambda *a, **k: None) bench.run() assert events == [ - "before_model_load", - "factory", "baseline", - "inputs", + "factory", + "after_load", "compile", + "model_ready", "after_compile", + "inputs", "inference", "after_inference", ] diff --git a/tests/unit/session/test_runtime_session.py b/tests/unit/session/test_runtime_session.py index e93d0c389..897cac4ec 100644 --- a/tests/unit/session/test_runtime_session.py +++ b/tests/unit/session/test_runtime_session.py @@ -271,9 +271,7 @@ def test_onnx_io_ranges_reach_tensor_comparison(tmp_path: Path) -> None: assert io_config["input_types"] == ["int64"] evaluator = object.__new__(TensorSimilarityEvaluator) evaluator.model = SimpleNamespace(io_config=io_config) - evaluator.config = SimpleNamespace( - input_data=None, dataset=SimpleNamespace(samples=2, seed=42) - ) + evaluator.config = SimpleNamespace(input_data=None, dataset=SimpleNamespace(samples=2, seed=42)) dataset = evaluator.prepare_data() @@ -286,8 +284,14 @@ def test_onnx_io_ranges_reach_tensor_comparison(tmp_path: Path) -> None: @pytest.mark.parametrize( "failure_point", - ["_load_mlir", "_load_onnx_on_cgc", "_load_onnx_on_ort", - "create_pipeline_builder", "build", "_stage_diagnostics"], + [ + "_load_mlir", + "_load_onnx_on_cgc", + "_load_onnx_on_ort", + "create_pipeline_builder", + "build", + "_stage_diagnostics", + ], ) def test_build_failure_releases_adapter(monkeypatch, mlir_target, failure_point): ep_device, adapter = mlir_target @@ -301,17 +305,24 @@ def test_build_failure_releases_adapter(monkeypatch, mlir_target, failure_point) session._is_mlir = False if failure_point == "_load_onnx_on_ort": session._backend = "ort" - monkeypatch.setattr(session, "_resolve_target", lambda *_args: SimpleNamespace( - adapter=adapter, execution_target=adapter.pointer, - provider_name=None, device_class="gpu", - )) + monkeypatch.setattr( + session, + "_resolve_target", + lambda *_args: SimpleNamespace( + adapter=adapter, + execution_target=adapter.pointer, + provider_name=None, + device_class="gpu", + ), + ) failure = RuntimeError("adapter cleanup probe") failing_call = Mock(side_effect=failure) if failure_point.startswith("_load_"): monkeypatch.setattr(session, failure_point, failing_call) elif failure_point == "_stage_diagnostics": monkeypatch.setattr( - "winml.modelkit.session.runtime_session._stage_diagnostics", failing_call, + "winml.modelkit.session.runtime_session._stage_diagnostics", + failing_call, ) elif failure_point == "build": monkeypatch.setattr(runtime.builder, failure_point, failing_call) @@ -376,6 +387,25 @@ def output(index): session.close() +def test_older_projection_materializes_outputs_without_request_api(monkeypatch, mlir_target): + stage = _Stage() + monkeypatch.delattr(_Stage, "request_output") + monkeypatch.setattr(stage, "output", lambda index: _Tensor(stage.bound[index].to_numpy())) + pipeline = _Pipeline() + wr = SimpleNamespace( + Runtime=lambda: _Runtime(stage, pipeline), NotSupportedError=_NotSupportedError + ) + monkeypatch.setattr("winml.modelkit.session.runtime_session.import_runtime", lambda: wr) + session = WinMLRuntimeSession("model.mlir", ep_device=mlir_target[0], backend="cgc") + try: + values = np.random.default_rng(42).normal(size=(2, 3, 8, 8)).astype(np.float32) + actual = session.run({"input_0": values}) + np.testing.assert_array_equal(actual["output_0"], values) + assert pipeline.runs == 1 + finally: + session.close() + + def test_mlir_session_builds_runs_and_resets( monkeypatch: pytest.MonkeyPatch, mlir_target: tuple[SimpleNamespace, _AdapterHandle], @@ -405,9 +435,7 @@ def test_mlir_session_builds_runs_and_resets( assert session.device == "gpu" assert session.io_config["input_names"] == ["input_0"] - outputs = session.run( - {"input_0": np.zeros((2, 3, 8, 8), dtype=np.float64)} - ) + outputs = session.run({"input_0": np.zeros((2, 3, 8, 8), dtype=np.float64)}) assert outputs["output_0"].shape == (2, 5) assert stage.bound[0].to_numpy().dtype == np.float32 assert pipeline.runs == 1 @@ -427,7 +455,9 @@ def test_mlir_session_builds_runs_and_resets( ids=["scalar", "vector", "matrix", "noncontiguous"], ) def test_prepare_inputs_preserves_shape( - dtype: type, shape: tuple[int, ...], transpose: bool, + dtype: type, + shape: tuple[int, ...], + transpose: bool, ) -> None: values = np.random.default_rng(0).standard_normal(shape).astype(dtype) if transpose: @@ -559,9 +589,7 @@ def test_runtime_session_rejects_provider_options() -> None: def test_onnx_session_passes_resolved_ep_and_device_to_runtime( monkeypatch: pytest.MonkeyPatch, ) -> None: - monkeypatch.setattr( - "winml.modelkit.onnx.get_io_config", lambda _path: {"value_ranges": {}} - ) + monkeypatch.setattr("winml.modelkit.onnx.get_io_config", lambda _path: {"value_ranges": {}}) stage = _Stage() pipeline = _Pipeline() @@ -640,9 +668,7 @@ def create_ort_execution_target( def test_onnx_session_without_ep_device_uses_request_resolution( monkeypatch: pytest.MonkeyPatch, ) -> None: - monkeypatch.setattr( - "winml.modelkit.onnx.get_io_config", lambda _path: {"value_ranges": {}} - ) + monkeypatch.setattr("winml.modelkit.onnx.get_io_config", lambda _path: {"value_ranges": {}}) stage = _Stage() pipeline = _Pipeline() @@ -781,9 +807,7 @@ def run(self) -> None: outputs: dict[str, float] = {} def run(name: str, value: float) -> None: - result = session.run( - {"input_0": np.full((1, 3, 8, 8), value, dtype=np.float32)} - ) + result = session.run({"input_0": np.full((1, 3, 8, 8), value, dtype=np.float32)}) outputs[name] = float(result["output_0"][0, 0]) first = threading.Thread(target=run, args=("first", 1.0)) @@ -850,8 +874,6 @@ def close(self) -> None: def test_perf_rejects_monitor_without_loading_runtime() -> None: - session = WinMLRuntimeSession( - "model.mlir", ep_device=_mlir_ep_device(), backend="cgc" - ) + session = WinMLRuntimeSession("model.mlir", ep_device=_mlir_ep_device(), backend="cgc") with pytest.raises(click.ClickException, match="monitor"), session.perf(monitor=object()): pass