diff --git a/_local/system.env.example b/_local/system.env.example index 79e8c0a6..4a162710 100644 --- a/_local/system.env.example +++ b/_local/system.env.example @@ -319,6 +319,8 @@ OCLAW_WS_SEND_QUEUE_MAX_BYTES= OCLAW_WS_EVENT_REPLAY_MAX=256 # netx -> oclaw AP 分析接口共享 token(可选;设置后可用 Bearer 直接访问 /admin/api/ops-ai/analyze-sync) OCLAW_OPS_AI_SHARED_TOKEN= +# netx -> oclaw 告警 WSS(/ws/netx-bridge)独立 token,与 OPS_AI 分开 +OCLAW_NETX_BRIDGE_TOKEN= # ----------------------------------------------------------------------------- # 十五、模型请求:工具 JSON、Replay 策略、Agent 消息回放 diff --git a/interfaces/admin/routes.py b/interfaces/admin/routes.py index d9bdc65c..80cdc25b 100644 --- a/interfaces/admin/routes.py +++ b/interfaces/admin/routes.py @@ -1301,6 +1301,118 @@ def build_admin_router() -> APIRouter: ) return {"ok": True, "deleted": deleted} + @router.get("/admin/api/whatsapp/groups") + def api_whatsapp_groups( + tenant_id: str = Query(default="default"), + account_id: str = Query(default=""), + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:read") + tid = str(tenant_id or "default").strip() + aid = str(account_id or os.getenv("AIA_WHATSAPP_ACCOUNT_ID") or "wa-default").strip() + known = store.list_whatsapp_known_groups(tenant_id=tid, account_id=aid) + session_groups: list[dict[str, Any]] = [] + try: + with store._connect() as conn: + rows = conn.execute( + """ + SELECT DISTINCT external_chat_id, session_id + FROM channel_session_v2 + WHERE tenant_id = ? AND channel = 'whatsapp' AND account_id = ? AND external_chat_id LIKE '%@g.us' + """, + (tid, aid), + ).fetchall() + for r in rows: + jid = str(r[0] or "").strip() + if jid: + session_groups.append({"group_jid": jid, "session_id": str(r[1] or "")}) + except Exception: + session_groups = [] + merged: dict[str, dict[str, Any]] = {} + for g in known: + jid = str(g.get("group_jid") or "").strip() + if jid: + merged[jid] = g + for sg in session_groups: + jid = str(sg.get("group_jid") or "").strip() + if jid and jid not in merged: + merged[jid] = { + "tenant_id": tid, + "account_id": aid, + "group_jid": jid, + "group_name": "", + "last_seen_at": "", + "session_id": sg.get("session_id"), + } + binding = store.get_whatsapp_alert_binding(tenant_id=tid, account_id=aid) + return {"ok": True, "items": list(merged.values()), "binding": binding} + + @router.get("/admin/api/whatsapp/alert-binding") + def api_whatsapp_alert_binding_get( + tenant_id: str = Query(default="default"), + account_id: str = Query(default=""), + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:read") + tid = str(tenant_id or "default").strip() + aid = str(account_id or os.getenv("AIA_WHATSAPP_ACCOUNT_ID") or "wa-default").strip() + binding = store.get_whatsapp_alert_binding(tenant_id=tid, account_id=aid) + return {"ok": True, "binding": binding} + + @router.post("/admin/api/whatsapp/alert-binding") + def api_whatsapp_alert_binding_upsert( + payload: dict[str, Any] | None = Body(default=None), + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + payload = payload or {} + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:runtime:write") + tid = str(payload.get("tenant_id") or "default").strip() + aid = str(payload.get("account_id") or os.getenv("AIA_WHATSAPP_ACCOUNT_ID") or "wa-default").strip() + group_jid = str(payload.get("group_jid") or "").strip() + if not group_jid: + return {"ok": False, "error": "group_jid_required"} + binding = store.upsert_whatsapp_alert_binding( + tenant_id=tid, + account_id=aid, + group_jid=group_jid, + group_name=str(payload.get("group_name") or "").strip(), + enabled=bool(payload.get("enabled", True)), + ) + return {"ok": True, "binding": binding} + + @router.post("/admin/api/whatsapp/alert-binding/test") + def api_whatsapp_alert_binding_test( + payload: dict[str, Any] | None = Body(default=None), + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + payload = payload or {} + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:runtime:write") + tid = str(payload.get("tenant_id") or "default").strip() + aid = str(payload.get("account_id") or os.getenv("AIA_WHATSAPP_ACCOUNT_ID") or "wa-default").strip() + binding = store.get_whatsapp_alert_binding(tenant_id=tid, account_id=aid) + if not binding or not bool(binding.get("enabled")): + return {"ok": False, "error": "binding_missing"} + group_jid = str(binding.get("group_jid") or "").strip() + if not group_jid: + return {"ok": False, "error": "group_jid_missing"} + msg_id = store.enqueue_channel_outbound_message( + channel="whatsapp", + chat_id=group_jid, + text="[OClaw] WhatsApp 告警群绑定测试消息", + tenant_id=tid, + account_id=aid, + source="admin.test", + ) + return {"ok": True, "outbound_id": msg_id} + @router.get("/admin/api/users") def api_users( tenant_id: str, diff --git a/interfaces/admin/static/app.js b/interfaces/admin/static/app.js index 91dc2a14..d68694ec 100644 --- a/interfaces/admin/static/app.js +++ b/interfaces/admin/static/app.js @@ -1472,7 +1472,7 @@ function markPrewarmReminder(reason) { } async function renderStack() { - const [st, anomaliesResp, scanResp, prewarmStatusResp, prewarmPromptsResp, channelSpecResp, weixinDispatchResp, whatsappDispatchResp] = await Promise.all([ + const [st, anomaliesResp, scanResp, prewarmStatusResp, prewarmPromptsResp, channelSpecResp, weixinDispatchResp, whatsappDispatchResp, whatsappGroupsResp] = await Promise.all([ apiGet("/admin/api/stack/status"), apiGet("/admin/api/runtime/anomalies"), apiGet("/admin/api/runtime/scan-artifacts"), @@ -1481,6 +1481,7 @@ async function renderStack() { apiGet("/admin/api/chat/settings/specialist-flags"), apiGet("/admin/api/chat/settings/channel-dispatch/weixin"), apiGet("/admin/api/chat/settings/channel-dispatch/whatsapp"), + apiGet("/admin/api/whatsapp/groups?tenant_id=default"), ]); const requiredServices = ["gateway", "channel:wecom"]; const runningNames = new Set( @@ -1582,6 +1583,68 @@ async function renderStack() { }; const weixinDispatchCard = createChannelDispatchCard("weixin", "Weixin dispatch", weixinDispatchResp || {}); const whatsappDispatchCard = createChannelDispatchCard("whatsapp", "WhatsApp dispatch", whatsappDispatchResp || {}); + const waGroups = Array.isArray(whatsappGroupsResp && whatsappGroupsResp.items) ? whatsappGroupsResp.items : []; + const waBinding = (whatsappGroupsResp && whatsappGroupsResp.binding) || {}; + const waGroupSel = el( + "select", + { class: "input" }, + [ + el("option", { value: "", text: currentLang === "zh" ? "选择群…" : "Select group…" }), + ...waGroups.map((g) => { + const jid = String((g && g.group_jid) || "").trim(); + const name = String((g && g.group_name) || "").trim(); + const label = name ? `${name} (${jid})` : jid; + return el("option", { + value: jid, + text: label || jid, + selected: jid && jid === String(waBinding.group_jid || "") ? "selected" : undefined, + }); + }), + ], + ); + const waBindingStatus = el("div", { + class: "muted", + text: waBinding.group_jid + ? `binding=${waBinding.group_jid} enabled=${Boolean(waBinding.enabled)}` + : (currentLang === "zh" ? "尚未绑定告警群" : "No alert group bound"), + }); + const whatsappAlertBindingCard = el("div", { class: "card" }, [ + el("div", { class: "card__title", text: currentLang === "zh" ? "WhatsApp 告警群绑定" : "WhatsApp alert group" }), + el("div", { class: "muted", text: currentLang === "zh" ? "NetX 关键告警将推送到此群;与 Chat 会话删除无关。" : "NetX key alerts go to this group; independent of chat sessions." }), + el("div", { class: "row" }, [waGroupSel]), + el("div", { class: "row" }, [ + el("button", { + class: "btn btn--primary", + text: currentLang === "zh" ? "保存绑定" : "Save binding", + onclick: async () => { + const group_jid = String(waGroupSel.value || "").trim(); + if (!group_jid) return; + const picked = waGroups.find((g) => String((g && g.group_jid) || "").trim() === group_jid) || {}; + const resp = await apiPost("/admin/api/whatsapp/alert-binding", { + tenant_id: "default", + group_jid, + group_name: String((picked && picked.group_name) || ""), + enabled: true, + }); + const b = (resp && resp.binding) || {}; + waBindingStatus.textContent = b.group_jid + ? `binding=${b.group_jid} enabled=${Boolean(b.enabled)}` + : (currentLang === "zh" ? "绑定失败" : "Bind failed"); + }, + }), + el("button", { + class: "btn", + text: currentLang === "zh" ? "测试推送" : "Test push", + onclick: async () => { + const resp = await apiPost("/admin/api/whatsapp/alert-binding/test", { tenant_id: "default" }); + waBindingStatus.textContent = resp && resp.ok + ? `test ok outbound_id=${String(resp.outbound_id || "")}` + : `test failed: ${String((resp && resp.error) || "unknown")}`; + }, + }), + ]), + waBindingStatus, + ]); const cleanupStatus = el("div", { class: "muted", text: "" }); const btnCleanup = el("button", { class: "btn btn--danger", text: t("stack.cleanup"), onclick: async () => { const resp = await apiPost("/admin/api/runtime/cleanup", {}); @@ -1773,6 +1836,7 @@ async function renderStack() { ]), weixinDispatchCard, whatsappDispatchCard, + whatsappAlertBindingCard, el("div", { class: "card" }, [ el("div", { class: "card__title", text: currentLang === "zh" ? "提示词/工具预热" : "Prompt/Tool Prewarm" }), el("div", { class: "muted", text: prewarmSummary }), diff --git a/interfaces/http/fastapi_app.py b/interfaces/http/fastapi_app.py index 5f5dc753..a99e8510 100644 --- a/interfaces/http/fastapi_app.py +++ b/interfaces/http/fastapi_app.py @@ -20,6 +20,7 @@ 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 @@ -278,6 +279,29 @@ def create_app() -> FastAPI: payload["channel"] = str(channel or "").strip().lower() return await asyncio.to_thread(process_inbound_payload_usecase, payload) + @app.get("/whatsapp/outbound/pending") + def whatsapp_outbound_pending(account_id: str = "", limit: int = 20) -> dict[str, Any]: + from svc.persistence.assistant_store import get_assistant_store + + store = get_assistant_store() + aid = str(account_id or os.getenv("AIA_WHATSAPP_ACCOUNT_ID") or "wa-default").strip() + items = store.list_pending_channel_outbound_messages(channel="whatsapp", account_id=aid, limit=limit) + return {"ok": True, "items": items} + + @app.post("/whatsapp/outbound/ack") + def whatsapp_outbound_ack(payload: dict[str, Any]) -> dict[str, Any]: + from svc.persistence.assistant_store import get_assistant_store + + body = payload if isinstance(payload, dict) else {} + msg_id = str(body.get("id") or "").strip() + if not msg_id: + return {"ok": False, "error": "missing id"} + ok = bool(body.get("ok", True)) + err = str(body.get("error") or "").strip() + store = get_assistant_store() + changed = store.ack_channel_outbound_message(message_id=msg_id, ok=ok, error=err) + return {"ok": changed} + @app.post("/wecom/inbound") async def wecom_inbound(payload: dict[str, Any]) -> dict[str, Any]: payload = payload if isinstance(payload, dict) else {} @@ -295,6 +319,10 @@ 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 new file mode 100644 index 00000000..0d47639a --- /dev/null +++ b/interfaces/ws/netx_bridge.py @@ -0,0 +1,137 @@ +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_NETX_BRIDGE_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 "关键告警").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": "新增", "updated": "更新", "deleted": "清除"}.get(action, action or "告警") + lines = [ + f"[NetX {action_label}] {label}", + f"设备: {device or '-'}", + f"对象: {str(payload.get('object_name') or '-').strip()}", + f"级别: {str(payload.get('perceived_severity') or '-').strip()}", + f"原因: {str(payload.get('native_probable_cause') or '-').strip()}", + f"时间: {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/runtime/application/gateway/inbound_service.py b/runtime/application/gateway/inbound_service.py index c420fc64..663a7f31 100644 --- a/runtime/application/gateway/inbound_service.py +++ b/runtime/application/gateway/inbound_service.py @@ -669,6 +669,21 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]: if not account_id: raise ValueError(f"missing {channel_name} account_id") + if str(inbound.channel or "").strip().lower() == "whatsapp" and inbound.is_group: + import os + + from runtime.extensions.whatsapp.api import is_whatsapp_group_jid + + chat_id = str(inbound.external_chat_id or "").strip() + if is_whatsapp_group_jid(chat_id): + meta = inbound.metadata if isinstance(inbound.metadata, dict) else {} + store.upsert_whatsapp_known_group( + tenant_id=str(os.getenv("OCLAW_DEFAULT_TENANT_ID") or "default"), + account_id=account_id, + group_jid=chat_id, + group_name=str(meta.get("group_name") or "").strip(), + ) + account = store.find_user_by_channel_account(channel=inbound.channel, account_id=account_id) or {} from runtime.orchestration.group_ingest import ( build_group_sender_context, diff --git a/runtime/operations/whatsapp_bridge/baileys_runner.ts b/runtime/operations/whatsapp_bridge/baileys_runner.ts index 54dcd1da..af8bd071 100644 --- a/runtime/operations/whatsapp_bridge/baileys_runner.ts +++ b/runtime/operations/whatsapp_bridge/baileys_runner.ts @@ -696,6 +696,53 @@ async function postInbound(payload: Json): Promise { return text ? (JSON.parse(text) as Json) : {}; } +async function pollOutboundQueue(sock: ReturnType): Promise { + const url = `${LOCAL_BASE_URL.replace(/\/+$/, "")}/whatsapp/outbound/pending?account_id=${encodeURIComponent(ACCOUNT_ID)}`; + try { + const res = await fetch(url); + if (!res.ok) return; + const body = (await res.json()) as Json; + const items = Array.isArray(body.items) ? (body.items as Json[]) : []; + for (const item of items) { + if (!item || typeof item !== "object") continue; + const id = String((item as any).id || "").trim(); + const chatId = String((item as any).chat_id || "").trim(); + const text = String((item as any).text || "").trim(); + if (!id || !chatId || !text) continue; + let ok = true; + let err = ""; + try { + await sock.sendMessage(chatId, { text }); + log(`outbound sent id=${id} chat=${chatId}`); + } catch (sendErr) { + ok = false; + err = String(sendErr); + log(`outbound send failed id=${id} chat=${chatId} err=${err.slice(0, 160)}`); + } + try { + await fetch(`${LOCAL_BASE_URL.replace(/\/+$/, "")}/whatsapp/outbound/ack`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ id, ok, error: err }), + }); + } catch { + // best-effort ack + } + } + } catch (err) { + if (VERBOSE) log(`outbound poll error: ${String(err)}`); + } +} + +function startOutboundPoller(getSock: () => ReturnType | null): void { + const intervalMs = Number(process.env.OCLAW_WHATSAPP_OUTBOUND_POLL_MS || "2500") || 2500; + setInterval(() => { + const s = getSock(); + if (!s) return; + void pollOutboundQueue(s); + }, Math.max(1000, intervalMs)); +} + function sleep(ms: number): Promise { return new Promise((r) => setTimeout(r, ms)); } @@ -952,6 +999,7 @@ async function main(): Promise { log(`runner started local=${LOCAL_BASE_URL} stateDir=${STATE_DIR} verbose=${VERBOSE} node=${process.version}`); await logNetworkHints(); + startOutboundPoller(() => sock); await connectOnce(); } diff --git a/svc/persistence/sqlite_store.py b/svc/persistence/sqlite_store.py index d86ef595..8e8f8556 100644 --- a/svc/persistence/sqlite_store.py +++ b/svc/persistence/sqlite_store.py @@ -632,6 +632,51 @@ class SqliteStore: uca_cols = {row[1] for row in conn.execute("PRAGMA table_info(user_channel_account)").fetchall()} if "name" not in uca_cols: conn.execute("ALTER TABLE user_channel_account ADD COLUMN name TEXT NOT NULL DEFAULT ''") + conn.execute( + """ + CREATE TABLE IF NOT EXISTS whatsapp_known_group ( + tenant_id TEXT NOT NULL, + account_id TEXT NOT NULL, + group_jid TEXT NOT NULL, + group_name TEXT NOT NULL DEFAULT '', + last_seen_at TEXT NOT NULL, + PRIMARY KEY (tenant_id, account_id, group_jid) + ); + """ + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS whatsapp_alert_binding ( + tenant_id TEXT NOT NULL, + account_id TEXT NOT NULL, + group_jid TEXT NOT NULL, + group_name TEXT NOT NULL DEFAULT '', + enabled INTEGER NOT NULL DEFAULT 1, + updated_at TEXT NOT NULL, + PRIMARY KEY (tenant_id, account_id) + ); + """ + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS channel_outbound_message ( + id TEXT PRIMARY KEY, + tenant_id TEXT NOT NULL DEFAULT '', + channel TEXT NOT NULL, + account_id TEXT NOT NULL DEFAULT '', + chat_id TEXT NOT NULL, + text TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'pending', + source TEXT NOT NULL DEFAULT '', + created_at TEXT NOT NULL, + sent_at TEXT, + error TEXT NOT NULL DEFAULT '' + ); + """ + ) + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_channel_outbound_pending ON channel_outbound_message(channel, account_id, status, created_at)" + ) conn.execute( """ CREATE TABLE IF NOT EXISTS todo_item ( @@ -1176,6 +1221,51 @@ class SqliteStore: def _init_db_postgresql(self) -> None: """PostgreSQL: tables from Alembic migration; run seeds and housekeeping.""" with self._connect() as conn: + conn.execute( + """ + CREATE TABLE IF NOT EXISTS whatsapp_known_group ( + tenant_id TEXT NOT NULL, + account_id TEXT NOT NULL, + group_jid TEXT NOT NULL, + group_name TEXT NOT NULL DEFAULT '', + last_seen_at TEXT NOT NULL, + PRIMARY KEY (tenant_id, account_id, group_jid) + ) + """ + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS whatsapp_alert_binding ( + tenant_id TEXT NOT NULL, + account_id TEXT NOT NULL, + group_jid TEXT NOT NULL, + group_name TEXT NOT NULL DEFAULT '', + enabled INTEGER NOT NULL DEFAULT 1, + updated_at TEXT NOT NULL, + PRIMARY KEY (tenant_id, account_id) + ) + """ + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS channel_outbound_message ( + id TEXT PRIMARY KEY, + tenant_id TEXT NOT NULL DEFAULT '', + channel TEXT NOT NULL, + account_id TEXT NOT NULL DEFAULT '', + chat_id TEXT NOT NULL, + text TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'pending', + source TEXT NOT NULL DEFAULT '', + created_at TEXT NOT NULL, + sent_at TEXT, + error TEXT NOT NULL DEFAULT '' + ) + """ + ) + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_channel_outbound_pending ON channel_outbound_message(channel, account_id, status, created_at)" + ) self._seed_builtin_llm_profiles(conn) self._seed_default_permissions(conn) conn.execute( @@ -5200,3 +5290,190 @@ class SqliteStore: (str(assignee_user_id), ts, str(tenant_id), str(todo_id)), ) return bool(cur.rowcount and cur.rowcount > 0) + + def upsert_whatsapp_known_group( + self, + *, + tenant_id: str, + account_id: str, + group_jid: str, + group_name: str = "", + ) -> None: + ts = utc_now_iso() + with self._connect() as conn: + conn.execute( + """ + INSERT INTO whatsapp_known_group (tenant_id, account_id, group_jid, group_name, last_seen_at) + VALUES (?, ?, ?, ?, ?) + ON CONFLICT(tenant_id, account_id, group_jid) DO UPDATE SET + group_name = CASE WHEN excluded.group_name != '' THEN excluded.group_name ELSE whatsapp_known_group.group_name END, + last_seen_at = excluded.last_seen_at + """, + (str(tenant_id), str(account_id), str(group_jid), str(group_name or ""), ts), + ) + + def list_whatsapp_known_groups( + self, + *, + tenant_id: str | None = None, + account_id: str | None = None, + ) -> list[dict[str, Any]]: + clauses: list[str] = [] + params: list[Any] = [] + if tenant_id: + clauses.append("tenant_id = ?") + params.append(str(tenant_id)) + if account_id: + clauses.append("account_id = ?") + params.append(str(account_id)) + where = f"WHERE {' AND '.join(clauses)}" if clauses else "" + with self._connect() as conn: + rows = conn.execute( + f""" + SELECT tenant_id, account_id, group_jid, group_name, last_seen_at + FROM whatsapp_known_group + {where} + ORDER BY last_seen_at DESC + """, + tuple(params), + ).fetchall() + return [ + { + "tenant_id": str(r[0] or ""), + "account_id": str(r[1] or ""), + "group_jid": str(r[2] or ""), + "group_name": str(r[3] or ""), + "last_seen_at": str(r[4] or ""), + } + for r in rows + ] + + def get_whatsapp_alert_binding( + self, + *, + tenant_id: str, + account_id: str, + ) -> dict[str, Any] | None: + with self._connect() as conn: + row = conn.execute( + """ + SELECT tenant_id, account_id, group_jid, group_name, enabled, updated_at + FROM whatsapp_alert_binding + WHERE tenant_id = ? AND account_id = ? + """, + (str(tenant_id), str(account_id)), + ).fetchone() + if row is None: + return None + return { + "tenant_id": str(row[0] or ""), + "account_id": str(row[1] or ""), + "group_jid": str(row[2] or ""), + "group_name": str(row[3] or ""), + "enabled": bool(int(row[4] or 0)), + "updated_at": str(row[5] or ""), + } + + def upsert_whatsapp_alert_binding( + self, + *, + tenant_id: str, + account_id: str, + group_jid: str, + group_name: str = "", + enabled: bool = True, + ) -> dict[str, Any]: + ts = utc_now_iso() + with self._connect() as conn: + conn.execute( + """ + INSERT INTO whatsapp_alert_binding (tenant_id, account_id, group_jid, group_name, enabled, updated_at) + VALUES (?, ?, ?, ?, ?, ?) + ON CONFLICT(tenant_id, account_id) DO UPDATE SET + group_jid = excluded.group_jid, + group_name = excluded.group_name, + enabled = excluded.enabled, + updated_at = excluded.updated_at + """, + (str(tenant_id), str(account_id), str(group_jid), str(group_name or ""), 1 if enabled else 0, ts), + ) + out = self.get_whatsapp_alert_binding(tenant_id=tenant_id, account_id=account_id) + return out or {} + + def enqueue_channel_outbound_message( + self, + *, + channel: str, + chat_id: str, + text: str, + tenant_id: str = "", + account_id: str = "", + source: str = "", + ) -> str: + import uuid + + msg_id = uuid.uuid4().hex + ts = utc_now_iso() + with self._connect() as conn: + conn.execute( + """ + INSERT INTO channel_outbound_message + (id, tenant_id, channel, account_id, chat_id, text, status, source, created_at, error) + VALUES (?, ?, ?, ?, ?, ?, 'pending', ?, ?, '') + """, + (msg_id, str(tenant_id), str(channel), str(account_id), str(chat_id), str(text), str(source), ts), + ) + return msg_id + + def list_pending_channel_outbound_messages( + self, + *, + channel: str, + account_id: str, + limit: int = 20, + ) -> list[dict[str, Any]]: + lim = max(1, min(int(limit), 100)) + with self._connect() as conn: + rows = conn.execute( + """ + SELECT id, tenant_id, channel, account_id, chat_id, text, source, created_at + FROM channel_outbound_message + WHERE channel = ? AND account_id = ? AND status = 'pending' + ORDER BY created_at ASC + LIMIT ? + """, + (str(channel), str(account_id), lim), + ).fetchall() + return [ + { + "id": str(r[0] or ""), + "tenant_id": str(r[1] or ""), + "channel": str(r[2] or ""), + "account_id": str(r[3] or ""), + "chat_id": str(r[4] or ""), + "text": str(r[5] or ""), + "source": str(r[6] or ""), + "created_at": str(r[7] or ""), + } + for r in rows + ] + + def ack_channel_outbound_message( + self, + *, + message_id: str, + ok: bool, + error: str = "", + ) -> bool: + ts = utc_now_iso() + status = "sent" if ok else "failed" + with self._connect() as conn: + cur = conn.execute( + """ + UPDATE channel_outbound_message + SET status = ?, sent_at = ?, error = ? + WHERE id = ? AND status = 'pending' + """, + (status, ts, str(error or ""), str(message_id)), + ) + return bool(cur.rowcount and cur.rowcount > 0) diff --git a/tests/test_netx_bridge_whatsapp_outbound.py b/tests/test_netx_bridge_whatsapp_outbound.py new file mode 100644 index 00000000..973d1d49 --- /dev/null +++ b/tests/test_netx_bridge_whatsapp_outbound.py @@ -0,0 +1,63 @@ +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 svc.persistence.assistant_store import reset_assistant_store_singleton + + +class NetxBridgeTests(unittest.TestCase): + def setUp(self) -> None: + os.environ["OCLAW_NETX_BRIDGE_TOKEN"] = "test-bridge-token" + os.environ.pop("OCLAW_OPS_AI_SHARED_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_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()