From 4646170743312b160b52d206870e3a63f3533cf0 Mon Sep 17 00:00:00 2001 From: oliver Date: Mon, 10 Aug 2026 22:59:40 +0800 Subject: [PATCH] Push scheduled-job failure summaries to WhatsApp and clarify pending access. Enqueue a short English failure notice on job crashes, default scheduler lang to en, and tell unauthorized users their request id is pending admin YES. Co-authored-by: Cursor --- .../gateway/whatsapp_inbound_access.py | 2 +- runtime/extensions/whatsapp/access_control.py | 14 ++- runtime/scheduler/service.py | 2 +- runtime/scheduler/turn_text.py | 26 ++++- runtime/scheduler/worker_turn.py | 45 +++++++- tests/test_scheduled_failure_delivery.py | 106 ++++++++++++++++++ 6 files changed, 186 insertions(+), 9 deletions(-) create mode 100644 tests/test_scheduled_failure_delivery.py diff --git a/runtime/application/gateway/whatsapp_inbound_access.py b/runtime/application/gateway/whatsapp_inbound_access.py index de717321..3d1abeb1 100644 --- a/runtime/application/gateway/whatsapp_inbound_access.py +++ b/runtime/application/gateway/whatsapp_inbound_access.py @@ -500,7 +500,7 @@ def handle_whatsapp_access( { "channel": "whatsapp", "chat_id": inbound.external_chat_id, - "text": denied_reply_text(lang=lang), + "text": denied_reply_text(lang=lang, pending_id=str(pending_id or "")), "attachments": [], "metadata": reply_meta, } diff --git a/runtime/extensions/whatsapp/access_control.py b/runtime/extensions/whatsapp/access_control.py index 00f2fb11..54e96f55 100644 --- a/runtime/extensions/whatsapp/access_control.py +++ b/runtime/extensions/whatsapp/access_control.py @@ -343,9 +343,21 @@ def coerce_whatsapp_access_target(value: str) -> str: return normalize_whatsapp_target(normalize_whatsapp_phone(value)) -def denied_reply_text(*, lang: str) -> str: +def denied_reply_text(*, lang: str, pending_id: str = "") -> str: + pid = str(pending_id or "").strip() if str(lang or "").strip().lower().startswith("zh"): + if pid: + return ( + f"访问申请已提交(编号 {pid})。管理员同意后即可使用;" + "请稍候,无需重复发送相同请求。" + ) return "无权限:您尚未获得使用此助手的授权。请联系管理员。" + if pid: + return ( + f"Access pending (request {pid}): an administrator was notified. " + "You can use the assistant after they approve with YES. " + "No need to resend the same request." + ) return "Access denied: you are not authorized to use this assistant. Please contact an administrator." diff --git a/runtime/scheduler/service.py b/runtime/scheduler/service.py index aa9ef1a0..0278cdf6 100644 --- a/runtime/scheduler/service.py +++ b/runtime/scheduler/service.py @@ -76,7 +76,7 @@ def enqueue_scheduled_job_run( trace_id = uuid.uuid4().hex agent_run_id = uuid.uuid4().hex prompt_text = str(getattr(job, "prompt_text", "") or "").strip() - lang = str(getattr(job, "lang", "") or "zh") + lang = str(getattr(job, "lang", "") or "en") recipe = load_recipe_from_job(job) playbook = recipe_has_playbook(recipe) user_text = build_scheduled_turn_instruction( diff --git a/runtime/scheduler/turn_text.py b/runtime/scheduler/turn_text.py index 3058f148..aae71537 100644 --- a/runtime/scheduler/turn_text.py +++ b/runtime/scheduler/turn_text.py @@ -5,13 +5,34 @@ from typing import Any from runtime.scheduler.recipe import compile_playbook_instruction, recipe_has_playbook -def format_scheduled_user_reminder(prompt_text: str) -> str: +def format_scheduled_user_reminder(prompt_text: str, *, lang: str = "en") -> str: body = str(prompt_text or "").strip() if not body: return "" if body.startswith("⏰"): return body - return f"⏰ 提醒:{body}" + if str(lang or "").lower().startswith("zh"): + return f"⏰ 提醒:{body}" + return f"⏰ Reminder: {body}" + + +def format_scheduled_failure_summary( + *, + job_name: str = "", + job_id: str = "", + error: str = "", + lang: str = "en", +) -> str: + """Short user-facing failure notice for WhatsApp/WeChat scheduled delivery.""" + name = str(job_name or "").strip() or str(job_id or "").strip() or "scheduled job" + err = " ".join(str(error or "").strip().split()) + if len(err) > 240: + err = err[:237] + "..." + if not err: + err = "unknown error" + if str(lang or "").lower().startswith("zh"): + return f"[定时任务失败] {name}\n错误:{err}\n请检查任务配置或手动重试。" + return f"[Scheduled job failed] {name}\nError: {err}\nCheck the job config or retry manually." def build_scheduled_turn_instruction( @@ -71,6 +92,7 @@ def scheduled_turn_system_suffix(*, lang: str, playbook: bool = False) -> str: __all__ = [ "build_scheduled_turn_instruction", + "format_scheduled_failure_summary", "format_scheduled_user_reminder", "scheduled_turn_system_suffix", ] diff --git a/runtime/scheduler/worker_turn.py b/runtime/scheduler/worker_turn.py index 0ac37c8e..5d2e85cf 100644 --- a/runtime/scheduler/worker_turn.py +++ b/runtime/scheduler/worker_turn.py @@ -6,7 +6,7 @@ from typing import Any from runtime.orchestration.group_ingest import is_nonsend_channel_reply_text -from runtime.scheduler.turn_text import format_scheduled_user_reminder +from runtime.scheduler.turn_text import format_scheduled_failure_summary, format_scheduled_user_reminder def resolve_scheduled_outbound_text(*, payload: dict[str, Any], reply_text: str) -> str: @@ -15,7 +15,8 @@ def resolve_scheduled_outbound_text(*, payload: dict[str, Any], reply_text: str) return text prompt = str(payload.get("prompt_text") or "").strip() if prompt: - return format_scheduled_user_reminder(prompt) + lang = str(payload.get("lang") or "en") + return format_scheduled_user_reminder(prompt, lang=lang) return "" @@ -158,9 +159,43 @@ def finalize_scheduled_turn_failure( payload: dict[str, Any], error: str, ) -> None: + from runtime.scheduler.channel_delivery import 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 "") + err = str(error or "")[:500] + job = store.scheduled_job_get(job_id=job_id, tenant_id=tenant_id) if job_id else None + delivery_json = "" + if job: + delivery_json = str(getattr(job, "delivery_json", "") or "") + if not delivery_json: + delivery = payload.get("delivery") if isinstance(payload.get("delivery"), dict) else {} + delivery_json = json.dumps(delivery, ensure_ascii=False) if delivery else "" + + lang = str(payload.get("lang") or getattr(job, "lang", "") or "en") + summary = format_scheduled_failure_summary( + job_name=str(getattr(job, "name", "") or ""), + job_id=job_id, + error=err, + lang=lang, + ) + delivery_status: dict[str, Any] = {} + try: + delivery_status = deliver_scheduled_reply( + store, + tenant_id=tenant_id, + reply_text=summary, + delivery_json=delivery_json, + 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="", + ) + except Exception as exc: + delivery_status = {"ok": False, "error": f"{type(exc).__name__}: {exc}"} + if scheduled_run_id: store.scheduled_job_run_update( run_id=scheduled_run_id, @@ -168,11 +203,13 @@ def finalize_scheduled_turn_failure( patch={ "status": "failed", "finished_at": datetime.now(timezone.utc).isoformat(), - "error": str(error or "")[:500], + "error": err, + "reply_text": summary, + "delivery_status": delivery_status, + "session_id": str(payload.get("session_id") or ""), }, ) if job_id: - job = store.scheduled_job_get(job_id=job_id, tenant_id=tenant_id) pause_after = bool(job and str(getattr(job, "schedule_kind", "") or "") == "once") store.scheduled_job_mark_run( job_id=job_id, diff --git a/tests/test_scheduled_failure_delivery.py b/tests/test_scheduled_failure_delivery.py new file mode 100644 index 00000000..3b91e0c6 --- /dev/null +++ b/tests/test_scheduled_failure_delivery.py @@ -0,0 +1,106 @@ +from __future__ import annotations + +import json +import tempfile +import unittest +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 svc.persistence.sqlite_store import SqliteStore + + +class ScheduledFailureTextTests(unittest.TestCase): + def test_failure_summary_english_default(self) -> None: + text = format_scheduled_failure_summary( + job_name="License daily", + error='TimeoutError: CLI timed out after 60s', + lang="en", + ) + self.assertIn("[Scheduled job failed]", text) + self.assertIn("License daily", text) + self.assertIn("TimeoutError", text) + self.assertNotIn("定时", 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")) + + +class DeniedPendingTextTests(unittest.TestCase): + def test_denied_includes_pending_id_english(self) -> None: + text = denied_reply_text(lang="en", pending_id="pend-1") + self.assertIn("Access pending", text) + self.assertIn("pend-1", text) + self.assertIn("YES", text) + + +class ScheduledFailureDeliveryTests(unittest.TestCase): + def setUp(self) -> None: + self._tmp = tempfile.TemporaryDirectory(ignore_cleanup_errors=True) + self.db = Path(self._tmp.name) / "fail.sqlite" + self.store = SqliteStore(str(self.db)) + tenant = self.store.create_tenant("Team") + self.tenant_id = str(tenant["id"]) + user = self.store.create_user(tenant_id=self.tenant_id, display_name="ops", role="administrator") + self.user_id = str(user["id"]) + + def tearDown(self) -> None: + self._tmp.cleanup() + + def test_finalize_failure_enqueues_whatsapp_summary(self) -> None: + job = self.store.scheduled_job_create( + tenant_id=self.tenant_id, + created_by_user_id=self.user_id, + name="Bandwidth watch", + prompt_text="check congestion", + schedule_kind="cron", + schedule_expr="0 * * * *", + 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_failure( + self.store, + payload=payload, + error="RuntimeError: tool boom", + ) + 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 failed]", body) + self.assertIn("Bandwidth watch", body) + self.assertIn("tool boom", body) + 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") + + +if __name__ == "__main__": + unittest.main()