Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 12 additions & 6 deletions tests/functional/object_model/ovms_instance.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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__)
Expand Down Expand Up @@ -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
158 changes: 157 additions & 1 deletion tests/functional/object_model/resource_monitor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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__)
Expand Down Expand Up @@ -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']],
Expand Down Expand Up @@ -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