feat(ume): OClaw forwarder runtime task and i18n status codes

Register oclaw_alarm_forwarder in background tasks with pause/resume, emit rt:/ws:/fwd: codes for last_error, and translate them in the UI for English and Chinese.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-06-24 10:29:10 +08:00
parent 5a00736f23
commit ba1a40f725
10 changed files with 380 additions and 37 deletions

View file

@ -9,7 +9,7 @@ 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 .oclaw_alarm_forwarder import enqueue_alarm_forward, is_forwarder_operational
from .ume_sync_service import (
_derive_ne_id_from_alarm,
_pick,
@ -69,7 +69,7 @@ def maybe_forward_key_alert(
alarm_key: str,
action: str,
) -> bool:
if not is_forwarder_enabled():
if not is_forwarder_operational():
return False
rule = match_key_alert_rule(db, norm=norm, action=action)
if rule is None:

View file

@ -75,7 +75,29 @@ from .key_alert_matcher import (
rule_storage_key,
serialize_rule_ne_types,
)
from .oclaw_alarm_forwarder import forwarder_status, shutdown_oclaw_alarm_forwarder, start_oclaw_alarm_forwarder
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,
@ -117,6 +139,7 @@ _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": ""},
}
_UME_WS_STOP_EVENT: threading.Event | None = None
@ -226,6 +249,10 @@ 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"
@ -343,7 +370,7 @@ def _run_startup_alarm_sync_before_ws() -> None:
"alarms_current_auto_sync",
status="running",
last_run_at=datetime.now(timezone.utc),
last_error="启动:正在同步当前告警(完成后连接 WSS)…",
last_error=RT_STARTUP_ALARM_SYNC_BEFORE_WS,
)
db = SessionLocal()
try:
@ -930,7 +957,7 @@ def on_startup() -> None:
client.renew_token()
_set_runtime_task("token_keepalive", status="running", last_run_at=datetime.now(timezone.utc), last_error="")
except Exception:
_set_runtime_task("token_keepalive", status="error", last_run_at=datetime.now(timezone.utc), last_error="keepalive_failed")
_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)
@ -983,7 +1010,7 @@ def on_startup() -> None:
_refresh_runtime_task_idle(
"alarms_current_auto_sync",
"alarms_current",
last_error="启动 REST 全量同步未完成,定时 REST 与 WSS 均待命",
last_error=RT_STARTUP_GATE_WAITING,
)
time.sleep(10)
continue
@ -994,7 +1021,7 @@ def on_startup() -> None:
_refresh_runtime_task_idle(
"alarms_current_auto_sync",
"alarms_current",
last_error="WSS 实时接收中,已跳过 REST 同步",
last_error=RT_WSS_ACTIVE_SKIP_REST,
)
time.sleep(max(30, min(alarms_interval_s, 300)))
continue
@ -1012,7 +1039,7 @@ def on_startup() -> None:
"alarms_current_auto_sync",
status="running",
last_run_at=datetime.now(timezone.utc),
last_error="正在拉取 UME 当前告警…",
last_error=RT_PULLING_ALARMS_CURRENT,
)
db = SessionLocal()
try:
@ -1032,7 +1059,7 @@ def on_startup() -> None:
_refresh_runtime_task_idle(
"alarms_current_auto_sync",
"alarms_current",
last_error="另一条当前告警 REST 同步进行中,已跳过",
last_error=RT_ALARMS_SYNC_IN_PROGRESS_SKIP,
)
time.sleep(30)
else:
@ -1091,7 +1118,7 @@ def on_startup() -> None:
"inventory_auto_sync",
status="running",
last_run_at=datetime.now(timezone.utc),
last_error="正在拉取 UME 网元清单…",
last_error=RT_PULLING_INVENTORY,
)
db = SessionLocal()
try:
@ -1151,7 +1178,7 @@ def on_startup() -> None:
)
_schedule_log.info("started thread %s alive=%s", t_ws.name, t_ws.is_alive())
else:
_set_runtime_task("alarms_current_ws_consumer", status="paused", last_error="未启用或未配置 UME_BASE_URL")
_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)
_set_runtime_task(
@ -1161,11 +1188,47 @@ def on_startup() -> None:
last_error=f"startup_thread_init_failed: {str(exc)[:180]}",
)
try:
def _fwd_on_status(msg: str) -> None:
paused = _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"
_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: _runtime_is_paused("oclaw_alarm_forwarder"),
on_status=_fwd_on_status,
)
if is_forwarder_enabled():
_set_runtime_task("oclaw_alarm_forwarder", status="running", last_error="")
else:
_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)
_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]}",
)
@app.on_event("shutdown")
@ -1695,6 +1758,8 @@ 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()}
@ -1707,12 +1772,15 @@ def ume_runtime_task_resume(task: str) -> dict[str, Any]:
_runtime_resume_task(tid)
if tid in ("alarms_current_auto_sync", "inventory_auto_sync"):
_request_force_sync_after_resume(tid)
resume_hint = "已恢复:将跳过本轮周期等待并尽快同步"
resume_hint = RT_RESUMED_SYNC_SOON
elif tid == "alarms_current_ws_consumer":
request_ws_reconnect()
resume_hint = "已恢复:将尽快重连 WSS"
resume_hint = RT_RESUMED_WSS_RECONNECT
elif tid == "oclaw_alarm_forwarder":
request_forwarder_reconnect()
resume_hint = RT_RESUMED_OCLAW_WSS_RECONNECT
else:
resume_hint = "已恢复"
resume_hint = RT_RESUMED
_set_runtime_task(tid, status="running", last_error=resume_hint)
return {"ok": True, "task": tid, "runtime_tasks": _list_runtime_tasks()}
@ -2287,6 +2355,16 @@ def integrations_status(db: Session = Depends(get_db)) -> dict:
"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",

View file

@ -5,6 +5,7 @@ import logging
import queue
import threading
import time
from collections.abc import Callable
from datetime import datetime, timezone
from typing import Any
@ -20,6 +21,8 @@ _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,
@ -46,8 +49,50 @@ def is_forwarder_enabled() -> bool:
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_enabled():
if not is_forwarder_operational():
return False
try:
_OUTBOUND_Q.put_nowait(dict(payload))
@ -112,12 +157,18 @@ def _run_loop() -> None:
global _WS
backoff_s = 2.0
while not _STOP_EVENT.is_set():
if not is_forwarder_enabled():
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):
@ -127,7 +178,14 @@ def _run_loop() -> None:
_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:
@ -174,7 +232,9 @@ def _run_loop() -> None:
_OUTBOUND_Q.task_done()
except Exception as exc:
_CONNECTED.clear()
_log.warning("oclaw netx-bridge disconnected: %s", str(exc)[:200])
err = str(exc)[:200]
_log.warning("oclaw netx-bridge disconnected: %s", err)
_notify_status(err)
time.sleep(backoff_s)
backoff_s = min(backoff_s * 1.5, 60.0)
finally:
@ -192,9 +252,6 @@ 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()
@ -208,8 +265,12 @@ def shutdown_oclaw_alarm_forwarder() -> None:
def forwarder_status() -> dict[str, Any]:
with _STATS_LOCK:
stats = dict(_STATS)
paused = _forwarder_paused()
enabled = is_forwarder_enabled()
return {
"enabled": is_forwarder_enabled(),
"enabled": enabled,
"operational": bool(enabled and not paused),
"paused": paused,
"connected": _CONNECTED.is_set(),
"queue_size": int(_OUTBOUND_Q.qsize()),
"url": _bridge_url(),

View file

@ -0,0 +1,23 @@
"""Stable runtime-task status codes for frontend i18n (rt:* / ws:* / fwd:*)."""
RT_STARTUP_ALARM_SYNC_BEFORE_WS = "rt:startup_alarm_sync_before_ws"
RT_STARTUP_GATE_WAITING = "rt:startup_gate_waiting"
RT_WSS_ACTIVE_SKIP_REST = "rt:wss_active_skip_rest"
RT_PULLING_ALARMS_CURRENT = "rt:pulling_alarms_current"
RT_ALARMS_SYNC_IN_PROGRESS_SKIP = "rt:alarms_sync_in_progress_skip"
RT_PULLING_INVENTORY = "rt:pulling_inventory"
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()}"

View file

@ -249,7 +249,7 @@ def get_ws_connection_status() -> dict[str, Any]:
detail = str(_ws_connection_detail or "")
return {
"state": state,
"label": _WS_CONNECTION_LABELS.get(state, state),
"label": f"ws:{state}",
"detail": detail,
}
@ -284,7 +284,7 @@ def _notify_ws_connection(
dedup=False,
)
if on_status is not None:
on_status(label)
on_status(f"ws:{state}")
def _parse_ws_message(raw: str) -> dict[str, Any] | None:
@ -804,7 +804,7 @@ def start_ume_alarm_ws_consumer(
if not str(client.base_url or "").strip():
append_ws_log("consumer disabled: no UME base URL", level="warning")
if on_status is not None:
on_status("disabled:no_base_url")
on_status("ws:disabled_no_base_url")
return
append_ws_log("ws consumer thread started")
global _active_client