Split WebCRT and API startup; default collectors to worker process.

Extract channel/session modules and CLI guard, move UME sidebands out of main, and default NETX_RUN_INLINE_SCHEDULERS off so ops run python -m netx_api.worker.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-08-02 16:50:52 +08:00
parent 57b3faf9fb
commit 136f40cdae
16 changed files with 2269 additions and 2150 deletions

144
netx_api/app_startup.py Normal file
View file

@ -0,0 +1,144 @@
"""API process startup orchestration (schema, recovery, schedulers, UME sidebands)."""
from __future__ import annotations
import logging
from .auth_service import bootstrap_admin_if_needed
from .config import settings
from .db import Base, SessionLocal, engine
from .schema_patches import (
apply_all_legacy_startup_ddl,
apply_auth_schema_patches,
run_alembic_upgrade_to_head,
)
from .security_bootstrap import assert_secure_defaults_or_exit
from .ume_runtime import start_api_sideband_threads, start_device_schedulers
import netx_api.ume_support as ume_support
from .ume_alarm_ws import (
begin_startup_alarm_sync_gate,
complete_startup_alarm_sync_gate,
)
_log = logging.getLogger("netx.ume.schedule")
def _configure_ume_diag_logging() -> None:
fmt = logging.Formatter("%(asctime)s %(levelname)s %(name)s: %(message)s")
for name in ("netx.ume.schedule", "netx.ume.sync"):
lg = logging.getLogger(name)
if lg.handlers:
continue
h = logging.StreamHandler()
h.setFormatter(fmt)
lg.addHandler(h)
lg.setLevel(logging.INFO)
lg.propagate = False
def run_api_startup() -> None:
"""Full API boot sequence previously inlined in ``main.on_startup``."""
assert_secure_defaults_or_exit()
_configure_ume_diag_logging()
Base.metadata.create_all(bind=engine)
alembic_ok = True
if bool(getattr(settings, "alembic_upgrade_on_start", True)):
try:
run_alembic_upgrade_to_head()
except Exception:
alembic_ok = False
_log.exception("startup: alembic upgrade head failed")
skip_ddl = bool(getattr(settings, "skip_legacy_startup_ddl", True))
try:
with engine.begin() as conn:
apply_auth_schema_patches(conn)
except Exception:
_log.exception("startup: auth schema patches failed")
if skip_ddl and alembic_ok:
_log.info("startup: schema via Alembic (legacy inline DDL skipped)")
else:
if skip_ddl and not alembic_ok:
_log.warning("startup: Alembic failed — falling back to legacy schema patches")
try:
apply_all_legacy_startup_ddl(engine)
except Exception:
_log.exception("startup: legacy schema patches failed")
ume_support._reset_runtime_pause_flags()
ume_support._fail_stale_running_sync_jobs_on_startup()
try:
from .topology_service import bootstrap_topology_tree, reclaim_stale_discover_jobs
db_topo = SessionLocal()
try:
bootstrap_topology_tree(db_topo)
closed = reclaim_stale_discover_jobs(db_topo, force_all_open=True)
if closed:
_log.warning("startup: closed %s orphaned topology discover jobs", closed)
finally:
db_topo.close()
except Exception:
_log.exception("startup: topology discover job cleanup failed")
if ume_support._needs_startup_alarm_sync_before_ws():
begin_startup_alarm_sync_gate()
_log.info(
"startup: WSS blocked until initial REST current-alarm sync completes (delay=%ss)",
ume_support._startup_alarm_pull_delay_s(),
)
else:
complete_startup_alarm_sync_gate()
db = SessionLocal()
try:
try:
bootstrap_admin_if_needed(db)
except Exception:
_log.exception("startup: auth bootstrap admin failed")
from .collection_recovery import recover_collection_jobs_on_startup
resumed = recover_collection_jobs_on_startup(db)
if resumed:
_log.info("startup: resumed %s pending ne collection runs", resumed)
from .config_sync_recovery import recover_config_sync_on_startup
from .config_sync_service import ensure_policy
from .port_traffic_recovery import recover_port_traffic_on_startup
ensure_policy(db)
cfg_resumed = recover_config_sync_on_startup(db)
if cfg_resumed:
_log.info("startup: resumed %s config_sync task(s) from interrupted cycle", cfg_resumed)
try:
from .lldp_collect_service import ensure_policy as ensure_lldp_collect_policy
ensure_lldp_collect_policy(db)
except Exception:
_log.exception("startup: lldp_collect policy ensure failed")
pt_cleared = recover_port_traffic_on_startup(db)
if pt_cleared:
_log.info("startup: cleared %s port_traffic stuck collect_running flag(s)", pt_cleared)
try:
from .port_traffic_migrate import backfill_port_traffic_series
backfill_port_traffic_series(db)
except Exception:
_log.exception("startup: port_traffic series backfill failed")
except Exception:
_log.exception("startup: ne collection / config_sync recovery failed")
finally:
db.close()
if bool(getattr(settings, "run_inline_schedulers", False)):
try:
start_device_schedulers()
except Exception:
_log.exception("startup: device schedulers init failed")
else:
_log.info(
"startup: inline schedulers disabled — run `python -m netx_api.worker` for "
"config_sync / lldp_collect / port_traffic"
)
start_api_sideband_threads()

View file

@ -140,8 +140,9 @@ class Settings(BaseSettings):
alembic_upgrade_on_start: bool = True
# Optional dedicated SQLAlchemy URL for /v1/sql/* (read-only DB role recommended).
sql_readonly_database_url: str = ""
# When false, API skips config_sync / lldp / port_traffic schedulers (run `python -m netx_api.worker`).
run_inline_schedulers: bool = True
# When false (default), API skips config_sync / lldp / port_traffic schedulers —
# run `python -m netx_api.worker` alongside the API. Set true only for single-process lab.
run_inline_schedulers: bool = False
settings = Settings()

View file

@ -2,13 +2,16 @@
from __future__ import annotations
import time
from typing import Any
from fastapi import APIRouter, Depends
from sqlalchemy import text as sql_text
from sqlalchemy.orm import Session
from .config import settings
from .db import get_db
from .oclaw_alarm_forwarder import forwarder_status
router = APIRouter(tags=["health"])
@ -21,9 +24,87 @@ def health_live() -> dict[str, str]:
@router.get("/health/ready", status_code=200)
def health_ready(db: Session = Depends(get_db)) -> dict[str, Any]:
"""Readiness — verifies database connectivity."""
"""Readiness — DB plus scheduler deployment hint."""
out: dict[str, Any] = {"status": "ok", "probe": "ready"}
try:
db.execute(sql_text("select 1"))
return {"status": "ok", "probe": "ready", "db": "up"}
out["db"] = "up"
except Exception as exc:
return {"status": "down", "probe": "ready", "db": "down", "error": str(exc)[:240]}
return {
"status": "down",
"probe": "ready",
"db": "down",
"error": str(exc)[:240],
}
inline = bool(getattr(settings, "run_inline_schedulers", False))
out["schedulers"] = {
"inline": inline,
"mode": "inline" if inline else "external_worker",
"hint": None
if inline
else "run `python -m netx_api.worker` for config_sync / lldp_collect / port_traffic",
}
return out
@router.get("/v1/integrations/status")
def integrations_status(db: Session = Depends(get_db)) -> dict:
"""netx API + DB + oclaw bridge status."""
netx_api = {"status": "up"}
db_status: dict = {"status": "unknown"}
try:
t0 = time.monotonic()
db.execute(sql_text("select 1"))
db_status = {"status": "up", "latency_ms": int((time.monotonic() - t0) * 1000)}
except Exception as exc:
db_status = {"status": "down", "error": str(exc)[:240]}
oclaw_status: dict = {"status": "unknown"}
fwd = forwarder_status()
if not bool(fwd.get("enabled")):
oclaw_status = {
"status": "unknown",
"mode": "ws",
"enabled": False,
"connected": False,
"error_kind": "disabled",
"error": "NETX_OCLAW_ALARM_WS_ENABLED=false or missing token/url",
"forwarder": fwd,
}
elif bool(fwd.get("paused")):
oclaw_status = {
"status": "unknown",
"mode": "ws",
"enabled": True,
"connected": False,
"error_kind": "paused",
"error": "oclaw_alarm_forwarder runtime task paused",
"forwarder": fwd,
}
elif bool(fwd.get("connected")):
oclaw_status = {
"status": "up",
"mode": "ws",
"enabled": True,
"connected": True,
"queue_size": int(fwd.get("queue_size") or 0),
"published_ok": int(fwd.get("published_ok") or 0),
"published_fail": int(fwd.get("published_fail") or 0),
"url": str(fwd.get("url") or ""),
"forwarder": fwd,
}
else:
oclaw_status = {
"status": "down",
"mode": "ws",
"enabled": True,
"connected": False,
"error_kind": "ws_disconnected",
"error": "oclaw netx-bridge WebSocket not connected",
"queue_size": int(fwd.get("queue_size") or 0),
"url": str(fwd.get("url") or ""),
"forwarder": fwd,
}
return {"netx_api": netx_api, "db": db_status, "oclaw_bridge": oclaw_status}

View file

@ -1,44 +1,30 @@
"""FastAPI application entry — routers + thin lifecycle hooks."""
from __future__ import annotations
import csv
import json
import logging
from datetime import datetime, timezone
from io import StringIO
import time
import re
import threading
_schedule_log = logging.getLogger("netx.ume.schedule")
_BOOT_MONO = time.monotonic()
from fastapi import Depends, FastAPI, File, HTTPException, Query, UploadFile
from fastapi.responses import Response
from sqlalchemy import text as sql_text
from sqlalchemy.orm import Session
from typing import Any
import uvicorn
from fastapi import FastAPI
from .ap_client import analyze_with_oclaw, health_with_oclaw
from .auth_middleware import AuthAuditMiddleware
from .auth_router import router as auth_router
from .auth_service import bootstrap_admin_if_needed
from .config import settings
from .db import Base, SessionLocal, engine, get_db
from .collection_router import router as collection_router
from .alarms_router import router as alarms_router
from .cli_router import router as cli_router
from .collection_router import router as collection_router
from .config import settings
from .config_sync_router import router as config_sync_router
from .port_traffic_router import router as port_traffic_router
from .managed_ne_router import router as managed_ne_router
from .webcrt_router import router as webcrt_router
from .topology_router import router as topology_router
from .integrations_router import router as integrations_router
from .lldp_collect_router import router as lldp_collect_router
from .managed_ne_router import router as managed_ne_router
from .ops_router import router as ops_router
from .parser_config import load_parser_config
from .port_traffic_router import router as port_traffic_router
from .sql_router import router as sql_router
from .sql_router import sql_query, sql_ume_query # noqa: F401 — tests import from main
from .security_bootstrap import assert_secure_defaults_or_exit
from .integrations_router import router as integrations_router
from .topology_router import router as topology_router
from .ume_router import router as ume_router
from .alarms_router import router as alarms_router
from .ume_router import ( # noqa: F401 — tests import from main
_extract_ume_raw_group_field,
_serialize_ume_alarm_raw_row,
@ -49,89 +35,12 @@ from .ume_support import ( # noqa: F401 — tests import from main
_protocol_bucket_label,
)
import netx_api.ume_support as ume_support
from .ume_runtime import start_device_schedulers
from .importer import aggregate_alarms, import_alarm_excel, query_alarms
from .models import (
AiAnalyzeHistory,
AlarmBatch,
AlarmNorm,
ApiToken,
AppUser,
AuditLog,
ImportErrorRow,
ManagedNE,
NeCollectionJob,
NeCollectionRun,
UmeAlarmCurrent,
UmeAlarmHistory,
UmeInventoryNE,
UmeKeyAlertRule,
UmeKeyAlertForwardLog,
UmeSyncJob,
)
from .models import ImportJob
from .parser_config import load_parser_config
from .ume_client import UMEClient
from .ume_alarm_ws import (
begin_startup_alarm_sync_gate,
cancel_alarm_subscription_manual,
clear_local_alarm_subscription_manual,
complete_startup_alarm_sync_gate,
establish_alarm_subscription_manual,
get_alarms_coordination_status,
get_subscription_status,
get_ws_connection_status,
get_ws_logs,
is_startup_alarm_sync_pending,
is_wss_active_for_current_alarms,
load_persisted_subscription,
request_ws_reconnect,
shutdown_ws_consumer,
start_ume_alarm_ws_consumer,
)
from .ume_sync_service import sync_alarms_current, sync_alarms_history_full, sync_inventory_full
from .runtime_task_messages import (
RT_ALARMS_SYNC_IN_PROGRESS_SKIP,
RT_OCLAW_FWD_DISABLED,
RT_PULLING_ALARMS_CURRENT,
RT_PULLING_INVENTORY,
RT_RESUMED,
RT_RESUMED_OCLAW_WSS_RECONNECT,
RT_RESUMED_SYNC_SOON,
RT_RESUMED_WSS_RECONNECT,
RT_STARTUP_ALARM_SYNC_BEFORE_WS,
RT_STARTUP_GATE_WAITING,
RT_KEEPALIVE_FAILED,
RT_UME_WS_DISABLED_NO_BASE_URL,
RT_WSS_ACTIVE_SKIP_REST,
)
from .oclaw_alarm_forwarder import (
forwarder_status,
is_forwarder_enabled,
request_forwarder_reconnect,
configure_oclaw_alarm_forwarder,
shutdown_oclaw_alarm_forwarder,
start_oclaw_alarm_forwarder,
)
from .ume_token_store import (
clear_shared_token,
load_shared_token,
release_refresh_lock,
save_shared_token,
try_acquire_refresh_lock,
wait_for_token_update,
)
from .schemas import (
AlarmAggregateBucket,
AlarmAggregateResponse,
AiAnalyzeHistoryItem,
AiAnalyzeHistoryResponse,
AlarmItem,
AlarmQueryResponse,
BatchSummary,
ImportJobItem,
ImportJobListResponse,
)
from .oclaw_alarm_forwarder import shutdown_oclaw_alarm_forwarder
from .ume_alarm_ws import shutdown_ws_consumer
from .webcrt_router import router as webcrt_router
_schedule_log = logging.getLogger("netx.ume.schedule")
_BOOT_MONO = time.monotonic()
app = FastAPI(
title="netx ops tool",
@ -156,423 +65,13 @@ app.include_router(integrations_router)
app.include_router(ume_router)
app.include_router(alarms_router)
parser_cfg = load_parser_config()
def _configure_ume_diag_logging() -> None:
"""Emit netx.ume.* INFO to stderr so background scripts/.run/*.log and consoles show scheduler lines."""
fmt = logging.Formatter("%(asctime)s %(levelname)s %(name)s: %(message)s")
for name in ("netx.ume.schedule", "netx.ume.sync"):
lg = logging.getLogger(name)
if lg.handlers:
continue
h = logging.StreamHandler()
h.setFormatter(fmt)
lg.addHandler(h)
lg.setLevel(logging.INFO)
lg.propagate = False
@app.on_event("startup")
def on_startup() -> None:
assert_secure_defaults_or_exit()
_configure_ume_diag_logging()
Base.metadata.create_all(bind=engine)
from .schema_patches import (
apply_all_legacy_startup_ddl,
apply_auth_schema_patches,
run_alembic_upgrade_to_head,
)
from .app_startup import run_api_startup
alembic_ok = True
if bool(getattr(settings, "alembic_upgrade_on_start", True)):
try:
run_alembic_upgrade_to_head()
except Exception:
alembic_ok = False
_schedule_log.exception("startup: alembic upgrade head failed")
skip_ddl = bool(getattr(settings, "skip_legacy_startup_ddl", True))
# Auth columns must exist before bootstrap even when legacy DDL is skipped.
try:
with engine.begin() as conn:
apply_auth_schema_patches(conn)
except Exception:
_schedule_log.exception("startup: auth schema patches failed")
if skip_ddl and alembic_ok:
_schedule_log.info("startup: schema via Alembic (legacy inline DDL skipped)")
else:
if skip_ddl and not alembic_ok:
_schedule_log.warning(
"startup: Alembic failed — falling back to legacy schema patches"
)
try:
apply_all_legacy_startup_ddl(engine)
except Exception:
_schedule_log.exception("startup: legacy schema patches failed")
ume_support._reset_runtime_pause_flags()
ume_support._fail_stale_running_sync_jobs_on_startup()
try:
from .topology_service import bootstrap_topology_tree, reclaim_stale_discover_jobs
db_topo = SessionLocal()
try:
bootstrap_topology_tree(db_topo)
closed = reclaim_stale_discover_jobs(db_topo, force_all_open=True)
if closed:
_schedule_log.warning(
"startup: closed %s orphaned topology discover jobs", closed
)
finally:
db_topo.close()
except Exception:
_schedule_log.exception("startup: topology discover job cleanup failed")
if ume_support._needs_startup_alarm_sync_before_ws():
begin_startup_alarm_sync_gate()
_schedule_log.info(
"startup: WSS blocked until initial REST current-alarm sync completes (delay=%ss)",
ume_support._startup_alarm_pull_delay_s(),
)
else:
complete_startup_alarm_sync_gate()
db = SessionLocal()
try:
try:
bootstrap_admin_if_needed(db)
except Exception:
_schedule_log.exception("startup: auth bootstrap admin failed")
from .collection_recovery import recover_collection_jobs_on_startup
resumed = recover_collection_jobs_on_startup(db)
if resumed:
_schedule_log.info("startup: resumed %s pending ne collection runs", resumed)
from .config_sync_recovery import recover_config_sync_on_startup
from .config_sync_service import ensure_policy
from .port_traffic_recovery import recover_port_traffic_on_startup
ensure_policy(db)
cfg_resumed = recover_config_sync_on_startup(db)
if cfg_resumed:
_schedule_log.info("startup: resumed %s config_sync task(s) from interrupted cycle", cfg_resumed)
try:
from .lldp_collect_service import ensure_policy as ensure_lldp_collect_policy
ensure_lldp_collect_policy(db)
except Exception:
_schedule_log.exception("startup: lldp_collect policy ensure failed")
pt_cleared = recover_port_traffic_on_startup(db)
if pt_cleared:
_schedule_log.info("startup: cleared %s port_traffic stuck collect_running flag(s)", pt_cleared)
try:
from .port_traffic_migrate import backfill_port_traffic_series
backfill_port_traffic_series(db)
except Exception:
_schedule_log.exception("startup: port_traffic series backfill failed")
except Exception:
_schedule_log.exception("startup: ne collection / config_sync recovery failed")
finally:
db.close()
if bool(getattr(settings, "run_inline_schedulers", True)):
try:
start_device_schedulers()
except Exception:
_schedule_log.exception("startup: device schedulers init failed")
else:
_schedule_log.info(
"startup: inline schedulers disabled — run `python -m netx_api.worker` for "
"config_sync / lldp_collect / port_traffic"
)
try:
if bool(getattr(settings, "ume_keepalive_enabled", True)):
interval_keepalive_s = int(getattr(settings, "ume_keepalive_interval_s", 600) or 600)
interval_keepalive_s = max(30, min(interval_keepalive_s, 3600))
renew_before_s = int(getattr(settings, "ume_keepalive_renew_before_s", 900) or 900)
renew_before_s = max(30, min(renew_before_s, 86400))
def _keepalive_loop() -> None:
# Best-effort keepalive: if token exists, periodically handshake to extend TTL.
while True:
try:
if ume_support._runtime_is_paused("token_keepalive"):
time.sleep(1)
continue
client = ume_support._ume_client()
st = client.token_status()
expires_in = int(st.get("expires_in_s") or 0)
# Renew when missing/invalid TTL (0) or nearing expiry — previously 0 skipped renew forever.
if bool(st.get("has_token")) and (expires_in <= 0 or expires_in < renew_before_s):
client.renew_token()
ume_support._set_runtime_task("token_keepalive", status="running", last_run_at=datetime.now(timezone.utc), last_error="")
except Exception:
ume_support._set_runtime_task("token_keepalive", status="error", last_run_at=datetime.now(timezone.utc), last_error=RT_KEEPALIVE_FAILED)
time.sleep(interval_keepalive_s)
t = threading.Thread(target=_keepalive_loop, name="ume-token-keepalive", daemon=True)
t.start()
except Exception as exc:
_schedule_log.exception("startup: token_keepalive thread init failed: %s", exc)
ume_support._set_runtime_task(
"token_keepalive",
status="error",
last_run_at=datetime.now(timezone.utc),
last_error=f"startup_thread_init_failed: {str(exc)[:180]}",
)
try:
def _startup_alarm_sync_worker() -> None:
try:
ume_support._run_startup_alarm_sync_before_ws()
except Exception as exc:
_schedule_log.exception("startup: alarm sync before WSS failed: %s", exc)
complete_startup_alarm_sync_gate()
# Do not block HTTP /health on slow UME REST pull; WSS waits on startup_alarm_sync_gate.
t_startup_sync = threading.Thread(
target=_startup_alarm_sync_worker,
name="ume-startup-alarm-sync",
daemon=True,
)
t_startup_sync.start()
except Exception as exc:
_schedule_log.exception("startup: alarm sync thread init failed: %s", exc)
complete_startup_alarm_sync_gate()
try:
if bool(getattr(settings, "ume_sync_alarms_current_enabled", True)):
alarms_interval_s = int(getattr(settings, "ume_sync_alarms_current_interval_s", 18000) or 18000)
alarms_interval_s = max(30, min(alarms_interval_s, 86400))
def _alarms_current_sync_loop() -> None:
ume_support._refresh_runtime_task_idle("alarms_current_auto_sync", "alarms_current")
ume_support._wait_until_startup_alarm_pull_allowed("alarms_current_auto_sync")
while True:
try:
_schedule_log.info(
"alarms_current_auto_sync: loop tick paused=%s",
ume_support._runtime_is_paused("alarms_current_auto_sync"),
)
if ume_support._runtime_is_paused("alarms_current_auto_sync"):
time.sleep(1)
continue
if is_startup_alarm_sync_pending():
ume_support._refresh_runtime_task_idle(
"alarms_current_auto_sync",
"alarms_current",
last_error=RT_STARTUP_GATE_WAITING,
)
time.sleep(10)
continue
if (
bool(getattr(settings, "ume_sync_alarms_current_skip_when_ws", True))
and is_wss_active_for_current_alarms()
):
ume_support._refresh_runtime_task_idle(
"alarms_current_auto_sync",
"alarms_current",
last_error=RT_WSS_ACTIVE_SKIP_REST,
)
time.sleep(max(30, min(alarms_interval_s, 300)))
continue
ume_support._maybe_wait_for_sync_interval(
task_id="alarms_current_auto_sync",
domain="alarms_current",
interval_s=alarms_interval_s,
label="alarms_current_auto_sync",
)
_schedule_log.info(
"alarms_current_auto_sync: iteration start (interval=%ss)",
alarms_interval_s,
)
ume_support._set_runtime_task(
"alarms_current_auto_sync",
status="running",
last_run_at=datetime.now(timezone.utc),
last_error=RT_PULLING_ALARMS_CURRENT,
)
db = SessionLocal()
try:
client = ume_support._ume_client()
sync_alarms_current(db, client, trigger_mode="schedule")
_schedule_log.info("alarms_current_auto_sync: sync finished ok")
ume_support._set_runtime_task(
"alarms_current_auto_sync",
status="running",
last_run_at=datetime.now(timezone.utc),
last_error="",
)
finally:
db.close()
except RuntimeError as exc:
if str(exc) == "alarms_current_sync_busy":
ume_support._refresh_runtime_task_idle(
"alarms_current_auto_sync",
"alarms_current",
last_error=RT_ALARMS_SYNC_IN_PROGRESS_SKIP,
)
time.sleep(30)
else:
raise
except Exception as exc:
_schedule_log.exception("alarms_current_auto_sync: sync failed: %s", exc)
ume_support._set_runtime_task(
"alarms_current_auto_sync",
status="error",
last_run_at=datetime.now(timezone.utc),
last_error=str(exc)[:240],
)
t2 = threading.Thread(target=_alarms_current_sync_loop, name="ume-alarms-current-sync", daemon=True)
t2.start()
_schedule_log.info("started thread %s alive=%s", t2.name, t2.is_alive())
if not t2.is_alive():
_schedule_log.error("ume-alarms-current-sync thread exited immediately (check uncaught errors above)")
except Exception as exc:
_schedule_log.exception("startup: alarms_current_auto_sync thread init failed: %s", exc)
ume_support._set_runtime_task(
"alarms_current_auto_sync",
status="error",
last_run_at=datetime.now(timezone.utc),
last_error=f"startup_thread_init_failed: {str(exc)[:180]}",
)
try:
if bool(getattr(settings, "ume_sync_inventory_auto_enabled", True)):
hours = int(getattr(settings, "ume_sync_inventory_every_hours", 48) or 48)
hours = max(1, min(hours, 168))
inventory_interval_s = int(hours * 3600)
ume_support._refresh_runtime_task_idle("inventory_auto_sync", "inventory")
def _inventory_auto_sync_loop() -> None:
ume_support._refresh_runtime_task_idle("inventory_auto_sync", "inventory")
while True:
try:
_schedule_log.info(
"inventory_auto_sync: loop tick paused=%s",
ume_support._runtime_is_paused("inventory_auto_sync"),
)
if ume_support._runtime_is_paused("inventory_auto_sync"):
time.sleep(1)
continue
ume_support._maybe_wait_for_sync_interval(
task_id="inventory_auto_sync",
domain="inventory",
interval_s=inventory_interval_s,
label="inventory_auto_sync",
)
_schedule_log.info(
"inventory_auto_sync: iteration start (interval=%ss)",
inventory_interval_s,
)
ume_support._set_runtime_task(
"inventory_auto_sync",
status="running",
last_run_at=datetime.now(timezone.utc),
last_error=RT_PULLING_INVENTORY,
)
db = SessionLocal()
try:
client = ume_support._ume_client()
sync_inventory_full(db, client, trigger_mode="schedule")
_schedule_log.info("inventory_auto_sync: sync finished ok")
ume_support._set_runtime_task(
"inventory_auto_sync",
status="running",
last_run_at=datetime.now(timezone.utc),
last_error="",
)
finally:
db.close()
except Exception as exc:
_schedule_log.exception("inventory_auto_sync: sync failed: %s", exc)
ume_support._set_runtime_task(
"inventory_auto_sync",
status="error",
last_run_at=datetime.now(timezone.utc),
last_error=str(exc)[:240],
)
t3 = threading.Thread(target=_inventory_auto_sync_loop, name="ume-inventory-auto-sync", daemon=True)
t3.start()
_schedule_log.info("started thread %s alive=%s", t3.name, t3.is_alive())
if not t3.is_alive():
_schedule_log.error("ume-inventory-auto-sync thread exited immediately (check uncaught errors above)")
except Exception as exc:
_schedule_log.exception("startup: inventory_auto_sync thread init failed: %s", exc)
ume_support._set_runtime_task(
"inventory_auto_sync",
status="error",
last_run_at=datetime.now(timezone.utc),
last_error=f"startup_thread_init_failed: {str(exc)[:180]}",
)
try:
if bool(getattr(settings, "ume_alarm_ws_enabled", True)) and str(getattr(settings, "ume_base_url", "") or "").strip():
if load_persisted_subscription():
_schedule_log.info("startup: loaded persisted UME alarm subscription")
ume_support._UME_WS_STOP_EVENT = threading.Event()
def _ws_on_status(msg: str) -> None:
ume_support._set_runtime_task(
"alarms_current_ws_consumer",
status="running",
last_run_at=datetime.now(timezone.utc),
last_error=str(msg or "")[:240],
)
t_ws = start_ume_alarm_ws_consumer(
ume_support._ume_client(),
on_status=_ws_on_status,
stop_event=ume_support._UME_WS_STOP_EVENT,
is_paused=lambda: ume_support._runtime_is_paused("alarms_current_ws_consumer"),
)
_schedule_log.info("started thread %s alive=%s", t_ws.name, t_ws.is_alive())
else:
ume_support._set_runtime_task("alarms_current_ws_consumer", status="paused", last_error=RT_UME_WS_DISABLED_NO_BASE_URL)
except Exception as exc:
_schedule_log.exception("startup: alarms_current_ws_consumer thread init failed: %s", exc)
ume_support._set_runtime_task(
"alarms_current_ws_consumer",
status="error",
last_run_at=datetime.now(timezone.utc),
last_error=f"startup_thread_init_failed: {str(exc)[:180]}",
)
try:
def _fwd_on_status(msg: str) -> None:
paused = ume_support._runtime_is_paused("oclaw_alarm_forwarder")
fwd = forwarder_status()
if paused:
status = "paused"
elif not bool(fwd.get("enabled")):
status = "paused"
elif bool(fwd.get("connected")):
status = "running"
else:
status = "running"
ume_support._set_runtime_task(
"oclaw_alarm_forwarder",
status=status,
last_run_at=datetime.now(timezone.utc),
last_error=str(msg or "")[:240],
)
configure_oclaw_alarm_forwarder(
is_paused=lambda: ume_support._runtime_is_paused("oclaw_alarm_forwarder"),
on_status=_fwd_on_status,
)
if is_forwarder_enabled():
ume_support._set_runtime_task("oclaw_alarm_forwarder", status="running", last_error="")
else:
ume_support._set_runtime_task(
"oclaw_alarm_forwarder",
status="paused",
last_error=RT_OCLAW_FWD_DISABLED,
)
t_fwd = start_oclaw_alarm_forwarder()
if t_fwd is not None:
_schedule_log.info("started thread %s alive=%s", t_fwd.name, t_fwd.is_alive())
except Exception as exc:
_schedule_log.exception("startup: oclaw_alarm_forwarder thread init failed: %s", exc)
ume_support._set_runtime_task(
"oclaw_alarm_forwarder",
status="error",
last_run_at=datetime.now(timezone.utc),
last_error=f"startup_thread_init_failed: {str(exc)[:180]}",
)
run_api_startup()
@app.on_event("shutdown")
@ -588,70 +87,6 @@ def health() -> dict[str, str]:
return {"status": "ok"}
@app.get("/v1/integrations/status")
def integrations_status(db: Session = Depends(get_db)) -> dict:
# netx api is up if this handler executes; still verify DB + oclaw bridge separately.
netx_api = {"status": "up"}
db_status: dict = {"status": "unknown"}
try:
t0 = time.monotonic()
db.execute(sql_text("select 1"))
db_status = {"status": "up", "latency_ms": int((time.monotonic() - t0) * 1000)}
except Exception as exc:
db_status = {"status": "down", "error": str(exc)[:240]}
oclaw_status: dict = {"status": "unknown"}
fwd = forwarder_status()
if not bool(fwd.get("enabled")):
oclaw_status = {
"status": "unknown",
"mode": "ws",
"enabled": False,
"connected": False,
"error_kind": "disabled",
"error": "NETX_OCLAW_ALARM_WS_ENABLED=false or missing token/url",
"forwarder": fwd,
}
elif bool(fwd.get("paused")):
oclaw_status = {
"status": "unknown",
"mode": "ws",
"enabled": True,
"connected": False,
"error_kind": "paused",
"error": "oclaw_alarm_forwarder runtime task paused",
"forwarder": fwd,
}
elif bool(fwd.get("connected")):
oclaw_status = {
"status": "up",
"mode": "ws",
"enabled": True,
"connected": True,
"queue_size": int(fwd.get("queue_size") or 0),
"published_ok": int(fwd.get("published_ok") or 0),
"published_fail": int(fwd.get("published_fail") or 0),
"url": str(fwd.get("url") or ""),
"forwarder": fwd,
}
else:
oclaw_status = {
"status": "down",
"mode": "ws",
"enabled": True,
"connected": False,
"error_kind": "ws_disconnected",
"error": "oclaw netx-bridge WebSocket not connected",
"queue_size": int(fwd.get("queue_size") or 0),
"url": str(fwd.get("url") or ""),
"forwarder": fwd,
}
return {"netx_api": netx_api, "db": db_status, "oclaw_bridge": oclaw_status}
@app.get("/")
def root() -> dict:
return {

View file

@ -2,7 +2,6 @@
from __future__ import annotations
import re
from typing import Any
from fastapi import HTTPException
@ -12,72 +11,24 @@ from .cli_resolve import resolve_cli_target
from .config import settings
from .ne_collect_runner import _collect_on_device
from .ne_crypto import credentials_configured
from .ne_exec_guard import _validate_command, validate_ne_exec_command
_EXEC_MAX_COMMANDS_CAP = 50
_EXEC_MAX_OUTPUT = 32_000
_EXEC_READ_TIMEOUT_DEFAULT = 60
_EXEC_READ_TIMEOUT_MAX = 120
__all__ = [
"_validate_command",
"execute_managed_ne_commands",
"validate_ne_exec_command",
]
def _exec_max_commands() -> int:
raw = int(settings.ne_exec_max_commands or 5)
return max(1, min(_EXEC_MAX_COMMANDS_CAP, raw))
# Block obvious config-change / destructive patterns (case-insensitive).
_BLOCKED_RE = re.compile(
r"(?i)("
r"configure\s+terminal|conf\s+t\b|"
r"\bwrite\s+(memory|erase)|\bcopy\s+run|\bcopy\s+startup|"
r"\breload\b|\breboot\b|\berase\b|\bformat\b|\bdelete\b|"
r"\bcommit\b|\brollback\b|startup-config|"
r"\bsystem-view\b|\bip\s+address\b|\bvlan\s+\d"
r")"
)
# Read-only CLI: show/display plus ping/traceroute reachability checks.
_ALLOWED_PREFIX_RE = re.compile(
r"(?i)^(show\s|display\s|ping\s|ping6\s|traceroute\s|tracert\s|trace\s|trace6\s)"
)
# Unicode / C1 line separators that can smuggle a second CLI after a show prefix.
_FORBIDDEN_LINE_SEPARATORS = ("\u2028", "\u2029", "\x85", "\x0b", "\x0c")
# Pipe segments allowed after show/display (output filtering only).
_ALLOWED_PIPE_SEGMENT_RE = re.compile(
r"(?i)^(include|exclude|begin|section|count|match|grep|one-line|no-more)(\s|$)"
)
_BLOCKED_PIPE_SEGMENT_RE = re.compile(r"(?i)\b(redirect|append|tee|send)\b")
def _validate_pipe_segments(cmd: str) -> None:
if "|" not in cmd:
return
parts = [p.strip() for p in cmd.split("|")]
if len(parts) < 2 or not parts[0] or any(not p for p in parts[1:]):
raise HTTPException(status_code=400, detail="command_pipe_not_allowed")
for segment in parts[1:]:
if _BLOCKED_PIPE_SEGMENT_RE.search(segment):
raise HTTPException(status_code=400, detail="command_pipe_not_allowed")
if not _ALLOWED_PIPE_SEGMENT_RE.match(segment):
raise HTTPException(status_code=400, detail="command_pipe_not_allowed")
def _validate_command(command: str) -> None:
cmd = str(command or "").strip()
if not cmd:
raise HTTPException(status_code=400, detail="empty_command")
if len(cmd) > 500:
raise HTTPException(status_code=400, detail="command_too_long")
if any(ch in cmd for ch in (";", "\n", "\r", "`")):
raise HTTPException(status_code=400, detail="command_chars_not_allowed")
if any(sep in cmd for sep in _FORBIDDEN_LINE_SEPARATORS):
raise HTTPException(status_code=400, detail="command_chars_not_allowed")
if _BLOCKED_RE.search(cmd):
raise HTTPException(status_code=400, detail="command_blocked")
if not _ALLOWED_PREFIX_RE.match(cmd):
raise HTTPException(status_code=400, detail="command_not_allowed_prefix")
_validate_pipe_segments(cmd)
def _normalize_read_timeout(sec: int | None) -> int:
raw = int(sec if sec is not None else _EXEC_READ_TIMEOUT_DEFAULT)
@ -105,7 +56,7 @@ def execute_managed_ne_commands(
if len(cmds) > max_cmds:
raise HTTPException(status_code=400, detail=f"too_many_commands (max {max_cmds})")
for c in cmds:
_validate_command(c)
validate_ne_exec_command(c)
creds, device = resolve_cli_target(db, managed_ne_id=mid or None, ume_ne_id=uid or None)
read_timeout = _normalize_read_timeout(read_timeout_sec)

74
netx_api/ne_exec_guard.py Normal file
View file

@ -0,0 +1,74 @@
"""NE CLI command allow/deny gates (read-only exec for ops tools)."""
from __future__ import annotations
import re
from fastapi import HTTPException
# Block obvious config-change / destructive patterns (case-insensitive).
_BLOCKED_RE = re.compile(
r"(?i)("
r"configure\s+terminal|conf\s+t\b|"
r"\bwrite\s+(memory|erase)|\bcopy\s+run|\bcopy\s+startup|"
r"\breload\b|\breboot\b|\berase\b|\bformat\b|\bdelete\b|"
r"\bcommit\b|\brollback\b|startup-config|"
r"\bsystem-view\b|\bip\s+address\b|\bvlan\s+\d|"
# Extra vendor / destructive surface (avoid words that appear in show output filters)
r"\bclear\s+configuration\b|\breset\s+saved-configuration\b|"
r"\bundo\s+|\bsave\s*$|\bsave\s+\S|"
r"\bfile\s+delete\b|\bftp\s+put\b|\btftp\s+put\b|"
r"\bdebug\s+all\b|\bundebug\s+all\b|"
r"\brequest\s+system\s+(reboot|halt|power-off|zeroize)\b|"
r"\bset\s+system\s+reboot\b"
r")"
)
# Read-only CLI: show/display plus ping/traceroute reachability checks.
_ALLOWED_PREFIX_RE = re.compile(
r"(?i)^(show\s|display\s|ping\s|ping6\s|traceroute\s|tracert\s|trace\s|trace6\s)"
)
# Unicode / C1 line separators that can smuggle a second CLI after a show prefix.
_FORBIDDEN_LINE_SEPARATORS = ("\u2028", "\u2029", "\x85", "\x0b", "\x0c")
# Pipe segments allowed after show/display (output filtering only).
_ALLOWED_PIPE_SEGMENT_RE = re.compile(
r"(?i)^(include|exclude|begin|section|count|match|grep|one-line|no-more)(\s|$)"
)
_BLOCKED_PIPE_SEGMENT_RE = re.compile(r"(?i)\b(redirect|append|tee|send)\b")
def _validate_pipe_segments(cmd: str) -> None:
if "|" not in cmd:
return
parts = [p.strip() for p in cmd.split("|")]
if len(parts) < 2 or not parts[0] or any(not p for p in parts[1:]):
raise HTTPException(status_code=400, detail="command_pipe_not_allowed")
for segment in parts[1:]:
if _BLOCKED_PIPE_SEGMENT_RE.search(segment):
raise HTTPException(status_code=400, detail="command_pipe_not_allowed")
if not _ALLOWED_PIPE_SEGMENT_RE.match(segment):
raise HTTPException(status_code=400, detail="command_pipe_not_allowed")
def validate_ne_exec_command(command: str) -> None:
"""Raise HTTPException if command is empty, smuggled, blocked, or not allowlisted."""
cmd = str(command or "").strip()
if not cmd:
raise HTTPException(status_code=400, detail="empty_command")
if len(cmd) > 500:
raise HTTPException(status_code=400, detail="command_too_long")
if any(ch in cmd for ch in (";", "\n", "\r", "`")):
raise HTTPException(status_code=400, detail="command_chars_not_allowed")
if any(sep in cmd for sep in _FORBIDDEN_LINE_SEPARATORS):
raise HTTPException(status_code=400, detail="command_chars_not_allowed")
if _BLOCKED_RE.search(cmd):
raise HTTPException(status_code=400, detail="command_blocked")
if not _ALLOWED_PREFIX_RE.match(cmd):
raise HTTPException(status_code=400, detail="command_not_allowed_prefix")
_validate_pipe_segments(cmd)
# Back-compat alias used by tests / callers.
_validate_command = validate_ne_exec_command

View file

@ -1,15 +1,48 @@
"""UME / long-task runtime helpers shared by API and optional worker process.
The API process still owns UME token keepalive, alarm WSS, and oclaw forwarder.
Config-sync / LLDP / port-traffic tick loops can run inline (default) or via
``python -m netx_api.worker`` when ``NETX_RUN_INLINE_SCHEDULERS=false``.
Device collectors (config_sync / LLDP / port_traffic) run via ``start_device_schedulers``
(API inline when ``NETX_RUN_INLINE_SCHEDULERS=true``, otherwise ``python -m netx_api.worker``).
API process also owns UME keepalive, alarm WSS, current-alarm/inventory sync loops,
and oclaw forwarder via ``start_api_sideband_threads``.
"""
from __future__ import annotations
import logging
import threading
import time
from datetime import datetime, timezone
from .config import settings
from .db import SessionLocal
import netx_api.ume_support as ume_support
from .ume_alarm_ws import (
complete_startup_alarm_sync_gate,
is_startup_alarm_sync_pending,
is_wss_active_for_current_alarms,
load_persisted_subscription,
start_ume_alarm_ws_consumer,
)
from .ume_sync_service import sync_alarms_current, sync_inventory_full
from .runtime_task_messages import (
RT_ALARMS_SYNC_IN_PROGRESS_SKIP,
RT_KEEPALIVE_FAILED,
RT_OCLAW_FWD_DISABLED,
RT_PULLING_ALARMS_CURRENT,
RT_PULLING_INVENTORY,
RT_STARTUP_GATE_WAITING,
RT_UME_WS_DISABLED_NO_BASE_URL,
RT_WSS_ACTIVE_SKIP_REST,
)
from .oclaw_alarm_forwarder import (
configure_oclaw_alarm_forwarder,
forwarder_status,
is_forwarder_enabled,
start_oclaw_alarm_forwarder,
)
_log = logging.getLogger("netx.ume.runtime")
_schedule_log = logging.getLogger("netx.ume.schedule")
def start_device_schedulers() -> None:
@ -22,3 +55,303 @@ def start_device_schedulers() -> None:
start_lldp_collect_scheduler()
start_port_traffic_scheduler()
_log.info("device schedulers started")
def start_api_sideband_threads() -> None:
"""UME keepalive / alarm sync / inventory / WSS / oclaw forwarder (API process)."""
try:
if bool(getattr(settings, "ume_keepalive_enabled", True)):
interval_keepalive_s = int(getattr(settings, "ume_keepalive_interval_s", 600) or 600)
interval_keepalive_s = max(30, min(interval_keepalive_s, 3600))
renew_before_s = int(getattr(settings, "ume_keepalive_renew_before_s", 900) or 900)
renew_before_s = max(30, min(renew_before_s, 86400))
def _keepalive_loop() -> None:
# Best-effort keepalive: if token exists, periodically handshake to extend TTL.
while True:
try:
if ume_support._runtime_is_paused("token_keepalive"):
time.sleep(1)
continue
client = ume_support._ume_client()
st = client.token_status()
expires_in = int(st.get("expires_in_s") or 0)
# Renew when missing/invalid TTL (0) or nearing expiry — previously 0 skipped renew forever.
if bool(st.get("has_token")) and (expires_in <= 0 or expires_in < renew_before_s):
client.renew_token()
ume_support._set_runtime_task("token_keepalive", status="running", last_run_at=datetime.now(timezone.utc), last_error="")
except Exception:
ume_support._set_runtime_task("token_keepalive", status="error", last_run_at=datetime.now(timezone.utc), last_error=RT_KEEPALIVE_FAILED)
time.sleep(interval_keepalive_s)
t = threading.Thread(target=_keepalive_loop, name="ume-token-keepalive", daemon=True)
t.start()
except Exception as exc:
_schedule_log.exception("startup: token_keepalive thread init failed: %s", exc)
ume_support._set_runtime_task(
"token_keepalive",
status="error",
last_run_at=datetime.now(timezone.utc),
last_error=f"startup_thread_init_failed: {str(exc)[:180]}",
)
try:
def _startup_alarm_sync_worker() -> None:
try:
ume_support._run_startup_alarm_sync_before_ws()
except Exception as exc:
_schedule_log.exception("startup: alarm sync before WSS failed: %s", exc)
complete_startup_alarm_sync_gate()
# Do not block HTTP /health on slow UME REST pull; WSS waits on startup_alarm_sync_gate.
t_startup_sync = threading.Thread(
target=_startup_alarm_sync_worker,
name="ume-startup-alarm-sync",
daemon=True,
)
t_startup_sync.start()
except Exception as exc:
_schedule_log.exception("startup: alarm sync thread init failed: %s", exc)
complete_startup_alarm_sync_gate()
try:
if bool(getattr(settings, "ume_sync_alarms_current_enabled", True)):
alarms_interval_s = int(getattr(settings, "ume_sync_alarms_current_interval_s", 18000) or 18000)
alarms_interval_s = max(30, min(alarms_interval_s, 86400))
def _alarms_current_sync_loop() -> None:
ume_support._refresh_runtime_task_idle("alarms_current_auto_sync", "alarms_current")
ume_support._wait_until_startup_alarm_pull_allowed("alarms_current_auto_sync")
while True:
try:
_schedule_log.info(
"alarms_current_auto_sync: loop tick paused=%s",
ume_support._runtime_is_paused("alarms_current_auto_sync"),
)
if ume_support._runtime_is_paused("alarms_current_auto_sync"):
time.sleep(1)
continue
if is_startup_alarm_sync_pending():
ume_support._refresh_runtime_task_idle(
"alarms_current_auto_sync",
"alarms_current",
last_error=RT_STARTUP_GATE_WAITING,
)
time.sleep(10)
continue
if (
bool(getattr(settings, "ume_sync_alarms_current_skip_when_ws", True))
and is_wss_active_for_current_alarms()
):
ume_support._refresh_runtime_task_idle(
"alarms_current_auto_sync",
"alarms_current",
last_error=RT_WSS_ACTIVE_SKIP_REST,
)
time.sleep(max(30, min(alarms_interval_s, 300)))
continue
ume_support._maybe_wait_for_sync_interval(
task_id="alarms_current_auto_sync",
domain="alarms_current",
interval_s=alarms_interval_s,
label="alarms_current_auto_sync",
)
_schedule_log.info(
"alarms_current_auto_sync: iteration start (interval=%ss)",
alarms_interval_s,
)
ume_support._set_runtime_task(
"alarms_current_auto_sync",
status="running",
last_run_at=datetime.now(timezone.utc),
last_error=RT_PULLING_ALARMS_CURRENT,
)
db = SessionLocal()
try:
client = ume_support._ume_client()
sync_alarms_current(db, client, trigger_mode="schedule")
_schedule_log.info("alarms_current_auto_sync: sync finished ok")
ume_support._set_runtime_task(
"alarms_current_auto_sync",
status="running",
last_run_at=datetime.now(timezone.utc),
last_error="",
)
finally:
db.close()
except RuntimeError as exc:
if str(exc) == "alarms_current_sync_busy":
ume_support._refresh_runtime_task_idle(
"alarms_current_auto_sync",
"alarms_current",
last_error=RT_ALARMS_SYNC_IN_PROGRESS_SKIP,
)
time.sleep(30)
else:
raise
except Exception as exc:
_schedule_log.exception("alarms_current_auto_sync: sync failed: %s", exc)
ume_support._set_runtime_task(
"alarms_current_auto_sync",
status="error",
last_run_at=datetime.now(timezone.utc),
last_error=str(exc)[:240],
)
t2 = threading.Thread(target=_alarms_current_sync_loop, name="ume-alarms-current-sync", daemon=True)
t2.start()
_schedule_log.info("started thread %s alive=%s", t2.name, t2.is_alive())
if not t2.is_alive():
_schedule_log.error("ume-alarms-current-sync thread exited immediately (check uncaught errors above)")
except Exception as exc:
_schedule_log.exception("startup: alarms_current_auto_sync thread init failed: %s", exc)
ume_support._set_runtime_task(
"alarms_current_auto_sync",
status="error",
last_run_at=datetime.now(timezone.utc),
last_error=f"startup_thread_init_failed: {str(exc)[:180]}",
)
try:
if bool(getattr(settings, "ume_sync_inventory_auto_enabled", True)):
hours = int(getattr(settings, "ume_sync_inventory_every_hours", 48) or 48)
hours = max(1, min(hours, 168))
inventory_interval_s = int(hours * 3600)
ume_support._refresh_runtime_task_idle("inventory_auto_sync", "inventory")
def _inventory_auto_sync_loop() -> None:
ume_support._refresh_runtime_task_idle("inventory_auto_sync", "inventory")
while True:
try:
_schedule_log.info(
"inventory_auto_sync: loop tick paused=%s",
ume_support._runtime_is_paused("inventory_auto_sync"),
)
if ume_support._runtime_is_paused("inventory_auto_sync"):
time.sleep(1)
continue
ume_support._maybe_wait_for_sync_interval(
task_id="inventory_auto_sync",
domain="inventory",
interval_s=inventory_interval_s,
label="inventory_auto_sync",
)
_schedule_log.info(
"inventory_auto_sync: iteration start (interval=%ss)",
inventory_interval_s,
)
ume_support._set_runtime_task(
"inventory_auto_sync",
status="running",
last_run_at=datetime.now(timezone.utc),
last_error=RT_PULLING_INVENTORY,
)
db = SessionLocal()
try:
client = ume_support._ume_client()
sync_inventory_full(db, client, trigger_mode="schedule")
_schedule_log.info("inventory_auto_sync: sync finished ok")
ume_support._set_runtime_task(
"inventory_auto_sync",
status="running",
last_run_at=datetime.now(timezone.utc),
last_error="",
)
finally:
db.close()
except Exception as exc:
_schedule_log.exception("inventory_auto_sync: sync failed: %s", exc)
ume_support._set_runtime_task(
"inventory_auto_sync",
status="error",
last_run_at=datetime.now(timezone.utc),
last_error=str(exc)[:240],
)
t3 = threading.Thread(target=_inventory_auto_sync_loop, name="ume-inventory-auto-sync", daemon=True)
t3.start()
_schedule_log.info("started thread %s alive=%s", t3.name, t3.is_alive())
if not t3.is_alive():
_schedule_log.error("ume-inventory-auto-sync thread exited immediately (check uncaught errors above)")
except Exception as exc:
_schedule_log.exception("startup: inventory_auto_sync thread init failed: %s", exc)
ume_support._set_runtime_task(
"inventory_auto_sync",
status="error",
last_run_at=datetime.now(timezone.utc),
last_error=f"startup_thread_init_failed: {str(exc)[:180]}",
)
try:
if bool(getattr(settings, "ume_alarm_ws_enabled", True)) and str(getattr(settings, "ume_base_url", "") or "").strip():
if load_persisted_subscription():
_schedule_log.info("startup: loaded persisted UME alarm subscription")
ume_support._UME_WS_STOP_EVENT = threading.Event()
def _ws_on_status(msg: str) -> None:
ume_support._set_runtime_task(
"alarms_current_ws_consumer",
status="running",
last_run_at=datetime.now(timezone.utc),
last_error=str(msg or "")[:240],
)
t_ws = start_ume_alarm_ws_consumer(
ume_support._ume_client(),
on_status=_ws_on_status,
stop_event=ume_support._UME_WS_STOP_EVENT,
is_paused=lambda: ume_support._runtime_is_paused("alarms_current_ws_consumer"),
)
_schedule_log.info("started thread %s alive=%s", t_ws.name, t_ws.is_alive())
else:
ume_support._set_runtime_task("alarms_current_ws_consumer", status="paused", last_error=RT_UME_WS_DISABLED_NO_BASE_URL)
except Exception as exc:
_schedule_log.exception("startup: alarms_current_ws_consumer thread init failed: %s", exc)
ume_support._set_runtime_task(
"alarms_current_ws_consumer",
status="error",
last_run_at=datetime.now(timezone.utc),
last_error=f"startup_thread_init_failed: {str(exc)[:180]}",
)
try:
def _fwd_on_status(msg: str) -> None:
paused = ume_support._runtime_is_paused("oclaw_alarm_forwarder")
fwd = forwarder_status()
if paused:
status = "paused"
elif not bool(fwd.get("enabled")):
status = "paused"
elif bool(fwd.get("connected")):
status = "running"
else:
status = "running"
ume_support._set_runtime_task(
"oclaw_alarm_forwarder",
status=status,
last_run_at=datetime.now(timezone.utc),
last_error=str(msg or "")[:240],
)
configure_oclaw_alarm_forwarder(
is_paused=lambda: ume_support._runtime_is_paused("oclaw_alarm_forwarder"),
on_status=_fwd_on_status,
)
if is_forwarder_enabled():
ume_support._set_runtime_task("oclaw_alarm_forwarder", status="running", last_error="")
else:
ume_support._set_runtime_task(
"oclaw_alarm_forwarder",
status="paused",
last_error=RT_OCLAW_FWD_DISABLED,
)
t_fwd = start_oclaw_alarm_forwarder()
if t_fwd is not None:
_schedule_log.info("started thread %s alive=%s", t_fwd.name, t_fwd.is_alive())
except Exception as exc:
_schedule_log.exception("startup: oclaw_alarm_forwarder thread init failed: %s", exc)
ume_support._set_runtime_task(
"oclaw_alarm_forwarder",
status="error",
last_run_at=datetime.now(timezone.utc),
last_error=f"startup_thread_init_failed: {str(exc)[:180]}",
)

View file

@ -47,6 +47,7 @@ from .ume_token_store import (
)
_schedule_log = logging.getLogger("netx.ume.schedule")
_BOOT_MONO = time.monotonic()
_UME_CLIENT_SINGLETON = UMEClient(
token_loader=lambda: load_shared_token(),

425
netx_api/webcrt_channel.py Normal file
View file

@ -0,0 +1,425 @@
"""WebCRT channel helpers: keymap, prompt heuristics, encoding, queues."""
from __future__ import annotations
import io
import json
import logging
import queue
import re
import time
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
from netmiko import ConnectHandler
from .config import settings
_log = logging.getLogger("netx.webcrt")
_NETWORK_CLI_KEY_SEQS: tuple[tuple[str, str], ...] = (
("\x1b[1~", "\x01"), # Home -> Ctrl-A
("\x1b[3~", "\x04"), # Delete key -> Ctrl-D
("\x1b[4~", "\x05"), # End -> Ctrl-E
("\x1b[H", "\x01"),
("\x1b[F", "\x05"),
("\x1bOH", "\x01"),
("\x1bOF", "\x05"),
("\x1bOA", "\x1b[A"), # App Up -> CSI Up
("\x1bOB", "\x1b[B"),
("\x1bOC", "\x1b[C"),
("\x1bOD", "\x1b[D"), # App Left -> CSI Left
("\x7f", "\x08"), # DEL -> BS
)
def uses_network_cli_keymap(device_type: str = "", vendor: str = "") -> bool:
blob = f"{device_type} {vendor}".strip().lower()
if not blob:
return True
for token in ("linux", "ubuntu", "centos", "debian", "redhat", "unix", "generic_telnet", "generic"):
if token in blob:
return False
return True
def map_network_cli_keys(
data: str,
*,
device_type: str = "",
vendor: str = "",
protocol: str = "",
) -> str:
"""Rewrite xterm key sequences for network-device CLIs."""
del device_type, vendor, protocol # protocol kept for call-site compatibility
text = str(data or "")
if not text:
return text
out: list[str] = []
i = 0
n = len(text)
while i < n:
matched = False
for seq, repl in _NETWORK_CLI_KEY_SEQS:
if text.startswith(seq, i):
out.append(repl)
i += len(seq)
matched = True
break
if not matched:
out.append(text[i])
i += 1
return "".join(out)
def channel_return(conn: ConnectHandler | None) -> str:
"""Netmiko line ending for this session (SSH usually \\n, Telnet often \\r\\n)."""
if conn is None:
return "\n"
ret = getattr(conn, "RETURN", None)
if isinstance(ret, str) and ret:
return ret
return "\n"
def map_network_cli_enter(data: str, conn: ConnectHandler | None) -> str:
"""Map xterm Enter (\\r) to the device's Netmiko RETURN."""
text = str(data or "")
if not text:
return text
ret = channel_return(conn)
if ret == "\r":
return text
# Prefer replacing CRLF first so Telnet RETURN \\r\\n does not double-expand.
return text.replace("\r\n", ret).replace("\r", ret)
def _drain_channel(conn: ConnectHandler, *, rounds: int = 6, wait: float = 0.06) -> str:
"""Read whatever is already sitting on the channel after login."""
chunks: list[str] = []
empty_streak = 0
for _ in range(max(1, rounds)):
time.sleep(wait)
try:
part = conn.read_channel()
except Exception:
break
if part:
chunks.append(str(part))
empty_streak = 0
else:
empty_streak += 1
if empty_streak >= 2 and chunks:
break
return "".join(chunks)
def _session_log_text(buf: io.BytesIO | None) -> str:
"""Decode Netmiko session_log buffer into display text."""
if buf is None:
return ""
try:
raw = buf.getvalue()
except Exception:
return ""
if isinstance(raw, bytes):
return raw.decode("utf-8", errors="replace")
return str(raw or "")
def _looks_like_cli_prompt(text: str) -> bool:
s = str(text or "").rstrip()
if not s:
return False
# Buffer races can leave a stray ':' after Huawei ``<r1>`` (from prior ``[Y/N]:``).
if s.endswith(":") and ">" in s:
s = s[:-1].rstrip()
# Common network CLI prompts: <r1> [HUAWEI] Router# Router>
return bool(re.search(r"(?:[>\]]|#)\s*$", s)) or bool(re.search(r"<[^>\r\n]+>\s*$", s))
def _looks_like_login_prompt(text: str) -> bool:
"""True when the transcript ends at Username:/Login:/Password: (interactive auth)."""
s = str(text or "").replace("\r\n", "\n").replace("\r", "\n")
lines = [ln.strip() for ln in s.split("\n") if ln.strip()]
if not lines:
return False
last = lines[-1]
return bool(re.search(r"(?i)(user\s*name|login|password)\s*:\s*$", last))
def _looks_like_password_change_prompt(text: str) -> bool:
"""Huawei/VRP post-auth ``Change now? [Y/N]:`` (Netmiko already answers N)."""
s = str(text or "").replace("\r\n", "\n").replace("\r", "\n")
lines = [ln.strip() for ln in s.split("\n") if ln.strip()]
if not lines:
return False
last = lines[-1]
return bool(re.search(r"(?i)(change\s*now|please\s*choose|password\s+needs\s+to\s+be\s+changed).{0,80}:\s*$", last)) or bool(
re.search(r"\[Y/N\]\s*:\s*$", last, flags=re.I)
)
# Cisco/Netmiko often yields "R2#R2#" when a sync Enter is appended without a newline.
_GLUED_PROMPT_RE = re.compile(r"(?<=[#>])(?=(?:[A-Za-z0-9][\w.\-:]{0,62})[#>])")
def normalize_cli_transcript(text: str) -> str:
"""Normalize login transcript for xterm (convertEol) and un-glue prompts."""
s = str(text or "").replace("\r\n", "\n").replace("\r", "\n")
s = _GLUED_PROMPT_RE.sub("\n", s)
lines = s.split("\n")
while lines and not str(lines[-1]).strip():
lines.pop()
# Drop blank lines immediately before a final prompt (banner\n\nR2# -> banner\nR2#).
while len(lines) >= 2 and not str(lines[-2]).strip() and _looks_like_cli_prompt(lines[-1]):
lines.pop(-2)
# Collapse trailing duplicate prompt lines (slow VMs often echo R2# several times).
while len(lines) >= 2 and str(lines[-1]).strip() == str(lines[-2]).strip() and _looks_like_cli_prompt(lines[-1]):
lines.pop()
return "\n".join(lines)
def prepare_bootstrap_output(text: str) -> str:
"""Full login transcript for UI replay; keep final prompt, no trailing newline after it.
Trailing newline would leave the cursor on a blank line so the first typed line
looks wrong; cursor should sit after the prompt like a real CRT.
"""
s = normalize_cli_transcript(text)
# Drop a stray ':' glued onto Huawei ``<host>`` after ``[Y/N]:`` buffer races.
s = re.sub(r"(<[^\r\n>]+>):\s*$", r"\1", s)
return s
def _capture_raw_channel(conn: ConnectHandler, *, duration: float = 0.5) -> str:
"""Read leftover PTY bytes into text (banner/MOTD after SSH auth).
Interactive WebCRT skips Netmiko session_preparation, so the post-auth banner
often never lands in ``session_log`` and must be pulled from the live channel.
"""
chunks: list[str] = []
channel = getattr(conn, "remote_conn", None)
if channel is None:
try:
return _drain_channel(conn, rounds=max(2, int(duration / 0.05)), wait=0.05)
except Exception:
return ""
end = time.time() + max(0.1, float(duration))
while time.time() < end:
got = False
try:
# Paramiko SSH channel
if hasattr(channel, "recv_ready") and hasattr(channel, "recv"):
if channel.recv_ready():
raw = channel.recv(65535)
if raw:
got = True
if isinstance(raw, bytes):
chunks.append(raw.decode("utf-8", errors="replace"))
else:
chunks.append(str(raw))
# telnetlib-style
elif callable(getattr(channel, "read_very_eager", None)):
data = channel.read_very_eager()
if data:
got = True
if isinstance(data, bytes):
chunks.append(data.decode("utf-8", errors="replace"))
else:
chunks.append(str(data))
else:
part = conn.read_channel()
if part:
got = True
chunks.append(str(part))
except Exception:
break
if not got:
time.sleep(0.04)
return "".join(chunks)
def _drain_raw_channel(conn: ConnectHandler, *, duration: float = 0.5) -> None:
"""Discard leftover bytes on the live channel (SSH/Telnet) after login priming."""
_capture_raw_channel(conn, duration=duration)
def _prime_interactive_channel(conn: ConnectHandler, *, already_prompted: bool = False) -> str:
"""Sync interactive channel after login; return captured banner/prompt text.
Skip the sync Enter when the login transcript already ends with a CLI prompt —
otherwise slow Cisco VMs accumulate duplicate ``R2#`` lines in the bootstrap.
"""
parts: list[str] = []
try:
parts.append(_capture_raw_channel(conn, duration=0.25))
except Exception:
pass
if not already_prompted:
try:
conn.write_channel(channel_return(conn))
except Exception:
try:
conn.write_channel("\n")
except Exception:
return "".join(parts)
try:
parts.append(_drain_channel(conn, rounds=6, wait=0.08))
except Exception:
pass
try:
parts.append(_capture_raw_channel(conn, duration=0.35))
except Exception:
pass
return "".join(parts)
def _is_prompt_only_echo(text: str, prompt_hint: str = "") -> bool:
"""True when chunk is only whitespace / CR / a repeated prompt (safe to drop after bootstrap)."""
s = str(text or "").replace("\r\n", "\n").replace("\r", "\n").strip()
if not s:
return True
hint = str(prompt_hint or "").strip()
if hint and s == hint:
return True
# Single-line prompt echo only.
if "\n" not in s and _looks_like_cli_prompt(s):
return True
if hint and all(line.strip() in ("", hint) for line in s.split("\n")):
return True
return False
def _normalize_encoding(name: str) -> str:
enc = str(name or "utf-8").strip().lower().replace("_", "-")
if enc in ("gbk", "gb2312", "gb18030", "cp936"):
return "gbk"
return "utf-8"
def _decode_bytes(data: bytes, encoding: str) -> str:
enc = _normalize_encoding(encoding)
try:
return data.decode(enc, errors="replace")
except Exception:
return data.decode("utf-8", errors="replace")
def _encode_text(text: str, encoding: str) -> bytes:
enc = _normalize_encoding(encoding)
try:
return text.encode(enc, errors="replace")
except Exception:
return text.encode("utf-8", errors="replace")
class _BoundedByteQueue:
"""Thread-safe queue that drops oldest chunks when full (backpressure)."""
def __init__(self, maxsize: int = 2000) -> None:
self._q: queue.Queue[bytes | None] = queue.Queue()
self._max = max(8, int(maxsize or 2000))
self._cond = threading.Condition()
self.dropped = 0
self._reported = 0
def put(self, item: bytes | None) -> None:
with self._cond:
while self._q.qsize() >= self._max:
try:
self._q.get_nowait()
self.dropped += 1
except queue.Empty:
break
self._q.put(item)
self._cond.notify()
def put_nowait(self, item: bytes | None) -> None:
self.put(item)
def get_nowait(self) -> bytes | None:
with self._cond:
return self._q.get_nowait()
def get(self, timeout: float = 0.25) -> bytes | None:
"""Block until a chunk is available or timeout (raises queue.Empty)."""
deadline = time.time() + max(0.0, float(timeout))
with self._cond:
while self._q.empty():
remaining = deadline - time.time()
if remaining <= 0:
raise queue.Empty
self._cond.wait(timeout=remaining)
return self._q.get_nowait()
def qsize(self) -> int:
with self._cond:
return self._q.qsize()
def take_drop_delta(self) -> int:
"""Return newly dropped chunk count since last call (for client notice)."""
with self._cond:
delta = int(self.dropped) - int(self._reported)
if delta <= 0:
return 0
self._reported = int(self.dropped)
return delta
def _utc_now() -> datetime:
return datetime.now(timezone.utc)
def _utc_iso() -> str:
return _utc_now().isoformat()
def webcrt_data_root() -> Path:
root = Path(str(settings.webcrt_data_dir or "data/webcrt"))
root.mkdir(parents=True, exist_ok=True)
return root.resolve()
def _session_log_path(session_id: str) -> Path:
folder = webcrt_data_root() / "sessions"
folder.mkdir(parents=True, exist_ok=True)
return folder / f"{session_id}.log"
def read_session_log_tail(session_id: str, *, max_bytes: int = 49152) -> str:
"""Best-effort UTF-8 tail of the on-disk session transcript (for WS re-attach)."""
path = _session_log_path(session_id)
try:
if not path.is_file():
return ""
size = path.stat().st_size
take = max(1024, min(int(max_bytes or 49152), 256 * 1024))
with path.open("rb") as fh:
if size > take:
fh.seek(size - take)
raw = fh.read()
# Drop partial first line after seek.
nl = raw.find(b"\n")
if 0 <= nl < len(raw) - 1:
raw = raw[nl + 1 :]
else:
raw = fh.read()
text = raw.decode("utf-8", errors="replace")
# Strip header comment lines from the visible replay.
lines = [ln for ln in text.splitlines(keepends=True) if not ln.startswith("# session=")]
return "".join(lines)
except Exception:
_log.debug("webcrt session log tail failed session=%s", session_id, exc_info=True)
return ""
def _audit(event: str, **fields: Any) -> None:
record = {"ts": _utc_iso(), "event": event, **fields}
try:
path = webcrt_data_root() / "audit.jsonl"
with path.open("a", encoding="utf-8") as fh:
fh.write(json.dumps(record, ensure_ascii=False) + "\n")
except Exception:
_log.exception("webcrt audit write failed")
_log.info("webcrt.%s %s", event, {k: v for k, v in fields.items() if k != "detail"})

File diff suppressed because it is too large Load diff

1111
netx_api/webcrt_session.py Normal file

File diff suppressed because it is too large Load diff

View file

@ -1,6 +1,7 @@
"""Background worker process for long-running schedulers.
Run separately from the API when NETX_RUN_INLINE_SCHEDULERS=false:
Default deployment: API has ``NETX_RUN_INLINE_SCHEDULERS=false``; run this
alongside the API:
python -m netx_api.worker