diff --git a/interfaces/admin/routes.py b/interfaces/admin/routes.py index 58e9e9f8..f4ee276d 100644 --- a/interfaces/admin/routes.py +++ b/interfaces/admin/routes.py @@ -1468,6 +1468,14 @@ def build_admin_router() -> APIRouter: for source_tid in (tid, legacy_tid): if not source_tid: continue + try: + store.expire_stale_whatsapp_access_pending( + tenant_id=source_tid, + account_id=aid, + older_than_hours=168, + ) + except Exception: + pass for row in store.list_whatsapp_access_pending( tenant_id=source_tid, account_id=aid, status="pending", limit=50 ): @@ -1683,6 +1691,24 @@ def build_admin_router() -> APIRouter: status="approved", extra_tenant_ids=extra_tids, ) + try: + from runtime.application.gateway.whatsapp_inbound_access import ( + notify_whatsapp_access_decision, + ) + + cfg = store.get_whatsapp_access_config(tenant_id=tid, account_id=aid) or {} + notify_whatsapp_access_decision( + store, + tenant_id=tid, + account_id=aid, + lang=str(cfg.get("lang") or "en"), + external_user_id=external_user_id, + phone=phone_val, + pending_id=pending_id, + approved=True, + ) + except Exception: + pass return { "ok": bool(changed), "action": "approve", @@ -1713,6 +1739,24 @@ def build_admin_router() -> APIRouter: status="denied", extra_tenant_ids=extra_tids, ) + try: + from runtime.application.gateway.whatsapp_inbound_access import ( + notify_whatsapp_access_decision, + ) + + cfg = store.get_whatsapp_access_config(tenant_id=tid, account_id=aid) or {} + notify_whatsapp_access_decision( + store, + tenant_id=tid, + account_id=aid, + lang=str(cfg.get("lang") or "en"), + external_user_id=external_user_id, + phone=phone_val or "", + pending_id=pending_id, + approved=False, + ) + except Exception: + pass return {"ok": bool(changed), "action": "deny", "phone": phone_val or external_user_id} @router.delete("/admin/api/whatsapp/access/pending") diff --git a/runtime/application/gateway/inbound_service.py b/runtime/application/gateway/inbound_service.py index 5e637e2f..ea50f08a 100644 --- a/runtime/application/gateway/inbound_service.py +++ b/runtime/application/gateway/inbound_service.py @@ -1052,8 +1052,12 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]: ) channel_is_wa_bind = str(inbound.channel or "").strip().lower() == "whatsapp" if info: - guide = _menu_text(channel=str(inbound.channel or "")) - reply = ("Bound successfully.\n\n" + guide) if channel_is_wa_bind else ("绑定成功。\n\n" + guide) + if channel_is_wa_bind: + from runtime.extensions.whatsapp.access_control import access_granted_guide_text + + reply = access_granted_guide_text(lang="en") + else: + reply = "绑定成功。\n\n" + _menu_text(channel=str(inbound.channel or "")) else: reply = ( "Bind failed: invalid or already used code." diff --git a/runtime/application/gateway/whatsapp_inbound_access.py b/runtime/application/gateway/whatsapp_inbound_access.py index bbe20e1b..b45695e4 100644 --- a/runtime/application/gateway/whatsapp_inbound_access.py +++ b/runtime/application/gateway/whatsapp_inbound_access.py @@ -264,6 +264,30 @@ def _notify_requester_access_decision( ) +def notify_whatsapp_access_decision( + store: Any, + *, + tenant_id: str, + account_id: str, + lang: str = "en", + external_user_id: str, + phone: str = "", + pending_id: str = "", + approved: bool, +) -> None: + """Public wrapper for admin/API approve paths to DM the requester.""" + _notify_requester_access_decision( + store, + tenant_id=tenant_id, + account_id=account_id, + lang=lang, + external_user_id=external_user_id, + phone=phone, + pending_id=pending_id, + approved=approved, + ) + + def _upsert_whatsapp_contact_profile( store: Any, *, diff --git a/runtime/scheduler/turn_text.py b/runtime/scheduler/turn_text.py index aae71537..62d2a57e 100644 --- a/runtime/scheduler/turn_text.py +++ b/runtime/scheduler/turn_text.py @@ -16,6 +16,40 @@ def format_scheduled_user_reminder(prompt_text: str, *, lang: str = "en") -> str return f"⏰ Reminder: {body}" +def format_scheduled_success_summary( + *, + job_name: str = "", + job_id: str = "", + reply_text: str = "", + attachment_count: int = 0, + lang: str = "en", + max_body_chars: int = 1200, +) -> str: + """Short user-facing success notice; keeps body but caps length for WhatsApp.""" + name = str(job_name or "").strip() or str(job_id or "").strip() or "scheduled job" + body = str(reply_text or "").strip() + # Drop pure reminder fallbacks that just echo the job prompt. + if body.startswith("⏰"): + body = "" + att_n = max(0, int(attachment_count or 0)) + cap = max(200, int(max_body_chars or 1200)) + if body and len(body) > cap: + body = body[: cap - 3].rstrip() + "..." + if str(lang or "").lower().startswith("zh"): + head = f"[定时任务完成] {name}" + if att_n: + head += f"\n附件:{att_n} 个" + if body: + return f"{head}\n{body}" + return head + ("\n(本次无文字摘要,请查看附件。)" if att_n else "\n(本次无摘要内容。)") + head = f"[Scheduled job done] {name}" + if att_n: + head += f"\nAttachments: {att_n}" + if body: + return f"{head}\n{body}" + return head + ("\n(No text summary; see attachment(s).)" if att_n else "\n(No summary content.)") + + def format_scheduled_failure_summary( *, job_name: str = "", @@ -72,11 +106,14 @@ def scheduled_turn_system_suffix(*, lang: str, playbook: bool = False) -> str: "\n\n[Scheduled playbook mode] You are executing a recurring workflow for the user. " "Follow the playbook steps, use tools as needed, and deliver a useful update " "(including save_deliverable_attachment for generated files). " + "Lead the final reply with a short English summary (3–8 lines: what ran, key counts, " + "ok/failed highlights), then optional detail. " "Do not pretend the user just messaged you." ) return ( "\n\n【定时工作流模式】你正在执行周期性工作流。" "按 playbook 步骤完成任务,按需调用工具;若生成文件须 save_deliverable_attachment。" + "最终回复先给 3–8 行摘要(做了什么、关键计数、成败),再写细节。" "不要假装用户刚刚发了消息,不要只回一句空提醒。" ) if is_en: @@ -93,6 +130,7 @@ def scheduled_turn_system_suffix(*, lang: str, playbook: bool = False) -> str: __all__ = [ "build_scheduled_turn_instruction", "format_scheduled_failure_summary", + "format_scheduled_success_summary", "format_scheduled_user_reminder", "scheduled_turn_system_suffix", ] diff --git a/runtime/scheduler/worker_turn.py b/runtime/scheduler/worker_turn.py index 5d2e85cf..1c64a48e 100644 --- a/runtime/scheduler/worker_turn.py +++ b/runtime/scheduler/worker_turn.py @@ -6,7 +6,11 @@ from typing import Any from runtime.orchestration.group_ingest import is_nonsend_channel_reply_text -from runtime.scheduler.turn_text import format_scheduled_failure_summary, format_scheduled_user_reminder +from runtime.scheduler.turn_text import ( + format_scheduled_failure_summary, + format_scheduled_success_summary, + format_scheduled_user_reminder, +) def resolve_scheduled_outbound_text(*, payload: dict[str, Any], reply_text: str) -> str: @@ -101,11 +105,16 @@ def finalize_scheduled_turn_success( payload: dict[str, Any], base_result: dict[str, Any], ) -> None: - from runtime.scheduler.channel_delivery import deliver_scheduled_reply + from runtime.scheduler.channel_delivery import ( + _collect_scheduled_turn_attachments, + deliver_scheduled_reply, + ) tenant_id = str(payload.get("tenant_id") or "") job_id = str(payload.get("job_id") or "") scheduled_run_id = str(payload.get("run_id_scheduled") or "") + session_id = str(payload.get("session_id") or "") + turn_uuid = str(base_result.get("turn_uuid") or payload.get("run_id") or "") reply_text = resolve_scheduled_outbound_text(payload=payload, reply_text=str(base_result.get("reply_text") or "")) delivery = payload.get("delivery") if isinstance(payload.get("delivery"), dict) else {} delivery_json = json.dumps(delivery, ensure_ascii=False) @@ -114,10 +123,34 @@ def finalize_scheduled_turn_success( if job: delivery_json = str(getattr(job, "delivery_json", "") or delivery_json) + lang = str(payload.get("lang") or getattr(job, "lang", "") or "en") + att_count = 0 + try: + att_count = len( + _collect_scheduled_turn_attachments( + store=store, + session_id=session_id, + turn_uuid=turn_uuid, + ) + or [] + ) + except Exception: + att_count = 0 + + # Playbook / ops jobs: always wrap with a short success header for WA readability. + if job_id or att_count: + reply_text = format_scheduled_success_summary( + job_name=str(getattr(job, "name", "") or ""), + job_id=job_id, + reply_text=reply_text, + attachment_count=att_count, + lang=lang, + ) + _persist_scheduled_assistant_reply( store, - session_id=str(payload.get("session_id") or ""), - turn_uuid=str(base_result.get("turn_uuid") or payload.get("run_id") or ""), + session_id=session_id, + turn_uuid=turn_uuid, reply_text=reply_text, ) delivery_status = deliver_scheduled_reply( @@ -128,8 +161,8 @@ def finalize_scheduled_turn_success( resolved_channel=str(payload.get("resolved_channel") or ""), resolved_chat_id=str(payload.get("resolved_chat_id") or ""), resolved_account_id=str(payload.get("resolved_account_id") or ""), - session_id=str(payload.get("session_id") or ""), - turn_uuid=str(base_result.get("turn_uuid") or payload.get("run_id") or ""), + session_id=session_id, + turn_uuid=turn_uuid, ) if scheduled_run_id: store.scheduled_job_run_update( diff --git a/runtime/tools/tool_error_hints.py b/runtime/tools/tool_error_hints.py index 2e13961f..fee6de56 100644 --- a/runtime/tools/tool_error_hints.py +++ b/runtime/tools/tool_error_hints.py @@ -88,16 +88,38 @@ def enrich_mcp_scope_error(result: dict[str, Any]) -> dict[str, Any]: ] out["user_facing_hint"] = ( "SQL query is not enabled for this bot token. " - "I will use alarm aggregate/report tools instead, or an admin can grant sql:query." + "I will use alarm aggregate/report tools instead. " + "To enable SQL: ask a WhatsApp admin → Admin UI → netx MCP token scopes → grant sql:query." ) + out["next_steps"] = [ + "Use aggregateUmeAlarms / queryUmeAlarmsRaw / ume_alarm_xlsx_report instead of sqlQueryUme.", + "Ask a WhatsApp admin to open Admin → MCP / netx token and grant scope sql:query.", + "Do not retry sqlQueryUme until the scope is granted.", + ] + out["admin_action"] = { + "required_scope": "sql:query", + "where": "Admin UI → MCP server / netx API token scopes", + "ask": "WhatsApp access admin (whitelist contact with list_type=admin)", + } else: out["hint"] = ( f"Current netx token lacks scope {scope or '(unknown)'}. " "Do not retry the same tool; ask an admin to grant it, or use tools that do not need this scope." ) out["user_facing_hint"] = ( - f"Permission missing ({scope or 'scope'}). An admin needs to grant this on the netx API token." + f"Permission missing ({scope or 'scope'}). " + "Ask a WhatsApp admin to grant this scope on the netx API token " + "(Admin → MCP / netx token scopes)." ) + out["next_steps"] = [ + f"Ask a WhatsApp admin to grant scope {scope or '(unknown)'} on the netx MCP token.", + "Retry only after the scope is granted; do not blind-retry this tool.", + ] + out["admin_action"] = { + "required_scope": scope or "", + "where": "Admin UI → MCP server / netx API token scopes", + "ask": "WhatsApp access admin (whitelist contact with list_type=admin)", + } out["failure_class"] = "auth" out["retry_forbidden"] = True return out diff --git a/svc/persistence/sqlite_store.py b/svc/persistence/sqlite_store.py index c4051167..f8a9a761 100644 --- a/svc/persistence/sqlite_store.py +++ b/svc/persistence/sqlite_store.py @@ -6087,6 +6087,14 @@ class SqliteStore(ScheduledJobStoreMixin): ) if existing: return str(existing.get("id") or "") or None + try: + self.expire_stale_whatsapp_access_pending( + tenant_id=str(tenant_id), + account_id=str(account_id), + older_than_hours=168, + ) + except Exception: + pass pending_id = uuid.uuid4().hex ts = utc_now_iso() with self._connect() as conn: @@ -6273,6 +6281,56 @@ class SqliteStore(ScheduledJobStoreMixin): ) return bool(cur.rowcount and cur.rowcount > 0) + def expire_stale_whatsapp_access_pending( + self, + *, + tenant_id: str = "", + account_id: str = "", + older_than_hours: int = 168, + limit: int = 200, + ) -> int: + """Mark old open pending requests as dismissed (default 7 days).""" + from datetime import datetime, timedelta, timezone + + hours = max(1, min(int(older_than_hours or 168), 24 * 90)) + cutoff = (datetime.now(timezone.utc) - timedelta(hours=hours)).isoformat() + lim = max(1, min(int(limit or 200), 500)) + tid = str(tenant_id or "").strip() + aid = str(account_id or "").strip() + where = ["status = 'pending'", "created_at < ?"] + args: list[Any] = [cutoff] + if tid: + where.append("tenant_id = ?") + args.append(tid) + if aid: + where.append("account_id = ?") + args.append(aid) + args.append(lim) + ts = utc_now_iso() + with self._connect() as conn: + rows = conn.execute( + f""" + SELECT id FROM whatsapp_access_pending + WHERE {' AND '.join(where)} + ORDER BY created_at ASC + LIMIT ? + """, + tuple(args), + ).fetchall() + ids = [str(r["id"] or "") for r in rows if str(r["id"] or "").strip()] + if not ids: + return 0 + placeholders = ", ".join("?" for _ in ids) + cur = conn.execute( + f""" + UPDATE whatsapp_access_pending + SET status = 'dismissed', resolved_at = ?, resolved_by = 'system:expire' + WHERE id IN ({placeholders}) AND status = 'pending' + """, + (ts, *ids), + ) + return int(cur.rowcount or 0) + def delete_whatsapp_access_pending( self, *, pending_id: str, statuses: tuple[str, ...] = ("pending",) ) -> bool: diff --git a/tests/test_scheduled_failure_delivery.py b/tests/test_scheduled_failure_delivery.py index 3b91e0c6..fcfd6c49 100644 --- a/tests/test_scheduled_failure_delivery.py +++ b/tests/test_scheduled_failure_delivery.py @@ -7,8 +7,12 @@ from pathlib import Path from unittest.mock import patch from runtime.extensions.whatsapp.access_control import denied_reply_text -from runtime.scheduler.turn_text import format_scheduled_failure_summary, format_scheduled_user_reminder -from runtime.scheduler.worker_turn import finalize_scheduled_turn_failure +from runtime.scheduler.turn_text import ( + format_scheduled_failure_summary, + format_scheduled_success_summary, + format_scheduled_user_reminder, +) +from runtime.scheduler.worker_turn import finalize_scheduled_turn_failure, finalize_scheduled_turn_success from svc.persistence.sqlite_store import SqliteStore @@ -24,6 +28,18 @@ class ScheduledFailureTextTests(unittest.TestCase): self.assertIn("TimeoutError", text) self.assertNotIn("定时", text) + def test_success_summary_wraps_body_and_attachments(self) -> None: + text = format_scheduled_success_summary( + job_name="Bandwidth watch", + reply_text="Critical: 3\nMajor: 1", + attachment_count=1, + lang="en", + ) + self.assertIn("[Scheduled job done]", text) + self.assertIn("Bandwidth watch", text) + self.assertIn("Attachments: 1", text) + self.assertIn("Critical: 3", text) + def test_reminder_fallback_english(self) -> None: self.assertIn("Reminder", format_scheduled_user_reminder("stretch", lang="en")) self.assertIn("提醒", format_scheduled_user_reminder("活动一下", lang="zh")) @@ -101,6 +117,56 @@ class ScheduledFailureDeliveryTests(unittest.TestCase): updated = self.store.scheduled_job_run_get(run_id=str(run.id), tenant_id=self.tenant_id) self.assertEqual(str(getattr(updated, "status", "") or ""), "failed") + def test_finalize_success_enqueues_wrapped_summary(self) -> None: + job = self.store.scheduled_job_create( + tenant_id=self.tenant_id, + created_by_user_id=self.user_id, + name="License daily", + prompt_text="license report", + schedule_kind="cron", + schedule_expr="0 11 * * *", + lang="en", + delivery={ + "whatsapp": { + "enabled": True, + "target_type": "group", + "chat_id": "120363011111111111@g.us", + "account_id": "wa-default", + } + }, + timezone_name="UTC", + status="active", + ) + run = self.store.scheduled_job_run_create( + job_id=str(job.id), + tenant_id=self.tenant_id, + scheduled_at="2026-08-10T00:00:00+00:00", + ) + payload = { + "tenant_id": self.tenant_id, + "job_id": str(job.id), + "run_id_scheduled": str(run.id), + "lang": "en", + "resolved_channel": "whatsapp", + "resolved_chat_id": "120363011111111111@g.us", + "resolved_account_id": "wa-default", + "session_id": "", + } + finalize_scheduled_turn_success( + self.store, + task=None, + payload=payload, + base_result={"reply_text": "Found 2 license alarms.", "turn_uuid": "tu1"}, + ) + pending = self.store.list_pending_channel_outbound_messages( + channel="whatsapp", account_id="wa-default", limit=5 + ) + self.assertEqual(len(pending), 1) + body = str(pending[0].get("text") or "") + self.assertIn("[Scheduled job done]", body) + self.assertIn("License daily", body) + self.assertIn("Found 2 license alarms.", body) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_tool_error_hints.py b/tests/test_tool_error_hints.py index 1bb60532..a2d664d8 100644 --- a/tests/test_tool_error_hints.py +++ b/tests/test_tool_error_hints.py @@ -40,6 +40,9 @@ def test_enrich_mcp_scope_sql() -> None: assert out["required_scope"] == "sql:query" assert "fallback_tools" in out assert "ume_alarm_xlsx_report" in out["fallback_tools"] + assert out.get("next_steps") + assert out.get("admin_action", {}).get("required_scope") == "sql:query" + assert "Admin" in str(out.get("user_facing_hint") or "") def test_build_finalize_system_suffix() -> None: diff --git a/tests/test_whatsapp_ops_onboarding.py b/tests/test_whatsapp_ops_onboarding.py index 2c0b03cc..365f893f 100644 --- a/tests/test_whatsapp_ops_onboarding.py +++ b/tests/test_whatsapp_ops_onboarding.py @@ -20,6 +20,36 @@ def test_access_granted_guide_includes_ops_help() -> None: assert "YES" in text or "continue" in text.lower() +def test_expire_stale_whatsapp_access_pending(tmp_path) -> None: + from datetime import datetime, timedelta, timezone + + from svc.persistence.sqlite_store import SqliteStore + + store = SqliteStore(str(tmp_path / "exp.sqlite")) + tenant = store.create_tenant("T") + tid = str(tenant["id"]) + aid = "wa-default" + old_id = store.create_whatsapp_access_pending( + tenant_id=tid, + account_id=aid, + external_user_id="8611111111111@s.whatsapp.net", + phone="8611111111111", + request_text="hi", + ) + assert old_id + # Backdate created_at beyond 7 days. + old_ts = (datetime.now(timezone.utc) - timedelta(days=10)).isoformat() + with store._connect() as conn: # noqa: SLF001 + conn.execute( + "UPDATE whatsapp_access_pending SET created_at = ? WHERE id = ?", + (old_ts, old_id), + ) + n = store.expire_stale_whatsapp_access_pending(tenant_id=tid, account_id=aid, older_than_hours=168) + assert n == 1 + row = store.get_whatsapp_access_pending_by_id(pending_id=str(old_id)) + assert row and str(row.get("status") or "") == "dismissed" + + def test_menu_text_whatsapp_not_productivity_chinese() -> None: text = _menu_text(channel="whatsapp") assert "记待办" not in text