diff --git a/_local/system.env.example b/_local/system.env.example index 28c5fd51..79e8c0a6 100644 --- a/_local/system.env.example +++ b/_local/system.env.example @@ -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= # ----------------------------------------------------------------------------- diff --git a/interfaces/http/fastapi_app.py b/interfaces/http/fastapi_app.py index 615daf29..d5fac4c8 100644 --- a/interfaces/http/fastapi_app.py +++ b/interfaces/http/fastapi_app.py @@ -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 diff --git a/interfaces/ws/netx_bridge.py b/interfaces/ws/netx_bridge.py deleted file mode 100644 index 84e1dc2d..00000000 --- a/interfaces/ws/netx_bridge.py +++ /dev/null @@ -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 diff --git a/tests/test_netx_bridge_whatsapp_outbound.py b/tests/test_netx_bridge_whatsapp_outbound.py deleted file mode 100644 index b3e776df..00000000 --- a/tests/test_netx_bridge_whatsapp_outbound.py +++ /dev/null @@ -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()