Remove unused NetX alarm WebSocket bridge.

Alarm delivery moved to NetX DSH hub; drop /ws/netx-bridge and its tests.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-09-15 23:10:32 +08:00
parent ad0cc98b38
commit 050b99426f
4 changed files with 1 additions and 230 deletions

View file

@ -317,7 +317,7 @@ OCLAW_WS_RATE_LIMIT_USER_PER_WINDOW=360
OCLAW_WS_SEND_QUEUE_MAX_MESSAGES=256
OCLAW_WS_SEND_QUEUE_MAX_BYTES=
OCLAW_WS_EVENT_REPLAY_MAX=256
# netx -> oclaw AP 分析接口共享 token(可选;设置后可用 Bearer 直接访问 /admin/api/ops-ai/analyze-sync,亦用于 /ws/netx-bridge 告警 WSS)
# netx -> oclaw AP 分析接口共享 token(可选;设置后可用 Bearer 直接访问 /admin/api/ops-ai/analyze-sync)
OCLAW_OPS_AI_SHARED_TOKEN=
# -----------------------------------------------------------------------------

View file

@ -20,7 +20,6 @@ from fastapi.staticfiles import StaticFiles
from runtime.application.gateway import process_inbound_payload_usecase
from interfaces.gateway.http_adapter import dispatch_gateway_http_method
from interfaces.ws import ws_gateway_loop
from interfaces.ws.netx_bridge import netx_bridge_loop
from interfaces.ws.common import MAX_PAYLOAD_BYTES
from interfaces.admin.routes import admin_static_dir, build_admin_router
from runtime.agents.agent_scope import list_agent_ids, resolve_agent_workspace_dir, resolve_default_agent_id
@ -363,10 +362,6 @@ def create_app() -> FastAPI:
async def ws_endpoint(ws: WebSocket) -> None:
await ws_gateway_loop(ws)
@app.websocket("/ws/netx-bridge")
async def netx_bridge_endpoint(ws: WebSocket) -> None:
await netx_bridge_loop(ws)
return app

View file

@ -1,141 +0,0 @@
from __future__ import annotations
import hmac
import json
import os
from typing import Any
from fastapi import WebSocket, WebSocketDisconnect
from svc.persistence.assistant_store import get_assistant_store
def _bridge_token_expected() -> str:
return str(os.getenv("OCLAW_OPS_AI_SHARED_TOKEN") or "").strip()
def _verify_bridge_token(token: str) -> bool:
expected = _bridge_token_expected()
if not expected:
return False
got = str(token or "").strip()
if not got:
return False
return hmac.compare_digest(expected, got)
def _default_whatsapp_account_id() -> str:
return str(os.getenv("AIA_WHATSAPP_ACCOUNT_ID") or "wa-default").strip()
def _default_tenant_id() -> str:
return str(os.getenv("OCLAW_DEFAULT_TENANT_ID") or "default").strip()
def _format_alarm_text(payload: dict[str, Any]) -> str:
action = str(payload.get("action") or "").strip().lower()
label = str(payload.get("rule_label") or payload.get("native_probable_cause") or "Key alarm").strip()
ne = payload.get("ne") if isinstance(payload.get("ne"), dict) else {}
host = str(ne.get("host_name") or payload.get("host_name") or "").strip()
ip = str(ne.get("ip_address") or "").strip()
ne_name = str(ne.get("ne_name") or ne.get("user_label") or "").strip()
device = host or ne_name or str(payload.get("ne_id") or "").strip()
if ip:
device = f"{device} ({ip})" if device else ip
action_label = {
"inserted": "Alarm Raised",
"updated": "Alarm Updated",
"deleted": "Alarm Cleared",
}.get(action, action or "Alarm")
lines = [
f"[UME {action_label}] {label}",
f"Device: {device or '-'}",
f"Object: {str(payload.get('object_name') or '-').strip()}",
f"Severity: {str(payload.get('perceived_severity') or '-').strip()}",
f"Cause: {str(payload.get('native_probable_cause') or '-').strip()}",
f"Time: {str(payload.get('time_created') or '-').strip()}",
f"notificationId: {str(payload.get('notification_id') or '-').strip()}",
]
return "\n".join(lines)
def handle_netx_alarm_event(payload: dict[str, Any]) -> dict[str, Any]:
store = get_assistant_store()
tenant_id = _default_tenant_id()
account_id = _default_whatsapp_account_id()
binding = store.get_whatsapp_alert_binding(tenant_id=tenant_id, account_id=account_id)
if not binding or not bool(binding.get("enabled")):
return {"ok": False, "error": "whatsapp_alert_binding_missing"}
group_jid = str(binding.get("group_jid") or "").strip()
if not group_jid:
return {"ok": False, "error": "group_jid_missing"}
text = _format_alarm_text(payload)
msg_id = store.enqueue_channel_outbound_message(
channel="whatsapp",
chat_id=group_jid,
text=text,
tenant_id=tenant_id,
account_id=account_id,
source="netx.alarm",
)
return {"ok": True, "outbound_id": msg_id, "chat_id": group_jid}
async def netx_bridge_loop(ws: WebSocket) -> None:
await ws.accept()
authed = False
try:
while True:
raw = await ws.receive_text()
try:
msg = json.loads(raw)
except Exception:
await ws.send_json({"type": "error", "error": "invalid_json"})
continue
if not isinstance(msg, dict):
await ws.send_json({"type": "error", "error": "invalid_message"})
continue
mtype = str(msg.get("type") or "").strip().lower()
if not authed:
if mtype != "auth":
await ws.send_json({"type": "auth-fail", "error": "auth_required"})
await ws.close(code=4401)
return
token = str(msg.get("token") or "").strip()
if not _verify_bridge_token(token):
await ws.send_json({"type": "auth-fail", "error": "invalid_token"})
await ws.close(code=4401)
return
authed = True
await ws.send_json({"type": "auth-ok"})
continue
if mtype == "ping":
await ws.send_json({"type": "pong"})
continue
if mtype == "event" and str(msg.get("event") or "").strip() == "netx.alarm":
payload = msg.get("payload") if isinstance(msg.get("payload"), dict) else {}
alarm_key = str(payload.get("alarm_key") or "").strip()
try:
out = handle_netx_alarm_event(payload)
await ws.send_json(
{
"type": "ack",
"alarm_key": alarm_key,
"ok": bool(out.get("ok")),
"error": str(out.get("error") or ""),
"outbound_id": str(out.get("outbound_id") or ""),
}
)
except Exception as exc:
await ws.send_json(
{
"type": "ack",
"alarm_key": alarm_key,
"ok": False,
"error": f"{type(exc).__name__}: {exc}",
}
)
continue
await ws.send_json({"type": "error", "error": f"unknown_type:{mtype}"})
except WebSocketDisconnect:
return

View file

@ -1,83 +0,0 @@
from __future__ import annotations
import asyncio
import os
import unittest
from unittest.mock import patch
from fastapi.testclient import TestClient
from interfaces.http.fastapi_app import create_app
from interfaces.ws.netx_bridge import _format_alarm_text
from svc.persistence.assistant_store import reset_assistant_store_singleton
class NetxBridgeTests(unittest.TestCase):
def setUp(self) -> None:
os.environ["OCLAW_OPS_AI_SHARED_TOKEN"] = "test-bridge-token"
os.environ.pop("OCLAW_NETX_BRIDGE_TOKEN", None)
os.environ["AIA_WHATSAPP_ACCOUNT_ID"] = "wa-default"
reset_assistant_store_singleton()
self.client = TestClient(create_app())
def tearDown(self) -> None:
reset_assistant_store_singleton()
def test_format_alarm_text_english(self) -> None:
text = _format_alarm_text(
{
"action": "inserted",
"rule_label": "Fan",
"object_name": "ME{abc},FAN={/module=0}",
"perceived_severity": "major",
"native_probable_cause": "Fan The fan speed level is abnormally high",
"time_created": "2026-06-22T19:29:41.086+07:00",
"notification_id": "1680996323029",
"ne": {"host_name": "LPG-BKM-AN1-ZM3SP", "ip_address": "114.0.24.178"},
}
)
self.assertIn("[UME Alarm Raised] Fan", text)
self.assertIn("Device: LPG-BKM-AN1-ZM3SP (114.0.24.178)", text)
self.assertIn("Severity: major", text)
self.assertNotIn("设备", text)
self.assertNotIn("NetX", text)
def test_netx_bridge_auth_and_alarm_ack(self) -> None:
with self.client.websocket_connect("/ws/netx-bridge") as ws:
ws.send_json({"type": "auth", "token": "bad"})
msg = ws.receive_json()
self.assertEqual(msg.get("type"), "auth-fail")
with self.client.websocket_connect("/ws/netx-bridge") as ws:
ws.send_json({"type": "auth", "token": "test-bridge-token"})
msg = ws.receive_json()
self.assertEqual(msg.get("type"), "auth-ok")
ws.send_json(
{
"type": "event",
"event": "netx.alarm",
"payload": {
"action": "inserted",
"alarm_key": "AK-1",
"notification_id": "NID-1",
"native_probable_cause": "link down",
"perceived_severity": "critical",
"ne": {"host_name": "host-a", "ip_address": "10.0.0.1"},
},
}
)
ack = ws.receive_json()
self.assertEqual(ack.get("type"), "ack")
self.assertEqual(ack.get("alarm_key"), "AK-1")
self.assertFalse(ack.get("ok"))
def test_whatsapp_outbound_pending_empty(self) -> None:
resp = self.client.get("/whatsapp/outbound/pending?account_id=wa-default")
self.assertEqual(resp.status_code, 200)
body = resp.json()
self.assertTrue(body.get("ok"))
self.assertEqual(body.get("items"), [])
if __name__ == "__main__":
unittest.main()