diff --git a/.env.example b/.env.example index 1d7cabb..6b8bbe9 100644 --- a/.env.example +++ b/.env.example @@ -81,4 +81,8 @@ NETX_UME_NOTIFICATION_TOPIC=ALARM # NETX_AUDIT_QUEUE_MAX=5000 # NETX_OCLAW_FORWARD_QUEUE_MAX=5000 # NETX_OCLAW_FORWARD_MAX_RETRIES=3 +# NETX_UME_RAW_JSON_MAX_BYTES=65536 +# NETX_NE_COLLECTION_KEEP_DAYS=14 # Heavier fleets: raise CLI/DB together; also ensure Postgres max_connections and bastion session limits. +# One-click start keeps collectors inline: .\scripts\start_netx.ps1 -Background -WithWeb +# (do not require a separate worker process unless you opt out of inline schedulers) diff --git a/netx_api/app_startup.py b/netx_api/app_startup.py index 56a5ccf..650ad04 100644 --- a/netx_api/app_startup.py +++ b/netx_api/app_startup.py @@ -125,11 +125,20 @@ def run_api_startup() -> None: backfill_port_traffic_series(db) except Exception: _log.exception("startup: port_traffic series backfill failed") + try: + from .ne_collection_paths import prune_old_collection_dirs + + pruned = prune_old_collection_dirs() + if pruned: + _log.info("startup: pruned %s old ne_collection job dir(s)", pruned) + except Exception: + _log.exception("startup: ne_collection prune failed") except Exception: _log.exception("startup: ne collection / config_sync recovery failed") finally: db.close() + # One-click start (scripts/start_netx.ps1) keeps collectors inline with the API. if bool(getattr(settings, "run_inline_schedulers", True)): try: start_device_schedulers() diff --git a/netx_api/config.py b/netx_api/config.py index e0de37c..d779f2f 100644 --- a/netx_api/config.py +++ b/netx_api/config.py @@ -167,6 +167,10 @@ class Settings(BaseSettings): # oclaw alarm forwarder: requeue attempts before drop on send failure. oclaw_forward_max_retries: int = 3 oclaw_forward_queue_max: int = 5000 + # Cap persisted UME raw_json blobs (inventory / alarms / sync batch meta). + ume_raw_json_max_bytes: int = 64 * 1024 + # Delete NE collection job dirs older than N days (0 disables). + ne_collection_keep_days: int = 14 settings = Settings() diff --git a/netx_api/ne_collection_paths.py b/netx_api/ne_collection_paths.py index 7bb7e3b..c15b88e 100644 --- a/netx_api/ne_collection_paths.py +++ b/netx_api/ne_collection_paths.py @@ -1,9 +1,14 @@ from __future__ import annotations +import logging +import shutil +import time from pathlib import Path from .config import settings +_log = logging.getLogger("netx.ne.collection.paths") + def collection_data_root() -> Path: root = Path(str(settings.ne_collection_data_dir or "data/ne_collections")) @@ -24,3 +29,42 @@ def clear_run_output_files(job_id: str, run_id: str) -> None: for path in out_dir.iterdir(): if path.is_file(): path.unlink(missing_ok=True) + + +def prune_old_collection_dirs(*, keep_days: int | None = None) -> int: + """Remove top-level job directories older than ``keep_days``. + + Returns number of job directories removed. Safe no-op when keep_days <= 0. + """ + days = int( + keep_days + if keep_days is not None + else (getattr(settings, "ne_collection_keep_days", 14) or 0) + ) + if days <= 0: + return 0 + root = collection_data_root() + cutoff = time.time() - (days * 86400) + removed = 0 + try: + children = list(root.iterdir()) + except Exception: # noqa: BLE001 + _log.exception("list collection root failed path=%s", root) + return 0 + for path in children: + if not path.is_dir(): + continue + try: + mtime = path.stat().st_mtime + except Exception: # noqa: BLE001 + continue + if mtime >= cutoff: + continue + try: + shutil.rmtree(path, ignore_errors=False) + removed += 1 + except Exception: # noqa: BLE001 + _log.warning("prune collection dir failed path=%s", path, exc_info=True) + if removed: + _log.info("pruned %s ne_collection job dir(s) older than %s day(s)", removed, days) + return removed diff --git a/netx_api/runtime_budget.py b/netx_api/runtime_budget.py index 672a785..9b003e3 100644 --- a/netx_api/runtime_budget.py +++ b/netx_api/runtime_budget.py @@ -48,11 +48,11 @@ def log_runtime_budget(*, role: str = "api") -> None: http_reserve, ) host = str(getattr(settings, "host", "") or "").strip().lower() - if host not in {"127.0.0.1", "localhost", "::1"} and bool( - getattr(settings, "run_inline_schedulers", True) - ): - _log.warning( - "Non-loopback bind (%s) with inline schedulers — for multi-user production prefer " - "NETX_RUN_INLINE_SCHEDULERS=false and `python -m netx_api.worker`.", - host, + inline = bool(getattr(settings, "run_inline_schedulers", True)) + if inline: + _log.info( + "runtime budget: inline schedulers ON (one-click start via scripts/start_netx.ps1). " + "Optional split: NETX_RUN_INLINE_SCHEDULERS=false + python -m netx_api.worker" ) + elif host not in {"127.0.0.1", "localhost", "::1"}: + _log.info("runtime budget: external worker mode on bind=%s", host) diff --git a/netx_api/ume_alarm_apply.py b/netx_api/ume_alarm_apply.py index c01f77d..16da957 100644 --- a/netx_api/ume_alarm_apply.py +++ b/netx_api/ume_alarm_apply.py @@ -15,6 +15,7 @@ from sqlalchemy.orm import Session from .config import settings from .models import UmeAlarmCurrent +from .ume_raw import dumps_ume_raw from .ume_sync_common import _lookup_host_name, _pick, _s, _utc_now_naive def _alarm_key(alarm: dict[str, Any]) -> str: @@ -178,7 +179,7 @@ def _alarm_row_from_norm(key: str, norm: dict[str, Any], *, touch_ts: datetime, "notification_id": notification_id_from_norm(norm), "first_seen_at": first_seen_at, "last_seen_at": touch_ts, - "raw_json": json.dumps(norm, ensure_ascii=False, default=str), + "raw_json": dumps_ume_raw(norm), } @@ -197,7 +198,7 @@ def _apply_row_to_model(db: Session, existing: UmeAlarmCurrent, norm: dict[str, prev_seen = existing.last_seen_at if prev_seen is None or touch_ts >= prev_seen: existing.last_seen_at = touch_ts - existing.raw_json = json.dumps(norm, ensure_ascii=False, default=str) + existing.raw_json = dumps_ume_raw(norm) def _upsert_alarm_current(db: Session, key: str, norm: dict[str, Any], *, touch_ts: datetime) -> tuple[str, bool]: diff --git a/netx_api/ume_raw.py b/netx_api/ume_raw.py new file mode 100644 index 0000000..3e74ea8 --- /dev/null +++ b/netx_api/ume_raw.py @@ -0,0 +1,29 @@ +"""Bounded JSON persistence for UME raw payloads.""" +from __future__ import annotations + +import json +from typing import Any + +from .config import settings + + +def dumps_ume_raw(obj: Any, *, max_bytes: int | None = None) -> str: + """Serialize UME payloads with a hard byte cap to protect DB/RSS. + + Oversized payloads are truncated as UTF-8 bytes and annotated so operators + can tell the raw blob is incomplete. + """ + text = json.dumps(obj, ensure_ascii=False, default=str) + limit = int( + max_bytes + if max_bytes is not None + else (getattr(settings, "ume_raw_json_max_bytes", 65536) or 65536) + ) + if limit <= 0: + return text + raw = text.encode("utf-8", errors="replace") + if len(raw) <= limit: + return text + keep = max(64, limit - 96) + truncated = raw[:keep].decode("utf-8", errors="replace") + return f"{truncated}\n/* truncated raw_json {limit} bytes */" diff --git a/netx_api/ume_sync_pull.py b/netx_api/ume_sync_pull.py index 236d38f..e113055 100644 --- a/netx_api/ume_sync_pull.py +++ b/netx_api/ume_sync_pull.py @@ -28,6 +28,7 @@ from .ume_alarm_apply import ( normalize_yang_alarm, notification_id_from_norm, ) +from .ume_raw import dumps_ume_raw from .ume_client import UMEClient from .ume_sync_common import ( _backfill_alarm_host_names, @@ -205,7 +206,7 @@ def sync_inventory_full(db: Session, client: UMEClient, *, trigger_mode: str = " existing.creator = _s(_pick(row, "creator")) existing.vendor = _s(_pick(row, "vendor-name")) or "ZTE" existing.last_seen_at = now - existing.raw_json = json.dumps(row, ensure_ascii=False, default=str) + existing.raw_json = dumps_ume_raw(row) _propagate_host_name_to_alarms(db, ne_id, existing.host_name) db.flush() @@ -352,7 +353,7 @@ def _sync_alarms_common( ) existing.notification_id = notification_id_from_norm(alarm) existing.last_seen_at = touch_ts - existing.raw_json = json.dumps(alarm, ensure_ascii=False, default=str) + existing.raw_json = dumps_ume_raw(alarm) iterator_500_as_end = bool(getattr(settings, "ume_iterator_500_as_end", True)) pages, meta = _collect_marker_pages( @@ -408,7 +409,7 @@ def _sync_alarms_common( batch.failed_rows = max(0, int(pulled) - int(inserted + updated)) batch.status = "done" batch.ended_at = _utc_now_naive() - batch.raw_json = json.dumps( + batch.raw_json = dumps_ume_raw( { "pulled": pulled, "inserted": inserted, @@ -425,8 +426,7 @@ def _sync_alarms_common( "reconcile_mode": reconcile_mode if not is_uncleared else "", "wss_active_during_sync": bool(wss_active) if not is_uncleared else False, "seen_keys_count": len(seen_keys) if not is_uncleared else 0, - }, - ensure_ascii=False, + } ) job.status = "done" diff --git a/scripts/start_netx.ps1 b/scripts/start_netx.ps1 index 6320121..7e966b1 100644 --- a/scripts/start_netx.ps1 +++ b/scripts/start_netx.ps1 @@ -73,6 +73,8 @@ Write-Host "" Write-Host "==> netx API URL" Write-Host "Base: $baseUrl/" Write-Host "Health: $baseUrl/health" +Write-Host "Ready: $baseUrl/health/ready" +Write-Host "Metrics: $baseUrl/metrics" Write-Host "Integrations: $baseUrl/v1/integrations/status" if ($WithWeb) { Write-Host "" diff --git a/tests/test_stability_hardening.py b/tests/test_stability_hardening.py index f44b46a..d8c53b5 100644 --- a/tests/test_stability_hardening.py +++ b/tests/test_stability_hardening.py @@ -30,6 +30,9 @@ class StabilityHardeningTests(unittest.TestCase): self.assertEqual(s.oclaw_forward_queue_max, 5000) self.assertEqual(s.ne_collect_max_output_bytes, 8 * 1024 * 1024) self.assertEqual(s.webcrt_session_log_max_bytes, 4 * 1024 * 1024) + self.assertEqual(s.ume_raw_json_max_bytes, 64 * 1024) + self.assertEqual(s.ne_collection_keep_days, 14) + self.assertTrue(s.run_inline_schedulers) # 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) diff --git a/tests/test_ume_raw_and_collection_prune.py b/tests/test_ume_raw_and_collection_prune.py new file mode 100644 index 0000000..519691e --- /dev/null +++ b/tests/test_ume_raw_and_collection_prune.py @@ -0,0 +1,57 @@ +"""Tests for UME raw_json capping and NE collection dir prune.""" + +from __future__ import annotations + +import json +import tempfile +import time +import unittest +from pathlib import Path +from unittest.mock import patch + +from netx_api.config import Settings +from netx_api.ne_collection_paths import prune_old_collection_dirs +from netx_api.ume_raw import dumps_ume_raw + + +class UmeRawAndCollectionPruneTests(unittest.TestCase): + def test_dumps_ume_raw_truncates(self) -> None: + payload = {"x": "a" * 10000} + out = dumps_ume_raw(payload, max_bytes=200) + self.assertLessEqual(len(out.encode("utf-8")), 200 + 32) + self.assertIn("truncated raw_json", out) + + def test_dumps_ume_raw_unlimited(self) -> None: + payload = {"ok": True, "n": 1} + out = dumps_ume_raw(payload, max_bytes=0) + self.assertEqual(json.loads(out)["ok"], True) + + def test_production_defaults_include_budgets(self) -> None: + s = Settings(_env_file=None) + self.assertEqual(s.ume_raw_json_max_bytes, 64 * 1024) + self.assertEqual(s.ne_collection_keep_days, 14) + self.assertTrue(s.run_inline_schedulers) + + def test_prune_old_collection_dirs(self) -> None: + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + old = root / "job_old" + new = root / "job_new" + old.mkdir() + new.mkdir() + (old / "a.txt").write_text("x", encoding="utf-8") + (new / "b.txt").write_text("y", encoding="utf-8") + # Make old directory appear stale. + old_mtime = time.time() - (20 * 86400) + import os + + os.utime(old, (old_mtime, old_mtime)) + with patch("netx_api.ne_collection_paths.collection_data_root", return_value=root): + removed = prune_old_collection_dirs(keep_days=14) + self.assertEqual(removed, 1) + self.assertFalse(old.exists()) + self.assertTrue(new.exists()) + + +if __name__ == "__main__": + unittest.main()