mirror of
https://github.com/hansjone/netx.git
synced 2026-10-09 03:10:46 +08:00
Harden retention under one-click start: UME raw caps and collection TTL.
Keep inline schedulers as the start_netx.ps1 default while bounding UME raw_json and pruning old NE collection dirs on startup. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
parent
e0503df8a5
commit
96be90bd1a
11 changed files with 167 additions and 14 deletions
|
|
@ -81,4 +81,8 @@ NETX_UME_NOTIFICATION_TOPIC=ALARM
|
||||||
# NETX_AUDIT_QUEUE_MAX=5000
|
# NETX_AUDIT_QUEUE_MAX=5000
|
||||||
# NETX_OCLAW_FORWARD_QUEUE_MAX=5000
|
# NETX_OCLAW_FORWARD_QUEUE_MAX=5000
|
||||||
# NETX_OCLAW_FORWARD_MAX_RETRIES=3
|
# 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.
|
# 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)
|
||||||
|
|
|
||||||
|
|
@ -125,11 +125,20 @@ def run_api_startup() -> None:
|
||||||
backfill_port_traffic_series(db)
|
backfill_port_traffic_series(db)
|
||||||
except Exception:
|
except Exception:
|
||||||
_log.exception("startup: port_traffic series backfill failed")
|
_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:
|
except Exception:
|
||||||
_log.exception("startup: ne collection / config_sync recovery failed")
|
_log.exception("startup: ne collection / config_sync recovery failed")
|
||||||
finally:
|
finally:
|
||||||
db.close()
|
db.close()
|
||||||
|
|
||||||
|
# One-click start (scripts/start_netx.ps1) keeps collectors inline with the API.
|
||||||
if bool(getattr(settings, "run_inline_schedulers", True)):
|
if bool(getattr(settings, "run_inline_schedulers", True)):
|
||||||
try:
|
try:
|
||||||
start_device_schedulers()
|
start_device_schedulers()
|
||||||
|
|
|
||||||
|
|
@ -167,6 +167,10 @@ class Settings(BaseSettings):
|
||||||
# oclaw alarm forwarder: requeue attempts before drop on send failure.
|
# oclaw alarm forwarder: requeue attempts before drop on send failure.
|
||||||
oclaw_forward_max_retries: int = 3
|
oclaw_forward_max_retries: int = 3
|
||||||
oclaw_forward_queue_max: int = 5000
|
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()
|
settings = Settings()
|
||||||
|
|
|
||||||
|
|
@ -1,9 +1,14 @@
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
|
import shutil
|
||||||
|
import time
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
from .config import settings
|
from .config import settings
|
||||||
|
|
||||||
|
_log = logging.getLogger("netx.ne.collection.paths")
|
||||||
|
|
||||||
|
|
||||||
def collection_data_root() -> Path:
|
def collection_data_root() -> Path:
|
||||||
root = Path(str(settings.ne_collection_data_dir or "data/ne_collections"))
|
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():
|
for path in out_dir.iterdir():
|
||||||
if path.is_file():
|
if path.is_file():
|
||||||
path.unlink(missing_ok=True)
|
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
|
||||||
|
|
|
||||||
|
|
@ -48,11 +48,11 @@ def log_runtime_budget(*, role: str = "api") -> None:
|
||||||
http_reserve,
|
http_reserve,
|
||||||
)
|
)
|
||||||
host = str(getattr(settings, "host", "") or "").strip().lower()
|
host = str(getattr(settings, "host", "") or "").strip().lower()
|
||||||
if host not in {"127.0.0.1", "localhost", "::1"} and bool(
|
inline = bool(getattr(settings, "run_inline_schedulers", True))
|
||||||
getattr(settings, "run_inline_schedulers", True)
|
if inline:
|
||||||
):
|
_log.info(
|
||||||
_log.warning(
|
"runtime budget: inline schedulers ON (one-click start via scripts/start_netx.ps1). "
|
||||||
"Non-loopback bind (%s) with inline schedulers — for multi-user production prefer "
|
"Optional split: NETX_RUN_INLINE_SCHEDULERS=false + python -m netx_api.worker"
|
||||||
"NETX_RUN_INLINE_SCHEDULERS=false and `python -m netx_api.worker`.",
|
|
||||||
host,
|
|
||||||
)
|
)
|
||||||
|
elif host not in {"127.0.0.1", "localhost", "::1"}:
|
||||||
|
_log.info("runtime budget: external worker mode on bind=%s", host)
|
||||||
|
|
|
||||||
|
|
@ -15,6 +15,7 @@ from sqlalchemy.orm import Session
|
||||||
|
|
||||||
from .config import settings
|
from .config import settings
|
||||||
from .models import UmeAlarmCurrent
|
from .models import UmeAlarmCurrent
|
||||||
|
from .ume_raw import dumps_ume_raw
|
||||||
from .ume_sync_common import _lookup_host_name, _pick, _s, _utc_now_naive
|
from .ume_sync_common import _lookup_host_name, _pick, _s, _utc_now_naive
|
||||||
|
|
||||||
def _alarm_key(alarm: dict[str, Any]) -> str:
|
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),
|
"notification_id": notification_id_from_norm(norm),
|
||||||
"first_seen_at": first_seen_at,
|
"first_seen_at": first_seen_at,
|
||||||
"last_seen_at": touch_ts,
|
"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
|
prev_seen = existing.last_seen_at
|
||||||
if prev_seen is None or touch_ts >= prev_seen:
|
if prev_seen is None or touch_ts >= prev_seen:
|
||||||
existing.last_seen_at = touch_ts
|
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]:
|
def _upsert_alarm_current(db: Session, key: str, norm: dict[str, Any], *, touch_ts: datetime) -> tuple[str, bool]:
|
||||||
|
|
|
||||||
29
netx_api/ume_raw.py
Normal file
29
netx_api/ume_raw.py
Normal file
|
|
@ -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 */"
|
||||||
|
|
@ -28,6 +28,7 @@ from .ume_alarm_apply import (
|
||||||
normalize_yang_alarm,
|
normalize_yang_alarm,
|
||||||
notification_id_from_norm,
|
notification_id_from_norm,
|
||||||
)
|
)
|
||||||
|
from .ume_raw import dumps_ume_raw
|
||||||
from .ume_client import UMEClient
|
from .ume_client import UMEClient
|
||||||
from .ume_sync_common import (
|
from .ume_sync_common import (
|
||||||
_backfill_alarm_host_names,
|
_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.creator = _s(_pick(row, "creator"))
|
||||||
existing.vendor = _s(_pick(row, "vendor-name")) or "ZTE"
|
existing.vendor = _s(_pick(row, "vendor-name")) or "ZTE"
|
||||||
existing.last_seen_at = now
|
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)
|
_propagate_host_name_to_alarms(db, ne_id, existing.host_name)
|
||||||
|
|
||||||
db.flush()
|
db.flush()
|
||||||
|
|
@ -352,7 +353,7 @@ def _sync_alarms_common(
|
||||||
)
|
)
|
||||||
existing.notification_id = notification_id_from_norm(alarm)
|
existing.notification_id = notification_id_from_norm(alarm)
|
||||||
existing.last_seen_at = touch_ts
|
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))
|
iterator_500_as_end = bool(getattr(settings, "ume_iterator_500_as_end", True))
|
||||||
pages, meta = _collect_marker_pages(
|
pages, meta = _collect_marker_pages(
|
||||||
|
|
@ -408,7 +409,7 @@ def _sync_alarms_common(
|
||||||
batch.failed_rows = max(0, int(pulled) - int(inserted + updated))
|
batch.failed_rows = max(0, int(pulled) - int(inserted + updated))
|
||||||
batch.status = "done"
|
batch.status = "done"
|
||||||
batch.ended_at = _utc_now_naive()
|
batch.ended_at = _utc_now_naive()
|
||||||
batch.raw_json = json.dumps(
|
batch.raw_json = dumps_ume_raw(
|
||||||
{
|
{
|
||||||
"pulled": pulled,
|
"pulled": pulled,
|
||||||
"inserted": inserted,
|
"inserted": inserted,
|
||||||
|
|
@ -425,8 +426,7 @@ def _sync_alarms_common(
|
||||||
"reconcile_mode": reconcile_mode if not is_uncleared else "",
|
"reconcile_mode": reconcile_mode if not is_uncleared else "",
|
||||||
"wss_active_during_sync": bool(wss_active) if not is_uncleared else False,
|
"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,
|
"seen_keys_count": len(seen_keys) if not is_uncleared else 0,
|
||||||
},
|
}
|
||||||
ensure_ascii=False,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
job.status = "done"
|
job.status = "done"
|
||||||
|
|
|
||||||
|
|
@ -73,6 +73,8 @@ Write-Host ""
|
||||||
Write-Host "==> netx API URL"
|
Write-Host "==> netx API URL"
|
||||||
Write-Host "Base: $baseUrl/"
|
Write-Host "Base: $baseUrl/"
|
||||||
Write-Host "Health: $baseUrl/health"
|
Write-Host "Health: $baseUrl/health"
|
||||||
|
Write-Host "Ready: $baseUrl/health/ready"
|
||||||
|
Write-Host "Metrics: $baseUrl/metrics"
|
||||||
Write-Host "Integrations: $baseUrl/v1/integrations/status"
|
Write-Host "Integrations: $baseUrl/v1/integrations/status"
|
||||||
if ($WithWeb) {
|
if ($WithWeb) {
|
||||||
Write-Host ""
|
Write-Host ""
|
||||||
|
|
|
||||||
|
|
@ -30,6 +30,9 @@ class StabilityHardeningTests(unittest.TestCase):
|
||||||
self.assertEqual(s.oclaw_forward_queue_max, 5000)
|
self.assertEqual(s.oclaw_forward_queue_max, 5000)
|
||||||
self.assertEqual(s.ne_collect_max_output_bytes, 8 * 1024 * 1024)
|
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.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.
|
# 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)
|
||||||
|
|
||||||
|
|
|
||||||
57
tests/test_ume_raw_and_collection_prune.py
Normal file
57
tests/test_ume_raw_and_collection_prune.py
Normal file
|
|
@ -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()
|
||||||
Loading…
Add table
Add a link
Reference in a new issue