mirror of
https://github.com/hansjone/netx.git
synced 2026-10-09 02:00:46 +08:00
Adopt production-safe runtime defaults and forwarder retry limits.
Lower CLI/WebCRT/audit caps, log startup capacity warnings, and bound oclaw requeues so reconnect storms cannot thrash memory. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
parent
d56020d84c
commit
fa577aa3bd
22 changed files with 238 additions and 69 deletions
34
.env.example
34
.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.
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
|
|
@ -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))),
|
||||
|
|
|
|||
|
|
@ -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":
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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] = []
|
||||
|
|
|
|||
|
|
@ -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)}")
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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)),
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
|
|
|
|||
57
netx_api/runtime_budget.py
Normal file
57
netx_api/runtime_budget.py
Normal file
|
|
@ -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,
|
||||
)
|
||||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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():
|
||||
|
|
|
|||
|
|
@ -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__":
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue