diff --git a/tests/functional/object_model/ovms_instance.py b/tests/functional/object_model/ovms_instance.py index 3eee7096c1..f3c1d4b2b2 100644 --- a/tests/functional/object_model/ovms_instance.py +++ b/tests/functional/object_model/ovms_instance.py @@ -41,7 +41,7 @@ from tests.functional.utils.core import get_children_from_module from tests.functional.utils.inference.communication import GRPC, REST from tests.functional.utils.logger import get_logger -from tests.functional.constants.os_type import OsType +from tests.functional.constants.os_type import OsType, get_host_os from tests.functional.utils.port_manager import PortManager from tests.functional.utils.process import Process from tests.functional.utils.test_framework import change_dir_permissions, is_single_threaded @@ -66,7 +66,7 @@ from tests.functional.object_model.mediapipe_calculators import MediaPipeCalculator from tests.functional.object_model.ovms_config import OvmsConfig from tests.functional.object_model.package_manager import PackageManager -from tests.functional.object_model.resource_monitor import DockerResourceMonitor +from tests.functional.object_model.resource_monitor import DockerResourceMonitor, WindowsResourceMonitor from tests.functional.object_model.test_environment import TestEnvironment logger = get_logger(__name__) @@ -539,7 +539,13 @@ def attach_context(self, context): def attach_resource_monitor(self, context, start=True): if hasattr(self.ovms, "container"): self.resource_monitor = DockerResourceMonitor(self.ovms.container) - if start: - self.resource_monitor.start() - context.test_objects.append(self.resource_monitor) - return self.resource_monitor + elif get_host_os() == OsType.Windows or getattr(context, "base_os", None) == OsType.Windows: + ovms_pid = self.ovms._dmesg_log.ovms_pid + assert ovms_pid is not None, "Cannot attach Windows resource monitor: ovms_pid is not available" + self.resource_monitor = WindowsResourceMonitor(ovms_pid) + else: + return None + if start: + self.resource_monitor.start() + context.test_objects.append(self.resource_monitor) + return self.resource_monitor diff --git a/tests/functional/object_model/resource_monitor.py b/tests/functional/object_model/resource_monitor.py index 5787b69b1f..d0d6790d92 100644 --- a/tests/functional/object_model/resource_monitor.py +++ b/tests/functional/object_model/resource_monitor.py @@ -17,12 +17,14 @@ import csv import threading from abc import ABC, abstractmethod +from datetime import datetime from pathlib import Path import numpy as np from dateutil import parser from tests.functional.utils.logger import get_logger +from tests.functional.utils.process import Process from tests.functional.config import artifacts_dir logger = get_logger(__name__) @@ -55,13 +57,33 @@ def save_data(self): pass +def _cgroup_cache_bytes(cgroup_memory_stats): + if "cache" in cgroup_memory_stats: + return float(cgroup_memory_stats.get("cache", 0)) + return float(cgroup_memory_stats.get("file", 0)) + + class DockerResourceMonitor(ResourceMonitor): MEMORY_USAGE = "MEMORY_USAGE" - FIELDS = ["DATE", "PIDS_COUNT", MEMORY_USAGE] # + ["CPU_USAGE"] # Enable in further releases + PRIVATE_MEMORY = "PRIVATE_MEMORY" + MEMORY_CACHE = "MEMORY_CACHE" + FIELDS = ["DATE", "PIDS_COUNT", MEMORY_USAGE, PRIVATE_MEMORY, MEMORY_CACHE] # + ["CPU_USAGE"] # Enable in further releases + VALIDATED_FIELDS = [MEMORY_USAGE, PRIVATE_MEMORY] + LOGGED_MEMORY_FIELDS = [MEMORY_CACHE] + COUNTER_FIELDS = ["PIDS_COUNT"] + LOGGED_FIELDS = LOGGED_MEMORY_FIELDS + COUNTER_FIELDS FIELDS_TO_STATS = { "DATE": lambda x: x["read"], "PIDS_COUNT": lambda x: int(x["pids_stats"].get("current", "0")), MEMORY_USAGE: lambda x: "{:.2f}M".format(float(x["memory_stats"].get("usage", "0.0")) / (2**20)), + PRIVATE_MEMORY: lambda x: "{:.2f}M".format( + float(x["memory_stats"].get("stats", {}).get( + "anon", x["memory_stats"].get("stats", {}).get("rss", 0) + )) / (2**20) + ), + MEMORY_CACHE: lambda x: "{:.2f}M".format( + _cgroup_cache_bytes(x["memory_stats"].get("stats", {})) / (2**20) + ), # Enable after debug & fixing # "CPU_USAGE": lambda x: # [cpu / x['cpu_stats']['cpu_usage']['total_usage'] for cpu in x['cpu_stats']['cpu_usage']['percpu_usage']], @@ -146,3 +168,137 @@ def get_stats_by_field(self, field): result = self._get_resource_data() self._docker_stats_data_raw.append(result) return self.get_field_data(field, result) + + def sample_all(self): + """Read one stats snapshot and return all tracked metrics as floats (MB / counts). + + A single snapshot keeps every metric in the returned sample mutually + consistent (same instant) and avoids one docker stats call per metric. + """ + stats = self._get_resource_data() + self._docker_stats_data_raw.append(stats) + return { + field: float(str(self.get_field_data(field, stats)).replace("M", "")) + for field in self.get_validated_metric_names() + self.get_logged_metric_names() + } + + @classmethod + def get_validated_metric_names(cls): + return cls.VALIDATED_FIELDS + + @classmethod + def get_logged_metric_names(cls): + return cls.LOGGED_FIELDS + + @classmethod + def get_memory_metric_names(cls): + return cls.VALIDATED_FIELDS + cls.LOGGED_MEMORY_FIELDS + + @classmethod + def get_counter_metric_names(cls): + return cls.COUNTER_FIELDS + + +class WindowsResourceMonitor(ResourceMonitor): + WORKING_SET_SIZE = "WORKING_SET_SIZE" + PRIVATE_BYTES = "PRIVATE_BYTES" + PAGE_FILE_USAGE = "PAGE_FILE_USAGE" + PAGE_FAULTS = "PAGE_FAULTS" + + MEMORY_USAGE = WORKING_SET_SIZE + + FIELDS = ["DATE", WORKING_SET_SIZE, PRIVATE_BYTES, PAGE_FILE_USAGE, PAGE_FAULTS] + VALIDATED_FIELDS = [WORKING_SET_SIZE, PRIVATE_BYTES] + LOGGED_MEMORY_FIELDS = [PAGE_FILE_USAGE] + COUNTER_FIELDS = [PAGE_FAULTS] + LOGGED_FIELDS = LOGGED_MEMORY_FIELDS + COUNTER_FIELDS + SAMPLE_INTERVAL_SEC = 1.0 + + PS_COMMAND_TEMPLATE = ( + "powershell -NoProfile -Command \"" + "$p = Get-Process -Id {pid}; " + "Write-Output $p.WorkingSet64; " + "Write-Output $p.PrivateMemorySize64; " + "Write-Output $p.PagedMemorySize64; " + "Write-Output (Get-CimInstance Win32_Process -Filter 'ProcessId={pid}').PageFaults\"" + ) + # Optional callback invoked after save_data with (log_path). + on_data_saved = None + + def __init__(self, ovms_pid, proc=None): + super().__init__() + self.ovms_pid = ovms_pid + self.proc = proc if proc is not None else Process() + self._stats_data_raw = [] + + def cleanup(self): + if not self._stop_event.is_set(): + if self.is_alive(): + self.stop() + self.save_data() + + def _get_resource_data(self): + stats = {"DATE": datetime.now().isoformat()} + cmd = self.PS_COMMAND_TEMPLATE.format(pid=self.ovms_pid) + _, stdout, stderr = self.proc.run_and_check_return_all(cmd) + lines = [line.strip() for line in stdout.strip().splitlines() if line.strip()] + if len(lines) < 4: + raise AssertionError( + f"Unexpected PowerShell output while collecting resource data for " + f"pid {self.ovms_pid}: expected at least 4 non-empty lines, got " + f"{len(lines)}. stdout={stdout!r}, stderr={stderr!r}" + ) + stats[self.WORKING_SET_SIZE] = float(lines[0]) / (1024 * 1024) + stats[self.PRIVATE_BYTES] = float(lines[1]) / (1024 * 1024) + stats[self.PAGE_FILE_USAGE] = float(lines[2]) / (1024 * 1024) + stats[self.PAGE_FAULTS] = int(lines[3]) + return stats + + def check_resources(self): + result = self._get_resource_data() + self._stats_data_raw.append(result) + self._stop_event.wait(self.SAMPLE_INTERVAL_SEC) + + def save_data(self): + self.rows = list(self._stats_data_raw) + log_path = Path(artifacts_dir, f"windows_stats_pid_{self.ovms_pid}.log") + with log_path.open("w") as csvfile: + writer = csv.DictWriter(csvfile, fieldnames=self.FIELDS) + writer.writeheader() + writer.writerows(self.rows) + if WindowsResourceMonitor.on_data_saved: + WindowsResourceMonitor.on_data_saved(log_path) + return log_path + + def get_stats_by_field(self, field): + result = self._get_resource_data() + self._stats_data_raw.append(result) + value = result[field] + if field in (self.WORKING_SET_SIZE, self.PRIVATE_BYTES, self.PAGE_FILE_USAGE): + return f"{value:.2f}M" + return str(value) + + def sample_all(self): + """Read one process snapshot and return all tracked metrics as floats (MB / counts).""" + stats = self._get_resource_data() + self._stats_data_raw.append(stats) + return { + field: float(stats[field]) + for field in self.get_validated_metric_names() + self.get_logged_metric_names() + } + + @classmethod + def get_validated_metric_names(cls): + return cls.VALIDATED_FIELDS + + @classmethod + def get_logged_metric_names(cls): + return cls.LOGGED_FIELDS + + @classmethod + def get_memory_metric_names(cls): + return cls.VALIDATED_FIELDS + cls.LOGGED_MEMORY_FIELDS + + @classmethod + def get_counter_metric_names(cls): + return cls.COUNTER_FIELDS