feat(ume): key alarm forward to OClaw via WSS and AI monitor UI

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-06-22 11:17:31 +08:00
parent 9a39ddfcc7
commit 1f1b42e8bf
18 changed files with 1079 additions and 32 deletions

View file

@ -20,6 +20,9 @@ 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"
oclaw_alarm_ws_token: str = ""
# UME RESTCONF integration
ume_base_url: str = ""
ume_username: str = ""

View file

@ -0,0 +1,156 @@
from __future__ import annotations
import logging
from datetime import datetime, timezone
from typing import Any
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from .key_alert_matcher import match_key_alert_rule
from .models import UmeInventoryNE, UmeKeyAlertForwardLog
from .oclaw_alarm_forwarder import enqueue_alarm_forward, is_forwarder_enabled
from .ume_sync_service import (
_derive_ne_id_from_alarm,
_pick,
_s,
notification_id_from_norm,
)
_log = logging.getLogger("netx.key_alert.forward")
def _utc_now_naive() -> datetime:
return datetime.now(timezone.utc).replace(tzinfo=None)
def _ne_payload(db: Session, ne_id: str) -> dict[str, str]:
row = db.get(UmeInventoryNE, ne_id) if ne_id else None
if row is None:
return {}
return {
"ne_id": str(row.ne_id or ""),
"ne_name": str(row.ne_name or ""),
"user_label": str(row.user_label or ""),
"host_name": str(row.host_name or ""),
"ip_address": str(row.ip_address or ""),
"ne_type": str(row.ne_type or ""),
"device_level": str(row.device_level or ""),
}
def _build_forward_payload(
*,
norm: dict[str, Any],
alarm_key: str,
action: str,
rule_label: str,
) -> dict[str, Any]:
ne_id = _s(_derive_ne_id_from_alarm(norm))
return {
"action": str(action or ""),
"alarm_key": str(alarm_key or ""),
"notification_id": notification_id_from_norm(norm),
"rule_label": str(rule_label or ""),
"event_type": _s(_pick(norm, "eventType", "event-type")),
"native_probable_cause": _s(_pick(norm, "nativeProbableCause", "native-probable-cause")),
"perceived_severity": _s(_pick(norm, "perceivedSeverity", "perceived-severity")),
"is_cleared": _s(_pick(norm, "isCleared", "is-cleared")),
"time_created": _s(_pick(norm, "timeCreated", "time-created")),
"object_name": _s(_pick(norm, "objectName", "object-name")),
"ne_id": ne_id,
}
def maybe_forward_key_alert(
db: Session,
*,
norm: dict[str, Any],
alarm_key: str,
action: str,
) -> bool:
if not is_forwarder_enabled():
return False
rule = match_key_alert_rule(db, norm=norm, action=action)
if rule is None:
return False
act = str(action or "").strip().lower()
existing = (
db.query(UmeKeyAlertForwardLog)
.filter(
UmeKeyAlertForwardLog.alarm_key == str(alarm_key or ""),
UmeKeyAlertForwardLog.action == act,
)
.first()
)
if existing is not None and int(existing.oclaw_ok or 0) == 1:
return False
ne_id = _s(_derive_ne_id_from_alarm(norm))
payload = _build_forward_payload(
norm=norm,
alarm_key=alarm_key,
action=action,
rule_label=str(rule.label or ""),
)
payload["ne"] = _ne_payload(db, ne_id)
queued = enqueue_alarm_forward(payload)
if not queued:
return False
row = existing
if row is None:
row = UmeKeyAlertForwardLog(
alarm_key=str(alarm_key or ""),
action=act,
notification_id=notification_id_from_norm(norm),
forwarded_at=_utc_now_naive(),
oclaw_ok=0,
error="queued",
)
db.add(row)
else:
row.notification_id = notification_id_from_norm(norm)
row.forwarded_at = _utc_now_naive()
row.oclaw_ok = 0
row.error = "queued"
try:
db.commit()
except IntegrityError:
db.rollback()
return True
def record_forward_result(*, alarm_key: str, action: str, ok: bool, error: 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,
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:
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()

View file

@ -0,0 +1,65 @@
from __future__ import annotations
import threading
import time
from typing import Any
from sqlalchemy.orm import Session
from .models import UmeKeyAlertRule
from .ume_sync_service import _is_alarm_cleared, notification_id_from_norm
_RULE_CACHE_LOCK = threading.Lock()
_RULE_CACHE: dict[str, UmeKeyAlertRule] = {}
_RULE_CACHE_LOADED_AT = 0.0
_RULE_CACHE_TTL_S = 30.0
def invalidate_key_alert_rule_cache() -> None:
global _RULE_CACHE_LOADED_AT
with _RULE_CACHE_LOCK:
_RULE_CACHE.clear()
_RULE_CACHE_LOADED_AT = 0.0
def _load_enabled_rules(db: Session) -> dict[str, UmeKeyAlertRule]:
global _RULE_CACHE_LOADED_AT
now = time.time()
with _RULE_CACHE_LOCK:
if _RULE_CACHE and (now - _RULE_CACHE_LOADED_AT) < _RULE_CACHE_TTL_S:
return dict(_RULE_CACHE)
rows = (
db.query(UmeKeyAlertRule)
.filter(UmeKeyAlertRule.enabled == 1)
.all()
)
loaded = {str(row.notification_id or "").strip(): row for row in rows if str(row.notification_id or "").strip()}
with _RULE_CACHE_LOCK:
_RULE_CACHE.clear()
_RULE_CACHE.update(loaded)
_RULE_CACHE_LOADED_AT = now
return dict(loaded)
def match_key_alert_rule(
db: Session,
*,
norm: dict[str, Any],
action: str,
) -> UmeKeyAlertRule | None:
notification_id = notification_id_from_norm(norm)
if not notification_id:
return None
rules = _load_enabled_rules(db)
rule = rules.get(notification_id)
if rule is None:
return None
act = str(action or "").strip().lower()
if act in {"inserted", "updated"}:
if _is_alarm_cleared(norm):
return None
return rule
if act == "deleted":
if int(getattr(rule, "forward_on_clear", 0) or 0) == 1:
return rule
return None

View file

@ -35,6 +35,8 @@ from .models import (
UmeAlarmCurrent,
UmeAlarmHistory,
UmeInventoryNE,
UmeKeyAlertRule,
UmeKeyAlertForwardLog,
UmeSyncJob,
)
from .models import ImportJob
@ -58,6 +60,8 @@ from .ume_alarm_ws import (
start_ume_alarm_ws_consumer,
)
from .ume_sync_service import sync_alarms_current, sync_alarms_history_full, sync_inventory_full
from .key_alert_matcher import invalidate_key_alert_rule_cache
from .oclaw_alarm_forwarder import forwarder_status, shutdown_oclaw_alarm_forwarder, start_oclaw_alarm_forwarder
from .ume_token_store import (
clear_shared_token,
load_shared_token,
@ -755,6 +759,18 @@ def on_startup() -> None:
conn.exec_driver_sql("ALTER TABLE ume_alarms_current ALTER COLUMN root_cause_alarm_indication TYPE TEXT")
conn.exec_driver_sql("ALTER TABLE ume_alarms_current ADD COLUMN IF NOT EXISTS host_name VARCHAR(256) DEFAULT ''")
conn.exec_driver_sql("ALTER TABLE ume_alarms_history ADD COLUMN IF NOT EXISTS host_name VARCHAR(256) DEFAULT ''")
conn.exec_driver_sql(
"ALTER TABLE ume_alarms_current ADD COLUMN IF NOT EXISTS notification_id VARCHAR(128) DEFAULT ''"
)
conn.exec_driver_sql(
"ALTER TABLE ume_alarms_history ADD COLUMN IF NOT EXISTS notification_id VARCHAR(128) DEFAULT ''"
)
conn.exec_driver_sql(
"CREATE INDEX IF NOT EXISTS ix_ume_alarms_current_notification_id ON ume_alarms_current (notification_id)"
)
conn.exec_driver_sql(
"CREATE INDEX IF NOT EXISTS ix_ume_alarms_history_notification_id ON ume_alarms_history (notification_id)"
)
conn.exec_driver_sql(
"CREATE INDEX IF NOT EXISTS ix_ume_alarms_current_host_name ON ume_alarms_current (host_name)"
)
@ -1068,11 +1084,18 @@ def on_startup() -> None:
last_run_at=datetime.now(timezone.utc),
last_error=f"startup_thread_init_failed: {str(exc)[:180]}",
)
try:
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)
@app.on_event("shutdown")
def on_shutdown() -> None:
global _UME_WS_STOP_EVENT
shutdown_oclaw_alarm_forwarder()
if _UME_WS_STOP_EVENT is not None:
_UME_WS_STOP_EVENT.set()
shutdown_ws_consumer()
@ -1177,6 +1200,126 @@ def ume_alarm_subscription_clear_local(db: Session = Depends(get_db)) -> dict[st
raise HTTPException(status_code=502, detail=msg) from exc
@app.get("/v1/ume/key-alert-rules")
def ume_list_key_alert_rules(db: Session = Depends(get_db)) -> dict[str, Any]:
from sqlalchemy import func
rows = db.query(UmeKeyAlertRule).order_by(UmeKeyAlertRule.notification_id.asc()).all()
stat_rows = (
db.query(
UmeKeyAlertForwardLog.notification_id,
func.count(UmeKeyAlertForwardLog.id).label("attempts"),
func.sum(UmeKeyAlertForwardLog.oclaw_ok).label("published_ok"),
func.max(UmeKeyAlertForwardLog.forwarded_at).label("last_forwarded_at"),
)
.group_by(UmeKeyAlertForwardLog.notification_id)
.all()
)
stat_map = {
str(nid or ""): {
"attempts": int(attempts or 0),
"published_ok": int(published_ok or 0),
"last_forwarded_at": (_ensure_utc(last_at) or datetime.now(timezone.utc)).isoformat() if last_at else "",
}
for nid, attempts, published_ok, last_at in stat_rows
if str(nid or "").strip()
}
items = [
{
"notification_id": str(row.notification_id or ""),
"enabled": bool(int(row.enabled or 0)),
"forward_on_clear": bool(int(row.forward_on_clear or 0)),
"label": str(row.label or ""),
"created_at": (_ensure_utc(row.created_at) or datetime.now(timezone.utc)).isoformat(),
"updated_at": (_ensure_utc(row.updated_at) or datetime.now(timezone.utc)).isoformat(),
"forward_stats": stat_map.get(str(row.notification_id or ""), {
"attempts": 0,
"published_ok": 0,
"last_forwarded_at": "",
}),
}
for row in rows
]
fwd = forwarder_status()
return {"items": items, "total": len(items), "forwarder": fwd}
@app.get("/v1/ume/key-alert-monitor")
def ume_key_alert_monitor(db: Session = Depends(get_db)) -> dict[str, Any]:
base = ume_list_key_alert_rules(db)
return {
"ok": True,
"rules": base.get("items") or [],
"forwarder": base.get("forwarder") or forwarder_status(),
}
@app.post("/v1/ume/key-alert-rules")
def ume_upsert_key_alert_rule(payload: dict[str, Any], db: Session = Depends(get_db)) -> dict[str, Any]:
notification_id = str(payload.get("notification_id") or "").strip()
if not notification_id:
raise HTTPException(status_code=400, detail="notification_id_required")
now = datetime.now(timezone.utc).replace(tzinfo=None)
row = db.get(UmeKeyAlertRule, notification_id)
if row is None:
row = UmeKeyAlertRule(notification_id=notification_id, created_at=now, updated_at=now)
db.add(row)
row.enabled = 1 if bool(payload.get("enabled", True)) else 0
row.forward_on_clear = 1 if bool(payload.get("forward_on_clear", False)) else 0
row.label = str(payload.get("label") or "").strip()
row.updated_at = now
db.commit()
invalidate_key_alert_rule_cache()
return {
"ok": True,
"notification_id": notification_id,
"enabled": bool(row.enabled),
"forward_on_clear": bool(row.forward_on_clear),
"label": row.label,
}
@app.delete("/v1/ume/key-alert-rules/{notification_id}")
def ume_delete_key_alert_rule(notification_id: str, db: Session = Depends(get_db)) -> dict[str, Any]:
nid = str(notification_id or "").strip()
row = db.get(UmeKeyAlertRule, nid)
if row is None:
raise HTTPException(status_code=404, detail="rule_not_found")
db.delete(row)
db.commit()
invalidate_key_alert_rule_cache()
return {"ok": True, "deleted": nid}
@app.get("/v1/ume/notification-ids")
def ume_list_notification_ids(
limit: int = Query(default=200, ge=1, le=2000),
db: Session = Depends(get_db),
) -> dict[str, Any]:
from sqlalchemy import func
rows = (
db.query(
UmeAlarmCurrent.notification_id,
func.max(UmeAlarmCurrent.native_probable_cause).label("cause_sample"),
)
.filter(UmeAlarmCurrent.notification_id != "")
.group_by(UmeAlarmCurrent.notification_id)
.order_by(UmeAlarmCurrent.notification_id.asc())
.limit(limit)
.all()
)
items = [
{
"notification_id": str(nid or ""),
"native_probable_cause_sample": str(cause or ""),
}
for nid, cause in rows
if str(nid or "").strip()
]
return {"items": items, "total": len(items), "forwarder": forwarder_status()}
@app.post("/v1/ume/sync")
def ume_sync(payload: dict[str, Any] | None = None, db: Session = Depends(get_db)) -> dict[str, Any]:
body = payload or {}
@ -1464,6 +1607,7 @@ def ume_list_alarms(
UmeAlarmCurrent.alarm_key.contains(kw)
| UmeAlarmCurrent.object_name.contains(kw)
| UmeAlarmCurrent.native_probable_cause.contains(kw)
| UmeAlarmCurrent.notification_id.contains(kw)
| UmeAlarmCurrent.host_name.contains(kw)
| UmeInventoryNE.ne_name.contains(kw)
| UmeInventoryNE.user_label.contains(kw)
@ -1492,6 +1636,7 @@ def ume_list_alarms(
"object_name": str(alarm.object_name or ""),
"event_type": str(alarm.event_type or ""),
"native_probable_cause": str(alarm.native_probable_cause or ""),
"notification_id": str(alarm.notification_id or ""),
"perceived_severity": str(alarm.perceived_severity or ""),
"is_cleared": str(alarm.is_cleared or ""),
"time_created": str(alarm.time_created or ""),
@ -1884,32 +2029,41 @@ def integrations_status(db: Session = Depends(get_db)) -> dict:
db_status = {"status": "down", "error": str(exc)[:240]}
oclaw_status: dict = {"status": "unknown"}
try:
t0 = time.monotonic()
data = health_with_oclaw()
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("connected")):
oclaw_status = {
"status": "up",
"latency_ms": int((time.monotonic() - t0) * 1000),
"http_status": int(data.get("status_code") or 200),
"detail": data.get("data") or {},
"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,
}
except Exception as exc:
msg = str(exc)
http_status = None
kind = "unknown"
if " 401 " in msg or "401" in msg:
kind = "auth"
http_status = 401
elif " 404 " in msg or "404" in msg:
kind = "not_found"
http_status = 404
elif "timeout" in msg.lower():
kind = "timeout"
elif "connect" in msg.lower():
kind = "connect"
else:
kind = "other"
oclaw_status = {"status": "down", "error_kind": kind, "http_status": http_status, "error": msg[:240]}
return {"netx_api": netx_api, "db": db_status, "oclaw_bridge": oclaw_status}

View file

@ -176,6 +176,7 @@ class UmeAlarmCurrent(Base):
is_cleared: Mapped[str] = mapped_column(Text, default="", index=True)
time_created: Mapped[str] = mapped_column(Text, default="", index=True)
root_cause_alarm_indication: Mapped[str] = mapped_column(Text, default="")
notification_id: Mapped[str] = mapped_column(String(128), default="", index=True)
first_seen_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
last_seen_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True)
raw_json: Mapped[str] = mapped_column(Text, default="{}")
@ -194,11 +195,37 @@ class UmeAlarmHistory(Base):
is_cleared: Mapped[str] = mapped_column(Text, default="", index=True)
time_created: Mapped[str] = mapped_column(Text, default="", index=True)
root_cause_alarm_indication: Mapped[str] = mapped_column(Text, default="")
notification_id: Mapped[str] = mapped_column(String(128), default="", index=True)
first_seen_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
last_seen_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True)
raw_json: Mapped[str] = mapped_column(Text, default="{}")
class UmeKeyAlertRule(Base):
"""Key alert rule matched by UME notificationId."""
__tablename__ = "ume_key_alert_rule"
notification_id: Mapped[str] = mapped_column(String(128), primary_key=True)
enabled: Mapped[int] = mapped_column(Integer, default=1)
forward_on_clear: Mapped[int] = mapped_column(Integer, default=0)
label: Mapped[str] = mapped_column(String(256), default="")
created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True)
class UmeKeyAlertForwardLog(Base):
__tablename__ = "ume_key_alert_forward_log"
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
alarm_key: Mapped[str] = mapped_column(Text, index=True)
action: Mapped[str] = mapped_column(String(32), default="", index=True)
notification_id: Mapped[str] = mapped_column(String(128), default="", index=True)
forwarded_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True)
oclaw_ok: Mapped[int] = mapped_column(Integer, default=0)
error: Mapped[str] = mapped_column(String(512), default="")
class UmeAlarmSubscription(Base):
"""Persisted UME ALARM notification subscription (manual establish/cancel)."""

View file

@ -0,0 +1,218 @@
from __future__ import annotations
import json
import logging
import queue
import threading
import time
from datetime import datetime, timezone
from typing import Any
import websocket
from .config import settings
_log = logging.getLogger("netx.oclaw.alarm_forwarder")
_OUTBOUND_Q: "queue.Queue[dict[str, Any]]" = queue.Queue(maxsize=5000)
_STOP_EVENT = threading.Event()
_THREAD: threading.Thread | None = None
_CONN_LOCK = threading.Lock()
_WS: Any | None = None
_CONNECTED = threading.Event()
_STATS_LOCK = threading.Lock()
_STATS: dict[str, int] = {
"published_ok": 0,
"published_fail": 0,
"queued": 0,
}
def _utc_now_iso() -> str:
return datetime.now(timezone.utc).isoformat()
def _bridge_token() -> str:
return str(getattr(settings, "oclaw_alarm_ws_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 enqueue_alarm_forward(payload: dict[str, Any]) -> bool:
if not is_forwarder_enabled():
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:
_log.warning("oclaw alarm forward queue full; dropping alarm_key=%s", payload.get("alarm_key"))
return False
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():
time.sleep(2.0)
continue
url = _bridge_url()
ws = None
try:
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])
while not _STOP_EVENT.is_set():
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,
)
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],
)
except Exception as exc:
_log.warning("oclaw forward failed alarm_key=%s err=%s", payload.get("alarm_key"), str(exc)[:120])
try:
_OUTBOUND_Q.put_nowait(payload)
except queue.Full:
pass
raise
finally:
_OUTBOUND_Q.task_done()
except Exception as exc:
_CONNECTED.clear()
_log.warning("oclaw netx-bridge disconnected: %s", str(exc)[:200])
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
if not is_forwarder_enabled():
_log.info("oclaw alarm forwarder disabled")
return None
_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)
return {
"enabled": is_forwarder_enabled(),
"connected": _CONNECTED.is_set(),
"queue_size": int(_OUTBOUND_Q.qsize()),
"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)),
}

View file

@ -582,6 +582,16 @@ def _run_ws_session(
label = _WS_ALARM_ACTION_LABEL.get(action, action)
status_msg = f"{label} key={alarm_key}"
append_ws_log(status_msg, subscription_id=subscription_id, dedup=False)
try:
from .key_alert_forward import maybe_forward_key_alert
maybe_forward_key_alert(db, norm=norm, alarm_key=alarm_key, action=action)
except Exception as fwd_exc:
append_ws_log(
f"key alert forward failed: {str(fwd_exc)[:160]}",
level="error",
subscription_id=subscription_id,
)
except Exception as exc:
db.rollback()
append_ws_log(f"alarm apply failed: {str(exc)[:200]}", level="error", subscription_id=subscription_id)

View file

@ -237,6 +237,10 @@ def _is_alarm_cleared_tombstone(alarm_key: str) -> bool:
return True
def notification_id_from_norm(norm: dict[str, Any]) -> str:
return _s(_pick(norm, "notificationId", "notification-id"))
def _alarm_row_from_norm(key: str, norm: dict[str, Any], *, touch_ts: datetime, first_seen_at: datetime) -> dict[str, Any]:
return {
"alarm_key": key,
@ -251,6 +255,7 @@ def _alarm_row_from_norm(key: str, norm: dict[str, Any], *, touch_ts: datetime,
"root_cause_alarm_indication": _s(
_pick(norm, "rootCauseAlarmIndication", "root-cause-alarm-indication")
),
"notification_id": notification_id_from_norm(norm),
"first_seen_at": first_seen_at,
"last_seen_at": touch_ts,
"raw_json": json.dumps(norm, ensure_ascii=False, default=str),
@ -268,6 +273,7 @@ def _apply_row_to_model(db: Session, existing: UmeAlarmCurrent, norm: dict[str,
existing.root_cause_alarm_indication = _s(
_pick(norm, "rootCauseAlarmIndication", "root-cause-alarm-indication")
)
existing.notification_id = notification_id_from_norm(norm)
prev_seen = existing.last_seen_at
if prev_seen is None or touch_ts >= prev_seen:
existing.last_seen_at = touch_ts
@ -300,6 +306,7 @@ def _upsert_alarm_current(db: Session, key: str, norm: dict[str, Any], *, touch_
"is_cleared": excluded.is_cleared,
"time_created": excluded.time_created,
"root_cause_alarm_indication": excluded.root_cause_alarm_indication,
"notification_id": excluded.notification_id,
"last_seen_at": func.greatest(UmeAlarmCurrent.last_seen_at, excluded.last_seen_at),
"raw_json": excluded.raw_json,
},
@ -681,6 +688,7 @@ def _sync_alarms_common(
existing.root_cause_alarm_indication = _s(
_pick(alarm, "rootCauseAlarmIndication", "root-cause-alarm-indication")
)
existing.notification_id = notification_id_from_norm(alarm)
existing.last_seen_at = touch_ts
existing.raw_json = json.dumps(alarm, ensure_ascii=False, default=str)