mirror of
https://github.com/hansjone/netx.git
synced 2026-10-09 03:10:46 +08:00
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 <cursoragent@cursor.com>
This commit is contained in:
parent
7b62824c2c
commit
f64835749a
12 changed files with 366 additions and 13 deletions
|
|
@ -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
|
# 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).
|
# 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`
|
# 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_RUN_INLINE_SCHEDULERS=false
|
||||||
|
# NETX_SCHEDULER_HEARTBEAT_PATH=data/runtime/scheduler_heartbeat.json
|
||||||
# --- Multi-user shared-server capacity (defaults in Settings already match these) ---
|
# --- Multi-user shared-server capacity (defaults in Settings already match these) ---
|
||||||
# NETX_DB_POOL_SIZE=40
|
# NETX_DB_POOL_SIZE=40
|
||||||
# NETX_DB_MAX_OVERFLOW=40
|
# NETX_DB_MAX_OVERFLOW=40
|
||||||
|
|
|
||||||
|
|
@ -10,6 +10,13 @@ def shutdown_runtime(*, reason: str = "lifespan") -> None:
|
||||||
"""Best-effort stop of schedulers, sidebands, pools, and sessions."""
|
"""Best-effort stop of schedulers, sidebands, pools, and sessions."""
|
||||||
_log.info("shutdown_runtime begin reason=%s", reason)
|
_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:
|
try:
|
||||||
from .config_sync_scheduler import stop_config_sync_scheduler
|
from .config_sync_scheduler import stop_config_sync_scheduler
|
||||||
|
|
||||||
|
|
@ -67,6 +74,13 @@ def shutdown_runtime(*, reason: str = "lifespan") -> None:
|
||||||
except Exception: # noqa: BLE001
|
except Exception: # noqa: BLE001
|
||||||
_log.exception("shutdown_cli_timeout_pool failed")
|
_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:
|
try:
|
||||||
from .audit_async import shutdown_audit_worker
|
from .audit_async import shutdown_audit_worker
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -144,6 +144,8 @@ class Settings(BaseSettings):
|
||||||
# When true (default), API also runs config_sync / lldp / port_traffic schedulers.
|
# 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.
|
# Production split: set false and run `python -m netx_api.worker` beside the API.
|
||||||
run_inline_schedulers: bool = True
|
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.
|
# SQLAlchemy QueuePool for multi-user API + collectors + UME WS.
|
||||||
# Rule of thumb: pool_size + max_overflow >= HTTP/WS peak + cli_max_concurrent + sidebands.
|
# Rule of thumb: pool_size + max_overflow >= HTTP/WS peak + cli_max_concurrent + sidebands.
|
||||||
db_pool_size: int = 40
|
db_pool_size: int = 40
|
||||||
|
|
|
||||||
|
|
@ -23,6 +23,7 @@ _log = logging.getLogger("netx.config_sync.scheduler")
|
||||||
_stop = threading.Event()
|
_stop = threading.Event()
|
||||||
_thread: threading.Thread | None = None
|
_thread: threading.Thread | None = None
|
||||||
_BOOT_MONO = time.monotonic()
|
_BOOT_MONO = time.monotonic()
|
||||||
|
_last_tick_mono: float = 0.0
|
||||||
|
|
||||||
|
|
||||||
def _utcnow() -> datetime:
|
def _utcnow() -> datetime:
|
||||||
|
|
@ -114,11 +115,13 @@ def try_start_scheduled_cycle() -> str | None:
|
||||||
|
|
||||||
|
|
||||||
def _loop() -> None:
|
def _loop() -> None:
|
||||||
|
global _last_tick_mono
|
||||||
tick = max(15, int(settings.config_sync_scheduler_tick_sec or 60))
|
tick = max(15, int(settings.config_sync_scheduler_tick_sec or 60))
|
||||||
grace = max(0, int(settings.config_sync_startup_grace_sec or 0))
|
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)
|
_log.info("config_sync scheduler started tick=%ss startup_grace=%ss", tick, grace)
|
||||||
while not _stop.is_set():
|
while not _stop.is_set():
|
||||||
try:
|
try:
|
||||||
|
_last_tick_mono = time.monotonic()
|
||||||
if bool(settings.config_sync_scheduler_enabled):
|
if bool(settings.config_sync_scheduler_enabled):
|
||||||
try_start_scheduled_cycle()
|
try_start_scheduled_cycle()
|
||||||
except Exception:
|
except Exception:
|
||||||
|
|
@ -142,3 +145,12 @@ def start_config_sync_scheduler() -> None:
|
||||||
|
|
||||||
def stop_config_sync_scheduler() -> None:
|
def stop_config_sync_scheduler() -> None:
|
||||||
_stop.set()
|
_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),
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -42,13 +42,31 @@ def health_ready(db: Session = Depends(get_db)) -> dict[str, Any]:
|
||||||
out["db_pool"] = db_pool_status()
|
out["db_pool"] = db_pool_status()
|
||||||
out["cli_budget"] = cli_budget_status()
|
out["cli_budget"] = cli_budget_status()
|
||||||
inline = bool(getattr(settings, "run_inline_schedulers", True))
|
inline = bool(getattr(settings, "run_inline_schedulers", True))
|
||||||
out["schedulers"] = {
|
sched_block: dict[str, Any] = {
|
||||||
"inline": inline,
|
"inline": inline,
|
||||||
"mode": "inline" if inline else "external_worker",
|
"mode": "inline" if inline else "external_worker",
|
||||||
"hint": None
|
"hint": None
|
||||||
if inline
|
if inline
|
||||||
else "run `python -m netx_api.worker` for config_sync / lldp_collect / port_traffic",
|
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
|
return out
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -20,6 +20,7 @@ _log = logging.getLogger("netx.lldp_collect.scheduler")
|
||||||
_stop = threading.Event()
|
_stop = threading.Event()
|
||||||
_thread: threading.Thread | None = None
|
_thread: threading.Thread | None = None
|
||||||
_BOOT_MONO = time.monotonic()
|
_BOOT_MONO = time.monotonic()
|
||||||
|
_last_tick_mono: float = 0.0
|
||||||
|
|
||||||
|
|
||||||
def _utcnow() -> datetime:
|
def _utcnow() -> datetime:
|
||||||
|
|
@ -68,11 +69,13 @@ def try_start_scheduled_collect() -> str | None:
|
||||||
|
|
||||||
|
|
||||||
def _loop() -> None:
|
def _loop() -> None:
|
||||||
|
global _last_tick_mono
|
||||||
tick = max(15, int(getattr(settings, "lldp_collect_scheduler_tick_sec", 60) or 60))
|
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))
|
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)
|
_log.info("lldp_collect scheduler started tick=%ss startup_grace=%ss", tick, grace)
|
||||||
while not _stop.is_set():
|
while not _stop.is_set():
|
||||||
try:
|
try:
|
||||||
|
_last_tick_mono = time.monotonic()
|
||||||
try_start_scheduled_collect()
|
try_start_scheduled_collect()
|
||||||
except Exception:
|
except Exception:
|
||||||
_log.exception("lldp_collect scheduler tick failed")
|
_log.exception("lldp_collect scheduler tick failed")
|
||||||
|
|
@ -94,3 +97,12 @@ def start_lldp_collect_scheduler() -> None:
|
||||||
|
|
||||||
def stop_lldp_collect_scheduler() -> None:
|
def stop_lldp_collect_scheduler() -> None:
|
||||||
_stop.set()
|
_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),
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -39,11 +39,21 @@ def collect_runtime_metrics() -> dict[str, Any]:
|
||||||
except Exception: # noqa: BLE001
|
except Exception: # noqa: BLE001
|
||||||
pass
|
pass
|
||||||
try:
|
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
|
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:
|
try:
|
||||||
from .webcrt_session_registry import active_session_count, list_sessions
|
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:
|
if key in fwd:
|
||||||
lines.append(f"{prom} {int(fwd.get(key) or 0)}")
|
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 {}
|
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:
|
if pt.get("last_purge_age_sec") is not None:
|
||||||
lines.append(f'netx_port_traffic_purge_age_seconds {pt["last_purge_age_sec"]}')
|
lines.append(f'netx_port_traffic_purge_age_seconds {pt["last_purge_age_sec"]}')
|
||||||
web = metrics.get("webcrt") or {}
|
web = metrics.get("webcrt") or {}
|
||||||
|
|
|
||||||
188
netx_api/scheduler_heartbeat.py
Normal file
188
netx_api/scheduler_heartbeat.py
Normal file
|
|
@ -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()
|
||||||
|
|
@ -51,10 +51,14 @@ def start_device_schedulers() -> None:
|
||||||
from .config_sync_scheduler import start_config_sync_scheduler
|
from .config_sync_scheduler import start_config_sync_scheduler
|
||||||
from .lldp_collect_scheduler import start_lldp_collect_scheduler
|
from .lldp_collect_scheduler import start_lldp_collect_scheduler
|
||||||
from .port_traffic_scheduler import start_port_traffic_scheduler
|
from .port_traffic_scheduler import start_port_traffic_scheduler
|
||||||
|
from .scheduler_heartbeat import start_scheduler_heartbeat_publisher
|
||||||
|
|
||||||
start_config_sync_scheduler()
|
start_config_sync_scheduler()
|
||||||
start_lldp_collect_scheduler()
|
start_lldp_collect_scheduler()
|
||||||
start_port_traffic_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")
|
_log.info("device schedulers started")
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
42
netx_api/webcrt_io.py
Normal file
42
netx_api/webcrt_io.py
Normal file
|
|
@ -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
|
||||||
|
|
@ -17,6 +17,7 @@ from .auth_deps import AuthContext, require_user, resolve_user_from_token
|
||||||
from .auth_scopes import SCOPE_WEBCRT, has_scope
|
from .auth_scopes import SCOPE_WEBCRT, has_scope
|
||||||
from .config import settings
|
from .config import settings
|
||||||
from .webcrt_tickets import consume_ws_ticket, issue_ws_ticket
|
from .webcrt_tickets import consume_ws_ticket, issue_ws_ticket
|
||||||
|
from .webcrt_io import webcrt_io_executor
|
||||||
from .webcrt_service import (
|
from .webcrt_service import (
|
||||||
close_session,
|
close_session,
|
||||||
create_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))
|
slice_timeout = min(1.0, max(0.2, remaining))
|
||||||
try:
|
try:
|
||||||
await loop.run_in_executor(
|
await loop.run_in_executor(
|
||||||
None,
|
webcrt_io_executor(),
|
||||||
lambda t=slice_timeout: wait_session_ready(session_id, timeout=t),
|
lambda t=slice_timeout: wait_session_ready(session_id, timeout=t),
|
||||||
)
|
)
|
||||||
sess = get_session(session_id) or cur
|
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)
|
data = "".join(stdin_buf)
|
||||||
stdin_buf = []
|
stdin_buf = []
|
||||||
try:
|
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:
|
except Exception as exc:
|
||||||
await websocket.send_json(
|
await websocket.send_json(
|
||||||
{"type": "status", "state": "error", "message": f"write_failed:{exc}"}
|
{"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():
|
while not stop.is_set():
|
||||||
# Longer block is cheap now (Condition wait); cuts executor churn when idle.
|
# Longer block is cheap now (Condition wait); cuts executor churn when idle.
|
||||||
chunk = await loop.run_in_executor(
|
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":
|
if chunk == "stale":
|
||||||
break
|
break
|
||||||
|
|
@ -626,7 +627,7 @@ async def websocket_session(websocket: WebSocket, session_id: str) -> None:
|
||||||
if sess.needs_live_prompt:
|
if sess.needs_live_prompt:
|
||||||
sess.needs_live_prompt = False
|
sess.needs_live_prompt = False
|
||||||
try:
|
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:
|
except Exception:
|
||||||
_log.debug("webcrt live prompt sync failed session=%s", session_id, exc_info=True)
|
_log.debug("webcrt live prompt sync failed session=%s", session_id, exc_info=True)
|
||||||
try:
|
try:
|
||||||
|
|
@ -667,10 +668,10 @@ async def websocket_session(websocket: WebSocket, session_id: str) -> None:
|
||||||
elif mtype == "resize":
|
elif mtype == "resize":
|
||||||
cols = int(msg.get("cols") or sess.cols)
|
cols = int(msg.get("cols") or sess.cols)
|
||||||
rows = int(msg.get("rows") or sess.rows)
|
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":
|
elif mtype == "break":
|
||||||
try:
|
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:
|
except Exception as exc:
|
||||||
await websocket.send_json(
|
await websocket.send_json(
|
||||||
{"type": "status", "state": "error", "message": f"break_failed:{exc}"}
|
{"type": "status", "state": "error", "message": f"break_failed:{exc}"}
|
||||||
|
|
|
||||||
|
|
@ -13,7 +13,13 @@ from netx_api.runtime_budget import log_runtime_budget
|
||||||
|
|
||||||
class StabilityHardeningTests(unittest.TestCase):
|
class StabilityHardeningTests(unittest.TestCase):
|
||||||
def test_production_defaults(self) -> None:
|
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_pool_size, 40)
|
||||||
self.assertEqual(s.db_max_overflow, 40)
|
self.assertEqual(s.db_max_overflow, 40)
|
||||||
self.assertEqual(s.cli_max_concurrent, 24)
|
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.ume_raw_json_max_bytes, 64 * 1024)
|
||||||
self.assertEqual(s.ne_collection_keep_days, 14)
|
self.assertEqual(s.ne_collection_keep_days, 14)
|
||||||
self.assertTrue(s.run_inline_schedulers)
|
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.
|
# 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)
|
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_thread_count", body)
|
||||||
self.assertIn("netx_cli_budget_limit", body)
|
self.assertIn("netx_cli_budget_limit", body)
|
||||||
self.assertIn("netx_oclaw_forwarder_dropped", 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:
|
def test_log_runtime_budget_does_not_raise(self) -> None:
|
||||||
log_runtime_budget(role="test")
|
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__":
|
if __name__ == "__main__":
|
||||||
unittest.main()
|
unittest.main()
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue