mirror of
https://github.com/hansjone/oclaw.git
synced 2026-10-09 00:40:45 +08:00
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 <cursoragent@cursor.com>
This commit is contained in:
parent
d3f621b58b
commit
4646170743
6 changed files with 186 additions and 9 deletions
|
|
@ -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,
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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."
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
]
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
106
tests/test_scheduled_failure_delivery.py
Normal file
106
tests/test_scheduled_failure_delivery.py
Normal file
|
|
@ -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()
|
||||
Loading…
Add table
Add a link
Reference in a new issue