from __future__ import annotations import logging import os import threading import time from datetime import datetime, timezone from pathlib import Path from typing import Any, Callable logger = logging.getLogger(__name__) POLL_INTERVAL_SECONDS = 5.0 def isoformat_utc(value: datetime) -> str: return value.astimezone(timezone.utc).isoformat(timespec="milliseconds").replace("+00:00", "Z") def _clamp_percent(value: float) -> float: return min(100.0, max(0.0, value)) def _parse_cpu_values(fields: list[str], label: str) -> tuple[int, int]: if len(fields) < 4: raise ValueError(f"/proc/stat {label} 行无效") values = [int(value) for value in fields] total = sum(values) idle = values[3] + (values[4] if len(values) > 4 else 0) if total <= 0: raise ValueError(f"/proc/stat {label} 总计无效") return total, idle def _parse_cpu_snapshot(text: str) -> tuple[tuple[int, int], tuple[tuple[str, int, int], ...]]: aggregate: tuple[int, int] | None = None cores: list[tuple[str, int, int]] = [] for line in text.splitlines(): parts = line.split() if not parts: continue label = parts[0] if label == "cpu": aggregate = _parse_cpu_values(parts[1:], label) elif label.startswith("cpu") and label[3:].isdigit(): total, idle = _parse_cpu_values(parts[1:], label) cores.append((label, total, idle)) if aggregate is None or not cores: raise ValueError("/proc/stat 缺少 CPU 汇总或逻辑核心行") cores.sort(key=lambda item: int(item[0][3:])) return aggregate, tuple(cores) def _parse_process_cpu_ticks(text: str) -> int: closing_paren = text.rfind(")") if closing_paren < 0: raise ValueError("/proc/self/stat 进程名无效") fields = text[closing_paren + 1:].split() if len(fields) < 13: raise ValueError("/proc/self/stat 字段不足") return int(fields[11]) + int(fields[12]) def _parse_rss_bytes(text: str, page_size: int) -> int: fields = text.split() if len(fields) < 2: raise ValueError("/proc/self/statm 字段不足") rss_pages = int(fields[1]) if rss_pages < 0: raise ValueError("/proc/self/statm RSS 无效") return rss_pages * page_size def _parse_memory_bytes(text: str) -> tuple[int, int]: values: dict[str, int] = {} for line in text.splitlines(): key, separator, remainder = line.partition(":") if not separator: continue parts = remainder.split() if parts and parts[0].isdigit(): values[key] = int(parts[0]) * 1024 total = values.get("MemTotal", 0) available = values.get("MemAvailable", -1) if total <= 0 or not 0 <= available <= total: raise ValueError("/proc/meminfo 内存汇总无效") return total, available class ResourceMonitor: def __init__( self, proc_root: Path | str = "/proc", poll_interval: float = POLL_INTERVAL_SECONDS, page_size: int | None = None, now: Callable[[], datetime] | None = None, ) -> None: self._proc_root = Path(proc_root) self._poll_interval = poll_interval if page_size is None: sysconf = getattr(os, "sysconf", None) page_size = int(sysconf("SC_PAGE_SIZE")) if sysconf is not None else 4096 self._page_size = page_size self._now = now or (lambda: datetime.now(timezone.utc)) self._state_lock = threading.RLock() self._stop = threading.Event() self._thread: threading.Thread | None = None self._previous_cpu: tuple[ int, int, int, tuple[tuple[str, int, int], ...], ] | None = None self._snapshot = self._empty_snapshot("starting", None, None) @staticmethod def _empty_snapshot(status: str, sampled_at: str | None, error_code: str | None) -> dict[str, Any]: return { "status": status, "sampled_at": sampled_at, "cpu": { "application_percent": None, "total_percent": None, "logical_cpu_count": None, "application_core_equivalent": None, "cores_percent": [], }, "memory": { "application_bytes": None, "application_percent": None, "total_percent": None, }, "error_code": error_code, } def start(self) -> None: with self._state_lock: if self._thread is not None and self._thread.is_alive(): return self._stop.clear() self._thread = threading.Thread( target=self._run, name="system-resource-monitor", daemon=True, ) self._thread.start() def close(self) -> None: self._stop.set() with self._state_lock: thread = self._thread if thread is not None and thread is not threading.current_thread(): thread.join(timeout=max(2.0, self._poll_interval + 1.0)) def get_status(self) -> dict[str, Any]: with self._state_lock: snapshot = self._snapshot return { **snapshot, "cpu": { **snapshot["cpu"], "cores_percent": list(snapshot["cpu"]["cores_percent"]), }, "memory": dict(snapshot["memory"]), } def sample_now(self) -> dict[str, Any]: sampled_at = isoformat_utc(self._now()) try: (cpu_total, cpu_idle), cpu_cores = _parse_cpu_snapshot( (self._proc_root / "stat").read_text(encoding="ascii") ) process_ticks = _parse_process_cpu_ticks( (self._proc_root / "self" / "stat").read_text(encoding="ascii") ) rss_bytes = _parse_rss_bytes( (self._proc_root / "self" / "statm").read_text(encoding="ascii"), self._page_size, ) memory_total, memory_available = _parse_memory_bytes( (self._proc_root / "meminfo").read_text(encoding="ascii") ) application_cpu: float | None = None application_cores: float | None = None total_cpu: float | None = None cores_percent: list[float | None] = [None] * len(cpu_cores) previous = self._previous_cpu if previous is not None: previous_total, previous_idle, previous_process, previous_cores = previous current_labels = tuple(item[0] for item in cpu_cores) previous_labels = tuple(item[0] for item in previous_cores) if current_labels == previous_labels: total_delta = cpu_total - previous_total idle_delta = cpu_idle - previous_idle process_delta = process_ticks - previous_process if total_delta > 0 and idle_delta >= 0 and process_delta >= 0: total_cpu = _clamp_percent( (total_delta - idle_delta) * 100.0 / total_delta ) application_cpu = _clamp_percent(process_delta * 100.0 / total_delta) application_cores = min( float(len(cpu_cores)), max(0.0, process_delta * len(cpu_cores) / total_delta), ) for index, (current, old) in enumerate(zip(cpu_cores, previous_cores)): _, current_total, current_idle = current _, old_total, old_idle = old core_delta = current_total - old_total core_idle_delta = current_idle - old_idle if core_delta > 0 and core_idle_delta >= 0: cores_percent[index] = _clamp_percent( (core_delta - core_idle_delta) * 100.0 / core_delta ) self._previous_cpu = (cpu_total, cpu_idle, process_ticks, cpu_cores) snapshot = { "status": "ok", "sampled_at": sampled_at, "cpu": { "application_percent": application_cpu, "total_percent": total_cpu, "logical_cpu_count": len(cpu_cores), "application_core_equivalent": application_cores, "cores_percent": cores_percent, }, "memory": { "application_bytes": rss_bytes, "application_percent": _clamp_percent(rss_bytes * 100.0 / memory_total), "total_percent": _clamp_percent( (memory_total - memory_available) * 100.0 / memory_total ), }, "error_code": None, } except (OSError, UnicodeError, ValueError, IndexError): snapshot = self._empty_snapshot("error", sampled_at, "proc_read_error") except Exception: logger.exception("Unexpected system resource monitor failure") snapshot = self._empty_snapshot("error", sampled_at, "unexpected_error") with self._state_lock: previous_status = self._snapshot["status"] self._snapshot = snapshot if previous_status != snapshot["status"]: if snapshot["status"] == "ok": logger.info("System resource monitor is available") else: logger.warning("System resource monitor is unavailable: %s", snapshot["error_code"]) return self.get_status() def _run(self) -> None: while not self._stop.is_set(): started = time.monotonic() self.sample_now() remaining = max(0.0, self._poll_interval - (time.monotonic() - started)) self._stop.wait(remaining)