From f64835749ada4b29fa92d11547d7c44198411e53 Mon Sep 17 00:00:00 2001 From: oliver Date: Sun, 2 Aug 2026 18:37:06 +0800 Subject: [PATCH] Expose worker scheduler health via heartbeat and isolate WebCRT IO. API /metrics and /health/ready now read a worker heartbeat when collectors are split out; WebCRT blocking I/O uses a dedicated executor so session pumps do not starve the default pool. Co-authored-by: Cursor --- .env.example | 2 + netx_api/app_shutdown.py | 14 +++ netx_api/config.py | 2 + netx_api/config_sync_scheduler.py | 12 ++ netx_api/integrations_router.py | 20 ++- netx_api/lldp_collect_scheduler.py | 12 ++ netx_api/metrics_router.py | 33 ++++- netx_api/scheduler_heartbeat.py | 188 +++++++++++++++++++++++++++++ netx_api/ume_runtime.py | 4 + netx_api/webcrt_io.py | 42 +++++++ netx_api/webcrt_router.py | 13 +- tests/test_stability_hardening.py | 37 +++++- 12 files changed, 366 insertions(+), 13 deletions(-) create mode 100644 netx_api/scheduler_heartbeat.py create mode 100644 netx_api/webcrt_io.py diff --git a/.env.example b/.env.example index 6bb4706..efabeee 100644 --- a/.env.example +++ b/.env.example @@ -56,7 +56,9 @@ NETX_UME_NOTIFICATION_TOPIC=ALARM # NETX_SQL_READONLY_DATABASE_URL=postgresql+psycopg://netx_ro:xxx@127.0.0.1:5432/netx # Device collectors run inline with the API by default (frontend+backend start is enough). # Production split only: NETX_RUN_INLINE_SCHEDULERS=false and run `python -m netx_api.worker` +# (start_netx.ps1/.sh do this automatically). Worker writes heartbeat for API /metrics. # NETX_RUN_INLINE_SCHEDULERS=false +# NETX_SCHEDULER_HEARTBEAT_PATH=data/runtime/scheduler_heartbeat.json # --- Multi-user shared-server capacity (defaults in Settings already match these) --- # NETX_DB_POOL_SIZE=40 # NETX_DB_MAX_OVERFLOW=40 diff --git a/netx_api/app_shutdown.py b/netx_api/app_shutdown.py index 17fb593..733c1cc 100644 --- a/netx_api/app_shutdown.py +++ b/netx_api/app_shutdown.py @@ -10,6 +10,13 @@ def shutdown_runtime(*, reason: str = "lifespan") -> None: """Best-effort stop of schedulers, sidebands, pools, and sessions.""" _log.info("shutdown_runtime begin reason=%s", reason) + try: + from .scheduler_heartbeat import stop_scheduler_heartbeat_publisher + + stop_scheduler_heartbeat_publisher() + except Exception: # noqa: BLE001 + _log.exception("stop_scheduler_heartbeat_publisher failed") + try: from .config_sync_scheduler import stop_config_sync_scheduler @@ -67,6 +74,13 @@ def shutdown_runtime(*, reason: str = "lifespan") -> None: except Exception: # noqa: BLE001 _log.exception("shutdown_cli_timeout_pool failed") + try: + from .webcrt_io import shutdown_webcrt_io_executor + + shutdown_webcrt_io_executor(wait=False) + except Exception: # noqa: BLE001 + _log.exception("shutdown_webcrt_io_executor failed") + try: from .audit_async import shutdown_audit_worker diff --git a/netx_api/config.py b/netx_api/config.py index d779f2f..83151e1 100644 --- a/netx_api/config.py +++ b/netx_api/config.py @@ -144,6 +144,8 @@ class Settings(BaseSettings): # When true (default), API also runs config_sync / lldp / port_traffic schedulers. # Production split: set false and run `python -m netx_api.worker` beside the API. run_inline_schedulers: bool = True + # Worker→API heartbeat file (used when run_inline_schedulers=false). + scheduler_heartbeat_path: str = "data/runtime/scheduler_heartbeat.json" # SQLAlchemy QueuePool for multi-user API + collectors + UME WS. # Rule of thumb: pool_size + max_overflow >= HTTP/WS peak + cli_max_concurrent + sidebands. db_pool_size: int = 40 diff --git a/netx_api/config_sync_scheduler.py b/netx_api/config_sync_scheduler.py index 10f947b..f574e76 100644 --- a/netx_api/config_sync_scheduler.py +++ b/netx_api/config_sync_scheduler.py @@ -23,6 +23,7 @@ _log = logging.getLogger("netx.config_sync.scheduler") _stop = threading.Event() _thread: threading.Thread | None = None _BOOT_MONO = time.monotonic() +_last_tick_mono: float = 0.0 def _utcnow() -> datetime: @@ -114,11 +115,13 @@ def try_start_scheduled_cycle() -> str | None: def _loop() -> None: + global _last_tick_mono tick = max(15, int(settings.config_sync_scheduler_tick_sec or 60)) grace = max(0, int(settings.config_sync_startup_grace_sec or 0)) _log.info("config_sync scheduler started tick=%ss startup_grace=%ss", tick, grace) while not _stop.is_set(): try: + _last_tick_mono = time.monotonic() if bool(settings.config_sync_scheduler_enabled): try_start_scheduled_cycle() except Exception: @@ -142,3 +145,12 @@ def start_config_sync_scheduler() -> None: def stop_config_sync_scheduler() -> None: _stop.set() + + +def config_sync_scheduler_status() -> dict: + now = time.monotonic() + return { + "running": bool(_thread and _thread.is_alive()), + "last_tick_age_sec": (now - _last_tick_mono) if _last_tick_mono else None, + "startup_grace_remaining_sec": round(startup_grace_remaining_sec(), 1), + } diff --git a/netx_api/integrations_router.py b/netx_api/integrations_router.py index a0d7174..120454f 100644 --- a/netx_api/integrations_router.py +++ b/netx_api/integrations_router.py @@ -42,13 +42,31 @@ def health_ready(db: Session = Depends(get_db)) -> dict[str, Any]: out["db_pool"] = db_pool_status() out["cli_budget"] = cli_budget_status() inline = bool(getattr(settings, "run_inline_schedulers", True)) - out["schedulers"] = { + sched_block: dict[str, Any] = { "inline": inline, "mode": "inline" if inline else "external_worker", "hint": None if inline else "run `python -m netx_api.worker` for config_sync / lldp_collect / port_traffic", } + try: + from .scheduler_heartbeat import resolve_device_scheduler_metrics + + resolved = resolve_device_scheduler_metrics() + sched_block["source"] = resolved.get("source") + sched_block["stale"] = bool(resolved.get("stale")) + sched_block["age_sec"] = resolved.get("age_sec") + sched_block["worker_pid"] = resolved.get("pid") + for key in ("config_sync", "lldp_collect", "port_traffic"): + block = resolved.get(key) or {} + sched_block[key] = {"running": bool(block.get("running"))} + if resolved.get("hint"): + sched_block["hint"] = resolved["hint"] + if not inline and resolved.get("stale"): + out["status"] = "degraded" + except Exception: # noqa: BLE001 + pass + out["schedulers"] = sched_block return out diff --git a/netx_api/lldp_collect_scheduler.py b/netx_api/lldp_collect_scheduler.py index 109fa0f..1d34f80 100644 --- a/netx_api/lldp_collect_scheduler.py +++ b/netx_api/lldp_collect_scheduler.py @@ -20,6 +20,7 @@ _log = logging.getLogger("netx.lldp_collect.scheduler") _stop = threading.Event() _thread: threading.Thread | None = None _BOOT_MONO = time.monotonic() +_last_tick_mono: float = 0.0 def _utcnow() -> datetime: @@ -68,11 +69,13 @@ def try_start_scheduled_collect() -> str | None: def _loop() -> None: + global _last_tick_mono tick = max(15, int(getattr(settings, "lldp_collect_scheduler_tick_sec", 60) or 60)) grace = max(0, int(getattr(settings, "lldp_collect_startup_grace_sec", 3600) or 0)) _log.info("lldp_collect scheduler started tick=%ss startup_grace=%ss", tick, grace) while not _stop.is_set(): try: + _last_tick_mono = time.monotonic() try_start_scheduled_collect() except Exception: _log.exception("lldp_collect scheduler tick failed") @@ -94,3 +97,12 @@ def start_lldp_collect_scheduler() -> None: def stop_lldp_collect_scheduler() -> None: _stop.set() + + +def lldp_collect_scheduler_status() -> dict: + now = time.monotonic() + return { + "running": bool(_thread is not None and _thread.is_alive()), + "last_tick_age_sec": (now - _last_tick_mono) if _last_tick_mono else None, + "startup_grace_remaining_sec": round(startup_grace_remaining_sec(), 1), + } diff --git a/netx_api/metrics_router.py b/netx_api/metrics_router.py index 56e4ee1..51c712f 100644 --- a/netx_api/metrics_router.py +++ b/netx_api/metrics_router.py @@ -39,11 +39,21 @@ def collect_runtime_metrics() -> dict[str, Any]: except Exception: # noqa: BLE001 pass try: - from .port_traffic_scheduler import port_traffic_scheduler_status + from .scheduler_heartbeat import resolve_device_scheduler_metrics - out["port_traffic"] = port_traffic_scheduler_status() + sched = resolve_device_scheduler_metrics() + out["device_schedulers"] = sched + # Convenience aliases (prefer heartbeat when split; local when inline). + out["port_traffic"] = sched.get("port_traffic") or {"running": False} + out["config_sync"] = sched.get("config_sync") or {"running": False} + out["lldp_collect"] = sched.get("lldp_collect") or {"running": False} except Exception: # noqa: BLE001 - pass + try: + from .port_traffic_scheduler import port_traffic_scheduler_status + + out["port_traffic"] = port_traffic_scheduler_status() + except Exception: # noqa: BLE001 + pass try: from .webcrt_session_registry import active_session_count, list_sessions @@ -89,9 +99,22 @@ def _prom_lines(metrics: dict[str, Any]) -> str: ): if key in fwd: lines.append(f"{prom} {int(fwd.get(key) or 0)}") + sched = metrics.get("device_schedulers") or {} + if "stale" in sched: + lines.append(f'netx_device_schedulers_stale {1 if sched.get("stale") else 0}') + if sched.get("age_sec") is not None: + lines.append(f'netx_device_schedulers_heartbeat_age_seconds {sched["age_sec"]}') + for name, key in ( + ("config_sync", "netx_config_sync_scheduler_running"), + ("lldp_collect", "netx_lldp_collect_scheduler_running"), + ("port_traffic", "netx_port_traffic_scheduler_running"), + ): + block = sched.get(name) or metrics.get(name) or {} + if "running" in block: + lines.append(f'{key} {1 if block.get("running") else 0}') + if block.get("last_tick_age_sec") is not None: + lines.append(f"netx_{name}_tick_age_seconds {block['last_tick_age_sec']}") pt = metrics.get("port_traffic") or {} - if pt.get("last_tick_age_sec") is not None: - lines.append(f'netx_port_traffic_tick_age_seconds {pt["last_tick_age_sec"]}') if pt.get("last_purge_age_sec") is not None: lines.append(f'netx_port_traffic_purge_age_seconds {pt["last_purge_age_sec"]}') web = metrics.get("webcrt") or {} diff --git a/netx_api/scheduler_heartbeat.py b/netx_api/scheduler_heartbeat.py new file mode 100644 index 0000000..6a9907f --- /dev/null +++ b/netx_api/scheduler_heartbeat.py @@ -0,0 +1,188 @@ +"""Cross-process device-scheduler heartbeat for API /metrics when worker is split. + +When ``NETX_RUN_INLINE_SCHEDULERS=false``, collectors live in ``python -m netx_api.worker``. +The API process cannot see their threads; the worker publishes a small JSON heartbeat +that ``/metrics`` and ``/health/ready`` can read. +""" + +from __future__ import annotations + +import json +import logging +import os +import tempfile +import threading +import time +from pathlib import Path +from typing import Any + +from .config import settings + +_log = logging.getLogger("netx.scheduler.heartbeat") + +_HB_LOCK = threading.Lock() +_publisher_stop = threading.Event() +_publisher_thread: threading.Thread | None = None + +# Consider worker gone if heartbeat older than this (worker publishes every ~5s). +DEFAULT_STALE_SEC = 45.0 + + +def heartbeat_path() -> Path: + raw = str(getattr(settings, "scheduler_heartbeat_path", "") or "").strip() + if raw: + return Path(raw) + return Path("data") / "runtime" / "scheduler_heartbeat.json" + + +def local_device_scheduler_status(*, role: str = "unknown") -> dict[str, Any]: + """In-process scheduler thread status (API inline or worker).""" + out: dict[str, Any] = { + "pid": os.getpid(), + "role": role, + "updated_at": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()), + "updated_mono": time.monotonic(), + } + try: + from .config_sync_scheduler import config_sync_scheduler_status + + out["config_sync"] = config_sync_scheduler_status() + except Exception: # noqa: BLE001 + out["config_sync"] = {"running": False, "error": "unavailable"} + try: + from .lldp_collect_scheduler import lldp_collect_scheduler_status + + out["lldp_collect"] = lldp_collect_scheduler_status() + except Exception: # noqa: BLE001 + out["lldp_collect"] = {"running": False, "error": "unavailable"} + try: + from .port_traffic_scheduler import port_traffic_scheduler_status + + out["port_traffic"] = port_traffic_scheduler_status() + except Exception: # noqa: BLE001 + out["port_traffic"] = {"running": False, "error": "unavailable"} + return out + + +def publish_scheduler_heartbeat(*, role: str = "worker") -> Path: + """Atomically write local scheduler status for the API process to read.""" + path = heartbeat_path() + path.parent.mkdir(parents=True, exist_ok=True) + payload = local_device_scheduler_status(role=role) + # Prefer wall-clock for cross-process age; drop mono (not comparable across processes). + payload.pop("updated_mono", None) + payload["updated_at_epoch"] = time.time() + text = json.dumps(payload, ensure_ascii=False, separators=(",", ":")) + fd, tmp_name = tempfile.mkstemp(prefix=".hb-", suffix=".json", dir=str(path.parent)) + try: + with os.fdopen(fd, "w", encoding="utf-8") as fh: + fh.write(text) + fh.flush() + os.fsync(fh.fileno()) + os.replace(tmp_name, path) + except Exception: + try: + os.unlink(tmp_name) + except OSError: + pass + raise + return path + + +def read_scheduler_heartbeat(*, max_age_sec: float = DEFAULT_STALE_SEC) -> dict[str, Any] | None: + path = heartbeat_path() + try: + raw = path.read_text(encoding="utf-8") + data = json.loads(raw) + except FileNotFoundError: + return None + except Exception: # noqa: BLE001 + _log.debug("scheduler heartbeat read failed path=%s", path, exc_info=True) + return None + if not isinstance(data, dict): + return None + epoch = float(data.get("updated_at_epoch") or 0) + age = (time.time() - epoch) if epoch > 0 else None + data["age_sec"] = round(age, 1) if age is not None else None + data["stale"] = bool(age is None or age > float(max_age_sec)) + return data + + +def resolve_device_scheduler_metrics() -> dict[str, Any]: + """API-facing view: local threads when inline, else worker heartbeat file.""" + inline = bool(getattr(settings, "run_inline_schedulers", True)) + if inline: + local = local_device_scheduler_status(role="api_inline") + local.pop("updated_mono", None) + return { + "mode": "inline", + "source": "local", + "stale": False, + "hint": None, + **local, + } + + hb = read_scheduler_heartbeat() + if hb is None: + return { + "mode": "external_worker", + "source": "missing", + "stale": True, + "hint": "run `python -m netx_api.worker` (start_netx scripts do this by default)", + "pid": None, + "role": None, + "config_sync": {"running": False}, + "lldp_collect": {"running": False}, + "port_traffic": {"running": False}, + } + + return { + "mode": "external_worker", + "source": "heartbeat", + "stale": bool(hb.get("stale")), + "hint": "worker heartbeat stale — check worker.pid / restart start_netx" + if hb.get("stale") + else None, + "pid": hb.get("pid"), + "role": hb.get("role"), + "updated_at": hb.get("updated_at"), + "age_sec": hb.get("age_sec"), + "config_sync": hb.get("config_sync") or {"running": False}, + "lldp_collect": hb.get("lldp_collect") or {"running": False}, + "port_traffic": hb.get("port_traffic") or {"running": False}, + } + + +def _publisher_loop(*, role: str, interval_sec: float) -> None: + _log.info("scheduler heartbeat publisher started role=%s interval=%ss path=%s", role, interval_sec, heartbeat_path()) + while not _publisher_stop.is_set(): + try: + publish_scheduler_heartbeat(role=role) + except Exception: # noqa: BLE001 + _log.exception("scheduler heartbeat publish failed") + _publisher_stop.wait(max(1.0, float(interval_sec))) + _log.info("scheduler heartbeat publisher stopped") + + +def start_scheduler_heartbeat_publisher(*, role: str = "worker", interval_sec: float = 5.0) -> None: + """Daemon thread: keep heartbeat fresh while this process owns device schedulers.""" + global _publisher_thread + with _HB_LOCK: + if _publisher_thread and _publisher_thread.is_alive(): + return + _publisher_stop.clear() + _publisher_thread = threading.Thread( + target=_publisher_loop, + kwargs={"role": role, "interval_sec": interval_sec}, + name="scheduler-heartbeat", + daemon=True, + ) + _publisher_thread.start() + try: + publish_scheduler_heartbeat(role=role) + except Exception: # noqa: BLE001 + _log.exception("initial scheduler heartbeat publish failed") + + +def stop_scheduler_heartbeat_publisher() -> None: + _publisher_stop.set() diff --git a/netx_api/ume_runtime.py b/netx_api/ume_runtime.py index 36b1d7a..3ab0a47 100644 --- a/netx_api/ume_runtime.py +++ b/netx_api/ume_runtime.py @@ -51,10 +51,14 @@ def start_device_schedulers() -> None: from .config_sync_scheduler import start_config_sync_scheduler from .lldp_collect_scheduler import start_lldp_collect_scheduler from .port_traffic_scheduler import start_port_traffic_scheduler + from .scheduler_heartbeat import start_scheduler_heartbeat_publisher start_config_sync_scheduler() start_lldp_collect_scheduler() start_port_traffic_scheduler() + # Publish status so API /metrics can see collectors when run in a split worker. + role = "api_inline" if bool(getattr(settings, "run_inline_schedulers", True)) else "worker" + start_scheduler_heartbeat_publisher(role=role) _log.info("device schedulers started") diff --git a/netx_api/webcrt_io.py b/netx_api/webcrt_io.py new file mode 100644 index 0000000..b2c38ef --- /dev/null +++ b/netx_api/webcrt_io.py @@ -0,0 +1,42 @@ +"""Dedicated thread pool for WebCRT blocking I/O (stdout take / stdin / resize). + +Keeps session pumps off the default asyncio executor so HTTP handlers and other +``run_in_executor(None, ...)`` callers are not starved under many CRT sessions. +""" + +from __future__ import annotations + +import logging +import threading +from concurrent.futures import ThreadPoolExecutor + +from .config import settings + +_log = logging.getLogger("netx.webcrt.io") +_lock = threading.Lock() +_executor: ThreadPoolExecutor | None = None + + +def webcrt_io_executor() -> ThreadPoolExecutor: + global _executor + with _lock: + if _executor is None: + # Cap workers: each attached WS holds one blocking take_stdout wait. + n = max(4, min(48, int(getattr(settings, "webcrt_max_sessions", 40) or 40))) + _executor = ThreadPoolExecutor(max_workers=n, thread_name_prefix="webcrt-io") + _log.info("webcrt io executor started workers=%s", n) + return _executor + + +def shutdown_webcrt_io_executor(*, wait: bool = False) -> None: + global _executor + with _lock: + if _executor is None: + return + try: + _executor.shutdown(wait=wait, cancel_futures=True) + except TypeError: + _executor.shutdown(wait=wait) + except Exception: # noqa: BLE001 + _log.exception("webcrt io executor shutdown failed") + _executor = None diff --git a/netx_api/webcrt_router.py b/netx_api/webcrt_router.py index 6715fad..826d347 100644 --- a/netx_api/webcrt_router.py +++ b/netx_api/webcrt_router.py @@ -17,6 +17,7 @@ from .auth_deps import AuthContext, require_user, resolve_user_from_token from .auth_scopes import SCOPE_WEBCRT, has_scope from .config import settings from .webcrt_tickets import consume_ws_ticket, issue_ws_ticket +from .webcrt_io import webcrt_io_executor from .webcrt_service import ( close_session, create_session, @@ -450,7 +451,7 @@ async def websocket_session(websocket: WebSocket, session_id: str) -> None: slice_timeout = min(1.0, max(0.2, remaining)) try: await loop.run_in_executor( - None, + webcrt_io_executor(), lambda t=slice_timeout: wait_session_ready(session_id, timeout=t), ) sess = get_session(session_id) or cur @@ -533,7 +534,7 @@ async def websocket_session(websocket: WebSocket, session_id: str) -> None: data = "".join(stdin_buf) stdin_buf = [] try: - await asyncio.get_running_loop().run_in_executor(None, sess.write_stdin, data) + await asyncio.get_running_loop().run_in_executor(webcrt_io_executor(), sess.write_stdin, data) except Exception as exc: await websocket.send_json( {"type": "status", "state": "error", "message": f"write_failed:{exc}"} @@ -589,7 +590,7 @@ async def websocket_session(websocket: WebSocket, session_id: str) -> None: while not stop.is_set(): # Longer block is cheap now (Condition wait); cuts executor churn when idle. chunk = await loop.run_in_executor( - None, lambda: sess.take_stdout(attach_gen, timeout=0.2) + webcrt_io_executor(), lambda: sess.take_stdout(attach_gen, timeout=0.2) ) if chunk == "stale": break @@ -626,7 +627,7 @@ async def websocket_session(websocket: WebSocket, session_id: str) -> None: if sess.needs_live_prompt: sess.needs_live_prompt = False try: - await asyncio.get_running_loop().run_in_executor(None, sess.write_stdin, "\r") + await asyncio.get_running_loop().run_in_executor(webcrt_io_executor(), sess.write_stdin, "\r") except Exception: _log.debug("webcrt live prompt sync failed session=%s", session_id, exc_info=True) try: @@ -667,10 +668,10 @@ async def websocket_session(websocket: WebSocket, session_id: str) -> None: elif mtype == "resize": cols = int(msg.get("cols") or sess.cols) rows = int(msg.get("rows") or sess.rows) - await asyncio.get_running_loop().run_in_executor(None, sess.resize, cols, rows) + await asyncio.get_running_loop().run_in_executor(webcrt_io_executor(), sess.resize, cols, rows) elif mtype == "break": try: - await asyncio.get_running_loop().run_in_executor(None, sess.send_break) + await asyncio.get_running_loop().run_in_executor(webcrt_io_executor(), sess.send_break) except Exception as exc: await websocket.send_json( {"type": "status", "state": "error", "message": f"break_failed:{exc}"} diff --git a/tests/test_stability_hardening.py b/tests/test_stability_hardening.py index d8c53b5..f608001 100644 --- a/tests/test_stability_hardening.py +++ b/tests/test_stability_hardening.py @@ -13,7 +13,13 @@ from netx_api.runtime_budget import log_runtime_budget class StabilityHardeningTests(unittest.TestCase): def test_production_defaults(self) -> None: - s = Settings(_env_file=None) + import os + from unittest.mock import patch + + # Ignore process env (one-click start may set NETX_RUN_INLINE_SCHEDULERS=false). + clean = {k: v for k, v in os.environ.items() if not k.startswith("NETX_")} + with patch.dict(os.environ, clean, clear=True): + s = Settings(_env_file=None) self.assertEqual(s.db_pool_size, 40) self.assertEqual(s.db_max_overflow, 40) self.assertEqual(s.cli_max_concurrent, 24) @@ -33,6 +39,7 @@ class StabilityHardeningTests(unittest.TestCase): self.assertEqual(s.ume_raw_json_max_bytes, 64 * 1024) self.assertEqual(s.ne_collection_keep_days, 14) self.assertTrue(s.run_inline_schedulers) + self.assertEqual(s.scheduler_heartbeat_path, "data/runtime/scheduler_heartbeat.json") # Pool should cover CLI + multi-user HTTP/WS reserve under defaults. self.assertGreaterEqual(s.db_pool_size + s.db_max_overflow, s.cli_max_concurrent + 24) @@ -65,10 +72,38 @@ class StabilityHardeningTests(unittest.TestCase): self.assertIn("netx_thread_count", body) self.assertIn("netx_cli_budget_limit", body) self.assertIn("netx_oclaw_forwarder_dropped", body) + self.assertIn("netx_device_schedulers_stale", body) def test_log_runtime_budget_does_not_raise(self) -> None: log_runtime_budget(role="test") + def test_scheduler_heartbeat_roundtrip(self) -> None: + import tempfile + from pathlib import Path + from unittest.mock import patch + + from netx_api.scheduler_heartbeat import ( + publish_scheduler_heartbeat, + read_scheduler_heartbeat, + resolve_device_scheduler_metrics, + ) + + with tempfile.TemporaryDirectory() as td: + path = Path(td) / "hb.json" + with patch("netx_api.scheduler_heartbeat.heartbeat_path", return_value=path): + publish_scheduler_heartbeat(role="worker") + hb = read_scheduler_heartbeat(max_age_sec=60) + self.assertIsNotNone(hb) + assert hb is not None + self.assertFalse(hb["stale"]) + self.assertEqual(hb.get("role"), "worker") + with patch.object(settings, "run_inline_schedulers", False): + resolved = resolve_device_scheduler_metrics() + self.assertEqual(resolved["mode"], "external_worker") + self.assertEqual(resolved["source"], "heartbeat") + self.assertFalse(resolved["stale"]) + self.assertIn("port_traffic", resolved) + if __name__ == "__main__": unittest.main()