mirror of
https://github.com/hansjone/netx.git
synced 2026-10-10 09:50:44 +08:00
Remove legacy OClaw alarm WSS and fix key-alert rules table layout.
Key alerts now deliver only via the DSH hub; Modal Body no longer flex-shrinks the rules rows to empty. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
parent
f937efa0f9
commit
0a4561f385
25 changed files with 36 additions and 729 deletions
|
|
@ -67,13 +67,6 @@ def shutdown_runtime(*, reason: str = "lifespan") -> None:
|
|||
except Exception: # noqa: BLE001
|
||||
_log.exception("shutdown_ws_consumer failed")
|
||||
|
||||
try:
|
||||
from .oclaw_alarm_forwarder import shutdown_oclaw_alarm_forwarder
|
||||
|
||||
shutdown_oclaw_alarm_forwarder()
|
||||
except Exception: # noqa: BLE001
|
||||
_log.exception("shutdown_oclaw_alarm_forwarder failed")
|
||||
|
||||
try:
|
||||
from .webcrt_session_registry import close_all_sessions
|
||||
|
||||
|
|
|
|||
|
|
@ -20,8 +20,6 @@ class Settings(BaseSettings):
|
|||
oclaw_connect_timeout_sec: float = 15.0
|
||||
oclaw_analyze_read_timeout_sec: float = 180.0
|
||||
oclaw_health_timeout_sec: float = 8.0
|
||||
oclaw_alarm_ws_enabled: bool = False
|
||||
oclaw_alarm_ws_url: str = "ws://127.0.0.1:8787/ws/netx-bridge"
|
||||
# UME RESTCONF integration
|
||||
ume_base_url: str = ""
|
||||
ume_username: str = ""
|
||||
|
|
|
|||
|
|
@ -12,7 +12,6 @@ from sqlalchemy.orm import Session
|
|||
from .config import settings
|
||||
from .db import get_db
|
||||
from .dsh_alarm_hub import dsh_alarm_ws_loop, hub_status
|
||||
from .oclaw_alarm_forwarder import forwarder_status
|
||||
|
||||
router = APIRouter(tags=["health"])
|
||||
|
||||
|
|
@ -79,7 +78,7 @@ def health_ready(db: Session = Depends(get_db)) -> dict[str, Any]:
|
|||
|
||||
@router.get("/v1/integrations/status")
|
||||
def integrations_status(db: Session = Depends(get_db)) -> dict:
|
||||
"""netx API + DB + oclaw bridge status."""
|
||||
"""netx API + DB + DSH alarm hub status."""
|
||||
netx_api = {"status": "up"}
|
||||
|
||||
db_status: dict = {"status": "unknown"}
|
||||
|
|
@ -90,56 +89,8 @@ def integrations_status(db: Session = Depends(get_db)) -> dict:
|
|||
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,
|
||||
"dsh_alarm_hub": hub_status(),
|
||||
}
|
||||
|
|
|
|||
|
|
@ -7,10 +7,9 @@ from typing import Any
|
|||
from sqlalchemy.exc import IntegrityError
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from .dsh_alarm_hub import publish_alarm
|
||||
from .key_alert_matcher import match_key_alert_rule
|
||||
from .models import UmeInventoryNE, UmeKeyAlertForwardLog
|
||||
from .oclaw_alarm_forwarder import enqueue_alarm_forward, is_forwarder_operational
|
||||
from .dsh_alarm_hub import publish_alarm
|
||||
from .ume_sync_service import (
|
||||
_derive_ne_id_from_alarm,
|
||||
_pick,
|
||||
|
|
@ -82,6 +81,7 @@ def maybe_forward_key_alert(
|
|||
)
|
||||
.first()
|
||||
)
|
||||
# Column name is historical (oclaw_ok); now means DSH hub delivery succeeded.
|
||||
if existing is not None and int(existing.oclaw_ok or 0) == 1:
|
||||
return False
|
||||
|
||||
|
|
@ -96,22 +96,10 @@ def maybe_forward_key_alert(
|
|||
payload["rule_key"] = str(rule.notification_id or "")
|
||||
|
||||
hub_sent = publish_alarm(payload)
|
||||
oclaw_queued = False
|
||||
if is_forwarder_operational():
|
||||
oclaw_queued = bool(enqueue_alarm_forward(payload))
|
||||
if hub_sent <= 0 and not oclaw_queued:
|
||||
if hub_sent <= 0:
|
||||
return False
|
||||
|
||||
# Hub-only delivery already reached DSH clients — mark ok for dedup.
|
||||
# Oclaw path stays pending until the bridge records a result.
|
||||
delivered_ok = hub_sent > 0 and not oclaw_queued
|
||||
status = []
|
||||
if hub_sent > 0:
|
||||
status.append(f"dsh_hub:{hub_sent}")
|
||||
if oclaw_queued:
|
||||
status.append("oclaw_queued")
|
||||
status_text = ",".join(status) if status else "queued"
|
||||
|
||||
status_text = f"dsh_hub:{hub_sent}"
|
||||
row = existing
|
||||
if row is None:
|
||||
row = UmeKeyAlertForwardLog(
|
||||
|
|
@ -120,62 +108,21 @@ def maybe_forward_key_alert(
|
|||
rule_key=str(rule.notification_id or ""),
|
||||
notification_id=notification_id_from_norm(norm),
|
||||
forwarded_at=_utc_now_naive(),
|
||||
oclaw_ok=1 if delivered_ok else 0,
|
||||
error="" if delivered_ok else status_text,
|
||||
oclaw_ok=1,
|
||||
error="",
|
||||
)
|
||||
db.add(row)
|
||||
else:
|
||||
row.notification_id = notification_id_from_norm(norm)
|
||||
row.rule_key = str(rule.notification_id or "")
|
||||
row.forwarded_at = _utc_now_naive()
|
||||
row.oclaw_ok = 1 if delivered_ok else 0
|
||||
row.error = "" if delivered_ok else status_text
|
||||
row.oclaw_ok = 1
|
||||
row.error = ""
|
||||
try:
|
||||
db.commit()
|
||||
except IntegrityError:
|
||||
db.rollback()
|
||||
_log.debug("key_alert forward log race alarm_key=%s action=%s", alarm_key, act)
|
||||
else:
|
||||
_log.debug("key_alert forwarded via DSH hub=%s status=%s", hub_sent, status_text)
|
||||
return True
|
||||
|
||||
|
||||
def record_forward_result(
|
||||
*,
|
||||
alarm_key: str,
|
||||
action: str,
|
||||
ok: bool,
|
||||
error: str = "",
|
||||
rule_key: str = "",
|
||||
) -> None:
|
||||
from .db import SessionLocal
|
||||
|
||||
key = str(alarm_key or "").strip()
|
||||
act = str(action or "").strip().lower()
|
||||
if not key or not act:
|
||||
return
|
||||
db = SessionLocal()
|
||||
try:
|
||||
row = (
|
||||
db.query(UmeKeyAlertForwardLog)
|
||||
.filter(UmeKeyAlertForwardLog.alarm_key == key, UmeKeyAlertForwardLog.action == act)
|
||||
.first()
|
||||
)
|
||||
if row is None:
|
||||
row = UmeKeyAlertForwardLog(
|
||||
alarm_key=key,
|
||||
action=act,
|
||||
rule_key=str(rule_key or "").strip(),
|
||||
forwarded_at=_utc_now_naive(),
|
||||
oclaw_ok=1 if ok else 0,
|
||||
error="" if ok else str(error or "forward_failed")[:240],
|
||||
)
|
||||
db.add(row)
|
||||
else:
|
||||
if rule_key and not str(row.rule_key or "").strip():
|
||||
row.rule_key = str(rule_key or "").strip()
|
||||
row.forwarded_at = _utc_now_naive()
|
||||
row.oclaw_ok = 1 if ok else 0
|
||||
row.error = "" if ok else str(error or "forward_failed")[:240]
|
||||
db.commit()
|
||||
except Exception:
|
||||
db.rollback()
|
||||
finally:
|
||||
db.close()
|
||||
|
|
|
|||
|
|
@ -14,7 +14,6 @@ from .config import settings
|
|||
from .db import db_pool_status
|
||||
from .db_storage_metrics import collect_db_storage_metrics
|
||||
from .host_metrics import collect_host_metrics
|
||||
from .oclaw_alarm_forwarder import forwarder_status
|
||||
|
||||
router = APIRouter(tags=["metrics"])
|
||||
|
||||
|
|
@ -30,7 +29,6 @@ def collect_runtime_metrics() -> dict[str, Any]:
|
|||
"db_storage": collect_db_storage_metrics(),
|
||||
"cli_budget": cli_budget_status(),
|
||||
"audit_queue": audit_queue_status(),
|
||||
"oclaw_forwarder": forwarder_status(),
|
||||
"schedulers_inline": bool(getattr(settings, "run_inline_schedulers", True)),
|
||||
"host": collect_host_metrics(),
|
||||
}
|
||||
|
|
@ -92,17 +90,6 @@ def _prom_lines(metrics: dict[str, Any]) -> str:
|
|||
for key in ("depth", "dropped", "maxsize"):
|
||||
if key in audit:
|
||||
lines.append(f"netx_audit_queue_{key} {audit[key]}")
|
||||
fwd = metrics.get("oclaw_forwarder") or {}
|
||||
for key, prom in (
|
||||
("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)}")
|
||||
sched = metrics.get("device_schedulers") or {}
|
||||
if "stale" in sched:
|
||||
lines.append(f'netx_device_schedulers_stale {1 if sched.get("stale") else 0}')
|
||||
|
|
|
|||
|
|
@ -1,350 +0,0 @@
|
|||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import queue
|
||||
import threading
|
||||
import time
|
||||
from collections.abc import Callable
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any
|
||||
|
||||
import websocket
|
||||
|
||||
from .config import settings
|
||||
from .runtime_task_messages import fwd_state_code
|
||||
|
||||
_log = logging.getLogger("netx.oclaw.alarm_forwarder")
|
||||
|
||||
|
||||
def _normalize_bridge_error(exc: BaseException | str) -> str:
|
||||
"""Map low-level WS errors to stable fwd:* codes for UI i18n."""
|
||||
raw = str(exc or "").strip()
|
||||
low = raw.lower()
|
||||
name = type(exc).__name__ if isinstance(exc, BaseException) else ""
|
||||
if (
|
||||
"10061" in raw
|
||||
or "actively refused" in low
|
||||
or "积极拒绝" in raw
|
||||
or "connection refused" in low
|
||||
or name == "ConnectionRefusedError"
|
||||
):
|
||||
return fwd_state_code("connect_refused")
|
||||
if "timed out" in low or "timeout" in low or name in {"TimeoutError", "socket.timeout"}:
|
||||
return fwd_state_code("connect_timeout")
|
||||
if "auth failed" in low or "auth-fail" in low or "invalid_token" in low:
|
||||
return fwd_state_code("auth_failed")
|
||||
if raw:
|
||||
return raw[:200]
|
||||
return fwd_state_code("disconnected")
|
||||
|
||||
|
||||
_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()
|
||||
_WS: Any | None = None
|
||||
_CONNECTED = threading.Event()
|
||||
_IS_PAUSED: Callable[[], bool] | None = None
|
||||
_ON_STATUS: Callable[[str], None] | None = None
|
||||
_STATS_LOCK = threading.Lock()
|
||||
_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()
|
||||
|
||||
|
||||
def _bridge_token() -> str:
|
||||
return str(getattr(settings, "oclaw_analyze_token", "") or "").strip()
|
||||
|
||||
|
||||
def _bridge_url() -> str:
|
||||
return str(getattr(settings, "oclaw_alarm_ws_url", "") or "").strip()
|
||||
|
||||
|
||||
def is_forwarder_enabled() -> bool:
|
||||
if not bool(getattr(settings, "oclaw_alarm_ws_enabled", False)):
|
||||
return False
|
||||
return bool(_bridge_url()) and bool(_bridge_token())
|
||||
|
||||
|
||||
def configure_oclaw_alarm_forwarder(
|
||||
*,
|
||||
is_paused: Callable[[], bool] | None = None,
|
||||
on_status: Callable[[str], None] | None = None,
|
||||
) -> None:
|
||||
global _IS_PAUSED, _ON_STATUS
|
||||
_IS_PAUSED = is_paused
|
||||
_ON_STATUS = on_status
|
||||
|
||||
|
||||
def _forwarder_paused() -> bool:
|
||||
if _IS_PAUSED is None:
|
||||
return False
|
||||
try:
|
||||
return bool(_IS_PAUSED())
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
|
||||
def is_forwarder_operational() -> bool:
|
||||
return is_forwarder_enabled() and not _forwarder_paused()
|
||||
|
||||
|
||||
def _notify_status(msg: str) -> None:
|
||||
if _ON_STATUS is None:
|
||||
return
|
||||
try:
|
||||
_ON_STATUS(str(msg or "")[:240])
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def request_forwarder_reconnect() -> None:
|
||||
with _CONN_LOCK:
|
||||
ws = _WS
|
||||
if ws is not None:
|
||||
try:
|
||||
ws.close()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def enqueue_alarm_forward(payload: dict[str, Any]) -> bool:
|
||||
if not is_forwarder_operational():
|
||||
return False
|
||||
try:
|
||||
_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))
|
||||
deadline = time.time() + 10.0
|
||||
while time.time() < deadline:
|
||||
raw = ws.recv()
|
||||
if not raw:
|
||||
continue
|
||||
try:
|
||||
msg = json.loads(raw)
|
||||
except Exception:
|
||||
continue
|
||||
if not isinstance(msg, dict):
|
||||
continue
|
||||
if str(msg.get("type") or "").strip() == "auth-ok":
|
||||
return True
|
||||
if str(msg.get("type") or "").strip() == "auth-fail":
|
||||
return False
|
||||
return False
|
||||
|
||||
|
||||
def _dispatch_one(ws: Any, payload: dict[str, Any]) -> dict[str, Any]:
|
||||
envelope = {
|
||||
"type": "event",
|
||||
"event": "netx.alarm",
|
||||
"payload": payload,
|
||||
}
|
||||
ws.send(json.dumps(envelope, ensure_ascii=False, default=str))
|
||||
deadline = time.time() + 15.0
|
||||
alarm_key = str(payload.get("alarm_key") or "")
|
||||
while time.time() < deadline:
|
||||
raw = ws.recv()
|
||||
if not raw:
|
||||
continue
|
||||
try:
|
||||
msg = json.loads(raw)
|
||||
except Exception:
|
||||
continue
|
||||
if not isinstance(msg, dict):
|
||||
continue
|
||||
if str(msg.get("type") or "").strip() == "pong":
|
||||
continue
|
||||
if str(msg.get("type") or "").strip() == "ack":
|
||||
if alarm_key and str(msg.get("alarm_key") or "").strip() not in {"", alarm_key}:
|
||||
continue
|
||||
return msg
|
||||
return {"type": "ack", "alarm_key": alarm_key, "ok": False, "error": "ack_timeout"}
|
||||
|
||||
|
||||
def _run_loop() -> None:
|
||||
global _WS
|
||||
backoff_s = 2.0
|
||||
while not _STOP_EVENT.is_set():
|
||||
if not is_forwarder_enabled() or _forwarder_paused():
|
||||
_CONNECTED.clear()
|
||||
if _forwarder_paused():
|
||||
_notify_status("fwd:paused")
|
||||
else:
|
||||
_notify_status("fwd:disabled")
|
||||
time.sleep(2.0)
|
||||
continue
|
||||
url = _bridge_url()
|
||||
ws = None
|
||||
try:
|
||||
_notify_status("fwd:connecting")
|
||||
ws = websocket.create_connection(url, timeout=20)
|
||||
ws.settimeout(30)
|
||||
if not _send_auth(ws):
|
||||
raise RuntimeError("oclaw netx-bridge auth failed")
|
||||
with _CONN_LOCK:
|
||||
_WS = ws
|
||||
_CONNECTED.set()
|
||||
backoff_s = 2.0
|
||||
_log.info("oclaw netx-bridge connected url=%s", url[:120])
|
||||
_notify_status("fwd:connected")
|
||||
while not _STOP_EVENT.is_set():
|
||||
if _forwarder_paused():
|
||||
_notify_status("fwd:paused")
|
||||
raise RuntimeError("forwarder paused")
|
||||
if not is_forwarder_enabled():
|
||||
_notify_status("fwd:disabled")
|
||||
raise RuntimeError("forwarder disabled")
|
||||
try:
|
||||
payload = _outbound_q().get(timeout=1.0)
|
||||
except queue.Empty:
|
||||
try:
|
||||
ws.send(json.dumps({"type": "ping", "ts": _utc_now_iso()}, ensure_ascii=False))
|
||||
except Exception:
|
||||
raise
|
||||
continue
|
||||
try:
|
||||
ack = _dispatch_one(ws, payload)
|
||||
ok = bool(ack.get("ok"))
|
||||
err = str(ack.get("error") or "")[:240]
|
||||
with _STATS_LOCK:
|
||||
if ok:
|
||||
_STATS["published_ok"] = int(_STATS.get("published_ok", 0)) + 1
|
||||
else:
|
||||
_STATS["published_fail"] = int(_STATS.get("published_fail", 0)) + 1
|
||||
try:
|
||||
from .key_alert_forward import record_forward_result
|
||||
|
||||
record_forward_result(
|
||||
alarm_key=str(payload.get("alarm_key") or ""),
|
||||
action=str(payload.get("action") or ""),
|
||||
ok=ok,
|
||||
error=err,
|
||||
rule_key=str(payload.get("rule_key") or ""),
|
||||
)
|
||||
except Exception as rec_exc:
|
||||
_log.warning("forward result record failed: %s", str(rec_exc)[:120])
|
||||
if not ok:
|
||||
_log.warning(
|
||||
"oclaw ack failed alarm_key=%s error=%s",
|
||||
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])
|
||||
_requeue_or_drop(payload, reason=str(exc)[:120])
|
||||
raise
|
||||
finally:
|
||||
_outbound_q().task_done()
|
||||
except Exception as exc:
|
||||
_CONNECTED.clear()
|
||||
err = _normalize_bridge_error(exc)
|
||||
_log.warning("oclaw netx-bridge disconnected: %s", str(exc)[:200])
|
||||
_notify_status(err)
|
||||
time.sleep(backoff_s)
|
||||
backoff_s = min(backoff_s * 1.5, 60.0)
|
||||
finally:
|
||||
with _CONN_LOCK:
|
||||
_WS = None
|
||||
_CONNECTED.clear()
|
||||
if ws is not None:
|
||||
try:
|
||||
ws.close()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def start_oclaw_alarm_forwarder() -> threading.Thread | None:
|
||||
global _THREAD
|
||||
if _THREAD is not None and _THREAD.is_alive():
|
||||
return _THREAD
|
||||
_STOP_EVENT.clear()
|
||||
_THREAD = threading.Thread(target=_run_loop, name="oclaw-alarm-forwarder", daemon=True)
|
||||
_THREAD.start()
|
||||
return _THREAD
|
||||
|
||||
|
||||
def shutdown_oclaw_alarm_forwarder() -> None:
|
||||
_STOP_EVENT.set()
|
||||
|
||||
|
||||
def forwarder_status() -> dict[str, Any]:
|
||||
with _STATS_LOCK:
|
||||
stats = dict(_STATS)
|
||||
paused = _forwarder_paused()
|
||||
enabled = is_forwarder_enabled()
|
||||
return {
|
||||
"enabled": enabled,
|
||||
"operational": bool(enabled and not paused),
|
||||
"paused": paused,
|
||||
"connected": _CONNECTED.is_set(),
|
||||
"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)),
|
||||
}
|
||||
|
|
@ -336,13 +336,6 @@ def _ume_runtime_items() -> list[dict[str, Any]]:
|
|||
from .ume_support import _list_runtime_tasks
|
||||
except Exception:
|
||||
return []
|
||||
fwd: dict[str, Any] | None = None
|
||||
try:
|
||||
from .oclaw_alarm_forwarder import forwarder_status
|
||||
|
||||
fwd = forwarder_status()
|
||||
except Exception:
|
||||
fwd = None
|
||||
ws_conn: dict[str, Any] | None = None
|
||||
try:
|
||||
from .ume_alarm_ws import get_ws_connection_status
|
||||
|
|
@ -356,24 +349,7 @@ def _ume_runtime_items() -> list[dict[str, Any]]:
|
|||
status = str(row.get("status") or "unknown")
|
||||
progress = str(row.get("interval_label") or "")
|
||||
detail = str(row.get("last_error") or "")[:240]
|
||||
if task == "oclaw_alarm_forwarder" and isinstance(fwd, dict):
|
||||
bits: list[str] = []
|
||||
if not bool(fwd.get("enabled")):
|
||||
bits.append("disabled")
|
||||
elif bool(fwd.get("paused")):
|
||||
bits.append("paused")
|
||||
elif bool(fwd.get("connected")):
|
||||
bits.append("connected")
|
||||
else:
|
||||
bits.append("disconnected")
|
||||
q = int(fwd.get("queue_size") or 0)
|
||||
if q > 0:
|
||||
bits.append(f"q={q}")
|
||||
pub_ok = int(fwd.get("published_ok") or 0)
|
||||
if pub_ok > 0:
|
||||
bits.append(f"pub={pub_ok}")
|
||||
progress = " · ".join([p for p in (progress, *bits) if p])
|
||||
elif task == "alarms_current_ws_consumer" and isinstance(ws_conn, dict):
|
||||
if task == "alarms_current_ws_consumer" and isinstance(ws_conn, dict):
|
||||
state = str(ws_conn.get("state") or ws_conn.get("status") or "").strip()
|
||||
if state:
|
||||
progress = " · ".join([p for p in (progress, state) if p])
|
||||
|
|
|
|||
|
|
@ -8,17 +8,11 @@ RT_ALARMS_SYNC_IN_PROGRESS_SKIP = "rt:alarms_sync_in_progress_skip"
|
|||
RT_PULLING_INVENTORY = "rt:pulling_inventory"
|
||||
RT_PULLING_TOPOLOGY = "rt:pulling_topology"
|
||||
RT_UME_WS_DISABLED_NO_BASE_URL = "rt:ume_ws_disabled_no_base_url"
|
||||
RT_OCLAW_FWD_DISABLED = "rt:oclaw_fwd_disabled"
|
||||
RT_RESUMED_SYNC_SOON = "rt:resumed_sync_soon"
|
||||
RT_RESUMED_WSS_RECONNECT = "rt:resumed_wss_reconnect"
|
||||
RT_RESUMED_OCLAW_WSS_RECONNECT = "rt:resumed_oclaw_wss_reconnect"
|
||||
RT_RESUMED = "rt:resumed"
|
||||
RT_KEEPALIVE_FAILED = "rt:keepalive_failed"
|
||||
|
||||
|
||||
def ws_state_code(state: str) -> str:
|
||||
return f"ws:{str(state or '').strip()}"
|
||||
|
||||
|
||||
def fwd_state_code(state: str) -> str:
|
||||
return f"fwd:{str(state or '').strip()}"
|
||||
return f"ws:{str(state or '').strip()}"
|
||||
|
|
@ -34,10 +34,6 @@ from .models import (
|
|||
UmeKeyAlertRule,
|
||||
UmeSyncJob,
|
||||
)
|
||||
from .oclaw_alarm_forwarder import (
|
||||
forwarder_status,
|
||||
request_forwarder_reconnect,
|
||||
)
|
||||
from .ume_alarm_ws import (
|
||||
cancel_alarm_subscription_manual,
|
||||
clear_local_alarm_subscription_manual,
|
||||
|
|
|
|||
|
|
@ -34,10 +34,6 @@ from .models import (
|
|||
UmeKeyAlertRule,
|
||||
UmeSyncJob,
|
||||
)
|
||||
from .oclaw_alarm_forwarder import (
|
||||
forwarder_status,
|
||||
request_forwarder_reconnect,
|
||||
)
|
||||
from .ume_alarm_ws import (
|
||||
cancel_alarm_subscription_manual,
|
||||
clear_local_alarm_subscription_manual,
|
||||
|
|
|
|||
|
|
@ -35,10 +35,6 @@ from .models import (
|
|||
UmeSyncJob,
|
||||
)
|
||||
from .dsh_alarm_hub import hub_status as dsh_alarm_hub_status
|
||||
from .oclaw_alarm_forwarder import (
|
||||
forwarder_status,
|
||||
request_forwarder_reconnect,
|
||||
)
|
||||
from .ume_alarm_ws import (
|
||||
cancel_alarm_subscription_manual,
|
||||
clear_local_alarm_subscription_manual,
|
||||
|
|
@ -146,14 +142,12 @@ def ume_list_key_alert_rules(
|
|||
}
|
||||
for row in rows
|
||||
]
|
||||
fwd = forwarder_status()
|
||||
hub = dsh_alarm_hub_status()
|
||||
return {
|
||||
"items": items,
|
||||
"total": total,
|
||||
"page": page,
|
||||
"page_size": page_size,
|
||||
"forwarder": fwd,
|
||||
"dsh_alarm_hub": hub,
|
||||
}
|
||||
|
||||
|
|
@ -183,8 +177,6 @@ def ume_key_alert_monitor(
|
|||
"page_size": int(base.get("page_size") or page_size),
|
||||
"config": get_key_alert_monitor_config(db),
|
||||
"dsh_alarm_hub": base.get("dsh_alarm_hub") or dsh_alarm_hub_status(),
|
||||
# Legacy OClaw outbound bridge (optional; demoted in UI).
|
||||
"forwarder": base.get("forwarder") or forwarder_status(),
|
||||
}
|
||||
|
||||
|
||||
|
|
@ -345,6 +337,6 @@ def ume_list_notification_ids(
|
|||
for nid, cause in rows
|
||||
if str(nid or "").strip()
|
||||
]
|
||||
return {"items": items, "total": len(items), "forwarder": forwarder_status()}
|
||||
return {"items": items, "total": len(items)}
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -4,7 +4,7 @@ Device collectors (config_sync / LLDP / ne_collect / port_traffic) run via ``sta
|
|||
(API inline by default; set ``NETX_RUN_INLINE_SCHEDULERS=false`` and run
|
||||
``python -m netx_api.worker`` for a split process).
|
||||
API process also owns UME keepalive, alarm WSS, current-alarm/inventory sync loops,
|
||||
and oclaw forwarder via ``start_api_sideband_threads``.
|
||||
and DSH alarm hub via ``start_api_sideband_threads``.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
|
@ -29,7 +29,6 @@ from .ume_sync_topology import fail_stale_topology_running_jobs
|
|||
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_PULLING_TOPOLOGY,
|
||||
|
|
@ -37,12 +36,6 @@ from .runtime_task_messages import (
|
|||
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")
|
||||
|
|
@ -420,48 +413,6 @@ def start_api_sideband_threads() -> None:
|
|||
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]}",
|
||||
)
|
||||
|
||||
# Dock already has topology but Fabric/World empty (or last apply partial): heal without waiting 24h.
|
||||
try:
|
||||
|
|
|
|||
|
|
@ -18,11 +18,9 @@ from .timeutil import utcnow_naive
|
|||
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_RESUMED,
|
||||
RT_RESUMED_OCLAW_WSS_RECONNECT,
|
||||
RT_RESUMED_SYNC_SOON,
|
||||
RT_RESUMED_WSS_RECONNECT,
|
||||
RT_STARTUP_ALARM_SYNC_BEFORE_WS,
|
||||
|
|
@ -30,7 +28,6 @@ from .runtime_task_messages import (
|
|||
RT_UME_WS_DISABLED_NO_BASE_URL,
|
||||
RT_WSS_ACTIVE_SKIP_REST,
|
||||
)
|
||||
from .oclaw_alarm_forwarder import is_forwarder_enabled
|
||||
from .ume_alarm_ws import (
|
||||
begin_startup_alarm_sync_gate,
|
||||
complete_startup_alarm_sync_gate,
|
||||
|
|
@ -64,7 +61,6 @@ _UME_RUNTIME_TASKS: dict[str, dict[str, Any]] = {
|
|||
"token_keepalive": {"task": "token_keepalive", "status": "init", "last_run_at": None, "last_error": ""},
|
||||
"alarms_current_auto_sync": {"task": "alarms_current_auto_sync", "status": "init", "last_run_at": None, "last_error": ""},
|
||||
"alarms_current_ws_consumer": {"task": "alarms_current_ws_consumer", "status": "init", "last_run_at": None, "last_error": ""},
|
||||
"oclaw_alarm_forwarder": {"task": "oclaw_alarm_forwarder", "status": "init", "last_run_at": None, "last_error": ""},
|
||||
"inventory_auto_sync": {"task": "inventory_auto_sync", "status": "init", "last_run_at": None, "last_error": ""},
|
||||
"topology_auto_sync": {"task": "topology_auto_sync", "status": "init", "last_run_at": None, "last_error": ""},
|
||||
}
|
||||
|
|
@ -175,10 +171,6 @@ def _runtime_task_interval_fields(task_id: str) -> tuple[int | None, str]:
|
|||
if not bool(getattr(settings, "ume_alarm_ws_enabled", True)):
|
||||
return None, "disabled"
|
||||
return None, "realtime"
|
||||
if task_id == "oclaw_alarm_forwarder":
|
||||
if not is_forwarder_enabled():
|
||||
return None, "disabled"
|
||||
return None, "realtime"
|
||||
if task_id == "inventory_auto_sync":
|
||||
if not bool(getattr(settings, "ume_sync_inventory_auto_enabled", True)):
|
||||
return None, "disabled"
|
||||
|
|
|
|||
|
|
@ -11,10 +11,8 @@ from sqlalchemy.orm import Session
|
|||
|
||||
from .db import get_db
|
||||
from .models import UmeSyncJob
|
||||
from .oclaw_alarm_forwarder import request_forwarder_reconnect
|
||||
from .runtime_task_messages import (
|
||||
RT_RESUMED,
|
||||
RT_RESUMED_OCLAW_WSS_RECONNECT,
|
||||
RT_RESUMED_SYNC_SOON,
|
||||
RT_RESUMED_WSS_RECONNECT,
|
||||
)
|
||||
|
|
@ -236,8 +234,6 @@ def ume_runtime_task_pause(task: str) -> dict[str, Any]:
|
|||
_clear_force_resume_hints(tid)
|
||||
if tid == "alarms_current_ws_consumer":
|
||||
request_ws_reconnect()
|
||||
if tid == "oclaw_alarm_forwarder":
|
||||
request_forwarder_reconnect()
|
||||
_set_runtime_task(tid, status="paused", last_error="")
|
||||
return {"ok": True, "task": tid, "runtime_tasks": _list_runtime_tasks()}
|
||||
|
||||
|
|
@ -254,9 +250,6 @@ def ume_runtime_task_resume(task: str) -> dict[str, Any]:
|
|||
elif tid == "alarms_current_ws_consumer":
|
||||
request_ws_reconnect()
|
||||
resume_hint = RT_RESUMED_WSS_RECONNECT
|
||||
elif tid == "oclaw_alarm_forwarder":
|
||||
request_forwarder_reconnect()
|
||||
resume_hint = RT_RESUMED_OCLAW_WSS_RECONNECT
|
||||
else:
|
||||
resume_hint = RT_RESUMED
|
||||
_set_runtime_task(tid, status="running", last_error=resume_hint)
|
||||
|
|
|
|||
|
|
@ -34,10 +34,6 @@ from .models import (
|
|||
UmeKeyAlertRule,
|
||||
UmeSyncJob,
|
||||
)
|
||||
from .oclaw_alarm_forwarder import (
|
||||
forwarder_status,
|
||||
request_forwarder_reconnect,
|
||||
)
|
||||
from .ume_alarm_ws import (
|
||||
cancel_alarm_subscription_manual,
|
||||
clear_local_alarm_subscription_manual,
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue