diff --git a/.env.example b/.env.example index f46e737..4a19cb0 100644 --- a/.env.example +++ b/.env.example @@ -57,17 +57,27 @@ NETX_UME_NOTIFICATION_TOPIC=ALARM # 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` # NETX_RUN_INLINE_SCHEDULERS=false -# NETX_AUDIT_ASYNC=true -# NETX_AUDIT_SAMPLE_N=1 -# DB pool (QueuePool). Aim: pool_size + max_overflow >= HTTP peak + NETX_CLI_MAX_CONCURRENT + UME WS burst. -# NETX_DB_POOL_SIZE=20 -# NETX_DB_MAX_OVERFLOW=20 +# --- Production-recommended capacity (defaults in Settings already match these) --- +# NETX_DB_POOL_SIZE=25 +# NETX_DB_MAX_OVERFLOW=15 # NETX_DB_POOL_RECYCLE_SEC=1800 # NETX_DB_POOL_TIMEOUT_SEC=30 -# Global SSH/Netmiko concurrency across discover / collect / config_sync / port_traffic. -# NETX_CLI_MAX_CONCURRENT=20 -# NETX_CLI_TIMEOUT_POOL_WORKERS=16 -# NETX_PORT_TRAFFIC_DISPATCH_WORKERS=4 -# NETX_AUDIT_QUEUE_MAX=5000 -# NETX_NE_COLLECT_MAX_OUTPUT_BYTES=8388608 -# NETX_WEBCRT_SESSION_LOG_MAX_BYTES=4194304 +# NETX_CLI_MAX_CONCURRENT=12 +# NETX_CLI_FEATURE_HARD_CAP=16 +# NETX_CLI_TIMEOUT_POOL_WORKERS=12 +# NETX_PORT_TRAFFIC_DISPATCH_WORKERS=3 +# NETX_NE_CONNECT_MAX_WORKERS=4 +# NETX_NE_COLLECT_MAX_WORKERS=4 +# NETX_NE_COLLECT_RUN_TIMEOUT_CAP_SEC=480 +# NETX_NE_COLLECT_MAX_OUTPUT_BYTES=4194304 +# NETX_WEBCRT_MAX_SESSIONS=12 +# NETX_WEBCRT_KEEPALIVE_SEC=30 +# NETX_WEBCRT_OUT_QUEUE_MAX=1500 +# NETX_WEBCRT_SESSION_LOG_MAX_BYTES=2097152 +# NETX_WEBCRT_SFTP_MAX_FILE_BYTES=268435456 +# NETX_AUDIT_ASYNC=true +# NETX_AUDIT_SAMPLE_N=10 +# NETX_AUDIT_QUEUE_MAX=2000 +# NETX_OCLAW_FORWARD_QUEUE_MAX=2000 +# NETX_OCLAW_FORWARD_MAX_RETRIES=3 +# Lab-only raises (optional): higher CLI/DB when bastion and Postgres are sized for it. diff --git a/netx_api/app_startup.py b/netx_api/app_startup.py index a8f19f2..56a5ccf 100644 --- a/netx_api/app_startup.py +++ b/netx_api/app_startup.py @@ -142,3 +142,9 @@ def run_api_startup() -> None: ) start_api_sideband_threads() + try: + from .runtime_budget import log_runtime_budget + + log_runtime_budget(role="api") + except Exception: + _log.exception("startup: runtime budget log failed") diff --git a/netx_api/cli_budget.py b/netx_api/cli_budget.py index 18c35bf..cabcf8e 100644 --- a/netx_api/cli_budget.py +++ b/netx_api/cli_budget.py @@ -17,7 +17,7 @@ _in_use_lock = threading.Lock() def _ensure_sem() -> threading.BoundedSemaphore: global _sem, _limit with _lock: - want = max(1, int(getattr(settings, "cli_max_concurrent", 20) or 20)) + want = max(1, int(getattr(settings, "cli_max_concurrent", 12) or 12)) if _sem is None or want != _limit: _sem = threading.BoundedSemaphore(want) _limit = want @@ -32,10 +32,20 @@ def cli_budget_status() -> dict[str, int]: return {"limit": _limit, "in_use": used, "available": max(0, _limit - used)} -def clamp_cli_workers(requested: int, *, hard_cap: int) -> int: +def clamp_cli_workers(requested: int, *, hard_cap: int | None = None) -> int: """Clamp a feature concurrency against the global CLI budget and a hard cap.""" - budget = max(1, int(getattr(settings, "cli_max_concurrent", 20) or 20)) - return max(1, min(int(hard_cap), int(requested or 1), budget)) + budget = max(1, int(getattr(settings, "cli_max_concurrent", 12) or 12)) + feature_cap = int( + hard_cap + if hard_cap is not None + else (getattr(settings, "cli_feature_hard_cap", 16) or 16) + ) + feature_cap = max(1, feature_cap) + return max(1, min(int(feature_cap), int(requested or 1), budget)) + + +def feature_hard_cap() -> int: + return max(1, int(getattr(settings, "cli_feature_hard_cap", 16) or 16)) @contextmanager diff --git a/netx_api/config.py b/netx_api/config.py index ca4852a..b189b04 100644 --- a/netx_api/config.py +++ b/netx_api/config.py @@ -64,13 +64,14 @@ class Settings(BaseSettings): ume_sync_alarms_history_every_hours: int = 24 # Managed NE credentials (Fernet key; generate with cryptography.fernet.Fernet.generate_key()) credential_secret_key: str = "" - ne_connect_max_workers: int = 5 + # Production-oriented worker caps (raise only when bastion/DB capacity allows). + ne_connect_max_workers: int = 4 ne_connect_timeout_sec: int = 30 - ne_collect_max_workers: int = 5 + ne_collect_max_workers: int = 4 ne_collect_read_timeout_sec: int = 120 ne_collect_stale_run_sec: int = 900 ne_collect_pending_stale_sec: int = 180 - ne_collect_run_timeout_cap_sec: int = 600 + ne_collect_run_timeout_cap_sec: int = 480 ne_collection_data_dir: str = "data/ne_collections" # Config sync (periodic running-config backup into DB) config_sync_scheduler_enabled: bool = True @@ -91,29 +92,29 @@ class Settings(BaseSettings): # Managed NE exec: max CLI commands per request (lab can raise; hard-capped in ne_exec). ne_exec_max_commands: int = 5 # WebCRT interactive terminal sessions - webcrt_max_sessions: int = 20 + webcrt_max_sessions: int = 12 webcrt_idle_timeout_sec: int = 1800 webcrt_connect_timeout_sec: int = 90 webcrt_attach_timeout_sec: int = 60 # Keep device PTY after WS drop so the UI can re-attach (reconnect / remount). webcrt_detach_grace_sec: int = 120 webcrt_data_dir: str = "data/webcrt" - # SSH transport keepalive interval (seconds); 0 disables (default off). - webcrt_keepalive_sec: int = 0 + # SSH transport keepalive (seconds); 30s recommended behind NAT/firewall. + webcrt_keepalive_sec: int = 30 # Device anti-idle CLI nudge (0 = off). Keep off: NEs close idle VTY themselves. webcrt_anti_idle_sec: int = 0 webcrt_anti_idle_payload: str = " " # Cap stdout queue depth (drop oldest when full) to protect memory. - webcrt_out_queue_max: int = 2000 + webcrt_out_queue_max: int = 1500 # Persist per-session transcripts under webcrt_data_dir/sessions/. webcrt_session_log_enabled: bool = True # Reader: short blocking wait instead of fixed 40ms spin (seconds). webcrt_reader_poll_sec: float = 0.01 - # WebCRT SFTP transfer limits (streamed; default 512 MiB per file). - webcrt_sftp_max_file_bytes: int = 512 * 1024 * 1024 + # WebCRT SFTP transfer limits (streamed; default 256 MiB per file). + webcrt_sftp_max_file_bytes: int = 256 * 1024 * 1024 webcrt_sftp_chunk_bytes: int = 64 * 1024 # Cap directory listings so huge folders cannot pin the API/UI. - webcrt_sftp_list_max_entries: int = 5000 + webcrt_sftp_list_max_entries: int = 2000 webcrt_sftp_list_timeout_sec: float = 30.0 # Local app login / audit (lab defaults; override in production) auth_enabled: bool = True @@ -132,7 +133,7 @@ class Settings(BaseSettings): allow_insecure_defaults: bool = False # Async audit writer; sample_n>1 keeps 1/N of generic http.* events. audit_async: bool = True - audit_sample_n: int = 1 + audit_sample_n: int = 10 # Prefer Alembic on API start; brownfield patches live in schema_patches + revisions. # Auth column ensures still run as a safety net before bootstrap. skip_legacy_startup_ddl: bool = True @@ -145,22 +146,27 @@ class Settings(BaseSettings): run_inline_schedulers: bool = True # SQLAlchemy QueuePool (API + background workers share one engine). # Rule of thumb: pool_size + max_overflow >= HTTP peak + cli_max_concurrent + UME WS burst. - db_pool_size: int = 20 - db_max_overflow: int = 20 + db_pool_size: int = 25 + db_max_overflow: int = 15 db_pool_recycle_sec: int = 1800 db_pool_timeout_sec: int = 30 # Global Netmiko/SSH concurrency across discover / collect / config_sync / port_traffic. - cli_max_concurrent: int = 20 + cli_max_concurrent: int = 12 + # Per-feature concurrency hard ceiling (API body / policy cannot exceed this). + cli_feature_hard_cap: int = 16 # Shared timeout watchdog pool (not per-task executors). - cli_timeout_pool_workers: int = 16 + cli_timeout_pool_workers: int = 12 # Port-traffic: how many devices may collect in parallel on the scheduler tick. - port_traffic_dispatch_workers: int = 4 + port_traffic_dispatch_workers: int = 3 # Bound async audit queue; drop oldest when full to protect RSS. - audit_queue_max: int = 5000 + audit_queue_max: int = 2000 # Cap NE collection output files (bytes); 0 = unlimited (not recommended). - ne_collect_max_output_bytes: int = 8 * 1024 * 1024 + ne_collect_max_output_bytes: int = 4 * 1024 * 1024 # WebCRT session transcript rotate size (bytes); 0 disables rotate. - webcrt_session_log_max_bytes: int = 4 * 1024 * 1024 + webcrt_session_log_max_bytes: int = 2 * 1024 * 1024 + # oclaw alarm forwarder: requeue attempts before drop on send failure. + oclaw_forward_max_retries: int = 3 + oclaw_forward_queue_max: int = 2000 settings = Settings() diff --git a/netx_api/config_sync_common.py b/netx_api/config_sync_common.py index fa5cba9..ef5c2ab 100644 --- a/netx_api/config_sync_common.py +++ b/netx_api/config_sync_common.py @@ -84,7 +84,7 @@ def policy_to_out(row: ConfigSyncPolicy) -> ConfigSyncPolicyOut: return ConfigSyncPolicyOut( enabled=bool(row.enabled), interval_days=max(1, int(row.interval_days or 3)), - concurrency=max(1, min(30, int(row.concurrency or 5))), + concurrency=max(1, min(16, int(row.concurrency or 5))), scope_mode=str(row.scope_mode or "all"), selected_targets=_targets_from_json(row.selected_targets), history_keep=max(0, min(30, int(row.history_keep if row.history_keep is not None else 3))), diff --git a/netx_api/config_sync_cycles.py b/netx_api/config_sync_cycles.py index e9a703b..0602a54 100644 --- a/netx_api/config_sync_cycles.py +++ b/netx_api/config_sync_cycles.py @@ -50,7 +50,7 @@ def update_policy(db: Session, body: ConfigSyncPolicyUpdate) -> ConfigSyncPolicy if "interval_days" in data and data["interval_days"] is not None: row.interval_days = int(data["interval_days"]) if "concurrency" in data and data["concurrency"] is not None: - row.concurrency = max(1, min(30, int(data["concurrency"]))) + row.concurrency = max(1, min(16, int(data["concurrency"]))) if "scope_mode" in data and data["scope_mode"] is not None: row.scope_mode = str(data["scope_mode"]) if "selected_targets" in data and data["selected_targets"] is not None: @@ -183,7 +183,7 @@ def create_cycle(db: Session, body: ConfigSyncCycleCreate) -> ConfigSyncCycleOut policy = ensure_policy(db) mode = str(body.mode or "full").strip().lower() trigger = "retry_failed" if mode == "retry_failed" else "manual" - concurrency = max(1, min(30, int(policy.concurrency or 5))) + concurrency = max(1, min(16, int(policy.concurrency or 5))) targets: list[dict[str, str]] = [] if mode == "retry_failed": diff --git a/netx_api/config_sync_runner.py b/netx_api/config_sync_runner.py index 3a064dd..aa354d6 100644 --- a/netx_api/config_sync_runner.py +++ b/netx_api/config_sync_runner.py @@ -43,7 +43,7 @@ def _pool_for_cycle(cycle_id: str, concurrency: int) -> ThreadPoolExecutor: with _pools_lock: pool = _pools.get(cycle_id) if pool is None: - workers = clamp_cli_workers(int(concurrency or 5), hard_cap=30) + workers = clamp_cli_workers(int(concurrency or 5)) pool = ThreadPoolExecutor(max_workers=workers, thread_name_prefix=f"cfg-sync-{cycle_id[:8]}") _pools[cycle_id] = pool return pool @@ -462,7 +462,7 @@ def dispatch_cycle(cycle_id: str) -> int: .all() ) task_ids = [str(t.id) for t in pending] - concurrency = max(1, min(30, int(cycle.concurrency or 5))) + concurrency = max(1, min(16, int(cycle.concurrency or 5))) finally: db.close() return schedule_cycle_tasks(cycle_id, task_ids, concurrency) diff --git a/netx_api/config_sync_scheduler.py b/netx_api/config_sync_scheduler.py index 4f052b6..3a6c4d1 100644 --- a/netx_api/config_sync_scheduler.py +++ b/netx_api/config_sync_scheduler.py @@ -74,7 +74,7 @@ def try_start_scheduled_cycle() -> str | None: if not targets: _log.info("config_sync schedule skip: no targets") return None - concurrency = max(1, min(30, int(policy.concurrency or 5))) + concurrency = max(1, min(16, int(policy.concurrency or 5))) cycle = ConfigSyncCycle( id=uuid4().hex, trigger_mode="schedule", diff --git a/netx_api/config_sync_schemas.py b/netx_api/config_sync_schemas.py index 67418dc..419696c 100644 --- a/netx_api/config_sync_schemas.py +++ b/netx_api/config_sync_schemas.py @@ -27,7 +27,7 @@ class ConfigSyncPolicyOut(BaseModel): class ConfigSyncPolicyUpdate(BaseModel): enabled: bool | None = None interval_days: int | None = Field(default=None, ge=1, le=365) - concurrency: int | None = Field(default=None, ge=1, le=30) + concurrency: int | None = Field(default=None, ge=1, le=16) scope_mode: Literal["all", "selected"] | None = None selected_targets: list[ConfigSyncTargetRef] | None = None history_keep: int | None = Field(default=None, ge=0, le=30) diff --git a/netx_api/lldp_collect_schemas.py b/netx_api/lldp_collect_schemas.py index 39e9e36..a7f72ef 100644 --- a/netx_api/lldp_collect_schemas.py +++ b/netx_api/lldp_collect_schemas.py @@ -29,7 +29,7 @@ class LldpCollectPolicyUpdate(BaseModel): enabled: bool | None = None interval_days: int | None = Field(default=None, ge=1, le=365) interval_hours: int | None = Field(default=None, ge=1, le=8760) - concurrency: int | None = Field(default=None, ge=1, le=32) + concurrency: int | None = Field(default=None, ge=1, le=16) scope_mode: str | None = None selected_targets: list[LldpCollectTargetRef] | None = None auto_add_unmatched: bool | None = None diff --git a/netx_api/lldp_collect_service.py b/netx_api/lldp_collect_service.py index 99330bd..f35a377 100644 --- a/netx_api/lldp_collect_service.py +++ b/netx_api/lldp_collect_service.py @@ -113,7 +113,7 @@ def update_policy(db: Session, body: LldpCollectPolicyUpdate) -> LldpCollectPoli row.interval_days = days row.interval_hours = days * 24 if "concurrency" in data and data["concurrency"] is not None: - row.concurrency = max(1, min(32, int(data["concurrency"]))) + row.concurrency = max(1, min(16, int(data["concurrency"]))) if "scope_mode" in data and data["scope_mode"] is not None: mode = str(data["scope_mode"] or "").strip().lower() if mode not in {"all", "selected"}: @@ -208,7 +208,7 @@ def next_due_at(db: Session, policy: LldpCollectPolicy) -> datetime | None: def build_discover_request(policy: LldpCollectPolicy) -> FabricDiscoverRequest: - concurrency = max(1, min(32, int(policy.concurrency or 4))) + concurrency = max(1, min(16, int(policy.concurrency or 4))) auto_add = bool(policy.auto_add_unmatched) if str(policy.scope_mode or "") == "selected": managed_ids: list[str] = [] diff --git a/netx_api/metrics_router.py b/netx_api/metrics_router.py index 05e430d..56e4ee1 100644 --- a/netx_api/metrics_router.py +++ b/netx_api/metrics_router.py @@ -83,6 +83,9 @@ def _prom_lines(metrics: dict[str, Any]) -> str: ("queue_size", "netx_oclaw_forwarder_queue_size"), ("published_ok", "netx_oclaw_forwarder_published_ok"), ("published_fail", "netx_oclaw_forwarder_published_fail"), + ("dropped", "netx_oclaw_forwarder_dropped"), + ("requeued", "netx_oclaw_forwarder_requeued"), + ("retry_exhausted", "netx_oclaw_forwarder_retry_exhausted"), ): if key in fwd: lines.append(f"{prom} {int(fwd.get(key) or 0)}") diff --git a/netx_api/ne_collect_runner.py b/netx_api/ne_collect_runner.py index 6d066bb..bdbf407 100644 --- a/netx_api/ne_collect_runner.py +++ b/netx_api/ne_collect_runner.py @@ -29,7 +29,7 @@ _executor: ThreadPoolExecutor | None = None def _executor_pool() -> ThreadPoolExecutor: global _executor if _executor is None: - workers = clamp_cli_workers(int(settings.ne_collect_max_workers or 5), hard_cap=32) + workers = clamp_cli_workers(int(settings.ne_collect_max_workers or 4)) _executor = ThreadPoolExecutor(max_workers=workers, thread_name_prefix="ne-collect") return _executor diff --git a/netx_api/ne_connect.py b/netx_api/ne_connect.py index 08287d2..94b7548 100644 --- a/netx_api/ne_connect.py +++ b/netx_api/ne_connect.py @@ -32,7 +32,7 @@ def _executor_pool() -> ThreadPoolExecutor: if _executor is None: from .cli_budget import clamp_cli_workers - workers = clamp_cli_workers(int(settings.ne_connect_max_workers or 5), hard_cap=32) + workers = clamp_cli_workers(int(settings.ne_connect_max_workers or 4)) _executor = ThreadPoolExecutor(max_workers=workers, thread_name_prefix="ne-connect") return _executor diff --git a/netx_api/oclaw_alarm_forwarder.py b/netx_api/oclaw_alarm_forwarder.py index f9ea655..164806b 100644 --- a/netx_api/oclaw_alarm_forwarder.py +++ b/netx_api/oclaw_alarm_forwarder.py @@ -15,7 +15,8 @@ from .config import settings _log = logging.getLogger("netx.oclaw.alarm_forwarder") -_OUTBOUND_Q: "queue.Queue[dict[str, Any]]" = queue.Queue(maxsize=5000) +_OUTBOUND_Q: "queue.Queue[dict[str, Any]] | None" = None +_Q_LOCK = threading.Lock() _STOP_EVENT = threading.Event() _THREAD: threading.Thread | None = None _CONN_LOCK = threading.Lock() @@ -28,9 +29,21 @@ _STATS: dict[str, int] = { "published_ok": 0, "published_fail": 0, "queued": 0, + "dropped": 0, + "requeued": 0, + "retry_exhausted": 0, } +def _outbound_q() -> "queue.Queue[dict[str, Any]]": + global _OUTBOUND_Q + with _Q_LOCK: + if _OUTBOUND_Q is None: + maxsize = max(100, int(getattr(settings, "oclaw_forward_queue_max", 2000) or 2000)) + _OUTBOUND_Q = queue.Queue(maxsize=maxsize) + return _OUTBOUND_Q + + def _utc_now_iso() -> str: return datetime.now(timezone.utc).isoformat() @@ -95,15 +108,46 @@ def enqueue_alarm_forward(payload: dict[str, Any]) -> bool: if not is_forwarder_operational(): return False try: - _OUTBOUND_Q.put_nowait(dict(payload)) + _outbound_q().put_nowait(dict(payload)) with _STATS_LOCK: _STATS["queued"] = int(_STATS.get("queued", 0)) + 1 return True except queue.Full: + with _STATS_LOCK: + _STATS["dropped"] = int(_STATS.get("dropped", 0)) + 1 _log.warning("oclaw alarm forward queue full; dropping alarm_key=%s", payload.get("alarm_key")) return False +def _requeue_or_drop(payload: dict[str, Any], *, reason: str) -> None: + max_retries = max(0, int(getattr(settings, "oclaw_forward_max_retries", 3) or 3)) + item = dict(payload) + attempts = int(item.get("_fwd_attempts") or 0) + 1 + item["_fwd_attempts"] = attempts + if attempts > max_retries: + with _STATS_LOCK: + _STATS["retry_exhausted"] = int(_STATS.get("retry_exhausted", 0)) + 1 + _STATS["dropped"] = int(_STATS.get("dropped", 0)) + 1 + _log.warning( + "oclaw forward drop after retries alarm_key=%s attempts=%s reason=%s", + item.get("alarm_key"), + attempts, + reason[:80], + ) + return + try: + _outbound_q().put_nowait(item) + with _STATS_LOCK: + _STATS["requeued"] = int(_STATS.get("requeued", 0)) + 1 + except queue.Full: + with _STATS_LOCK: + _STATS["dropped"] = int(_STATS.get("dropped", 0)) + 1 + _log.warning( + "oclaw forward requeue full; dropping alarm_key=%s", + item.get("alarm_key"), + ) + + def _send_auth(ws: Any) -> bool: token = _bridge_token() ws.send(json.dumps({"type": "auth", "token": token}, ensure_ascii=False)) @@ -187,7 +231,7 @@ def _run_loop() -> None: _notify_status("fwd:disabled") raise RuntimeError("forwarder disabled") try: - payload = _OUTBOUND_Q.get(timeout=1.0) + payload = _outbound_q().get(timeout=1.0) except queue.Empty: try: ws.send(json.dumps({"type": "ping", "ts": _utc_now_iso()}, ensure_ascii=False)) @@ -221,15 +265,13 @@ def _run_loop() -> None: payload.get("alarm_key"), str(ack.get("error") or "")[:120], ) + # Soft ack failure: do not infinite-requeue; count as fail only. except Exception as exc: _log.warning("oclaw forward failed alarm_key=%s err=%s", payload.get("alarm_key"), str(exc)[:120]) - try: - _OUTBOUND_Q.put_nowait(payload) - except queue.Full: - pass + _requeue_or_drop(payload, reason=str(exc)[:120]) raise finally: - _OUTBOUND_Q.task_done() + _outbound_q().task_done() except Exception as exc: _CONNECTED.clear() err = str(exc)[:200] @@ -272,9 +314,13 @@ def forwarder_status() -> dict[str, Any]: "operational": bool(enabled and not paused), "paused": paused, "connected": _CONNECTED.is_set(), - "queue_size": int(_OUTBOUND_Q.qsize()), + "queue_size": int(_outbound_q().qsize()), + "queue_max": max(100, int(getattr(settings, "oclaw_forward_queue_max", 2000) or 2000)), "url": _bridge_url(), "published_ok": int(stats.get("published_ok", 0)), "published_fail": int(stats.get("published_fail", 0)), "queued_total": int(stats.get("queued", 0)), + "dropped": int(stats.get("dropped", 0)), + "requeued": int(stats.get("requeued", 0)), + "retry_exhausted": int(stats.get("retry_exhausted", 0)), } diff --git a/netx_api/port_traffic_scheduler.py b/netx_api/port_traffic_scheduler.py index deaa031..ccc8fe1 100644 --- a/netx_api/port_traffic_scheduler.py +++ b/netx_api/port_traffic_scheduler.py @@ -33,8 +33,7 @@ def _dispatch_pool_get() -> ThreadPoolExecutor: with _pool_lock: if _dispatch_pool is None: workers = clamp_cli_workers( - int(getattr(settings, "port_traffic_dispatch_workers", 4) or 4), - hard_cap=16, + int(getattr(settings, "port_traffic_dispatch_workers", 3) or 3), ) _dispatch_pool = ThreadPoolExecutor( max_workers=workers, thread_name_prefix="pt-dispatch" diff --git a/netx_api/runtime_budget.py b/netx_api/runtime_budget.py new file mode 100644 index 0000000..3863a58 --- /dev/null +++ b/netx_api/runtime_budget.py @@ -0,0 +1,57 @@ +"""Startup capacity checks and operator-facing budget logging.""" +from __future__ import annotations + +import logging + +from .cli_budget import cli_budget_status, feature_hard_cap +from .config import settings +from .db import db_pool_status + +_log = logging.getLogger("netx.runtime.budget") + + +def log_runtime_budget(*, role: str = "api") -> None: + """Log effective pools/CLI caps and warn when capacity looks undersized.""" + pool_size = max(1, int(getattr(settings, "db_pool_size", 25) or 25)) + overflow = max(0, int(getattr(settings, "db_max_overflow", 15) or 15)) + pool_cap = pool_size + overflow + cli = max(1, int(getattr(settings, "cli_max_concurrent", 12) or 12)) + hard = feature_hard_cap() + http_reserve = 8 # rough floor for request handlers + UME WS bursts + recommended_pool = cli + http_reserve + + _log.info( + "runtime budget role=%s db_pool=%s+%s=%s cli_max=%s feature_hard_cap=%s " + "timeout_pool=%s pt_dispatch=%s webcrt_sessions=%s audit_sample_n=%s inline_schedulers=%s", + role, + pool_size, + overflow, + pool_cap, + cli, + hard, + int(getattr(settings, "cli_timeout_pool_workers", 12) or 12), + int(getattr(settings, "port_traffic_dispatch_workers", 3) or 3), + int(getattr(settings, "webcrt_max_sessions", 12) or 12), + int(getattr(settings, "audit_sample_n", 10) or 10), + bool(getattr(settings, "run_inline_schedulers", True)), + ) + _log.info("runtime budget snapshot pool=%s cli=%s", db_pool_status(), cli_budget_status()) + + if pool_cap < recommended_pool: + _log.warning( + "DB pool capacity %s < recommended %s (cli_max=%s + http_reserve=%s). " + "Raise NETX_DB_POOL_SIZE / NETX_DB_MAX_OVERFLOW or lower NETX_CLI_MAX_CONCURRENT.", + pool_cap, + recommended_pool, + cli, + 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 production prefer " + "NETX_RUN_INLINE_SCHEDULERS=false and `python -m netx_api.worker`.", + host, + ) diff --git a/netx_api/topology_discover_jobs.py b/netx_api/topology_discover_jobs.py index 873929e..5c0f037 100644 --- a/netx_api/topology_discover_jobs.py +++ b/netx_api/topology_discover_jobs.py @@ -59,7 +59,7 @@ def _run_discover_job(job_id: str, body: FabricDiscoverRequest) -> None: from .cli_budget import clamp_cli_workers - concurrency = clamp_cli_workers(int(body.concurrency or 4), hard_cap=32) + concurrency = clamp_cli_workers(int(body.concurrency or 4)) added = 0 updated = 0 stale = 0 diff --git a/netx_api/topology_schemas.py b/netx_api/topology_schemas.py index a564b3a..d4f7598 100644 --- a/netx_api/topology_schemas.py +++ b/netx_api/topology_schemas.py @@ -85,7 +85,7 @@ class FabricDiscoverRequest(BaseModel): default=True, description="Create SSH placeholder ManagedNEs for LLDP neighbors not in inventory", ) - concurrency: int = Field(default=4, ge=1, le=32) + concurrency: int = Field(default=4, ge=1, le=16) trigger_mode: str = Field(default="manual", description="manual | schedule | topology") diff --git a/netx_api/worker.py b/netx_api/worker.py index f2c97dc..fc6c66a 100644 --- a/netx_api/worker.py +++ b/netx_api/worker.py @@ -39,6 +39,12 @@ def main() -> None: from .ume_runtime import start_device_schedulers start_device_schedulers() + try: + from .runtime_budget import log_runtime_budget + + log_runtime_budget(role="worker") + except Exception: + _log.exception("worker runtime budget log failed") _log.info("netx worker schedulers started (config_sync, lldp_collect, port_traffic)") while not stop.is_set(): diff --git a/tests/test_stability_hardening.py b/tests/test_stability_hardening.py index bba42c3..a0731fb 100644 --- a/tests/test_stability_hardening.py +++ b/tests/test_stability_hardening.py @@ -1,25 +1,47 @@ -"""Smoke tests for DB pool / CLI budget / metrics wiring.""" +"""Smoke tests for DB pool / CLI budget / metrics / production defaults.""" from __future__ import annotations import unittest -from netx_api.cli_budget import acquire_cli_slot, clamp_cli_workers, cli_budget_status -from netx_api.config import settings +from netx_api.cli_budget import acquire_cli_slot, clamp_cli_workers, cli_budget_status, feature_hard_cap +from netx_api.config import Settings, settings from netx_api.db import db_pool_status from netx_api.metrics_router import _prom_lines, collect_runtime_metrics +from netx_api.runtime_budget import log_runtime_budget class StabilityHardeningTests(unittest.TestCase): + def test_production_defaults(self) -> None: + s = Settings(_env_file=None) + self.assertEqual(s.db_pool_size, 25) + self.assertEqual(s.db_max_overflow, 15) + self.assertEqual(s.cli_max_concurrent, 12) + self.assertEqual(s.cli_feature_hard_cap, 16) + self.assertEqual(s.cli_timeout_pool_workers, 12) + self.assertEqual(s.port_traffic_dispatch_workers, 3) + self.assertEqual(s.ne_connect_max_workers, 4) + self.assertEqual(s.ne_collect_max_workers, 4) + self.assertEqual(s.webcrt_max_sessions, 12) + self.assertEqual(s.webcrt_keepalive_sec, 30) + self.assertEqual(s.audit_sample_n, 10) + self.assertEqual(s.audit_queue_max, 2000) + self.assertEqual(s.oclaw_forward_max_retries, 3) + self.assertEqual(s.oclaw_forward_queue_max, 2000) + self.assertEqual(s.ne_collect_max_output_bytes, 4 * 1024 * 1024) + self.assertEqual(s.webcrt_session_log_max_bytes, 2 * 1024 * 1024) + # Pool should cover CLI + HTTP reserve under defaults. + self.assertGreaterEqual(s.db_pool_size + s.db_max_overflow, s.cli_max_concurrent + 8) + def test_settings_have_pool_and_budget_knobs(self) -> None: self.assertGreaterEqual(int(settings.db_pool_size), 1) self.assertGreaterEqual(int(settings.cli_max_concurrent), 1) self.assertGreaterEqual(int(settings.audit_queue_max), 100) def test_clamp_cli_workers(self) -> None: - n = clamp_cli_workers(999, hard_cap=32) + n = clamp_cli_workers(999) self.assertLessEqual(n, int(settings.cli_max_concurrent)) - self.assertLessEqual(n, 32) + self.assertLessEqual(n, feature_hard_cap()) self.assertGreaterEqual(n, 1) def test_cli_budget_acquire_release(self) -> None: @@ -39,6 +61,10 @@ class StabilityHardeningTests(unittest.TestCase): self.assertIn("netx_uptime_seconds", body) self.assertIn("netx_thread_count", body) self.assertIn("netx_cli_budget_limit", body) + self.assertIn("netx_oclaw_forwarder_dropped", body) + + def test_log_runtime_budget_does_not_raise(self) -> None: + log_runtime_budget(role="test") if __name__ == "__main__": diff --git a/tests/test_webcrt.py b/tests/test_webcrt.py index cb70437..fff29a8 100644 --- a/tests/test_webcrt.py +++ b/tests/test_webcrt.py @@ -304,8 +304,8 @@ class WebcrtServiceTests(unittest.TestCase): self.assertTrue(called_creds["hop_enabled"]) self.assertEqual(called_creds["hop_vendor"], "bastion") self.assertIn("session_log", mock_open.call_args.kwargs) - self.assertEqual(mock_open.call_args.kwargs.get("keepalive"), 0) - self.assertEqual(out.get("keepalive_sec"), 0) + self.assertEqual(mock_open.call_args.kwargs.get("keepalive"), 30) + self.assertEqual(out.get("keepalive_sec"), 30) fake.remote_conn.resize_pty.assert_called() self.assertEqual(out["ne_id"], "ne-hop") self.assertFalse(out.get("cli_hop")) # bastion hop is not vendor CLI hop guard @@ -867,8 +867,8 @@ class WebcrtServiceTests(unittest.TestCase): self.assertIn("filename=", dispo) self.assertIn("filename*=UTF-8''", dispo) self.assertNotIn("\n", dispo) - self.assertEqual(int(cfg.settings.webcrt_sftp_max_file_bytes), 512 * 1024 * 1024) - self.assertEqual(int(cfg.settings.webcrt_sftp_list_max_entries), 5000) + self.assertEqual(int(cfg.settings.webcrt_sftp_max_file_bytes), 256 * 1024 * 1024) + self.assertEqual(int(cfg.settings.webcrt_sftp_list_max_entries), 2000) def test_sftp_rename_rejects_bad_paths(self) -> None: from unittest.mock import MagicMock