diff --git a/.env.example b/.env.example index 4a19cb0..1d7cabb 100644 --- a/.env.example +++ b/.env.example @@ -57,27 +57,28 @@ 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 -# --- Production-recommended capacity (defaults in Settings already match these) --- -# NETX_DB_POOL_SIZE=25 -# NETX_DB_MAX_OVERFLOW=15 +# --- Multi-user shared-server capacity (defaults in Settings already match these) --- +# NETX_DB_POOL_SIZE=40 +# NETX_DB_MAX_OVERFLOW=40 # NETX_DB_POOL_RECYCLE_SEC=1800 # NETX_DB_POOL_TIMEOUT_SEC=30 -# 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_CLI_MAX_CONCURRENT=24 +# NETX_CLI_FEATURE_HARD_CAP=32 +# NETX_CLI_TIMEOUT_POOL_WORKERS=24 +# NETX_PORT_TRAFFIC_DISPATCH_WORKERS=6 +# NETX_NE_CONNECT_MAX_WORKERS=8 +# NETX_NE_COLLECT_MAX_WORKERS=8 +# NETX_NE_COLLECT_RUN_TIMEOUT_CAP_SEC=600 +# NETX_NE_COLLECT_MAX_OUTPUT_BYTES=8388608 +# NETX_WEBCRT_MAX_SESSIONS=40 # 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_WEBCRT_OUT_QUEUE_MAX=2000 +# NETX_WEBCRT_SESSION_LOG_MAX_BYTES=4194304 +# NETX_WEBCRT_SFTP_MAX_FILE_BYTES=536870912 +# NETX_WEBCRT_SFTP_LIST_MAX_ENTRIES=5000 # NETX_AUDIT_ASYNC=true -# NETX_AUDIT_SAMPLE_N=10 -# NETX_AUDIT_QUEUE_MAX=2000 -# NETX_OCLAW_FORWARD_QUEUE_MAX=2000 +# NETX_AUDIT_SAMPLE_N=5 +# NETX_AUDIT_QUEUE_MAX=5000 +# NETX_OCLAW_FORWARD_QUEUE_MAX=5000 # NETX_OCLAW_FORWARD_MAX_RETRIES=3 -# Lab-only raises (optional): higher CLI/DB when bastion and Postgres are sized for it. +# Heavier fleets: raise CLI/DB together; also ensure Postgres max_connections and bastion session limits. diff --git a/netx_api/cli_budget.py b/netx_api/cli_budget.py index cabcf8e..ee1c67d 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", 12) or 12)) + want = max(1, int(getattr(settings, "cli_max_concurrent", 24) or 24)) if _sem is None or want != _limit: _sem = threading.BoundedSemaphore(want) _limit = want @@ -34,18 +34,18 @@ def cli_budget_status() -> dict[str, 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", 12) or 12)) + budget = max(1, int(getattr(settings, "cli_max_concurrent", 24) or 24)) feature_cap = int( hard_cap if hard_cap is not None - else (getattr(settings, "cli_feature_hard_cap", 16) or 16) + else (getattr(settings, "cli_feature_hard_cap", 32) or 32) ) 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)) + return max(1, int(getattr(settings, "cli_feature_hard_cap", 32) or 32)) @contextmanager diff --git a/netx_api/config.py b/netx_api/config.py index b189b04..e0de37c 100644 --- a/netx_api/config.py +++ b/netx_api/config.py @@ -64,14 +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 = "" - # Production-oriented worker caps (raise only when bastion/DB capacity allows). - ne_connect_max_workers: int = 4 + # Shared-server worker caps (sized for multi-operator use; raise if bastion/DB allow). + ne_connect_max_workers: int = 8 ne_connect_timeout_sec: int = 30 - ne_collect_max_workers: int = 4 + ne_collect_max_workers: int = 8 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 = 480 + ne_collect_run_timeout_cap_sec: int = 600 ne_collection_data_dir: str = "data/ne_collections" # Config sync (periodic running-config backup into DB) config_sync_scheduler_enabled: bool = True @@ -91,8 +91,8 @@ class Settings(BaseSettings): port_traffic_scheduler_tick_sec: int = 15 # 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 = 12 + # WebCRT interactive terminal sessions (multi-operator concurrent terminals). + webcrt_max_sessions: int = 40 webcrt_idle_timeout_sec: int = 1800 webcrt_connect_timeout_sec: int = 90 webcrt_attach_timeout_sec: int = 60 @@ -105,16 +105,16 @@ class Settings(BaseSettings): 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 = 1500 + webcrt_out_queue_max: int = 2000 # 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 256 MiB per file). - webcrt_sftp_max_file_bytes: int = 256 * 1024 * 1024 + # WebCRT SFTP transfer limits (streamed; default 512 MiB per file). + webcrt_sftp_max_file_bytes: int = 512 * 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 = 2000 + webcrt_sftp_list_max_entries: int = 5000 webcrt_sftp_list_timeout_sec: float = 30.0 # Local app login / audit (lab defaults; override in production) auth_enabled: bool = True @@ -133,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 = 10 + audit_sample_n: int = 5 # 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 @@ -144,29 +144,29 @@ 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 - # 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 = 25 - db_max_overflow: int = 15 + # 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 + db_max_overflow: int = 40 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 = 12 + cli_max_concurrent: int = 24 # Per-feature concurrency hard ceiling (API body / policy cannot exceed this). - cli_feature_hard_cap: int = 16 + cli_feature_hard_cap: int = 32 # Shared timeout watchdog pool (not per-task executors). - cli_timeout_pool_workers: int = 12 + cli_timeout_pool_workers: int = 24 # Port-traffic: how many devices may collect in parallel on the scheduler tick. - port_traffic_dispatch_workers: int = 3 + port_traffic_dispatch_workers: int = 6 # Bound async audit queue; drop oldest when full to protect RSS. - audit_queue_max: int = 2000 + audit_queue_max: int = 5000 # Cap NE collection output files (bytes); 0 = unlimited (not recommended). - ne_collect_max_output_bytes: int = 4 * 1024 * 1024 + ne_collect_max_output_bytes: int = 8 * 1024 * 1024 # WebCRT session transcript rotate size (bytes); 0 disables rotate. - webcrt_session_log_max_bytes: int = 2 * 1024 * 1024 + webcrt_session_log_max_bytes: int = 4 * 1024 * 1024 # oclaw alarm forwarder: requeue attempts before drop on send failure. oclaw_forward_max_retries: int = 3 - oclaw_forward_queue_max: int = 2000 + oclaw_forward_queue_max: int = 5000 settings = Settings() diff --git a/netx_api/config_sync_common.py b/netx_api/config_sync_common.py index ef5c2ab..555b0f7 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(16, int(row.concurrency or 5))), + concurrency=max(1, min(32, 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 0602a54..a02de81 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(16, int(data["concurrency"]))) + row.concurrency = max(1, min(32, 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(16, int(policy.concurrency or 5))) + concurrency = max(1, min(32, 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 aa354d6..1349186 100644 --- a/netx_api/config_sync_runner.py +++ b/netx_api/config_sync_runner.py @@ -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(16, int(cycle.concurrency or 5))) + concurrency = max(1, min(32, 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 3a6c4d1..10f947b 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(16, int(policy.concurrency or 5))) + concurrency = max(1, min(32, 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 419696c..e9216f9 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=16) + concurrency: int | None = Field(default=None, ge=1, le=32) 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 a7f72ef..39e9e36 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=16) + concurrency: int | None = Field(default=None, ge=1, le=32) 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 f35a377..99330bd 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(16, int(data["concurrency"]))) + row.concurrency = max(1, min(32, 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(16, int(policy.concurrency or 4))) + concurrency = max(1, min(32, 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/ne_collect_runner.py b/netx_api/ne_collect_runner.py index bdbf407..f4f8624 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 4)) + workers = clamp_cli_workers(int(settings.ne_collect_max_workers or 8)) _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 94b7548..4b4f682 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 4)) + workers = clamp_cli_workers(int(settings.ne_connect_max_workers or 8)) _executor = ThreadPoolExecutor(max_workers=workers, thread_name_prefix="ne-connect") return _executor diff --git a/netx_api/runtime_budget.py b/netx_api/runtime_budget.py index 3863a58..672a785 100644 --- a/netx_api/runtime_budget.py +++ b/netx_api/runtime_budget.py @@ -12,12 +12,13 @@ _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_size = max(1, int(getattr(settings, "db_pool_size", 40) or 40)) + overflow = max(0, int(getattr(settings, "db_max_overflow", 40) or 40)) pool_cap = pool_size + overflow - cli = max(1, int(getattr(settings, "cli_max_concurrent", 12) or 12)) + cli = max(1, int(getattr(settings, "cli_max_concurrent", 24) or 24)) hard = feature_hard_cap() - http_reserve = 8 # rough floor for request handlers + UME WS bursts + # Multi-user HTTP + WebCRT WS + UME alarm WS headroom. + http_reserve = 24 recommended_pool = cli + http_reserve _log.info( @@ -29,10 +30,10 @@ def log_runtime_budget(*, role: str = "api") -> None: 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), + int(getattr(settings, "cli_timeout_pool_workers", 24) or 24), + int(getattr(settings, "port_traffic_dispatch_workers", 6) or 6), + int(getattr(settings, "webcrt_max_sessions", 40) or 40), + int(getattr(settings, "audit_sample_n", 5) or 5), bool(getattr(settings, "run_inline_schedulers", True)), ) _log.info("runtime budget snapshot pool=%s cli=%s", db_pool_status(), cli_budget_status()) @@ -51,7 +52,7 @@ def log_runtime_budget(*, role: str = "api") -> None: getattr(settings, "run_inline_schedulers", True) ): _log.warning( - "Non-loopback bind (%s) with inline schedulers — for production prefer " + "Non-loopback bind (%s) with inline schedulers — for multi-user production prefer " "NETX_RUN_INLINE_SCHEDULERS=false and `python -m netx_api.worker`.", host, ) diff --git a/netx_api/topology_schemas.py b/netx_api/topology_schemas.py index d4f7598..a564b3a 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=16) + concurrency: int = Field(default=4, ge=1, le=32) trigger_mode: str = Field(default="manual", description="manual | schedule | topology") diff --git a/tests/test_stability_hardening.py b/tests/test_stability_hardening.py index a0731fb..f44b46a 100644 --- a/tests/test_stability_hardening.py +++ b/tests/test_stability_hardening.py @@ -14,24 +14,24 @@ 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.db_pool_size, 40) + self.assertEqual(s.db_max_overflow, 40) + self.assertEqual(s.cli_max_concurrent, 24) + self.assertEqual(s.cli_feature_hard_cap, 32) + self.assertEqual(s.cli_timeout_pool_workers, 24) + self.assertEqual(s.port_traffic_dispatch_workers, 6) + self.assertEqual(s.ne_connect_max_workers, 8) + self.assertEqual(s.ne_collect_max_workers, 8) + self.assertEqual(s.webcrt_max_sessions, 40) self.assertEqual(s.webcrt_keepalive_sec, 30) - self.assertEqual(s.audit_sample_n, 10) - self.assertEqual(s.audit_queue_max, 2000) + self.assertEqual(s.audit_sample_n, 5) + self.assertEqual(s.audit_queue_max, 5000) 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) + 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) + # 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) def test_settings_have_pool_and_budget_knobs(self) -> None: self.assertGreaterEqual(int(settings.db_pool_size), 1) diff --git a/tests/test_webcrt.py b/tests/test_webcrt.py index fff29a8..2c364f8 100644 --- a/tests/test_webcrt.py +++ b/tests/test_webcrt.py @@ -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), 256 * 1024 * 1024) - self.assertEqual(int(cfg.settings.webcrt_sftp_list_max_entries), 2000) + self.assertEqual(int(cfg.settings.webcrt_sftp_max_file_bytes), 512 * 1024 * 1024) + self.assertEqual(int(cfg.settings.webcrt_sftp_list_max_entries), 5000) def test_sftp_rename_rejects_bad_paths(self) -> None: from unittest.mock import MagicMock