mirror of
https://github.com/hansjone/oclaw.git
synced 2026-10-08 23:33:16 +08:00
feat(netx-bridge): WhatsApp proactive alerts from NetX key alarms
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
parent
9b6bf1d4eb
commit
7210293f42
9 changed files with 747 additions and 1 deletions
|
|
@ -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 消息回放
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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 }),
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
||||
|
|
|
|||
137
interfaces/ws/netx_bridge.py
Normal file
137
interfaces/ws/netx_bridge.py
Normal file
|
|
@ -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
|
||||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -696,6 +696,53 @@ async function postInbound(payload: Json): Promise<Json> {
|
|||
return text ? (JSON.parse(text) as Json) : {};
|
||||
}
|
||||
|
||||
async function pollOutboundQueue(sock: ReturnType<typeof makeWASocket>): Promise<void> {
|
||||
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<typeof makeWASocket> | 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<void> {
|
||||
return new Promise((r) => setTimeout(r, ms));
|
||||
}
|
||||
|
|
@ -952,6 +999,7 @@ async function main(): Promise<void> {
|
|||
|
||||
log(`runner started local=${LOCAL_BASE_URL} stateDir=${STATE_DIR} verbose=${VERBOSE} node=${process.version}`);
|
||||
await logNetworkHints();
|
||||
startOutboundPoller(() => sock);
|
||||
await connectOnce();
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
63
tests/test_netx_bridge_whatsapp_outbound.py
Normal file
63
tests/test_netx_bridge_whatsapp_outbound.py
Normal file
|
|
@ -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()
|
||||
Loading…
Add table
Add a link
Reference in a new issue