From 1ea13a205df7f0b1933e3ce551fbecb1b026f307 Mon Sep 17 00:00:00 2001 From: oliver Date: Wed, 1 Jul 2026 11:36:58 +0800 Subject: [PATCH] feat(whatsapp): add access control with quote-based admin approval Introduce whitelist/blacklist gating, pending requests, Admin UI management, and stanza-mapped quote approval so admins can approve via WhatsApp DM without blocking normal LLM chat. Co-authored-by: Cursor --- interfaces/admin/chat_api.py | 3 +- interfaces/admin/routes.py | 319 +++++++ interfaces/admin/static/app.js | 411 ++++++++- interfaces/http/fastapi_app.py | 8 +- .../application/gateway/inbound_service.py | 15 +- .../gateway/whatsapp_inbound_access.py | 512 ++++++++++++ runtime/extensions/whatsapp/access_control.py | 392 +++++++++ runtime/extensions/whatsapp/tenant.py | 14 + .../whatsapp_bridge/baileys_runner.ts | 71 +- svc/persistence/sqlite_store.py | 787 +++++++++++++++++- tests/test_whatsapp_access_control.py | 145 ++++ tests/test_whatsapp_inbound_access.py | 389 +++++++++ tests/test_whatsapp_ops_scripts.py | 12 + 13 files changed, 3048 insertions(+), 30 deletions(-) create mode 100644 runtime/application/gateway/whatsapp_inbound_access.py create mode 100644 runtime/extensions/whatsapp/access_control.py create mode 100644 runtime/extensions/whatsapp/tenant.py create mode 100644 tests/test_whatsapp_access_control.py create mode 100644 tests/test_whatsapp_inbound_access.py diff --git a/interfaces/admin/chat_api.py b/interfaces/admin/chat_api.py index dc403b92..b0703780 100644 --- a/interfaces/admin/chat_api.py +++ b/interfaces/admin/chat_api.py @@ -1587,6 +1587,7 @@ def include_chat_routes(router: APIRouter, *, resolve_auth: Callable[[SqliteStor ctx = resolve_auth(store, authorization) _require_administrator_chat_viewer(ctx) ch = _normalize_channel_dispatch_channel(channel) + default_lang = "en" if ch == "whatsapp" else "auto" interaction_mode = normalize_interaction_mode( store.get_setting(_channel_dispatch_interaction_key(ch)) or "expert" ) @@ -1594,7 +1595,7 @@ def include_chat_routes(router: APIRouter, *, resolve_auth: Callable[[SqliteStor store.get_setting(_channel_dispatch_specialist_key(ch)) or "generalist" ) lang = _normalize_channel_dispatch_lang( - store.get_setting(_channel_dispatch_lang_key(ch)) or "auto" + store.get_setting(_channel_dispatch_lang_key(ch)) or default_lang ) specialist = _apply_specialist_flags(store, specialist) return { diff --git a/interfaces/admin/routes.py b/interfaces/admin/routes.py index 9f8a0753..1515f443 100644 --- a/interfaces/admin/routes.py +++ b/interfaces/admin/routes.py @@ -1413,6 +1413,325 @@ def build_admin_router() -> APIRouter: ) return {"ok": True, "outbound_id": msg_id} + @router.get("/admin/api/whatsapp/access") + def api_whatsapp_access_get( + tenant_id: str = Query(default="default"), + account_id: str = Query(default=""), + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + from runtime.extensions.whatsapp.tenant import resolve_whatsapp_tenant_id + + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:read") + aid = str(account_id or os.getenv("AIA_WHATSAPP_ACCOUNT_ID") or "wa-default").strip() + tid = resolve_whatsapp_tenant_id(store, account_id=aid) + legacy_tid = str(tenant_id or "default").strip() + # Be defensive: older DBs/gateways may not have the new tables yet. + # Do not break the entire Runtime page; return empty defaults instead. + try: + cfg = store.get_whatsapp_access_config(tenant_id=tid, account_id=aid) or {} + if not cfg and legacy_tid and legacy_tid != tid: + cfg = store.get_whatsapp_access_config(tenant_id=legacy_tid, account_id=aid) or {} + except Exception: + cfg = {} + contacts_by_phone: dict[str, dict[str, Any]] = {} + try: + from runtime.extensions.whatsapp.access_control import contact_phone_key + + priority = {"admin": 4, "blacklist": 3, "whitelist": 2} + + def _contact_rank(row: dict[str, Any]) -> int: + lt = str(row.get("list_type") or "").strip().lower() + return int(priority.get(lt, 0)) + + for source_tid in (tid, legacy_tid): + if not source_tid: + continue + for row in store.list_whatsapp_contacts(tenant_id=source_tid, account_id=aid, limit=500): + phone_key = contact_phone_key(row) + if not phone_key: + continue + prev = contacts_by_phone.get(phone_key) + if not prev or _contact_rank(row) > _contact_rank(prev): + contacts_by_phone[phone_key] = row + contacts = list(contacts_by_phone.values()) + contacts = [ + row + for row in contacts + if str(row.get("list_type") or "").strip().lower() in {"admin", "whitelist", "blacklist"} + ] + except Exception: + contacts = [] + pending_by_id: dict[str, dict[str, Any]] = {} + try: + for source_tid in (tid, legacy_tid): + if not source_tid: + continue + for row in store.list_whatsapp_access_pending( + tenant_id=source_tid, account_id=aid, status="pending", limit=50 + ): + pid = str(row.get("id") or "").strip() + if pid: + pending_by_id[pid] = row + pending = list(pending_by_id.values()) + except Exception: + pending = [] + denied_by_id: dict[str, dict[str, Any]] = {} + try: + for source_tid in (tid, legacy_tid): + if not source_tid: + continue + for row in store.list_whatsapp_access_pending( + tenant_id=source_tid, account_id=aid, status="denied", limit=50 + ): + pid = str(row.get("id") or "").strip() + if pid: + denied_by_id[pid] = row + denied = list(denied_by_id.values()) + except Exception: + denied = [] + return { + "ok": True, + "config": { + "tenant_id": tid, + "account_id": aid, + "access_mode": str(cfg.get("access_mode") or "blacklist"), + "lang": str(cfg.get("lang") or "en"), + }, + "contacts": contacts, + "pending": pending, + "denied": denied, + } + + @router.post("/admin/api/whatsapp/access/config") + def api_whatsapp_access_config_upsert( + payload: dict[str, Any] | None = Body(default=None), + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + from runtime.extensions.whatsapp.tenant import resolve_whatsapp_tenant_id + + payload = payload or {} + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:runtime:write") + aid = str(payload.get("account_id") or os.getenv("AIA_WHATSAPP_ACCOUNT_ID") or "wa-default").strip() + tid = resolve_whatsapp_tenant_id(store, account_id=aid) + access_mode = str(payload.get("access_mode") or "blacklist").strip().lower() + lang = str(payload.get("lang") or "en").strip().lower() + cfg = store.upsert_whatsapp_access_config( + tenant_id=tid, + account_id=aid, + access_mode=access_mode, + lang=lang, + ) + return {"ok": True, "config": cfg} + + @router.post("/admin/api/whatsapp/access/contacts") + def api_whatsapp_access_contact_upsert( + payload: dict[str, Any] | None = Body(default=None), + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + from runtime.extensions.whatsapp.access_control import ( + normalize_whatsapp_phone, + resolve_sender_phone, + ) + from runtime.extensions.whatsapp.tenant import resolve_whatsapp_tenant_id + + payload = payload or {} + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:runtime:write") + aid = str(payload.get("account_id") or os.getenv("AIA_WHATSAPP_ACCOUNT_ID") or "wa-default").strip() + tid = resolve_whatsapp_tenant_id(store, account_id=aid) + target_raw = str(payload.get("phone") or payload.get("external_user_id") or "").strip() + if not target_raw: + return {"ok": False, "error": "phone_required"} + try: + phone_val = normalize_whatsapp_phone(target_raw) + except Exception: + return {"ok": False, "error": "invalid_phone"} + list_type = str(payload.get("list_type") or "").strip().lower() + if list_type not in {"admin", "whitelist", "blacklist"}: + return {"ok": False, "error": "invalid_list_type"} + contact = store.apply_whatsapp_contact_access( + tenant_id=tid, + account_id=aid, + external_user_id=phone_val, + push_name=str(payload.get("push_name") or ""), + phone=phone_val, + list_type=list_type, + notes=str(payload.get("notes") or ""), + ) + resolved_by = str(ctx.get("user_id") or "") + pending_status = "approved" if list_type in {"admin", "whitelist"} else "denied" + legacy_tid = str(payload.get("tenant_id") or "default").strip() + extra_tids = [legacy_tid] if legacy_tid and legacy_tid != tid else [] + store.resolve_whatsapp_pending_for_sender( + tenant_id=tid, + account_id=aid, + external_user_id=phone_val, + resolved_by=resolved_by, + status=pending_status, + extra_tenant_ids=extra_tids, + ) + return {"ok": True, "contact": contact} + + @router.delete("/admin/api/whatsapp/access/contacts") + def api_whatsapp_access_contact_delete( + tenant_id: str = Query(default="default"), + account_id: str = Query(default=""), + phone: str = Query(default=""), + external_user_id: str = Query(default=""), + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + from runtime.extensions.whatsapp.access_control import normalize_whatsapp_phone + from runtime.extensions.whatsapp.tenant import resolve_whatsapp_tenant_id + + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:runtime:write") + aid = str(account_id or os.getenv("AIA_WHATSAPP_ACCOUNT_ID") or "wa-default").strip() + tid = resolve_whatsapp_tenant_id(store, account_id=aid) + key = str(phone or external_user_id or "").strip() + if not key: + return {"ok": False, "error": "phone_required"} + try: + phone_val = normalize_whatsapp_phone(key) + except Exception: + return {"ok": False, "error": "invalid_phone"} + legacy_tid = str(tenant_id or "default").strip() + extra_tids = [legacy_tid] if legacy_tid and legacy_tid != tid else [] + deleted_count = store.delete_whatsapp_contact_aliases( + tenant_id=tid, + account_id=aid, + external_user_id=phone_val, + extra_tenant_ids=extra_tids, + ) + return {"ok": True, "deleted": deleted_count > 0, "deleted_count": deleted_count} + + @router.post("/admin/api/whatsapp/access/pending/resolve") + def api_whatsapp_access_pending_resolve( + payload: dict[str, Any] | None = Body(default=None), + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + from runtime.extensions.whatsapp.access_control import resolve_sender_phone + from runtime.extensions.whatsapp.tenant import resolve_whatsapp_tenant_id + + body = payload or {} + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:runtime:write") + aid = str(body.get("account_id") or os.getenv("AIA_WHATSAPP_ACCOUNT_ID") or "wa-default").strip() + tid = resolve_whatsapp_tenant_id(store, account_id=aid) + legacy_tid = str(body.get("tenant_id") or "default").strip() + extra_tids = [legacy_tid] if legacy_tid and legacy_tid != tid else [] + pending_id = str(body.get("pending_id") or "").strip() + action = str(body.get("action") or "").strip().lower() + if not pending_id: + return {"ok": False, "error": "pending_id_required"} + if action not in {"approve", "deny", "dismiss"}: + return {"ok": False, "error": "invalid_action"} + + items: list[dict[str, Any]] = [] + seen: set[str] = set() + for source_tid in (tid, *extra_tids): + for status in ("pending", "denied"): + for row in store.list_whatsapp_access_pending( + tenant_id=source_tid, account_id=aid, status=status, limit=200 + ): + pid = str(row.get("id") or "") + if pid and pid not in seen: + seen.add(pid) + items.append(row) + + row = next((x for x in items if str(x.get("id") or "") == pending_id), None) + if not row: + return {"ok": False, "error": "pending_not_found"} + + external_user_id = str(row.get("external_user_id") or "").strip() + push_name = str(row.get("push_name") or "").strip() + resolved_by = str(ctx.get("user_id") or "") + + if action == "dismiss": + deleted = store.delete_whatsapp_access_pending(pending_id=pending_id) + return {"ok": deleted, "action": "dismiss", "deleted": deleted} + + if action == "approve": + phone_val = str(row.get("phone") or "").strip() or resolve_sender_phone(external_user_id) + if not phone_val: + return {"ok": False, "error": "invalid_phone"} + store.apply_whatsapp_contact_access( + tenant_id=tid, + account_id=aid, + external_user_id=phone_val, + push_name=push_name, + phone=phone_val, + list_type="whitelist", + ) + changed = store.resolve_whatsapp_access_pending( + pending_id=pending_id, + status="approved", + resolved_by=resolved_by, + from_statuses=("pending", "denied"), + ) + store.resolve_whatsapp_pending_for_sender( + tenant_id=tid, + account_id=aid, + external_user_id=phone_val, + resolved_by=resolved_by, + status="approved", + extra_tenant_ids=extra_tids, + ) + return { + "ok": bool(changed), + "action": "approve", + "phone": phone_val, + } + + phone_val = str(row.get("phone") or "").strip() or resolve_sender_phone(external_user_id) + if phone_val: + store.apply_whatsapp_contact_access( + tenant_id=tid, + account_id=aid, + external_user_id=phone_val, + push_name=push_name, + phone=phone_val, + list_type="blacklist", + ) + changed = store.resolve_whatsapp_access_pending( + pending_id=pending_id, + status="denied", + resolved_by=resolved_by, + from_statuses=("pending",), + ) + store.resolve_whatsapp_pending_for_sender( + tenant_id=tid, + account_id=aid, + external_user_id=phone_val or external_user_id, + resolved_by=resolved_by, + status="denied", + extra_tenant_ids=extra_tids, + ) + return {"ok": bool(changed), "action": "deny", "phone": phone_val or external_user_id} + + @router.delete("/admin/api/whatsapp/access/pending") + def api_whatsapp_access_pending_delete( + pending_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:runtime:write") + pid = str(pending_id or "").strip() + if not pid: + return {"ok": False, "error": "pending_id_required"} + deleted = store.delete_whatsapp_access_pending( + pending_id=pid, + statuses=("pending", "denied"), + ) + return {"ok": deleted, "deleted": deleted} + @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 2c925119..ad688624 100644 --- a/interfaces/admin/static/app.js +++ b/interfaces/admin/static/app.js @@ -1166,6 +1166,30 @@ async function apiGet(path) { return await res.json(); } +async function apiGetOptional(path) { + try { + return await apiGet(path); + } catch (err) { + console.warn(`apiGetOptional failed path=${String(path || "")} err=${String(err)}`); + return null; + } +} + +async function apiGetNoHang(path) { + const url = resolveAdminApiUrl(path); + const token = getStoredAuthToken(); + const headers = { "accept": "application/json" }; + if (token) headers["authorization"] = `Bearer ${token}`; + const res = await fetch(url, { headers }); + if (res.status === 401) { + scheduleReauthAfter401(url); + // Returning null avoids Promise.all / await from hanging forever. + return null; + } + if (!res.ok) throw new Error(`GET ${url} ${res.status}`); + return await res.json(); +} + async function apiPost(path, body) { const url = resolveAdminApiUrl(path); const token = getStoredAuthToken(); @@ -1544,17 +1568,36 @@ function markPrewarmReminder(reason) { } async function renderStack() { - 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"), - apiGet("/admin/api/runtime/prewarm/status"), - apiGet("/admin/api/runtime/prewarm/prompts?role=manager"), - 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 results = await Promise.allSettled([ + apiGetNoHang("/admin/api/stack/status"), + apiGetNoHang("/admin/api/runtime/anomalies"), + apiGetNoHang("/admin/api/runtime/scan-artifacts"), + apiGetNoHang("/admin/api/runtime/prewarm/status"), + apiGetNoHang("/admin/api/runtime/prewarm/prompts?role=manager"), + apiGetNoHang("/admin/api/chat/settings/specialist-flags"), + apiGetNoHang("/admin/api/chat/settings/channel-dispatch/weixin"), + apiGetNoHang("/admin/api/chat/settings/channel-dispatch/whatsapp"), + apiGetNoHang("/admin/api/whatsapp/groups?tenant_id=default"), + apiGetNoHang("/admin/api/whatsapp/access?tenant_id=default"), ]); + const st = results[0].status === "fulfilled" ? results[0].value : null; + const anomaliesResp = results[1].status === "fulfilled" ? results[1].value : null; + const scanResp = results[2].status === "fulfilled" ? results[2].value : null; + const prewarmStatusResp = results[3].status === "fulfilled" ? results[3].value : null; + const prewarmPromptsResp = results[4].status === "fulfilled" ? results[4].value : null; + const channelSpecResp = results[5].status === "fulfilled" ? results[5].value : null; + const weixinDispatchResp = results[6].status === "fulfilled" ? results[6].value : null; + const whatsappDispatchResp = results[7].status === "fulfilled" ? results[7].value : null; + const whatsappGroupsResp = results[8].status === "fulfilled" ? results[8].value : null; + const whatsappAccessResp = results[9].status === "fulfilled" ? results[9].value : null; + + // If auth failed (401), show a gentle message instead of an infinite spinner. + if (!st) { + return el("div", { class: "card" }, [ + el("div", { class: "card__title", text: t("title.stack") || "Stack" }), + el("div", { class: "muted", text: currentLang === "zh" ? "未登录或会话已过期,请重新登录。" : "Not logged in or session expired. Please log in again." }), + ]); + } const requiredServices = ["gateway", "channel:wecom"]; const runningNames = new Set( (Array.isArray(st.items) ? st.items : []) @@ -1680,6 +1723,353 @@ async function renderStack() { ? `binding=${waBinding.group_jid} enabled=${Boolean(waBinding.enabled)}` : (currentLang === "zh" ? "尚未绑定告警群" : "No alert group bound"), }); + const waAccessCfg = (whatsappAccessResp && whatsappAccessResp.config) || {}; + const waContacts = Array.isArray(whatsappAccessResp && whatsappAccessResp.contacts) ? whatsappAccessResp.contacts : []; + const waPending = Array.isArray(whatsappAccessResp && whatsappAccessResp.pending) ? whatsappAccessResp.pending : []; + const waDenied = Array.isArray(whatsappAccessResp && whatsappAccessResp.denied) ? whatsappAccessResp.denied : []; + const waPhoneDisplay = (row) => { + const phone = String((row && row.phone) || "").trim(); + if (phone) return phone; + const jid = String((row && row.external_user_id) || "").trim(); + const base = jid.split("@")[0] || ""; + const digits = base.replace(/\D/g, ""); + return digits || base || "-"; + }; + const waAccessModeSel = el("select", { class: "input" }, [ + el("option", { + value: "blacklist", + text: currentLang === "zh" ? "黑名单模式(默认拒绝,白名单放行)" : "Blacklist mode (deny by default)", + selected: String(waAccessCfg.access_mode || "blacklist") === "blacklist" ? "selected" : undefined, + }), + el("option", { + value: "whitelist", + text: currentLang === "zh" ? "白名单模式(默认放行,黑名单拒绝)" : "Whitelist mode (allow by default)", + selected: String(waAccessCfg.access_mode || "") === "whitelist" ? "selected" : undefined, + }), + ]); + const waAccessLangSel = el("select", { class: "input" }, [ + el("option", { value: "en", text: "English", selected: String(waAccessCfg.lang || "en") === "en" ? "selected" : undefined }), + el("option", { value: "zh", text: currentLang === "zh" ? "中文" : "Chinese", selected: String(waAccessCfg.lang || "") === "zh" ? "selected" : undefined }), + ]); + const waCounts = { admin: 0, whitelist: 0, blacklist: 0 }; + waContacts.forEach((c) => { + const lt = String((c && c.list_type) || "").trim().toLowerCase(); + if (lt === "admin") waCounts.admin += 1; + else if (lt === "whitelist") waCounts.whitelist += 1; + else if (lt === "blacklist") waCounts.blacklist += 1; + }); + const waAccessStatus = el("div", { + class: "muted", + text: `mode=${String(waAccessCfg.access_mode || "blacklist")} lang=${String(waAccessCfg.lang || "en")} admin=${waCounts.admin} whitelist=${waCounts.whitelist} blacklist=${waCounts.blacklist} pending=${waPending.length} denied=${waDenied.length}`, + }); + const waContactFilterSel = el("select", { class: "input" }, [ + el("option", { value: "all", text: currentLang === "zh" ? "全部" : "All" }), + el("option", { value: "admin", text: currentLang === "zh" ? "管理员" : "Admin" }), + el("option", { value: "whitelist", text: currentLang === "zh" ? "白名单" : "Whitelist" }), + el("option", { value: "blacklist", text: currentLang === "zh" ? "黑名单" : "Blacklist" }), + ]); + const waContactsTbody = el("tbody", {}); + const waContactPhone = (row) => { + const phone = String((row && row.phone) || "").trim(); + if (phone) return phone; + return waPhoneDisplay(row); + }; + const waSaveContact = async (phone, pushName, listType) => { + const phoneVal = String(phone || "").trim(); + const list_type = String(listType || "").trim().toLowerCase(); + if (!phoneVal) return; + if (!list_type) { + window.alert(currentLang === "zh" ? "请选择类型" : "Please select a type"); + return; + } + try { + const resp = await apiPost("/admin/api/whatsapp/access/contacts", { + phone: phoneVal, + push_name: String(pushName || "").trim(), + list_type, + }); + if (!resp || resp.ok === false) { + window.alert(String((resp && resp.error) || (currentLang === "zh" ? "修改失败" : "Update failed"))); + return; + } + router(); + } catch (err) { + window.alert(String(err)); + } + }; + const renderWaContactRows = (filterValue) => { + const filter = String(filterValue || "all").trim().toLowerCase() || "all"; + const rows = waContacts + .filter((c) => { + const lt = String((c && c.list_type) || "").trim().toLowerCase(); + if (filter === "all") return true; + return lt === filter; + }) + .map((c) => { + const name = String((c && c.push_name) || ""); + const lt = String((c && c.list_type) || "").trim().toLowerCase(); + const phone = waContactPhone(c); + const typeSel = el("select", { class: "input" }, [ + el("option", { value: "admin", text: currentLang === "zh" ? "管理员" : "Admin" }), + el("option", { value: "whitelist", text: currentLang === "zh" ? "白名单" : "Whitelist" }), + el("option", { value: "blacklist", text: currentLang === "zh" ? "黑名单" : "Blacklist" }), + ]); + typeSel.value = lt || "whitelist"; + const actionBtns = [ + el("button", { + class: "btn", + text: currentLang === "zh" ? "修改" : "Update", + onclick: async () => { + await waSaveContact(phone, name, String(typeSel.value || "").trim()); + }, + }), + ]; + if (lt === "blacklist") { + actionBtns.push(el("span", { text: " " })); + actionBtns.push(el("button", { + class: "btn btn--primary", + text: currentLang === "zh" ? "加白名单" : "Whitelist", + onclick: async () => { + await waSaveContact(phone, name, "whitelist"); + }, + })); + } + actionBtns.push(el("span", { text: " " })); + actionBtns.push(el("button", { + class: "btn btn--danger", + text: currentLang === "zh" ? "删除" : "Remove", + onclick: async () => { + try { + const resp = await apiDeleteJson( + `/admin/api/whatsapp/access/contacts?phone=${encodeURIComponent(phone)}`, + ); + if (resp && resp.deleted === false) { + window.alert(currentLang === "zh" ? "未找到可删除的联系人" : "Contact not found"); + return; + } + router(); + } catch (err) { + window.alert(String(err)); + } + }, + })); + return el("tr", {}, [ + el("td", { "data-copy-disabled": "1" }, [typeSel]), + el("td", { text: name || "-" }), + el("td", { text: phone || "-" }), + el("td", { "data-copy-disabled": "1" }, actionBtns), + ]); + }); + waContactsTbody.replaceChildren( + ...(rows.length + ? rows + : [el("tr", {}, [el("td", { colspan: "4", text: currentLang === "zh" ? "暂无联系人" : "No contacts" })])]), + ); + }; + waContactFilterSel.addEventListener("change", () => { + renderWaContactRows(String(waContactFilterSel.value || "all")); + }); + renderWaContactRows("all"); + const waContactPhoneInput = el("input", { class: "input", placeholder: currentLang === "zh" ? "电话,如 +8615601877957" : "Phone, e.g. +8615601877957" }); + const waContactNameInput = el("input", { class: "input", placeholder: currentLang === "zh" ? "显示名(可选)" : "Display name (optional)" }); + const waContactTypeSel = el("select", { class: "input" }, [ + el("option", { value: "admin", text: currentLang === "zh" ? "管理员" : "Admin" }), + el("option", { value: "whitelist", text: currentLang === "zh" ? "白名单" : "Whitelist" }), + el("option", { value: "blacklist", text: currentLang === "zh" ? "黑名单" : "Blacklist" }), + ]); + const waPendingPickSel = el( + "select", + { class: "input" }, + [ + el("option", { value: "", text: currentLang === "zh" ? "从待审批选择联系人…" : "Pick from pending requests…" }), + ...waPending.map((p, idx) => { + const phone = waPhoneDisplay(p); + const name = String((p && p.push_name) || "").trim(); + const label = name ? `${name} (${phone})` : phone; + return el("option", { value: String(idx), text: label || phone }); + }), + ], + ); + waPendingPickSel.addEventListener("change", () => { + const idx = Number(String(waPendingPickSel.value || "").trim()); + if (!Number.isFinite(idx) || idx < 0 || idx >= waPending.length) return; + const picked = waPending[idx] || {}; + const phone = waPhoneDisplay(picked); + const name = String((picked && picked.push_name) || "").trim(); + if (phone && phone !== "-") waContactPhoneInput.value = phone; + waContactNameInput.value = name || ""; + }); + const waResolvePending = async (pendingId, action) => { + if (!pendingId) return; + try { + const resp = await apiPost("/admin/api/whatsapp/access/pending/resolve", { + pending_id: pendingId, + action, + }); + if (!resp || resp.ok === false) { + window.alert(String((resp && resp.error) || (currentLang === "zh" ? "操作失败" : "Request failed"))); + return; + } + router(); + } catch (err) { + window.alert(String(err)); + } + }; + const waDeletePending = async (pendingId) => { + if (!pendingId) return; + try { + const resp = await apiDeleteJson( + `/admin/api/whatsapp/access/pending?pending_id=${encodeURIComponent(pendingId)}`, + ); + if (!resp || resp.ok === false) { + window.alert(String((resp && resp.error) || (currentLang === "zh" ? "删除失败" : "Delete failed"))); + return; + } + router(); + } catch (err) { + window.alert(String(err)); + } + }; + const waDeniedRows = waDenied.map((p) => { + const pendingId = String((p && p.id) || "").trim(); + const phone = waPhoneDisplay(p); + const name = String((p && p.push_name) || "").trim(); + return el("tr", {}, [ + el("td", { text: name || "-" }), + el("td", { text: phone }), + el("td", { text: String((p && p.request_text) || "").slice(0, 80) }), + el("td", { text: String((p && p.resolved_at) || p.created_at || "") }), + el("td", {}, [ + el("button", { + class: "btn btn--primary", + text: currentLang === "zh" ? "加白名单" : "Whitelist", + onclick: async () => { await waResolvePending(pendingId, "approve"); }, + }), + el("span", { text: " " }), + el("button", { + class: "btn btn--danger", + text: currentLang === "zh" ? "删除" : "Delete", + onclick: async () => { await waDeletePending(pendingId); }, + }), + ]), + ]); + }); + const waPendingRows = waPending.map((p) => { + const pendingId = String((p && p.id) || "").trim(); + return el("tr", {}, [ + el("td", { text: String((p && p.push_name) || "-") }), + el("td", { text: waPhoneDisplay(p) }), + el("td", { text: String((p && p.external_user_id) || "") }), + el("td", { text: String((p && p.request_text) || "").slice(0, 80) }), + el("td", { text: String((p && p.created_at) || "") }), + el("td", {}, [ + el("button", { + class: "btn btn--primary", + text: currentLang === "zh" ? "通过" : "Approve", + onclick: async () => { await waResolvePending(pendingId, "approve"); }, + }), + el("span", { text: " " }), + el("button", { + class: "btn", + text: currentLang === "zh" ? "拒绝" : "Deny", + onclick: async () => { await waResolvePending(pendingId, "deny"); }, + }), + el("span", { text: " " }), + el("button", { + class: "btn btn--danger", + text: currentLang === "zh" ? "删除" : "Delete", + onclick: async () => { await waDeletePending(pendingId); }, + }), + ]), + ]); + }); + const whatsappAccessCard = el("div", { class: "card" }, [ + el("div", { class: "card__title", text: currentLang === "zh" ? "WhatsApp 访问控制" : "WhatsApp access control" }), + el("div", { class: "muted", text: currentLang === "zh" ? "默认黑名单模式:未授权用户进入待审批;拒绝后进入黑名单/已拒绝列表,可一键加白名单。" : "Default blacklist mode: unauthorized users go to Pending; denied users go to blacklist/denied list and can be whitelisted." }), + el("div", { class: "row" }, [ + el("label", { text: currentLang === "zh" ? "模式" : "Mode" }), + waAccessModeSel, + el("label", { text: currentLang === "zh" ? "提示语言" : "Message lang" }), + waAccessLangSel, + el("button", { + class: "btn btn--primary", + text: currentLang === "zh" ? "保存配置" : "Save config", + onclick: async () => { + const resp = await apiPost("/admin/api/whatsapp/access/config", { + tenant_id: "default", + access_mode: String(waAccessModeSel.value || "blacklist"), + lang: String(waAccessLangSel.value || "en"), + }); + const cfg = (resp && resp.config) || {}; + waAccessStatus.textContent = `mode=${String(cfg.access_mode || "blacklist")} lang=${String(cfg.lang || "en")}`; + }, + }), + ]), + waAccessStatus, + el("div", { class: "row" }, [ + el("label", { text: currentLang === "zh" ? "筛选" : "Filter" }), + waContactFilterSel, + ]), + el("div", { class: "row" }, [waPendingPickSel]), + el("div", { class: "row" }, [waContactPhoneInput, waContactNameInput, waContactTypeSel]), + el("div", { class: "row" }, [ + el("button", { + class: "btn", + text: currentLang === "zh" ? "添加联系人" : "Add contact", + onclick: async () => { + const phone = String(waContactPhoneInput.value || "").trim(); + if (!phone) return; + try { + const resp = await apiPost("/admin/api/whatsapp/access/contacts", { + phone, + push_name: String(waContactNameInput.value || "").trim(), + list_type: String(waContactTypeSel.value || "whitelist"), + }); + if (!resp || resp.ok === false) { + window.alert(String((resp && resp.error) || (currentLang === "zh" ? "添加失败" : "Add failed"))); + return; + } + router(); + } catch (err) { + window.alert(String(err)); + } + }, + }), + ]), + el("table", { class: "table" }, [ + el("thead", {}, [el("tr", {}, [ + el("th", { text: currentLang === "zh" ? "类型" : "Type" }), + el("th", { text: currentLang === "zh" ? "名称" : "Name" }), + el("th", { text: currentLang === "zh" ? "电话" : "Phone" }), + el("th", { text: currentLang === "zh" ? "操作" : "Action" }), + ])]), + waContactsTbody, + ]), + el("div", { class: "card__title", text: currentLang === "zh" ? "待审批请求" : "Pending requests" }), + el("table", { class: "table" }, [ + el("thead", {}, [el("tr", {}, [ + el("th", { text: currentLang === "zh" ? "名称" : "Name" }), + el("th", { text: currentLang === "zh" ? "电话" : "Phone" }), + el("th", { text: "JID" }), + el("th", { text: currentLang === "zh" ? "消息" : "Message" }), + el("th", { text: currentLang === "zh" ? "时间" : "Time" }), + el("th", { text: currentLang === "zh" ? "操作" : "Action" }), + ])]), + el("tbody", {}, waPendingRows.length ? waPendingRows : [el("tr", {}, [el("td", { colspan: "6", text: currentLang === "zh" ? "无待审批" : "None" })])]), + ]), + el("div", { class: "card__title", text: currentLang === "zh" ? "已拒绝 / 黑名单记录" : "Denied / blacklist records" }), + el("div", { class: "muted", text: currentLang === "zh" ? "点「拒绝」后用户会进入黑名单联系人,并保留在此列表;可一键加白名单恢复访问。" : "Denied users are blacklisted and listed here; use Whitelist to restore access." }), + el("table", { class: "table" }, [ + el("thead", {}, [el("tr", {}, [ + el("th", { text: currentLang === "zh" ? "名称" : "Name" }), + el("th", { text: currentLang === "zh" ? "电话" : "Phone" }), + el("th", { text: currentLang === "zh" ? "消息" : "Message" }), + el("th", { text: currentLang === "zh" ? "拒绝时间" : "Denied at" }), + el("th", { text: currentLang === "zh" ? "操作" : "Action" }), + ])]), + el("tbody", {}, waDeniedRows.length ? waDeniedRows : [el("tr", {}, [el("td", { colspan: "5", text: currentLang === "zh" ? "无已拒绝记录" : "None" })])]), + ]), + ]); 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." }), @@ -1908,6 +2298,7 @@ async function renderStack() { ]), weixinDispatchCard, whatsappDispatchCard, + whatsappAccessCard, whatsappAlertBindingCard, el("div", { class: "card" }, [ el("div", { class: "card__title", text: currentLang === "zh" ? "提示词/工具预热" : "Prompt/Tool Prewarm" }), diff --git a/interfaces/http/fastapi_app.py b/interfaces/http/fastapi_app.py index c83fd870..e5ef5cb4 100644 --- a/interfaces/http/fastapi_app.py +++ b/interfaces/http/fastapi_app.py @@ -307,8 +307,14 @@ def create_app() -> FastAPI: return {"ok": False, "error": "missing id"} ok = bool(body.get("ok", True)) err = str(body.get("error") or "").strip() + stanza_id = str(body.get("stanza_id") or "").strip() store = get_assistant_store() - changed = store.ack_channel_outbound_message(message_id=msg_id, ok=ok, error=err) + changed = store.ack_channel_outbound_message( + message_id=msg_id, + ok=ok, + error=err, + stanza_id=stanza_id, + ) return {"ok": changed} @app.get("/weixin/outbound/pending") diff --git a/runtime/application/gateway/inbound_service.py b/runtime/application/gateway/inbound_service.py index cde774df..1c493df2 100644 --- a/runtime/application/gateway/inbound_service.py +++ b/runtime/application/gateway/inbound_service.py @@ -58,7 +58,7 @@ def _resolve_channel_dispatch(store: Any, *, channel: str, account: dict[str, An _get_channel_dispatch_setting(store, _CHANNEL_DISPATCH_SPECIALIST_KEY_PREFIX, ch) or "generalist" ) lang = _normalize_channel_dispatch_lang( - _get_channel_dispatch_setting(store, _CHANNEL_DISPATCH_LANG_KEY_PREFIX, ch) or "auto" + _get_channel_dispatch_setting(store, _CHANNEL_DISPATCH_LANG_KEY_PREFIX, ch) or ("en" if ch == "whatsapp" else "auto") ) cfg = (account or {}).get("config") if isinstance(cfg, dict): @@ -820,6 +820,17 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]: text[:120], ) return {"ok": True, "replies": []} + if str(inbound.channel or "").strip().lower() == "whatsapp": + from runtime.application.gateway.whatsapp_inbound_access import handle_whatsapp_access + + access_out = handle_whatsapp_access( + store, + inbound=inbound, + account_id=account_id, + text=text, + ) + if access_out is not None: + return access_out reply = "" reply_attachments: list[dict[str, Any]] = [] ident = store.resolve_user_by_channel_identity_v2( @@ -829,7 +840,7 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]: ) if not ident: owner = _ensure_administrator_owner(store) - if owner: + if owner and str(inbound.channel or "").strip().lower() != "whatsapp": store.upsert_user_channel_account( tenant_id=str(owner.get("tenant_id") or ""), user_id=str(owner.get("user_id") or ""), diff --git a/runtime/application/gateway/whatsapp_inbound_access.py b/runtime/application/gateway/whatsapp_inbound_access.py new file mode 100644 index 00000000..de717321 --- /dev/null +++ b/runtime/application/gateway/whatsapp_inbound_access.py @@ -0,0 +1,512 @@ +from __future__ import annotations + +import json +import uuid +from typing import Any + +from runtime.extensions.whatsapp.access_control import ( + admin_approval_result_text, + admin_notify_text, + default_access_lang, + default_access_mode, + denied_reply_text, + extract_participant_alt, + extract_push_name, + extract_quote_context, + extract_remote_jid_alt, + is_access_allowed, + normalize_whatsapp_phone, + parse_admin_access_command, + parse_admin_approval_intent, + parse_pending_id_from_notify_text, + phone_from_jid, + resolve_sender_phone, + resolve_whatsapp_sender_jid, + whatsapp_sender_lookup_jids, +) +from runtime.extensions.whatsapp.api import is_whatsapp_user_target +from runtime.extensions.whatsapp.tenant import resolve_whatsapp_tenant_id + + +def _extract_contact_fields(inbound: Any) -> tuple[str, str, str, str]: + meta = inbound.metadata if isinstance(inbound.metadata, dict) else {} + push_name = extract_push_name(meta) + raw_jid = str(inbound.external_user_id or "").strip() + participant_alt = extract_participant_alt(meta) + remote_jid_alt = extract_remote_jid_alt(meta) + canonical_jid = resolve_whatsapp_sender_jid(raw_jid, meta) + phone = resolve_sender_phone(raw_jid, participant_alt, remote_jid_alt) + return push_name, phone, raw_jid, canonical_jid + + +def _link_identity_aliases( + store: Any, + *, + tenant_id: str, + account_id: str, + user_id: str, + jids: tuple[str, ...], +) -> None: + for jid in jids: + val = str(jid or "").strip() + if not val: + continue + store.upsert_channel_identity_v2( + tenant_id=tenant_id, + channel="whatsapp", + account_id=account_id, + external_user_id=val, + user_id=user_id, + ) + + +def _ensure_admin_identity( + store: Any, + *, + tenant_id: str, + account_id: str, + external_user_id: str, + participant_alt: str, + push_name: str, + phone: str, +) -> dict[str, Any] | None: + lookup_jids = whatsapp_sender_lookup_jids(external_user_id, participant_alt) + for jid in lookup_jids: + ident = store.resolve_user_by_channel_identity_v2( + channel="whatsapp", + account_id=account_id, + external_user_id=jid, + ) + if ident: + _link_identity_aliases( + store, + tenant_id=tenant_id, + account_id=account_id, + user_id=str(ident.get("user_id") or ""), + jids=lookup_jids, + ) + return ident + + admin_user = store.get_user_by_username(tenant_id=tenant_id, username="administrator") + user_id = str((admin_user or {}).get("id") or "") + if not user_id: + return _ensure_guest_identity( + store, + tenant_id=tenant_id, + account_id=account_id, + external_user_id=external_user_id, + participant_alt=participant_alt, + push_name=push_name, + phone=phone, + ) + _link_identity_aliases( + store, + tenant_id=tenant_id, + account_id=account_id, + user_id=user_id, + jids=lookup_jids, + ) + return store.resolve_user_by_channel_identity_v2( + channel="whatsapp", + account_id=account_id, + external_user_id=lookup_jids[0] if lookup_jids else external_user_id, + ) + + +def _ensure_guest_identity( + store: Any, + *, + tenant_id: str, + account_id: str, + external_user_id: str, + participant_alt: str, + push_name: str, + phone: str, +) -> dict[str, Any] | None: + lookup_jids = whatsapp_sender_lookup_jids(external_user_id, participant_alt) + for jid in lookup_jids: + ident = store.resolve_user_by_channel_identity_v2( + channel="whatsapp", + account_id=account_id, + external_user_id=jid, + ) + if ident: + _link_identity_aliases( + store, + tenant_id=tenant_id, + account_id=account_id, + user_id=str(ident.get("user_id") or ""), + jids=lookup_jids, + ) + return ident + + label = str(push_name or "").strip() or phone or external_user_id + username = f"wa_{phone or uuid.uuid4().hex[:10]}" + user = store.create_user_account( + tenant_id=tenant_id, + username=username, + display_name=label, + role="guest", + password_hash="", + is_active=True, + ) + user_id = str((user or {}).get("id") or "") + if not user_id: + return None + _link_identity_aliases( + store, + tenant_id=tenant_id, + account_id=account_id, + user_id=user_id, + jids=lookup_jids, + ) + return store.resolve_user_by_channel_identity_v2( + channel="whatsapp", + account_id=account_id, + external_user_id=lookup_jids[0] if lookup_jids else external_user_id, + ) + + +def _notify_admins( + store: Any, + *, + tenant_id: str, + account_id: str, + lang: str, + push_name: str, + external_user_id: str, + request_text: str, + pending_id: str, +) -> None: + admins = store.list_whatsapp_contacts( + tenant_id=tenant_id, + account_id=account_id, + list_type="admin", + ) + text = admin_notify_text( + lang=lang, + push_name=push_name, + external_user_id=external_user_id, + request_text=request_text, + pending_id=pending_id, + ) + for admin in admins: + chat_id = str(admin.get("external_user_id") or "").strip() + if not chat_id or not is_whatsapp_user_target(chat_id): + continue + store.enqueue_channel_outbound_message( + channel="whatsapp", + chat_id=chat_id, + text=text, + tenant_id=tenant_id, + account_id=account_id, + source=json.dumps( + { + "kind": "whatsapp_access_pending", + "pending_id": str(pending_id or ""), + }, + ensure_ascii=False, + ), + ) + + +def _upsert_whatsapp_contact_profile( + store: Any, + *, + tenant_id: str, + account_id: str, + raw_jid: str, + canonical_jid: str, + push_name: str, + phone: str, +) -> None: + store.upsert_whatsapp_contact( + tenant_id=tenant_id, + account_id=account_id, + external_user_id=raw_jid, + push_name=push_name, + phone=phone, + ) + if canonical_jid and canonical_jid != raw_jid: + store.upsert_whatsapp_contact( + tenant_id=tenant_id, + account_id=account_id, + external_user_id=canonical_jid, + push_name=push_name, + phone=phone_from_jid(canonical_jid), + ) + + +def _handle_admin_message( + store: Any, + *, + tenant_id: str, + account_id: str, + external_user_id: str, + text: str, + lang: str, + inbound: Any, +) -> dict[str, Any] | None: + cmd = parse_admin_access_command(text) + if cmd: + action = str(cmd.get("action") or "") + if action == "set_mode": + mode = str(cmd.get("access_mode") or default_access_mode()) + store.upsert_whatsapp_access_config( + tenant_id=tenant_id, + account_id=account_id, + access_mode=mode, + lang=lang, + ) + return {"text": f"Access mode set to {mode}.", "metadata": {}} + if action == "set_list": + list_type = str(cmd.get("list_type") or "") + target_raw = str(cmd.get("target") or "").strip() + try: + target_phone = normalize_whatsapp_phone(target_raw) + except Exception: + return {"text": f"Invalid target: {target_raw}", "metadata": {}} + store.apply_whatsapp_contact_access( + tenant_id=tenant_id, + account_id=account_id, + external_user_id=target_phone, + phone=target_phone, + list_type=list_type, + ) + return {"text": f"Set {target_phone} as {list_type}.", "metadata": {}} + + meta = inbound.metadata if isinstance(inbound.metadata, dict) else {} + quote = extract_quote_context(meta) + stanza_id = str(quote.get("stanza_id") or "").strip() + quoted_text = str(quote.get("quoted_text") or "").strip() + if not stanza_id and not quoted_text: + return None + + admin_chat_id = str(getattr(inbound, "external_chat_id", None) or external_user_id or "").strip() + pending_id = "" + if stanza_id: + found = store.find_whatsapp_pending_by_notify_stanza( + tenant_id=tenant_id, + account_id=account_id, + admin_chat_id=admin_chat_id, + notify_stanza_id=stanza_id, + ) + if found: + pending_id = str(found) + + if not pending_id and quoted_text: + fallback = parse_pending_id_from_notify_text(quoted_text) + if fallback: + pending_id = fallback + + if not pending_id: + return None + + item = store.get_whatsapp_access_pending_by_id(pending_id=pending_id) + if not item or str(item.get("status") or "").strip().lower() != "pending": + return None + + intent = parse_admin_approval_intent(text) + if not intent: + return None + + target_jid = str(item.get("external_user_id") or "").strip() + target_name = str(item.get("push_name") or "").strip() + target_phone = str(item.get("phone") or "").strip() or resolve_sender_phone(target_jid) + approved = intent == "approve" + if approved: + store.apply_whatsapp_contact_access( + tenant_id=tenant_id, + account_id=account_id, + external_user_id=target_phone, + push_name=target_name, + phone=target_phone, + list_type="whitelist", + ) + store.resolve_whatsapp_access_pending( + pending_id=pending_id, + status="approved", + resolved_by=external_user_id, + ) + else: + if target_phone: + store.apply_whatsapp_contact_access( + tenant_id=tenant_id, + account_id=account_id, + external_user_id=target_phone, + push_name=target_name, + phone=target_phone, + list_type="blacklist", + ) + store.resolve_whatsapp_access_pending( + pending_id=pending_id, + status="denied", + resolved_by=external_user_id, + ) + deleter = getattr(store, "delete_whatsapp_pending_notify_for_pending", None) + if callable(deleter): + deleter(pending_id=pending_id) + + raw = meta.get("raw") if isinstance(meta.get("raw"), dict) else {} + reply_meta: dict[str, Any] = { + "quote_remote_jid": admin_chat_id, + "quote_stanza_id": str(raw.get("id") or "").strip(), + "quote_participant": str(raw.get("participant") or external_user_id or "").strip(), + "quote_text": str(text or "").strip(), + } + return { + "text": admin_approval_result_text( + lang=lang, + approved=approved, + push_name=target_name, + external_user_id=target_jid, + ), + "metadata": reply_meta, + } + + +def handle_whatsapp_access( + store: Any, + *, + inbound: Any, + account_id: str, + text: str, +) -> dict[str, Any] | None: + """Return an inbound response dict when access gate short-circuits normal processing.""" + tenant_id = resolve_whatsapp_tenant_id(store, account_id=account_id) + meta = inbound.metadata if isinstance(inbound.metadata, dict) else {} + participant_alt = extract_participant_alt(meta) + remote_jid_alt = extract_remote_jid_alt(meta) + push_name, phone, raw_jid, canonical_jid = _extract_contact_fields(inbound) + if not raw_jid: + return None + + cfg = store.get_whatsapp_access_config(tenant_id=tenant_id, account_id=account_id) + access_mode = str((cfg or {}).get("access_mode") or default_access_mode()) + lang = str((cfg or {}).get("lang") or default_access_lang()) + + matched = store.find_whatsapp_contact_for_sender( + tenant_id=tenant_id, + account_id=account_id, + external_user_id=raw_jid, + participant_alt=participant_alt, + remote_jid_alt=remote_jid_alt, + ) + list_type = str((matched or {}).get("list_type") or "").strip().lower() or None + + if list_type == "admin": + admin_reply = _handle_admin_message( + store, + tenant_id=tenant_id, + account_id=account_id, + external_user_id=raw_jid, + text=text, + lang=lang, + inbound=inbound, + ) + if admin_reply: + return { + "ok": True, + "replies": [ + { + "channel": "whatsapp", + "chat_id": inbound.external_chat_id, + "text": str(admin_reply.get("text") or ""), + "attachments": [], + "metadata": admin_reply.get("metadata") + if isinstance(admin_reply.get("metadata"), dict) + else {}, + } + ], + "whatsapp_access": "admin_command", + } + _upsert_whatsapp_contact_profile( + store, + tenant_id=tenant_id, + account_id=account_id, + raw_jid=raw_jid, + canonical_jid=canonical_jid, + push_name=push_name, + phone=phone, + ) + _ensure_admin_identity( + store, + tenant_id=tenant_id, + account_id=account_id, + external_user_id=raw_jid, + participant_alt=participant_alt, + push_name=push_name, + phone=phone, + ) + return None + + if is_access_allowed(access_mode=access_mode, list_type=list_type): + _upsert_whatsapp_contact_profile( + store, + tenant_id=tenant_id, + account_id=account_id, + raw_jid=raw_jid, + canonical_jid=canonical_jid, + push_name=push_name, + phone=phone, + ) + _ensure_guest_identity( + store, + tenant_id=tenant_id, + account_id=account_id, + external_user_id=raw_jid, + participant_alt=participant_alt, + push_name=push_name, + phone=phone, + ) + return None + + pending_id = store.create_whatsapp_access_pending( + tenant_id=tenant_id, + account_id=account_id, + external_user_id=raw_jid, + push_name=push_name, + phone=phone, + request_text=text, + ) + if pending_id: + _notify_admins( + store, + tenant_id=tenant_id, + account_id=account_id, + lang=lang, + push_name=push_name, + external_user_id=raw_jid, + request_text=text, + pending_id=pending_id, + ) + + reply_meta: dict[str, Any] = {} + if bool(getattr(inbound, "is_group", False)): + raw = meta.get("raw") if isinstance(meta.get("raw"), dict) else {} + stanza_id = str(raw.get("id") or "").strip() + quote_participant = str(raw.get("participant") or raw_jid or "").strip() + reply_meta = { + "mention_jids": [raw_jid], + "quote_remote_jid": str(inbound.external_chat_id or "").strip(), + "quote_stanza_id": stanza_id, + "quote_participant": quote_participant, + "quote_text": str(text or "").strip(), + } + + return { + "ok": True, + "replies": [ + { + "channel": "whatsapp", + "chat_id": inbound.external_chat_id, + "text": denied_reply_text(lang=lang), + "attachments": [], + "metadata": reply_meta, + } + ], + "whatsapp_access": "denied", + } + + +__all__ = ["handle_whatsapp_access"] diff --git a/runtime/extensions/whatsapp/access_control.py b/runtime/extensions/whatsapp/access_control.py new file mode 100644 index 00000000..00f2fb11 --- /dev/null +++ b/runtime/extensions/whatsapp/access_control.py @@ -0,0 +1,392 @@ +from __future__ import annotations + +import re +from typing import Any, Literal + +from runtime.extensions.whatsapp.api import normalize_whatsapp_target + +AccessMode = Literal["blacklist", "whitelist"] +ListType = Literal["admin", "whitelist", "blacklist"] +ApprovalIntent = Literal["approve", "deny"] + +_AFFIRMATIVE = frozenset( + { + "yes", + "y", + "approve", + "allow", + "add", + "ok", + "grant", + "whitelist", + "同意", + "可以", + "是", + "添加", + "通过", + } +) +_NEGATIVE = frozenset( + { + "no", + "n", + "deny", + "reject", + "ignore", + "拒绝", + "不", + "否", + "忽略", + } +) + + +def default_access_mode() -> AccessMode: + return "blacklist" + + +def default_access_lang() -> str: + return "en" + + +def is_access_allowed(*, access_mode: str, list_type: str | None) -> bool: + lt = str(list_type or "").strip().lower() + if lt == "admin": + return True + mode = str(access_mode or default_access_mode()).strip().lower() + if mode == "whitelist": + return lt != "blacklist" + return lt == "whitelist" + + +def extract_push_name(metadata: dict[str, Any] | None) -> str: + if not isinstance(metadata, dict): + return "" + for key in ("push_name", "display_name", "pushName"): + val = metadata.get(key) + if val is not None and str(val).strip(): + return str(val).strip() + raw = metadata.get("raw") + if isinstance(raw, dict): + for key in ("pushName", "push_name", "display_name"): + val = raw.get(key) + if val is not None and str(val).strip(): + return str(val).strip() + return "" + + +def is_whatsapp_lid_jid(jid: str) -> bool: + return str(jid or "").strip().lower().endswith("@lid") + + +def phone_from_jid(jid: str) -> str: + if is_whatsapp_lid_jid(jid): + return "" + base = str(jid or "").split("@", 1)[0].strip() + digits = re.sub(r"\D", "", base) + return digits or base + + +def normalize_whatsapp_phone(value: str) -> str: + raw = str(value or "").strip() + if not raw: + raise ValueError("phone is required") + digits = phone_from_jid(raw) if "@" in raw else re.sub(r"\D", "", raw.lstrip("+")) + if len(digits) < 6: + raise ValueError(f"invalid phone: {value}") + return digits + + +def resolve_sender_phone( + external_user_id: str, + participant_alt: str = "", + remote_jid_alt: str = "", +) -> str: + for alt in (str(participant_alt or "").strip(), str(remote_jid_alt or "").strip()): + if alt.lower().endswith("@s.whatsapp.net"): + alt_phone = phone_from_jid(alt) + if len(alt_phone) >= 6: + return alt_phone + canonical = resolve_whatsapp_sender_jid( + external_user_id, + { + "raw": { + "participantAlt": str(participant_alt or "").strip() or None, + "remoteJidAlt": str(remote_jid_alt or "").strip() or None, + } + }, + ) + if str(canonical or "").lower().endswith("@s.whatsapp.net"): + canonical_phone = phone_from_jid(canonical) + if len(canonical_phone) >= 6: + return canonical_phone + if not is_whatsapp_lid_jid(external_user_id): + raw_phone = phone_from_jid(external_user_id) + if len(raw_phone) >= 6: + return raw_phone + return "" + + +def contact_phone_key(row: dict[str, Any] | None) -> str: + if not isinstance(row, dict): + return "" + phone = str(row.get("phone") or "").strip() + if phone: + return phone + return phone_from_jid(str(row.get("external_user_id") or "")) + + +def whatsapp_phones_match(a: str, b: str) -> bool: + left = str(a or "").strip() + right = str(b or "").strip() + if not left or not right: + return False + try: + return normalize_whatsapp_phone(left) == normalize_whatsapp_phone(right) + except Exception: + left_phone = phone_from_jid(left) + right_phone = phone_from_jid(right) + return bool(len(left_phone) >= 6 and left_phone == right_phone) + + +def jid_local_part(jid: str) -> str: + return str(jid or "").split("@", 1)[0].split(":", 1)[0].strip().lower() + + +def extract_remote_jid_alt(metadata: dict[str, Any] | None) -> str: + if not isinstance(metadata, dict): + return "" + raw = metadata.get("raw") + if isinstance(raw, dict): + for key in ("remoteJidAlt", "remote_jid_alt"): + val = raw.get(key) + if val is not None and str(val).strip(): + return str(val).strip() + return "" + + +def extract_participant_alt(metadata: dict[str, Any] | None) -> str: + if not isinstance(metadata, dict): + return "" + raw = metadata.get("raw") + if isinstance(raw, dict): + for key in ("participantAlt", "participant_alt"): + val = raw.get(key) + if val is not None and str(val).strip(): + return str(val).strip() + return "" + + +def resolve_whatsapp_sender_jid(external_user_id: str, metadata: dict[str, Any] | None = None) -> str: + jid = str(external_user_id or "").strip() + participant_alt = extract_participant_alt(metadata) + remote_jid_alt = extract_remote_jid_alt(metadata) + low = jid.lower() + for alt in (participant_alt, remote_jid_alt): + if alt and low.endswith("@lid"): + try: + return normalize_whatsapp_target(alt) + except Exception: + pass + if low.endswith("@s.whatsapp.net"): + return low + if low.endswith("@lid"): + for alt in (participant_alt, remote_jid_alt): + if alt: + try: + return normalize_whatsapp_target(alt) + except Exception: + pass + return jid + try: + return normalize_whatsapp_target(jid) + except Exception: + return jid + + +def whatsapp_sender_lookup_jids(external_user_id: str, participant_alt: str = "") -> tuple[str, ...]: + out: list[str] = [] + seen: set[str] = set() + canonical = resolve_whatsapp_sender_jid( + external_user_id, + {"raw": {"participantAlt": participant_alt}} if participant_alt else None, + ) + for item in (external_user_id, participant_alt, canonical): + val = str(item or "").strip() + if val and val not in seen: + seen.add(val) + out.append(val) + low = str(external_user_id or "").lower() + if low.endswith("@lid"): + alias = f"{jid_local_part(external_user_id)}@s.whatsapp.net" + if alias not in seen: + seen.add(alias) + out.append(alias) + return tuple(out) + + +def whatsapp_users_match(a: str, b: str) -> bool: + left = str(a or "").strip() + right = str(b or "").strip() + if not left or not right: + return False + if whatsapp_phones_match(left, right): + return True + if left.lower() == right.lower(): + return True + if jid_local_part(left) == jid_local_part(right): + return True + return False + + +def extract_quote_context(metadata: dict[str, Any] | None) -> dict[str, str]: + out = { + "stanza_id": "", + "quoted_text": "", + "remote_jid": "", + "participant": "", + } + if not isinstance(metadata, dict): + return out + raw = metadata.get("raw") + if not isinstance(raw, dict): + return out + out["stanza_id"] = str( + raw.get("quotedStanzaId") or raw.get("quoted_stanza_id") or "" + ).strip() + out["quoted_text"] = str( + raw.get("quotedText") or raw.get("quoted_text") or "" + ).strip() + out["participant"] = str( + raw.get("quotedParticipant") + or raw.get("quoted_participant") + or raw.get("participant") + or "" + ).strip() + for key in ("quotedRemoteJid", "quoted_remote_jid", "remoteJid"): + val = raw.get(key) + if val is not None and str(val).strip(): + out["remote_jid"] = str(val).strip() + break + return out + + +_PENDING_ID_IN_NOTIFY_RE = re.compile( + r"(?:请求编号|Request)\s*[::]\s*([0-9a-f]{16,64})", + re.IGNORECASE, +) + + +def parse_pending_id_from_notify_text(text: str) -> str | None: + m = _PENDING_ID_IN_NOTIFY_RE.search(str(text or "")) + if not m: + return None + return str(m.group(1) or "").strip() or None + + +def admin_approval_quote_required_text(*, lang: str) -> str: + if str(lang or "").strip().lower().startswith("zh"): + return "请引用待审批通知消息后回复「同意」或「拒绝」。" + return "Reply YES or NO by quoting the pending access notification." + + +def admin_approval_unknown_notify_text(*, lang: str) -> str: + if str(lang or "").strip().lower().startswith("zh"): + return "无法识别该审批请求,请引用正确的通知消息。" + return "Cannot identify this approval request. Please quote the correct notification message." + + +def parse_admin_approval_intent(text: str) -> ApprovalIntent | None: + blob = str(text or "").strip().lower() + if not blob: + return None + first = re.split(r"[\s,,。.!!??]+", blob, maxsplit=1)[0].strip() + if first in _AFFIRMATIVE or any(tok in blob for tok in ("add to whitelist", "grant access", "添加白名单", "加入白名单")): + return "approve" + if first in _NEGATIVE or any(tok in blob for tok in ("do not add", "don't add", "不要添加", "不加")): + return "deny" + return None + + +def parse_admin_access_command(text: str) -> dict[str, Any] | None: + raw = str(text or "").strip() + if not raw: + return None + low = raw.lower() + + mode_match = re.match(r"^(?:access\s+mode|模式)\s+(blacklist|whitelist|黑名单|白名单)\b", low) + if mode_match: + token = mode_match.group(1) + if token in {"黑名单", "blacklist"}: + return {"action": "set_mode", "access_mode": "blacklist"} + return {"action": "set_mode", "access_mode": "whitelist"} + + for verb, list_type in ( + (r"^(?:whitelist|白名单)\s+(?:add\s+)?(.+)$", "whitelist"), + (r"^(?:blacklist|黑名单)\s+(?:add\s+)?(.+)$", "blacklist"), + (r"^(?:admin|管理员)\s+(?:add\s+)?(.+)$", "admin"), + ): + m = re.match(verb, raw, flags=re.IGNORECASE) + if m: + target = str(m.group(1) or "").strip() + if target: + return {"action": "set_list", "list_type": list_type, "target": target} + return None + + +def normalize_contact_jid(value: str) -> str: + return normalize_whatsapp_target(str(value or "").strip()) + + +def coerce_whatsapp_access_target(value: str) -> str: + """Normalize manual contact input to canonical WhatsApp user JID from phone.""" + return normalize_whatsapp_target(normalize_whatsapp_phone(value)) + + +def denied_reply_text(*, lang: str) -> str: + if str(lang or "").strip().lower().startswith("zh"): + return "无权限:您尚未获得使用此助手的授权。请联系管理员。" + return "Access denied: you are not authorized to use this assistant. Please contact an administrator." + + +def admin_notify_text( + *, + lang: str, + push_name: str, + external_user_id: str, + request_text: str, + pending_id: str, +) -> str: + phone = phone_from_jid(external_user_id) + name = str(push_name or "").strip() or phone or external_user_id + preview = str(request_text or "").strip() + if len(preview) > 160: + preview = preview[:157] + "..." + if str(lang or "").strip().lower().startswith("zh"): + return ( + f"[oclaw] 未授权用户请求访问\n" + f"用户: {name}\n" + f"ID: {external_user_id}\n" + f"消息: {preview or '(empty)'}\n" + f"请求编号: {pending_id}\n" + f"请引用本条消息回复「同意」或「拒绝」。" + ) + return ( + f"[oclaw] Unauthorized access request\n" + f"User: {name}\n" + f"ID: {external_user_id}\n" + f"Message: {preview or '(empty)'}\n" + f"Request: {pending_id}\n" + f"Reply YES or NO by quoting this message." + ) + + +def admin_approval_result_text(*, lang: str, approved: bool, push_name: str, external_user_id: str) -> str: + name = str(push_name or "").strip() or phone_from_jid(external_user_id) or external_user_id + if str(lang or "").strip().lower().startswith("zh"): + if approved: + return f"已将 {name} ({external_user_id}) 加入白名单。" + return f"已忽略 {name} ({external_user_id}) 的访问请求。" + if approved: + return f"Whitelisted {name} ({external_user_id})." + return f"Ignored access request from {name} ({external_user_id})." diff --git a/runtime/extensions/whatsapp/tenant.py b/runtime/extensions/whatsapp/tenant.py new file mode 100644 index 00000000..b1f587e5 --- /dev/null +++ b/runtime/extensions/whatsapp/tenant.py @@ -0,0 +1,14 @@ +from __future__ import annotations + +import os +from typing import Any + + +def resolve_whatsapp_tenant_id(store: Any, *, account_id: str) -> str: + account = store.find_user_by_channel_account(channel="whatsapp", account_id=str(account_id or "").strip()) + if account and str(account.get("tenant_id") or "").strip(): + return str(account["tenant_id"]) + return str(os.getenv("OCLAW_DEFAULT_TENANT_ID") or "default").strip() or "default" + + +__all__ = ["resolve_whatsapp_tenant_id"] diff --git a/runtime/operations/whatsapp_bridge/baileys_runner.ts b/runtime/operations/whatsapp_bridge/baileys_runner.ts index 2058a7da..ff4eced5 100644 --- a/runtime/operations/whatsapp_bridge/baileys_runner.ts +++ b/runtime/operations/whatsapp_bridge/baileys_runner.ts @@ -162,6 +162,36 @@ function resolveSenderJid(key: proto.IMessageKey): string { return ""; } +function resolvePnFromLid(sock: ReturnType | null, lidJid: string): string { + const lid = String(lidJid || "").trim(); + if (!lid || !sock) return ""; + const lidMapping = (sock as any)?.signalRepository?.lidMapping; + if (!lidMapping || typeof lidMapping.getPNForLID !== "function") return ""; + try { + const pn = lidMapping.getPNForLID(lid); + return pn ? jidNormalizedUser(String(pn)) : ""; + } catch { + return ""; + } +} + +function resolveDmUserJid(key: proto.IMessageKey, sock: ReturnType | null): string { + const remoteJid = jidNormalizedUser(String(key.remoteJid || "")); + const remoteJidAlt = String((key as any).remoteJidAlt || "").trim(); + if (remoteJidAlt) { + const alt = jidNormalizedUser(remoteJidAlt); + if (remoteJid.toLowerCase().endsWith("@lid") && alt.toLowerCase().includes("@s.whatsapp")) { + return alt; + } + if (alt.toLowerCase().includes("@s.whatsapp")) return alt; + } + if (remoteJid.toLowerCase().endsWith("@lid")) { + const pn = resolvePnFromLid(sock, remoteJid); + if (pn) return pn; + } + return remoteJid; +} + function jidBaseLocal(jid: string): string { const s = String(jid || "").trim().toLowerCase(); if (!s) return ""; @@ -727,21 +757,31 @@ async function pollOutboundQueue(sock: ReturnType): Promise let ok = true; let err = ""; try { - await sock.sendMessage(chatId, { text }); - log(`outbound sent id=${id} chat=${chatId}`); + const sent = await sock.sendMessage(chatId, { text }); + const stanzaId = String((sent as any)?.key?.id || "").trim(); + log(`outbound sent id=${id} chat=${chatId} stanza=${stanzaId || "?"}`); + try { + await fetch(`${LOCAL_BASE_URL.replace(/\/+$/, "")}/whatsapp/outbound/ack`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ id, ok, error: err, stanza_id: stanzaId }), + }); + } catch { + // best-effort ack + } } 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 + 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) { @@ -933,10 +973,12 @@ async function main(): Promise { const isGroup = remoteJid.endsWith("@g.us"); const participantRaw = String(key.participant || "").trim(); const participantAlt = String((key as any).participantAlt || "").trim(); + const remoteJidAlt = String((key as any).remoteJidAlt || "").trim(); + const remoteJidNorm = jidNormalizedUser(remoteJid); const userId = isGroup - ? resolveSenderJid(key) || jidNormalizedUser(remoteJid) - : jidNormalizedUser(remoteJid); - const chatId = jidNormalizedUser(remoteJid); + ? resolveSenderJid(key) || remoteJidNorm + : resolveDmUserJid(key, sock); + const chatId = remoteJidNorm; const mentions = extractMentionsFromUpsert(msg); const quote = extractQuoteContext(msg.message); const botJidRaw = sock?.user?.id ? String(sock.user.id).trim() : ""; @@ -960,6 +1002,7 @@ async function main(): Promise { const raw = { id, remoteJid, + remoteJidAlt: remoteJidAlt || null, participant: participantRaw || null, participantAlt: participantAlt || null, pushName: (msg as any).pushName || null, diff --git a/svc/persistence/sqlite_store.py b/svc/persistence/sqlite_store.py index 2043d9d3..560ab308 100644 --- a/svc/persistence/sqlite_store.py +++ b/svc/persistence/sqlite_store.py @@ -658,6 +658,72 @@ class SqliteStore(ScheduledJobStoreMixin): ); """ ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS whatsapp_access_config ( + tenant_id TEXT NOT NULL, + account_id TEXT NOT NULL, + access_mode TEXT NOT NULL DEFAULT 'blacklist', + lang TEXT NOT NULL DEFAULT 'en', + updated_at TEXT NOT NULL, + PRIMARY KEY (tenant_id, account_id) + ); + """ + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS whatsapp_contact ( + tenant_id TEXT NOT NULL, + account_id TEXT NOT NULL, + external_user_id TEXT NOT NULL, + push_name TEXT NOT NULL DEFAULT '', + phone TEXT NOT NULL DEFAULT '', + list_type TEXT, + notes TEXT NOT NULL DEFAULT '', + first_seen_at TEXT NOT NULL, + last_seen_at TEXT NOT NULL, + PRIMARY KEY (tenant_id, account_id, external_user_id) + ); + """ + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS whatsapp_access_pending ( + id TEXT PRIMARY KEY, + tenant_id TEXT NOT NULL, + account_id TEXT NOT NULL, + external_user_id TEXT NOT NULL, + push_name TEXT NOT NULL DEFAULT '', + phone TEXT NOT NULL DEFAULT '', + request_text TEXT NOT NULL DEFAULT '', + status TEXT NOT NULL DEFAULT 'pending', + created_at TEXT NOT NULL, + resolved_at TEXT, + resolved_by TEXT NOT NULL DEFAULT '' + ); + """ + ) + self._ensure_whatsapp_access_pending_phone_column(conn) + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_whatsapp_access_pending_status ON whatsapp_access_pending(tenant_id, account_id, status, created_at)" + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS whatsapp_access_pending_notify ( + tenant_id TEXT NOT NULL, + account_id TEXT NOT NULL, + admin_chat_id TEXT NOT NULL, + notify_stanza_id TEXT NOT NULL, + pending_id TEXT NOT NULL, + notify_text TEXT NOT NULL DEFAULT '', + created_at TEXT NOT NULL, + PRIMARY KEY (tenant_id, account_id, admin_chat_id, notify_stanza_id) + ); + """ + ) + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_whatsapp_pending_notify_pending ON whatsapp_access_pending_notify(pending_id)" + ) conn.execute( """ CREATE TABLE IF NOT EXISTS channel_outbound_message ( @@ -1248,6 +1314,72 @@ class SqliteStore(ScheduledJobStoreMixin): ) """ ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS whatsapp_access_config ( + tenant_id TEXT NOT NULL, + account_id TEXT NOT NULL, + access_mode TEXT NOT NULL DEFAULT 'blacklist', + lang TEXT NOT NULL DEFAULT 'en', + updated_at TEXT NOT NULL, + PRIMARY KEY (tenant_id, account_id) + ) + """ + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS whatsapp_contact ( + tenant_id TEXT NOT NULL, + account_id TEXT NOT NULL, + external_user_id TEXT NOT NULL, + push_name TEXT NOT NULL DEFAULT '', + phone TEXT NOT NULL DEFAULT '', + list_type TEXT, + notes TEXT NOT NULL DEFAULT '', + first_seen_at TEXT NOT NULL, + last_seen_at TEXT NOT NULL, + PRIMARY KEY (tenant_id, account_id, external_user_id) + ) + """ + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS whatsapp_access_pending ( + id TEXT PRIMARY KEY, + tenant_id TEXT NOT NULL, + account_id TEXT NOT NULL, + external_user_id TEXT NOT NULL, + push_name TEXT NOT NULL DEFAULT '', + phone TEXT NOT NULL DEFAULT '', + request_text TEXT NOT NULL DEFAULT '', + status TEXT NOT NULL DEFAULT 'pending', + created_at TEXT NOT NULL, + resolved_at TEXT, + resolved_by TEXT NOT NULL DEFAULT '' + ) + """ + ) + self._ensure_whatsapp_access_pending_phone_column(conn) + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_whatsapp_access_pending_status ON whatsapp_access_pending(tenant_id, account_id, status, created_at)" + ) + conn.execute( + """ + CREATE TABLE IF NOT EXISTS whatsapp_access_pending_notify ( + tenant_id TEXT NOT NULL, + account_id TEXT NOT NULL, + admin_chat_id TEXT NOT NULL, + notify_stanza_id TEXT NOT NULL, + pending_id TEXT NOT NULL, + notify_text TEXT NOT NULL DEFAULT '', + created_at TEXT NOT NULL, + PRIMARY KEY (tenant_id, account_id, admin_chat_id, notify_stanza_id) + ); + """ + ) + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_whatsapp_pending_notify_pending ON whatsapp_access_pending_notify(pending_id)" + ) conn.execute( """ CREATE TABLE IF NOT EXISTS channel_outbound_message ( @@ -5583,6 +5715,621 @@ class SqliteStore(ScheduledJobStoreMixin): out = self.get_whatsapp_alert_binding(tenant_id=tenant_id, account_id=account_id) return out or {} + def get_whatsapp_access_config( + self, + *, + tenant_id: str, + account_id: str, + ) -> dict[str, Any] | None: + with self._connect() as conn: + row = conn.execute( + """ + SELECT tenant_id, account_id, access_mode, lang, updated_at + FROM whatsapp_access_config + WHERE tenant_id = ? AND account_id = ? + """, + (str(tenant_id), str(account_id)), + ).fetchone() + if not row: + return None + return { + "tenant_id": str(row["tenant_id"] or ""), + "account_id": str(row["account_id"] or ""), + "access_mode": str(row["access_mode"] or "blacklist"), + "lang": str(row["lang"] or "en"), + "updated_at": str(row["updated_at"] or ""), + } + + def upsert_whatsapp_access_config( + self, + *, + tenant_id: str, + account_id: str, + access_mode: str, + lang: str = "en", + ) -> dict[str, Any]: + ts = utc_now_iso() + mode = str(access_mode or "blacklist").strip().lower() + if mode not in {"blacklist", "whitelist"}: + mode = "blacklist" + lang_val = str(lang or "en").strip().lower() or "en" + with self._connect() as conn: + conn.execute( + """ + INSERT INTO whatsapp_access_config (tenant_id, account_id, access_mode, lang, updated_at) + VALUES (?, ?, ?, ?, ?) + ON CONFLICT(tenant_id, account_id) DO UPDATE SET + access_mode = excluded.access_mode, + lang = excluded.lang, + updated_at = excluded.updated_at + """, + (str(tenant_id), str(account_id), mode, lang_val, ts), + ) + return self.get_whatsapp_access_config(tenant_id=tenant_id, account_id=account_id) or {} + + def upsert_whatsapp_contact( + self, + *, + tenant_id: str, + account_id: str, + external_user_id: str, + push_name: str = "", + phone: str = "", + list_type: str | None = None, + notes: str = "", + ) -> dict[str, Any]: + ts = utc_now_iso() + lt = str(list_type).strip().lower() if list_type is not None else None + if lt is not None and lt not in {"admin", "whitelist", "blacklist"}: + lt = None + with self._connect() as conn: + conn.execute( + """ + INSERT INTO whatsapp_contact + (tenant_id, account_id, external_user_id, push_name, phone, list_type, notes, first_seen_at, last_seen_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(tenant_id, account_id, external_user_id) DO UPDATE SET + push_name = CASE WHEN excluded.push_name != '' THEN excluded.push_name ELSE whatsapp_contact.push_name END, + phone = CASE WHEN excluded.phone != '' THEN excluded.phone ELSE whatsapp_contact.phone END, + list_type = COALESCE(excluded.list_type, whatsapp_contact.list_type), + notes = CASE WHEN excluded.notes != '' THEN excluded.notes ELSE whatsapp_contact.notes END, + last_seen_at = excluded.last_seen_at + """, + ( + str(tenant_id), + str(account_id), + str(external_user_id), + str(push_name or ""), + str(phone or ""), + lt, + str(notes or ""), + ts, + ts, + ), + ) + return self.get_whatsapp_contact( + tenant_id=tenant_id, + account_id=account_id, + external_user_id=external_user_id, + ) or {} + + def get_whatsapp_contact( + self, + *, + tenant_id: str, + account_id: str, + external_user_id: str, + ) -> dict[str, Any] | None: + with self._connect() as conn: + row = conn.execute( + """ + SELECT tenant_id, account_id, external_user_id, push_name, phone, list_type, notes, first_seen_at, last_seen_at + FROM whatsapp_contact + WHERE tenant_id = ? AND account_id = ? AND external_user_id = ? + """, + (str(tenant_id), str(account_id), str(external_user_id)), + ).fetchone() + if not row: + return None + return { + "tenant_id": str(row["tenant_id"] or ""), + "account_id": str(row["account_id"] or ""), + "external_user_id": str(row["external_user_id"] or ""), + "push_name": str(row["push_name"] or ""), + "phone": str(row["phone"] or ""), + "list_type": str(row["list_type"] or "") or None, + "notes": str(row["notes"] or ""), + "first_seen_at": str(row["first_seen_at"] or ""), + "last_seen_at": str(row["last_seen_at"] or ""), + } + + def list_whatsapp_contacts( + self, + *, + tenant_id: str, + account_id: str, + list_type: str | None = None, + limit: int = 200, + ) -> list[dict[str, Any]]: + lim = max(1, min(int(limit), 500)) + clauses = ["tenant_id = ?", "account_id = ?"] + params: list[Any] = [str(tenant_id), str(account_id)] + if list_type: + clauses.append("list_type = ?") + params.append(str(list_type)) + where = " AND ".join(clauses) + with self._connect() as conn: + rows = conn.execute( + f""" + SELECT tenant_id, account_id, external_user_id, push_name, phone, list_type, notes, first_seen_at, last_seen_at + FROM whatsapp_contact + WHERE {where} + ORDER BY last_seen_at DESC + LIMIT ? + """, + (*params, lim), + ).fetchall() + return [ + { + "tenant_id": str(r["tenant_id"] or ""), + "account_id": str(r["account_id"] or ""), + "external_user_id": str(r["external_user_id"] or ""), + "push_name": str(r["push_name"] or ""), + "phone": str(r["phone"] or ""), + "list_type": str(r["list_type"] or "") or None, + "notes": str(r["notes"] or ""), + "first_seen_at": str(r["first_seen_at"] or ""), + "last_seen_at": str(r["last_seen_at"] or ""), + } + for r in rows + ] + + def delete_whatsapp_contact( + self, + *, + tenant_id: str, + account_id: str, + external_user_id: str, + ) -> bool: + with self._connect() as conn: + cur = conn.execute( + """ + DELETE FROM whatsapp_contact + WHERE tenant_id = ? AND account_id = ? AND external_user_id = ? + """, + (str(tenant_id), str(account_id), str(external_user_id)), + ) + return bool(cur.rowcount and cur.rowcount > 0) + + def delete_whatsapp_contact_aliases( + self, + *, + tenant_id: str, + account_id: str, + external_user_id: str, + extra_tenant_ids: list[str] | None = None, + ) -> int: + from runtime.extensions.whatsapp.access_control import ( + contact_phone_key, + normalize_whatsapp_phone, + phone_from_jid, + whatsapp_users_match, + ) + + tenant_ids: list[str] = [] + for tid in [str(tenant_id), *(extra_tenant_ids or [])]: + val = str(tid or "").strip() + if val and val not in tenant_ids: + tenant_ids.append(val) + + try: + phone_val = normalize_whatsapp_phone(external_user_id) + except Exception: + phone_val = phone_from_jid(external_user_id) + + contacts: list[dict[str, Any]] = [] + for tid in tenant_ids: + contacts.extend(self.list_whatsapp_contacts(tenant_id=tid, account_id=account_id, limit=500)) + + targets: set[str] = set() + for row in contacts: + cj = str(row.get("external_user_id") or "").strip() + if not cj: + continue + if whatsapp_users_match(cj, external_user_id): + targets.add(cj) + elif phone_val and contact_phone_key(row) == phone_val: + targets.add(cj) + + deleted_count = 0 + with self._connect() as conn: + for tid in tenant_ids: + for jid in targets: + if not jid: + continue + cur = conn.execute( + """ + DELETE FROM whatsapp_contact + WHERE tenant_id = ? AND account_id = ? AND external_user_id = ? + """, + (str(tid), str(account_id), str(jid)), + ) + if cur.rowcount: + deleted_count += int(cur.rowcount) + return deleted_count + + def _ensure_whatsapp_access_pending_phone_column(self, conn: Any) -> None: + if self._use_pg: + conn.execute( + "ALTER TABLE whatsapp_access_pending ADD COLUMN IF NOT EXISTS phone TEXT NOT NULL DEFAULT ''" + ) + return + cols = {row[1] for row in conn.execute("PRAGMA table_info(whatsapp_access_pending)").fetchall()} + if "phone" not in cols: + conn.execute("ALTER TABLE whatsapp_access_pending ADD COLUMN phone TEXT NOT NULL DEFAULT ''") + + def create_whatsapp_access_pending( + self, + *, + tenant_id: str, + account_id: str, + external_user_id: str, + push_name: str = "", + phone: str = "", + request_text: str = "", + ) -> str | None: + import uuid + + from runtime.extensions.whatsapp.access_control import contact_phone_key, whatsapp_users_match + + pending = self.list_whatsapp_access_pending( + tenant_id=tenant_id, + account_id=account_id, + status="pending", + limit=200, + ) + phone_val = str(phone or "").strip() + for row in pending: + if whatsapp_users_match(str(row.get("external_user_id") or ""), external_user_id): + return str(row.get("id") or "") + if phone_val and contact_phone_key(row) == phone_val: + return str(row.get("id") or "") + pending_id = uuid.uuid4().hex + ts = utc_now_iso() + with self._connect() as conn: + conn.execute( + """ + INSERT INTO whatsapp_access_pending + (id, tenant_id, account_id, external_user_id, push_name, phone, request_text, status, created_at, resolved_at, resolved_by) + VALUES (?, ?, ?, ?, ?, ?, ?, 'pending', ?, NULL, '') + """, + ( + pending_id, + str(tenant_id), + str(account_id), + str(external_user_id), + str(push_name or ""), + str(phone or ""), + str(request_text or ""), + ts, + ), + ) + return pending_id + + def get_whatsapp_access_pending_by_id(self, *, pending_id: str) -> dict[str, Any] | None: + pid = str(pending_id or "").strip() + if not pid: + return None + with self._connect() as conn: + row = conn.execute( + """ + SELECT id, tenant_id, account_id, external_user_id, push_name, phone, request_text, status, created_at, resolved_at, resolved_by + FROM whatsapp_access_pending + WHERE id = ? + """, + (pid,), + ).fetchone() + if not row: + return None + from runtime.extensions.whatsapp.access_control import phone_from_jid + + return { + "id": str(row["id"] or ""), + "tenant_id": str(row["tenant_id"] or ""), + "account_id": str(row["account_id"] or ""), + "external_user_id": str(row["external_user_id"] or ""), + "push_name": str(row["push_name"] or ""), + "phone": str(row["phone"] or "") or phone_from_jid(str(row["external_user_id"] or "")), + "request_text": str(row["request_text"] or ""), + "status": str(row["status"] or ""), + "created_at": str(row["created_at"] or ""), + "resolved_at": str(row["resolved_at"] or "") or None, + "resolved_by": str(row["resolved_by"] or ""), + } + + def upsert_whatsapp_pending_notify( + self, + *, + tenant_id: str, + account_id: str, + admin_chat_id: str, + notify_stanza_id: str, + pending_id: str, + notify_text: str = "", + ) -> None: + ts = utc_now_iso() + with self._connect() as conn: + conn.execute( + """ + INSERT INTO whatsapp_access_pending_notify + (tenant_id, account_id, admin_chat_id, notify_stanza_id, pending_id, notify_text, created_at) + VALUES (?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(tenant_id, account_id, admin_chat_id, notify_stanza_id) DO UPDATE SET + pending_id = excluded.pending_id, + notify_text = CASE WHEN excluded.notify_text != '' THEN excluded.notify_text ELSE whatsapp_access_pending_notify.notify_text END, + created_at = excluded.created_at + """, + ( + str(tenant_id), + str(account_id), + str(admin_chat_id), + str(notify_stanza_id), + str(pending_id), + str(notify_text or ""), + ts, + ), + ) + + def find_whatsapp_pending_by_notify_stanza( + self, + *, + tenant_id: str, + account_id: str, + admin_chat_id: str, + notify_stanza_id: str, + ) -> str | None: + stanza = str(notify_stanza_id or "").strip() + chat = str(admin_chat_id or "").strip() + if not stanza or not chat: + return None + with self._connect() as conn: + row = conn.execute( + """ + SELECT pending_id + FROM whatsapp_access_pending_notify + WHERE tenant_id = ? AND account_id = ? AND admin_chat_id = ? AND notify_stanza_id = ? + """, + (str(tenant_id), str(account_id), chat, stanza), + ).fetchone() + return str(row["pending_id"] or "").strip() if row else None + + def delete_whatsapp_pending_notify_for_pending(self, *, pending_id: str) -> int: + pid = str(pending_id or "").strip() + if not pid: + return 0 + with self._connect() as conn: + cur = conn.execute( + "DELETE FROM whatsapp_access_pending_notify WHERE pending_id = ?", + (pid,), + ) + return int(cur.rowcount or 0) + + def list_whatsapp_access_pending( + self, + *, + tenant_id: str, + account_id: str, + status: str = "pending", + limit: int = 50, + ) -> list[dict[str, Any]]: + lim = max(1, min(int(limit), 200)) + from runtime.extensions.whatsapp.access_control import phone_from_jid + + with self._connect() as conn: + rows = conn.execute( + """ + SELECT id, tenant_id, account_id, external_user_id, push_name, phone, request_text, status, created_at, resolved_at, resolved_by + FROM whatsapp_access_pending + WHERE tenant_id = ? AND account_id = ? AND status = ? + ORDER BY created_at ASC + LIMIT ? + """, + (str(tenant_id), str(account_id), str(status), lim), + ).fetchall() + return [ + { + "id": str(r["id"] or ""), + "tenant_id": str(r["tenant_id"] or ""), + "account_id": str(r["account_id"] or ""), + "external_user_id": str(r["external_user_id"] or ""), + "push_name": str(r["push_name"] or ""), + "phone": str(r["phone"] or "") or phone_from_jid(str(r["external_user_id"] or "")), + "request_text": str(r["request_text"] or ""), + "status": str(r["status"] or ""), + "created_at": str(r["created_at"] or ""), + "resolved_at": str(r["resolved_at"] or "") or None, + "resolved_by": str(r["resolved_by"] or ""), + } + for r in rows + ] + + def resolve_whatsapp_access_pending( + self, + *, + pending_id: str, + status: str, + resolved_by: str = "", + from_statuses: tuple[str, ...] = ("pending",), + ) -> bool: + ts = utc_now_iso() + st = str(status or "").strip().lower() + if st not in {"approved", "denied", "dismissed"}: + return False + allowed = tuple(str(x or "").strip().lower() for x in from_statuses if str(x or "").strip()) + if not allowed: + allowed = ("pending",) + placeholders = ", ".join("?" for _ in allowed) + with self._connect() as conn: + cur = conn.execute( + f""" + UPDATE whatsapp_access_pending + SET status = ?, resolved_at = ?, resolved_by = ? + WHERE id = ? AND status IN ({placeholders}) + """, + (st, ts, str(resolved_by or ""), str(pending_id), *allowed), + ) + return bool(cur.rowcount and cur.rowcount > 0) + + def delete_whatsapp_access_pending( + self, *, pending_id: str, statuses: tuple[str, ...] = ("pending",) + ) -> bool: + allowed = tuple(str(x or "").strip().lower() for x in statuses if str(x or "").strip()) + if not allowed: + allowed = ("pending",) + placeholders = ", ".join("?" for _ in allowed) + with self._connect() as conn: + cur = conn.execute( + f"DELETE FROM whatsapp_access_pending WHERE id = ? AND status IN ({placeholders})", + (str(pending_id), *allowed), + ) + return bool(cur.rowcount and cur.rowcount > 0) + + def find_whatsapp_contact_for_sender( + self, + *, + tenant_id: str, + account_id: str, + external_user_id: str, + participant_alt: str = "", + remote_jid_alt: str = "", + ) -> dict[str, Any] | None: + from runtime.extensions.whatsapp.access_control import ( + contact_phone_key, + resolve_sender_phone, + whatsapp_users_match, + ) + + sender_phone = resolve_sender_phone(external_user_id, participant_alt, remote_jid_alt) + contacts = self.list_whatsapp_contacts(tenant_id=tenant_id, account_id=account_id, limit=500) + classified = [row for row in contacts if str(row.get("list_type") or "").strip()] + + if sender_phone: + for row in classified: + if contact_phone_key(row) == sender_phone: + return row + + for row in classified: + cj = str(row.get("external_user_id") or "") + if whatsapp_users_match(cj, external_user_id) or ( + participant_alt and whatsapp_users_match(cj, participant_alt) + ): + return row + return None + + def apply_whatsapp_contact_access( + self, + *, + tenant_id: str, + account_id: str, + external_user_id: str, + list_type: str, + push_name: str = "", + phone: str = "", + participant_alt: str = "", + notes: str = "", + ) -> dict[str, Any]: + from runtime.extensions.whatsapp.access_control import ( + contact_phone_key, + normalize_whatsapp_phone, + ) + from runtime.extensions.whatsapp.api import normalize_whatsapp_target + + phone_val = "" + if str(phone or "").strip(): + phone_val = normalize_whatsapp_phone(phone) + elif str(external_user_id or "").strip(): + phone_val = normalize_whatsapp_phone(external_user_id) + else: + return {} + canonical = normalize_whatsapp_target(phone_val) + primary = self.upsert_whatsapp_contact( + tenant_id=tenant_id, + account_id=account_id, + external_user_id=canonical, + push_name=push_name, + phone=phone_val, + list_type=list_type, + notes=notes, + ) + ts = utc_now_iso() + with self._connect() as conn: + for row in self.list_whatsapp_contacts(tenant_id=tenant_id, account_id=account_id, limit=500): + cj = str(row.get("external_user_id") or "") + if cj == canonical: + continue + if contact_phone_key(row) == phone_val: + conn.execute( + """ + UPDATE whatsapp_contact + SET list_type = ?, phone = ?, last_seen_at = ? + WHERE tenant_id = ? AND account_id = ? AND external_user_id = ? + """, + (str(list_type), phone_val, ts, str(tenant_id), str(account_id), cj), + ) + return primary or {} + + def resolve_whatsapp_pending_for_sender( + self, + *, + tenant_id: str, + account_id: str, + external_user_id: str, + participant_alt: str = "", + remote_jid_alt: str = "", + resolved_by: str = "", + status: str = "approved", + extra_tenant_ids: list[str] | None = None, + ) -> int: + from runtime.extensions.whatsapp.access_control import ( + contact_phone_key, + resolve_sender_phone, + whatsapp_users_match, + ) + + tenant_ids: list[str] = [] + for tid in [str(tenant_id), *(extra_tenant_ids or [])]: + val = str(tid or "").strip() + if val and val not in tenant_ids: + tenant_ids.append(val) + + sender_phone = resolve_sender_phone(external_user_id, participant_alt, remote_jid_alt) + resolved = 0 + seen_ids: set[str] = set() + for tid in tenant_ids: + pending = self.list_whatsapp_access_pending( + tenant_id=tid, + account_id=account_id, + status="pending", + limit=200, + ) + for item in pending: + pending_id = str(item.get("id") or "") + if not pending_id or pending_id in seen_ids: + continue + pending_jid = str(item.get("external_user_id") or "") + pending_phone = contact_phone_key(item) + if (sender_phone and pending_phone == sender_phone) or whatsapp_users_match( + pending_jid, external_user_id + ) or (participant_alt and whatsapp_users_match(pending_jid, participant_alt)): + if self.resolve_whatsapp_access_pending( + pending_id=pending_id, + status=status, + resolved_by=resolved_by, + ): + seen_ids.add(pending_id) + resolved += 1 + return resolved + def enqueue_channel_outbound_message( self, *, @@ -5697,16 +6444,52 @@ class SqliteStore(ScheduledJobStoreMixin): message_id: str, ok: bool, error: str = "", + stanza_id: str = "", ) -> bool: ts = utc_now_iso() status = "sent" if ok else "failed" + mid = str(message_id or "").strip() + if not mid: + return False with self._connect() as conn: + row = conn.execute( + """ + SELECT tenant_id, account_id, chat_id, text, source + FROM channel_outbound_message + WHERE id = ? AND status = 'pending' + """, + (mid,), + ).fetchone() + if not row: + return False 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)), + (status, ts, str(error or ""), mid), ) - return bool(cur.rowcount and cur.rowcount > 0) + changed = bool(cur.rowcount and cur.rowcount > 0) + if changed and ok and str(stanza_id or "").strip(): + source_raw = str(row["source"] or "").strip() + meta: dict[str, Any] = {} + if source_raw.startswith("{"): + try: + parsed = json.loads(source_raw) + if isinstance(parsed, dict): + meta = parsed + except Exception: + meta = {} + if str(meta.get("kind") or "") == "whatsapp_access_pending": + pending_id = str(meta.get("pending_id") or "").strip() + if pending_id: + self.upsert_whatsapp_pending_notify( + tenant_id=str(row["tenant_id"] or ""), + account_id=str(row["account_id"] or ""), + admin_chat_id=str(row["chat_id"] or ""), + notify_stanza_id=str(stanza_id or "").strip(), + pending_id=pending_id, + notify_text=str(row["text"] or ""), + ) + return changed diff --git a/tests/test_whatsapp_access_control.py b/tests/test_whatsapp_access_control.py new file mode 100644 index 00000000..6c68979b --- /dev/null +++ b/tests/test_whatsapp_access_control.py @@ -0,0 +1,145 @@ +from __future__ import annotations + +from runtime.extensions.whatsapp.access_control import ( + admin_notify_text, + extract_quote_context, + is_access_allowed, + normalize_whatsapp_phone, + parse_admin_access_command, + parse_admin_approval_intent, + parse_pending_id_from_notify_text, + resolve_sender_phone, + resolve_whatsapp_sender_jid, + whatsapp_phones_match, + whatsapp_users_match, +) + + +def test_blacklist_mode_denies_by_default() -> None: + assert is_access_allowed(access_mode="blacklist", list_type=None) is False + assert is_access_allowed(access_mode="blacklist", list_type="whitelist") is True + assert is_access_allowed(access_mode="blacklist", list_type="blacklist") is False + assert is_access_allowed(access_mode="blacklist", list_type="admin") is True + + +def test_whitelist_mode_allows_by_default() -> None: + assert is_access_allowed(access_mode="whitelist", list_type=None) is True + assert is_access_allowed(access_mode="whitelist", list_type="blacklist") is False + assert is_access_allowed(access_mode="whitelist", list_type="whitelist") is True + + +def test_admin_approval_intent() -> None: + assert parse_admin_approval_intent("yes") == "approve" + assert parse_admin_approval_intent("NO") == "deny" + assert parse_admin_approval_intent("add to whitelist") == "approve" + assert parse_admin_approval_intent("hello") is None + + +def test_whatsapp_users_match_lid_alias() -> None: + assert whatsapp_users_match("91010910658657@lid", "91010910658657@s.whatsapp.net") is True + + +def test_resolve_whatsapp_sender_jid_uses_participant_alt() -> None: + jid = resolve_whatsapp_sender_jid( + "91010910658657@lid", + {"raw": {"participantAlt": "8618142387786@s.whatsapp.net"}}, + ) + assert jid == "8618142387786@s.whatsapp.net" + + +def test_normalize_whatsapp_phone() -> None: + assert normalize_whatsapp_phone("+8615601877957") == "8615601877957" + assert normalize_whatsapp_phone("8615601877957@s.whatsapp.net") == "8615601877957" + + +def test_resolve_sender_phone_prefers_participant_alt() -> None: + phone = resolve_sender_phone( + "91010910658657@lid", + "8618142387786@s.whatsapp.net", + ) + assert phone == "8618142387786" + + +def test_resolve_sender_phone_uses_remote_jid_alt_for_dm() -> None: + phone = resolve_sender_phone( + "91010910658657@lid", + "", + "8615601877957@s.whatsapp.net", + ) + assert phone == "8615601877957" + + +def test_phone_from_jid_ignores_lid() -> None: + from runtime.extensions.whatsapp.access_control import phone_from_jid + + assert phone_from_jid("91010910658657@lid") == "" + + +def test_resolve_whatsapp_sender_jid_uses_remote_jid_alt() -> None: + jid = resolve_whatsapp_sender_jid( + "91010910658657@lid", + {"raw": {"remoteJidAlt": "8615601877957@s.whatsapp.net"}}, + ) + assert jid == "8615601877957@s.whatsapp.net" + + +def test_whatsapp_phones_match_across_formats() -> None: + assert whatsapp_phones_match("+8615601877957", "8615601877957@s.whatsapp.net") is True + + +def test_coerce_whatsapp_access_target_normalizes_phone() -> None: + from runtime.extensions.whatsapp.access_control import coerce_whatsapp_access_target + + assert coerce_whatsapp_access_target("+8618142387786") == "8618142387786@s.whatsapp.net" + assert coerce_whatsapp_access_target("8615601877957") == "8615601877957@s.whatsapp.net" + + +def test_admin_access_commands() -> None: + cmd = parse_admin_access_command("whitelist add +8615601877957") + assert cmd and cmd.get("action") == "set_list" and cmd.get("list_type") == "whitelist" + mode = parse_admin_access_command("access mode blacklist") + assert mode and mode.get("access_mode") == "blacklist" + + +def test_parse_pending_id_from_notify_text() -> None: + zh = "[oclaw] 未授权用户请求访问\n请求编号: abcdef0123456789\n" + en = "[oclaw] Unauthorized access request\nRequest: fedcba9876543210\n" + assert parse_pending_id_from_notify_text(zh) == "abcdef0123456789" + assert parse_pending_id_from_notify_text(en) == "fedcba9876543210" + assert parse_pending_id_from_notify_text("no id here") is None + + +def test_extract_quote_context() -> None: + ctx = extract_quote_context( + { + "raw": { + "quotedStanzaId": "stanza_notify_1", + "quotedText": "Request: abc123", + "participant": "111@s.whatsapp.net", + "remoteJid": "222@s.whatsapp.net", + } + } + ) + assert ctx["stanza_id"] == "stanza_notify_1" + assert ctx["quoted_text"] == "Request: abc123" + assert ctx["participant"] == "111@s.whatsapp.net" + assert ctx["remote_jid"] == "222@s.whatsapp.net" + + +def test_admin_notify_text_prompts_quote_reply() -> None: + zh = admin_notify_text( + lang="zh", + push_name="Bob", + external_user_id="8615601877957@s.whatsapp.net", + request_text="hi", + pending_id="pending123", + ) + en = admin_notify_text( + lang="en", + push_name="Bob", + external_user_id="8615601877957@s.whatsapp.net", + request_text="hi", + pending_id="pending123", + ) + assert "请引用本条消息" in zh + assert "quoting this message" in en diff --git a/tests/test_whatsapp_inbound_access.py b/tests/test_whatsapp_inbound_access.py new file mode 100644 index 00000000..193fbc11 --- /dev/null +++ b/tests/test_whatsapp_inbound_access.py @@ -0,0 +1,389 @@ +from __future__ import annotations + +from typing import Any + +from runtime.application.gateway.whatsapp_inbound_access import handle_whatsapp_access + + +class _AccessStore: + def __init__(self) -> None: + self.contacts: dict[str, dict[str, Any]] = {} + self.pending: list[dict[str, Any]] = [] + self.outbound: list[dict[str, Any]] = [] + self.config: dict[str, Any] = {"access_mode": "blacklist", "lang": "en"} + self.identities: dict[str, dict[str, Any]] = {} + self.users: list[dict[str, Any]] = [] + self.notify_map: dict[tuple[str, str], str] = {} + self._pending_seq = 0 + + def get_whatsapp_access_config(self, *, tenant_id: str, account_id: str) -> dict[str, Any]: + return dict(self.config) + + def upsert_whatsapp_access_config(self, **kwargs: Any) -> dict[str, Any]: + self.config.update({k: v for k, v in kwargs.items() if k in {"access_mode", "lang"}}) + return dict(self.config) + + def apply_whatsapp_contact_access(self, **kwargs: Any) -> dict[str, Any]: + return self.upsert_whatsapp_contact(**kwargs) + + def find_whatsapp_contact_for_sender(self, **kwargs: Any) -> dict[str, Any] | None: + from runtime.extensions.whatsapp.access_control import ( + contact_phone_key, + resolve_sender_phone, + whatsapp_users_match, + ) + + jid = str(kwargs.get("external_user_id") or "") + alt = str(kwargs.get("participant_alt") or "") + remote_alt = str(kwargs.get("remote_jid_alt") or "") + sender_phone = resolve_sender_phone(jid, alt, remote_alt) + for row in self.contacts.values(): + if not row.get("list_type"): + continue + if sender_phone and contact_phone_key(row) == sender_phone: + return row + cj = str(row.get("external_user_id") or "") + if whatsapp_users_match(cj, jid) or (alt and whatsapp_users_match(cj, alt)): + return row + return None + + def resolve_whatsapp_pending_for_sender(self, **kwargs: Any) -> int: + return 0 + + def upsert_whatsapp_contact(self, **kwargs: Any) -> dict[str, Any]: + jid = str(kwargs.get("external_user_id") or "") + row = dict(self.contacts.get(jid) or {}) + row.update({k: v for k, v in kwargs.items() if v is not None}) + if kwargs.get("list_type") is None and jid in self.contacts: + row["list_type"] = self.contacts[jid].get("list_type") + self.contacts[jid] = row + return row + + def list_whatsapp_contacts(self, *, tenant_id: str, account_id: str, list_type: str | None = None, limit: int = 200) -> list[dict[str, Any]]: + rows = list(self.contacts.values()) + if list_type: + rows = [r for r in rows if str(r.get("list_type") or "") == list_type] + return rows + + def create_whatsapp_access_pending(self, **kwargs: Any) -> str: + self._pending_seq += 1 + pending_id = f"pending{self._pending_seq}" + for row in self.pending: + if row.get("status") == "pending" and row.get("external_user_id") == kwargs.get("external_user_id"): + return str(row.get("id") or pending_id) + self.pending.append( + { + "id": pending_id, + "external_user_id": kwargs.get("external_user_id"), + "push_name": kwargs.get("push_name"), + "phone": kwargs.get("phone"), + "request_text": kwargs.get("request_text"), + "status": "pending", + } + ) + return pending_id + + def list_whatsapp_access_pending(self, **kwargs: Any) -> list[dict[str, Any]]: + return [r for r in self.pending if r.get("status") == kwargs.get("status", "pending")] + + def get_whatsapp_access_pending_by_id(self, *, pending_id: str) -> dict[str, Any] | None: + for row in self.pending: + if str(row.get("id") or "") == str(pending_id or ""): + return dict(row) + return None + + def find_whatsapp_pending_by_notify_stanza( + self, + *, + tenant_id: str, + account_id: str, + admin_chat_id: str, + notify_stanza_id: str, + ) -> str | None: + return self.notify_map.get((str(admin_chat_id), str(notify_stanza_id))) + + def delete_whatsapp_pending_notify_for_pending(self, *, pending_id: str) -> int: + removed = 0 + for key, pid in list(self.notify_map.items()): + if pid == pending_id: + del self.notify_map[key] + removed += 1 + return removed + + def resolve_whatsapp_access_pending(self, **kwargs: Any) -> bool: + for row in self.pending: + if row.get("id") == kwargs.get("pending_id"): + row["status"] = kwargs.get("status") + return True + return False + + def enqueue_channel_outbound_message(self, **kwargs: Any) -> str: + self.outbound.append(dict(kwargs)) + return "out1" + + def resolve_user_by_channel_identity_v2(self, **kwargs: Any) -> dict[str, Any] | None: + return self.identities.get(str(kwargs.get("external_user_id") or "")) + + def create_user_account(self, **kwargs: Any) -> dict[str, Any]: + user = {"id": "guest1", "role": kwargs.get("role", "guest")} + self.users.append(user) + return user + + def upsert_channel_identity_v2(self, **kwargs: Any) -> None: + jid = str(kwargs.get("external_user_id") or "") + self.identities[jid] = { + "tenant_id": kwargs.get("tenant_id"), + "user_id": kwargs.get("user_id"), + "role": "guest", + } + + def get_user_by_username(self, **kwargs: Any) -> dict[str, Any] | None: + return {"id": "admin1", "role": "owner"} + + +class _Inbound: + channel = "whatsapp" + external_user_id = "8615601877957@s.whatsapp.net" + external_chat_id = "8615601877957@s.whatsapp.net" + metadata = {"raw": {"pushName": "Bob"}} + + +class _InboundGroup: + channel = "whatsapp" + is_group = True + external_user_id = "91010910658657@lid" + external_chat_id = "999@g.us" + metadata = { + "raw": { + "pushName": "oliver", + "id": "stanza_1", + "participant": "91010910658657@lid", + "participantAlt": "8618142387786@s.whatsapp.net", + } + } + + +class _InboundAdmin: + channel = "whatsapp" + external_user_id = "111@s.whatsapp.net" + external_chat_id = "111@s.whatsapp.net" + metadata: dict[str, Any] = {"raw": {}} + + +def _setup_admin_store() -> _AccessStore: + store = _AccessStore() + store.contacts["111@s.whatsapp.net"] = { + "external_user_id": "111@s.whatsapp.net", + "phone": "111", + "list_type": "admin", + } + return store + + +def test_handle_whatsapp_access_admin_yes_without_quote_passes_to_llm(monkeypatch) -> None: + import runtime.application.gateway.whatsapp_inbound_access as mod + + monkeypatch.setattr(mod, "resolve_whatsapp_tenant_id", lambda store, account_id: "tenant1") + store = _setup_admin_store() + store.pending.append( + { + "id": "pending1", + "external_user_id": "8615601877957@s.whatsapp.net", + "push_name": "Bob", + "phone": "8615601877957", + "status": "pending", + } + ) + out = handle_whatsapp_access(store, inbound=_InboundAdmin(), account_id="wa-default", text="yes") + assert out is None + assert store.pending[0]["status"] == "pending" + + +def test_handle_whatsapp_access_admin_quoted_notify_without_intent_passes_to_llm(monkeypatch) -> None: + import runtime.application.gateway.whatsapp_inbound_access as mod + + monkeypatch.setattr(mod, "resolve_whatsapp_tenant_id", lambda store, account_id: "tenant1") + store = _setup_admin_store() + store.pending.append( + { + "id": "pending2", + "external_user_id": "8615601877957@s.whatsapp.net", + "push_name": "Bob", + "phone": "8615601877957", + "status": "pending", + } + ) + store.notify_map[("111@s.whatsapp.net", "notify_stanza_2")] = "pending2" + inbound = type("_InboundAdminQuoteHello", (), { + "channel": "whatsapp", + "external_user_id": "111@s.whatsapp.net", + "external_chat_id": "111@s.whatsapp.net", + "metadata": { + "raw": { + "quotedStanzaId": "notify_stanza_2", + "quotedText": "[oclaw] Unauthorized access request", + } + }, + })() + out = handle_whatsapp_access(store, inbound=inbound, account_id="wa-default", text="hello") + assert out is None + assert store.pending[0]["status"] == "pending" + + +def test_handle_whatsapp_access_admin_yes_with_stanza_mapping(monkeypatch) -> None: + import runtime.application.gateway.whatsapp_inbound_access as mod + + monkeypatch.setattr(mod, "resolve_whatsapp_tenant_id", lambda store, account_id: "tenant1") + store = _setup_admin_store() + store.pending.append( + { + "id": "pending2", + "external_user_id": "8615601877957@s.whatsapp.net", + "push_name": "Bob", + "phone": "8615601877957", + "status": "pending", + } + ) + store.notify_map[("111@s.whatsapp.net", "notify_stanza_2")] = "pending2" + inbound = type("_InboundAdminQuote", (), { + "channel": "whatsapp", + "external_user_id": "111@s.whatsapp.net", + "external_chat_id": "111@s.whatsapp.net", + "metadata": { + "raw": { + "id": "admin_reply_1", + "participant": "111@s.whatsapp.net", + "quotedStanzaId": "notify_stanza_2", + "quotedText": "[oclaw] Unauthorized access request", + } + }, + })() + out = handle_whatsapp_access(store, inbound=inbound, account_id="wa-default", text="yes") + assert out is not None + assert store.pending[0]["status"] == "approved" + assert store.contacts.get("8615601877957", {}).get("list_type") == "whitelist" + replies = out.get("replies") if isinstance(out.get("replies"), list) else [] + md = replies[0].get("metadata") if replies and isinstance(replies[0], dict) else {} + assert md.get("quote_stanza_id") == "admin_reply_1" + + +def test_handle_whatsapp_access_admin_yes_with_quoted_text_fallback(monkeypatch) -> None: + import runtime.application.gateway.whatsapp_inbound_access as mod + + monkeypatch.setattr(mod, "resolve_whatsapp_tenant_id", lambda store, account_id: "tenant1") + store = _setup_admin_store() + store.pending.append( + { + "id": "abcdef0123456789abcdef0123456789", + "external_user_id": "8615601877957@s.whatsapp.net", + "push_name": "Bob", + "phone": "8615601877957", + "status": "pending", + } + ) + inbound = type("_InboundAdminQuoteText", (), { + "channel": "whatsapp", + "external_user_id": "111@s.whatsapp.net", + "external_chat_id": "111@s.whatsapp.net", + "metadata": { + "raw": { + "id": "admin_reply_2", + "participant": "111@s.whatsapp.net", + "quotedStanzaId": "missing_stanza", + "quotedText": "[oclaw] Unauthorized access request\nRequest: abcdef0123456789abcdef0123456789\n", + } + }, + })() + out = handle_whatsapp_access(store, inbound=inbound, account_id="wa-default", text="no") + assert out is not None + assert store.pending[0]["status"] == "denied" + assert store.contacts.get("8615601877957", {}).get("list_type") == "blacklist" + + +def test_handle_whatsapp_access_denied_enqueues_notify_with_pending_source(monkeypatch) -> None: + import json + import runtime.application.gateway.whatsapp_inbound_access as mod + + monkeypatch.setattr(mod, "resolve_whatsapp_tenant_id", lambda store, account_id: "tenant1") + store = _AccessStore() + store.contacts["111@s.whatsapp.net"] = { + "external_user_id": "111@s.whatsapp.net", + "list_type": "admin", + } + out = handle_whatsapp_access(store, inbound=_Inbound(), account_id="wa-default", text="hello") + assert out is not None + assert store.outbound + source = store.outbound[0].get("source") + parsed = json.loads(str(source)) + assert parsed.get("kind") == "whatsapp_access_pending" + assert parsed.get("pending_id") + + +def test_handle_whatsapp_access_allows_admin_via_lid_alias(monkeypatch) -> None: + import runtime.application.gateway.whatsapp_inbound_access as mod + + monkeypatch.setattr(mod, "resolve_whatsapp_tenant_id", lambda store, account_id: "tenant1") + store = _AccessStore() + store.contacts["8618142387786@s.whatsapp.net"] = { + "external_user_id": "8618142387786@s.whatsapp.net", + "phone": "8618142387786", + "list_type": "admin", + } + out = handle_whatsapp_access(store, inbound=_InboundGroup(), account_id="wa-default", text="hello") + assert out is None + + +def test_handle_whatsapp_access_denied_unknown_user(monkeypatch) -> None: + import runtime.application.gateway.whatsapp_inbound_access as mod + + monkeypatch.setattr(mod, "resolve_whatsapp_tenant_id", lambda store, account_id: "tenant1") + store = _AccessStore() + store.contacts["111@s.whatsapp.net"] = {"external_user_id": "111@s.whatsapp.net", "list_type": "admin"} + out = handle_whatsapp_access(store, inbound=_Inbound(), account_id="wa-default", text="hello") + assert out is not None + assert out.get("whatsapp_access") == "denied" + assert "8615601877957@s.whatsapp.net" not in store.contacts + assert store.pending + assert store.pending[0].get("phone") == "8615601877957" + assert store.outbound + replies = out.get("replies") if isinstance(out.get("replies"), list) else [] + assert replies and "denied" in str(replies[0].get("text") or "").lower() + + +def test_handle_whatsapp_access_allows_whitelisted_user(monkeypatch) -> None: + import runtime.application.gateway.whatsapp_inbound_access as mod + + monkeypatch.setattr(mod, "resolve_whatsapp_tenant_id", lambda store, account_id: "tenant1") + store = _AccessStore() + store.contacts["8615601877957@s.whatsapp.net"] = { + "external_user_id": "8615601877957@s.whatsapp.net", + "phone": "8615601877957", + "list_type": "whitelist", + } + out = handle_whatsapp_access(store, inbound=_Inbound(), account_id="wa-default", text="hello") + assert out is None + assert store.identities.get("8615601877957@s.whatsapp.net") + + +def test_handle_whatsapp_access_denied_group_includes_mention_and_quote(monkeypatch) -> None: + import runtime.application.gateway.whatsapp_inbound_access as mod + + monkeypatch.setattr(mod, "resolve_whatsapp_tenant_id", lambda store, account_id: "tenant1") + store = _AccessStore() + inbound = type("_InboundGroupDenied", (), { + "channel": "whatsapp", + "is_group": True, + "external_user_id": "333@s.whatsapp.net", + "external_chat_id": "999@g.us", + "metadata": {"raw": {"pushName": "Bob", "id": "stanza_1", "participant": "333@s.whatsapp.net"}}, + })() + out = handle_whatsapp_access(store, inbound=inbound, account_id="wa-default", text="hello") + assert out is not None + assert out.get("whatsapp_access") == "denied" + replies = out.get("replies") if isinstance(out.get("replies"), list) else [] + assert replies and isinstance(replies[0], dict) + md = replies[0].get("metadata") if isinstance(replies[0].get("metadata"), dict) else {} + assert md.get("mention_jids") == ["333@s.whatsapp.net"] + assert md.get("quote_remote_jid") == "999@g.us" + assert md.get("quote_stanza_id") == "stanza_1" + assert md.get("quote_participant") == "333@s.whatsapp.net" diff --git a/tests/test_whatsapp_ops_scripts.py b/tests/test_whatsapp_ops_scripts.py index 1a0b69a5..e7d33e34 100644 --- a/tests/test_whatsapp_ops_scripts.py +++ b/tests/test_whatsapp_ops_scripts.py @@ -48,6 +48,12 @@ def test_whatsapp_runner_supports_reply_attachments_base64() -> None: assert "media_url" in text +def test_whatsapp_runner_resolves_dm_remote_jid_alt() -> None: + text = _read("runtime/operations/whatsapp_bridge/baileys_runner.ts") + assert "resolveDmUserJid" in text + assert "remoteJidAlt" in text + + def test_whatsapp_runner_supports_inbound_media_download() -> None: text = _read("runtime/operations/whatsapp_bridge/baileys_runner.ts") assert "downloadMediaMessage" in text @@ -58,3 +64,9 @@ def test_whatsapp_runner_supports_inbound_media_download() -> None: assert "buildOutboundMediaMessage" in text assert "video_base64" in text assert "audio_base64" in text + + +def test_whatsapp_runner_outbound_ack_passes_stanza_id() -> None: + text = _read("runtime/operations/whatsapp_bridge/baileys_runner.ts") + assert "stanza_id" in text + assert "sent as any)?.key?.id" in text