From 51eef09801ff9bbfb8e52d894fdb5fbc2c768d65 Mon Sep 17 00:00:00 2001 From: oliver Date: Mon, 10 Aug 2026 23:59:27 +0800 Subject: [PATCH] Synthesize playbooks from structured prompts and inject prior-run context. Legacy jobs with empty recipes (e.g. dying-gasp algorithms in prompt_text) now compile into durable playbooks on create/enqueue, and each run gets a short previous-run summary for continuity. Co-authored-by: Cursor --- interfaces/admin/routes.py | 12 ++ runtime/scheduler/recipe.py | 144 ++++++++++++++++++ runtime/scheduler/service.py | 88 ++++++++++- runtime/scheduler/turn_text.py | 77 ++++++++-- .../experts/productivity/schedule_tools.py | 36 +++-- tests/test_schedule_recipe.py | 63 ++++++++ tests/test_scheduler_playbook_backfill.py | 120 +++++++++++++++ tests/test_scheduler_viewer_username.py | 3 + 8 files changed, 516 insertions(+), 27 deletions(-) create mode 100644 tests/test_scheduler_playbook_backfill.py diff --git a/interfaces/admin/routes.py b/interfaces/admin/routes.py index f4ee276d..2836a80a 100644 --- a/interfaces/admin/routes.py +++ b/interfaces/admin/routes.py @@ -2437,6 +2437,7 @@ def build_admin_router() -> APIRouter: prompt_summary_from_recipe, recipe_has_playbook, resolve_ops_recipe_template, + synthesize_recipe_from_prompt, ) recipe = payload.get("recipe") if isinstance(payload.get("recipe"), dict) else None @@ -2457,6 +2458,17 @@ def build_admin_router() -> APIRouter: prompt_text = prompt_summary_from_recipe(recipe, fallback=name) if not prompt_text: return {"ok": False, "error": "name, prompt_text, schedule_expr are required"} + if not recipe_has_playbook(recipe): + synth = synthesize_recipe_from_prompt( + prompt_text, + session_id=str(payload.get("source_session_id") or ""), + ) + if recipe_has_playbook(synth): + recipe = synth + if not str((recipe.get("source") or {}).get("compiled_at") or "").strip(): + from datetime import datetime, timezone + + recipe.setdefault("source", {})["compiled_at"] = datetime.now(timezone.utc).isoformat() # WhatsApp field ops default to English when delivery targets WA and lang omitted. lang = str(payload.get("lang") or "").strip() if not lang: diff --git a/runtime/scheduler/recipe.py b/runtime/scheduler/recipe.py index 7e47707e..e093c80b 100644 --- a/runtime/scheduler/recipe.py +++ b/runtime/scheduler/recipe.py @@ -110,6 +110,7 @@ def normalize_recipe(raw: Any) -> dict[str, Any]: "source": { "session_id": str(source_raw.get("session_id") or source_raw.get("sessionId") or "").strip(), "compiled_at": str(source_raw.get("compiled_at") or source_raw.get("compiledAt") or "").strip(), + "compiled_from": str(source_raw.get("compiled_from") or source_raw.get("compiledFrom") or "").strip(), }, } return recipe @@ -166,6 +167,147 @@ def looks_like_complex_schedule_prompt(prompt_text: str, *, recipe: dict[str, An return False +_STEP_HEADER_RE = re.compile( + r"(?im)^\s*(?:" + r"step\s*\d+\s*[—\-–:.]?\s*" + r"|第[0-9一二三四五六七八九十百]+步\s*[—\-–:.]?\s*" + r"|\d+\s*[\.\)\-—–]\s+" + r").+$" +) + + +def _extract_numbered_steps(prompt_text: str) -> tuple[str, list[str]]: + """Split prompt into (preamble, step bodies) when Step N / 1. headers exist.""" + text = str(prompt_text or "").replace("\r\n", "\n").replace("\r", "\n").strip() + if not text: + return "", [] + matches = list(_STEP_HEADER_RE.finditer(text)) + if len(matches) < 2: + return text, [] + preamble = text[: matches[0].start()].strip() + steps: list[str] = [] + for i, m in enumerate(matches): + end = matches[i + 1].start() if i + 1 < len(matches) else len(text) + chunk = text[m.start() : end].strip() + # Drop leading "Step N —" / "1." so compile_playbook can re-number cleanly. + body = re.sub( + r"(?is)^\s*(?:step\s*\d+|第[0-9一二三四五六七八九十百]+步)\s*[—\-–:.]?\s*", + "", + chunk, + count=1, + ) + body = re.sub(r"(?is)^\s*\d+\s*[\.\)\-—–]\s+", "", body, count=1).strip() + if body: + steps.append(body) + elif chunk: + steps.append(chunk) + return preamble, steps + + +def synthesize_recipe_from_prompt(prompt_text: str, *, session_id: str = "") -> dict[str, Any] | None: + """ + Build a durable playbook recipe from a long/structured prompt. + + Used when field jobs stored algorithm text in prompt_text but left recipe_json empty. + Returns None when the prompt is too thin to treat as a playbook. + """ + text = str(prompt_text or "").strip() + if not text or not looks_like_complex_schedule_prompt(text): + return None + + preamble, steps = _extract_numbered_steps(text) + low = text.lower() + need_attachments = any( + tok in low or tok in text + for tok in ("xlsx", "csv", "pdf", "attachment", "附件", "save_deliverable", "write_xlsx") + ) + + if len(steps) >= 2: + goal = preamble.split("\n\n", 1)[0].strip() if preamble else "" + goal = re.sub(r"(?is)^\s*critical\s*[—\-–:]?\s*", "", goal).strip() + if not goal: + goal = steps[0][:240] + if len(goal) > 400: + goal = goal[:397].rstrip() + "..." + # Prefer keeping full algorithm fidelity: if preamble is short, fold remaining + # non-step prose into constraints rather than losing it. + constraints: list[str] = [] + if preamble and preamble != goal and len(preamble) > len(goal) + 20: + rest = preamble[len(goal) :].strip() if preamble.startswith(goal) else preamble + if rest: + constraints.append(rest[:800]) + success = [ + "Follow every step end-to-end with tools as needed.", + "Deliver a useful channel update reflecting completed work.", + ] + if need_attachments: + success.append("Generated files are saved via save_deliverable_attachment.") + recipe = normalize_recipe( + { + "version": 1, + "goal": goal, + "steps": steps, + "constraints": constraints, + "success_criteria": success, + "output": {"need_attachments": need_attachments, "style": "channel_update"}, + "source": { + "session_id": str(session_id or "").strip(), + "compiled_at": "", + "compiled_from": "prompt_text", + }, + } + ) + # Preserve compiled_from beyond normalize (normalize only keeps session_id/compiled_at). + recipe.setdefault("source", {})["compiled_from"] = "prompt_text" + return recipe + + # No clear Step N headers: still promote long ops prompts so they are not + # misclassified as short "reminder" turns. + if len(text) < 80 and text.count("\n") < 2: + return None + goal = text.split("\n\n", 1)[0].strip() + if len(goal) > 400: + goal = goal[:397].rstrip() + "..." + steps = [ + "Execute the full algorithm described in the goal/prompt end-to-end. Use tools as needed; do not reply with only a short reminder.", + "Deliver a useful channel update that reflects completed work" + + (" (call save_deliverable_attachment for any generated files)." if need_attachments else "."), + ] + if text != goal: + steps.insert(1, f"Full prompt/algorithm to follow:\n{text}") + recipe = normalize_recipe( + { + "version": 1, + "goal": goal, + "steps": steps, + "constraints": [], + "success_criteria": [ + "Workflow completed with tools as required.", + "Channel update delivered.", + ], + "output": {"need_attachments": need_attachments, "style": "channel_update"}, + "source": {"session_id": str(session_id or "").strip(), "compiled_at": ""}, + } + ) + recipe.setdefault("source", {})["compiled_from"] = "prompt_text" + return recipe + + +def resolve_effective_playbook_recipe( + *, + recipe: dict[str, Any] | None, + prompt_text: str, + session_id: str = "", +) -> dict[str, Any] | None: + """Return a playbook recipe from stored recipe or synthesized prompt; else None.""" + if recipe_has_playbook(recipe): + return normalize_recipe(recipe) + synth = synthesize_recipe_from_prompt(prompt_text, session_id=session_id) + if recipe_has_playbook(synth): + return synth + return None + + def prompt_summary_from_recipe(recipe: dict[str, Any] | None, *, fallback: str = "") -> str: norm = normalize_recipe(recipe or {}) goal = str(norm.get("goal") or "").strip() @@ -461,5 +603,7 @@ __all__ = [ "recipe_has_playbook", "recipe_is_empty", "recipe_missing_fields", + "resolve_effective_playbook_recipe", "resolve_ops_recipe_template", + "synthesize_recipe_from_prompt", ] diff --git a/runtime/scheduler/service.py b/runtime/scheduler/service.py index 0278cdf6..144060de 100644 --- a/runtime/scheduler/service.py +++ b/runtime/scheduler/service.py @@ -7,7 +7,12 @@ import uuid from datetime import datetime, timezone from typing import Any -from runtime.scheduler.recipe import load_recipe_from_job, recipe_has_playbook +from runtime.scheduler.recipe import ( + load_recipe_from_job, + recipe_has_playbook, + recipe_is_empty, + resolve_effective_playbook_recipe, +) from runtime.scheduler.session_resolver import resolve_scheduled_session, resolve_scheduled_viewer_username from runtime.scheduler.turn_text import build_scheduled_turn_instruction from runtime.worker import ensure_worker_started @@ -27,6 +32,70 @@ def _tick_interval_seconds() -> float: return 30.0 +def _load_previous_run_context( + store: Any, + *, + job_id: str, + tenant_id: str, + current_run_id: str = "", +) -> dict[str, Any] | None: + """Latest finished run for this job (excluding the run currently being queued).""" + lister = getattr(store, "scheduled_job_run_list", None) + if not callable(lister): + return None + try: + rows = lister(job_id=str(job_id), tenant_id=str(tenant_id), limit=8) or [] + except Exception: + return None + cur = str(current_run_id or "").strip() + for row in rows: + rid = str(getattr(row, "id", "") or "") + if cur and rid == cur: + continue + status = str(getattr(row, "status", "") or "").strip().lower() + if status in {"queued", "running", ""}: + continue + return { + "id": rid, + "status": status, + "finished_at": str(getattr(row, "finished_at", "") or "") or None, + "created_at": str(getattr(row, "created_at", "") or "") or None, + "reply_text": str(getattr(row, "reply_text", "") or ""), + "error": str(getattr(row, "error", "") or ""), + } + return None + + +def _maybe_backfill_job_recipe( + store: Any, + *, + job: Any, + recipe: dict[str, Any], +) -> None: + """Persist synthesized recipe onto jobs that only stored a long prompt_text.""" + if not recipe_has_playbook(recipe): + return + if not recipe_is_empty(load_recipe_from_job(job)): + return + updater = getattr(store, "scheduled_job_update", None) + if not callable(updater): + return + try: + stamped = dict(recipe) + src = dict(stamped.get("source") or {}) + if not str(src.get("compiled_at") or "").strip(): + src["compiled_at"] = datetime.now(timezone.utc).isoformat() + src["compiled_from"] = str(src.get("compiled_from") or "prompt_text") + stamped["source"] = src + updater( + job_id=str(getattr(job, "id") or ""), + tenant_id=str(getattr(job, "tenant_id") or ""), + patch={"recipe": stamped}, + ) + except Exception: + pass + + def enqueue_scheduled_job_run( store: Any, *, @@ -77,13 +146,27 @@ def enqueue_scheduled_job_run( agent_run_id = uuid.uuid4().hex prompt_text = str(getattr(job, "prompt_text", "") or "").strip() lang = str(getattr(job, "lang", "") or "en") - recipe = load_recipe_from_job(job) + stored_recipe = load_recipe_from_job(job) + recipe = resolve_effective_playbook_recipe( + recipe=stored_recipe, + prompt_text=prompt_text, + session_id=str(getattr(job, "source_session_id", "") or ""), + ) playbook = recipe_has_playbook(recipe) + if playbook and recipe is not None: + _maybe_backfill_job_recipe(store, job=job, recipe=recipe) + previous_run = _load_previous_run_context( + store, + job_id=job_id, + tenant_id=tenant_id, + current_run_id=str(getattr(run, "id", "") or ""), + ) user_text = build_scheduled_turn_instruction( prompt_text=prompt_text, mode=mode, lang=lang, recipe=recipe if playbook else None, + previous_run=previous_run, ) viewer_username = resolve_scheduled_viewer_username( store, @@ -105,6 +188,7 @@ def enqueue_scheduled_job_run( "text": user_text, "prompt_text": prompt_text, "recipe": recipe if playbook else {}, + "previous_run": previous_run or {}, "attachments": [], "metadata": { "scheduled_job_id": job_id, diff --git a/runtime/scheduler/turn_text.py b/runtime/scheduler/turn_text.py index 62d2a57e..1a51b7f4 100644 --- a/runtime/scheduler/turn_text.py +++ b/runtime/scheduler/turn_text.py @@ -75,27 +75,77 @@ def build_scheduled_turn_instruction( mode: str, lang: str, recipe: dict[str, Any] | None = None, + previous_run: dict[str, Any] | None = None, ) -> str: """Internal LLM instruction for proactive scheduled reminders/playbooks (not user-facing).""" _ = str(mode or "scheduled").strip() if recipe_has_playbook(recipe): - return compile_playbook_instruction(recipe=recipe or {}, lang=lang) + text = compile_playbook_instruction(recipe=recipe or {}, lang=lang) + else: + intent = str(prompt_text or "").strip() + is_en = str(lang or "").lower().startswith("en") + if is_en: + text = ( + "[Scheduled proactive reminder — internal instruction, not a user message]\n" + f"Reminder intent: {intent}\n" + "Write a short, friendly proactive reminder TO the user (second person). " + "Do not say you received a reminder or that you will remind someone; speak directly to the user." + ) + else: + text = ( + "【定时主动提醒·内部指令,不是用户发言】\n" + f"提醒意图:{intent}\n" + "请生成一条简短、自然、第二人称的主动提醒消息直接对用户说。" + "不要写「收到提醒」「好的我来提醒用户」等元对话;不要假装用户刚说了话。" + ) + return append_previous_run_context(text, previous_run=previous_run, lang=lang) - intent = str(prompt_text or "").strip() + +def append_previous_run_context( + instruction: str, + *, + previous_run: dict[str, Any] | None, + lang: str = "en", + max_body_chars: int = 800, +) -> str: + """Append a short prior-run note so recurring jobs can compare deltas.""" + base = str(instruction or "").rstrip() + if not previous_run or not isinstance(previous_run, dict): + return base + status = str(previous_run.get("status") or "").strip() or "unknown" + finished = str(previous_run.get("finished_at") or previous_run.get("created_at") or "").strip() + err = " ".join(str(previous_run.get("error") or "").split()) + body = " ".join(str(previous_run.get("reply_text") or "").split()) + # Prefer error text on failures; otherwise the outbound summary. + if status.lower() in {"failed", "error"} and err: + body = err + cap = max(120, int(max_body_chars or 800)) + if len(body) > cap: + body = body[: cap - 3].rstrip() + "..." + if not body and not finished: + return base is_en = str(lang or "").lower().startswith("en") if is_en: - return ( - "[Scheduled proactive reminder — internal instruction, not a user message]\n" - f"Reminder intent: {intent}\n" - "Write a short, friendly proactive reminder TO the user (second person). " - "Do not say you received a reminder or that you will remind someone; speak directly to the user." + lines = [ + "", + "[Previous run context — for continuity only; still execute this run fully]", + f"Status: {status}" + (f" | Finished: {finished}" if finished else ""), + ] + if body: + lines.append(f"Summary: {body}") + lines.append( + "Use this for deltas/comparisons when useful; do not skip work just because the prior run succeeded." ) - return ( - "【定时主动提醒·内部指令,不是用户发言】\n" - f"提醒意图:{intent}\n" - "请生成一条简短、自然、第二人称的主动提醒消息直接对用户说。" - "不要写「收到提醒」「好的我来提醒用户」等元对话;不要假装用户刚说了话。" - ) + else: + lines = [ + "", + "【上一轮运行摘要·仅供对照;本轮仍须完整执行】", + f"状态:{status}" + (f"|完成时间:{finished}" if finished else ""), + ] + if body: + lines.append(f"摘要:{body}") + lines.append("可参考做环比/差异,但不要因上轮成功而跳过本轮步骤。") + return base + "\n" + "\n".join(lines) def scheduled_turn_system_suffix(*, lang: str, playbook: bool = False) -> str: @@ -128,6 +178,7 @@ def scheduled_turn_system_suffix(*, lang: str, playbook: bool = False) -> str: __all__ = [ + "append_previous_run_context", "build_scheduled_turn_instruction", "format_scheduled_failure_summary", "format_scheduled_success_summary", diff --git a/runtime/tools/experts/productivity/schedule_tools.py b/runtime/tools/experts/productivity/schedule_tools.py index 957bec71..95da5a9c 100644 --- a/runtime/tools/experts/productivity/schedule_tools.py +++ b/runtime/tools/experts/productivity/schedule_tools.py @@ -23,6 +23,7 @@ from runtime.scheduler.recipe import ( recipe_has_playbook, recipe_missing_fields, resolve_ops_recipe_template, + synthesize_recipe_from_prompt, ) from runtime.scheduler.service import run_scheduled_job_now from runtime.scheduler.system_timezone import default_system_timezone @@ -255,18 +256,29 @@ def schedule_create_tool() -> ToolSpec: needs_recipe = looks_like_complex_schedule_prompt(prompt_text, recipe=recipe) if needs_recipe and not recipe_has_playbook(recipe): - missing = recipe_missing_fields(recipe) or ["goal", "steps", "success_criteria"] - return { - "ok": False, - "error": "recipe_required", - "missing_fields": missing, - "hint": ( - "This looks like a complex/multi-step job. Call schedule_propose first, " - "show the draft to the user, then schedule_create with a full recipe " - "(goal + >=2 steps + success_criteria). Do not store vague prompts like " - "'继续刚才那个'." - ), - } + # Prefer auto-compiling a durable recipe from a structured prompt + # (Step 1/2..., long algorithm text) over rejecting create. + synth = synthesize_recipe_from_prompt( + prompt_text, + session_id=str(args.get("session_id") or ""), + ) + if recipe_has_playbook(synth): + recipe = synth + if not str((recipe.get("source") or {}).get("compiled_at") or "").strip(): + recipe.setdefault("source", {})["compiled_at"] = datetime.now(timezone.utc).isoformat() + else: + missing = recipe_missing_fields(recipe) or ["goal", "steps", "success_criteria"] + return { + "ok": False, + "error": "recipe_required", + "missing_fields": missing, + "hint": ( + "This looks like a complex/multi-step job. Call schedule_propose first, " + "show the draft to the user, then schedule_create with a full recipe " + "(goal + >=2 steps + success_criteria). Do not store vague prompts like " + "'继续刚才那个'." + ), + } if recipe and not recipe_has_playbook(recipe) and recipe_missing_fields(recipe): # Explicit but incomplete recipe → reject rather than silently drop. return { diff --git a/tests/test_schedule_recipe.py b/tests/test_schedule_recipe.py index 6b21981b..62d61a83 100644 --- a/tests/test_schedule_recipe.py +++ b/tests/test_schedule_recipe.py @@ -104,6 +104,46 @@ class RecipeHelpersTests(unittest.TestCase): self.assertIn("工作流", scheduled_turn_system_suffix(lang="zh", playbook=True)) self.assertIn("提醒", scheduled_turn_system_suffix(lang="zh", playbook=False)) + def test_synthesize_step_headers_and_previous_run(self) -> None: + from runtime.scheduler.recipe import synthesize_recipe_from_prompt + from runtime.scheduler.turn_text import append_previous_run_context + + prompt = ( + 'Run "unmanaged by dying gasp" report.\n\n' + "CRITICAL — Follow this exact algorithm:\n\n" + "Step 1 — Query BN EMS failures:\n" + "- Query active alarms: native_probable_cause = \"BN EMS alarm NE communication failure\"\n\n" + "Step 2 — Query Remote dying gasp:\n" + "- Query active alarms: native_probable_cause LIKE \"%Remote dying gasp%\"\n\n" + "Step 3 — Join and deliver XLSX via save_deliverable_attachment.\n" + ) + recipe = synthesize_recipe_from_prompt(prompt) + assert recipe is not None + self.assertTrue(recipe_has_playbook(recipe)) + self.assertGreaterEqual(len(recipe["steps"]), 3) + self.assertEqual((recipe.get("source") or {}).get("compiled_from"), "prompt_text") + self.assertTrue((recipe.get("output") or {}).get("need_attachments")) + instr = build_scheduled_turn_instruction( + prompt_text=prompt, + mode="scheduled", + lang="en", + recipe=recipe, + previous_run={ + "status": "success", + "finished_at": "2026-08-10T06:00:00+00:00", + "reply_text": "Found 12 candidates; 3 confirmed.", + }, + ) + self.assertIn("Scheduled playbook", instr) + self.assertIn("Previous run context", instr) + self.assertIn("Found 12 candidates", instr) + with_prev = append_previous_run_context( + "base", + previous_run={"status": "failed", "error": "timeout on CLI"}, + lang="en", + ) + self.assertIn("timeout on CLI", with_prev) + class ScheduleRecipeToolTests(unittest.TestCase): def setUp(self) -> None: @@ -220,6 +260,29 @@ class ScheduleRecipeToolTests(unittest.TestCase): self.assertEqual((recipe.get("source") or {}).get("template_id"), "ume_alarm_tally_daily") self.assertIn("UME", str(job.get("prompt_text") or "")) + def test_create_auto_compiles_structured_prompt(self) -> None: + prompt = ( + "Daily unmanaged dying-gasp correlation.\n\n" + "Step 1 — Query BN EMS communication failures.\n" + "Step 2 — Query remote dying gasp alarms.\n" + "Step 3 — Correlate and attach XLSX report.\n" + ) + out = schedule_create_tool().handler( + { + "tenant_id": self.tenant_id, + "owner_user_id": self.user_id, + "name": "dying-gasp-daily", + "prompt_text": prompt, + "schedule_kind": "cron", + "schedule_expr": "0 14 * * *", + "lang": "en", + } + ) + self.assertTrue(out.get("ok"), out) + recipe = (out.get("job") or {}).get("recipe") or {} + self.assertTrue(recipe_has_playbook(recipe)) + self.assertGreaterEqual(len(recipe.get("steps") or []), 3) + def test_simple_reminder_still_works(self) -> None: out = schedule_create_tool().handler( { diff --git a/tests/test_scheduler_playbook_backfill.py b/tests/test_scheduler_playbook_backfill.py new file mode 100644 index 00000000..eccc653e --- /dev/null +++ b/tests/test_scheduler_playbook_backfill.py @@ -0,0 +1,120 @@ +from __future__ import annotations + +import os +import tempfile +import unittest +from pathlib import Path +from unittest.mock import MagicMock + +from runtime.scheduler.service import enqueue_scheduled_job_run +from svc.persistence.assistant_store import reset_assistant_store_singleton +from svc.persistence.sqlite_store import SqliteStore + + +class EnqueuePlaybookBackfillTests(unittest.TestCase): + def setUp(self) -> None: + self._tmp = tempfile.TemporaryDirectory(ignore_cleanup_errors=True) + self.db = Path(self._tmp.name) / "enqueue.sqlite" + os.environ["OPS_ASSISTANT_DB_PATH"] = str(self.db) + os.environ["AIA_ASSISTANT_DB_BACKEND"] = "sqlite" + reset_assistant_store_singleton() + self.store = SqliteStore(str(self.db)) + t = self.store.create_tenant("Team") + self.tenant_id = str(t["id"]) + user = self.store.create_user_account( + tenant_id=self.tenant_id, + username="administrator", + display_name="Admin", + role="owner", + password_hash="x", + is_active=True, + ) + self.user_id = str(user["id"]) + + def tearDown(self) -> None: + reset_assistant_store_singleton() + self._tmp.cleanup() + + def test_enqueue_synthesizes_recipe_and_injects_previous_run(self) -> None: + prompt = ( + "Hourly dying-gasp check.\n\n" + "Step 1 — Query BN EMS failures.\n" + "Step 2 — Query remote dying gasp.\n" + "Step 3 — Correlate and report.\n" + ) + job = self.store.scheduled_job_create( + tenant_id=self.tenant_id, + name="dying-gasp", + prompt_text=prompt, + schedule_kind="interval", + schedule_expr="3600", + timezone_name="UTC", + lang="en", + delivery={"channel": "admin_chat"}, + recipe={}, + created_by_user_id=self.user_id, + ) + prev = self.store.scheduled_job_run_create( + job_id=job.id, + tenant_id=self.tenant_id, + status="queued", + ) + self.store.scheduled_job_run_update( + run_id=prev.id, + tenant_id=self.tenant_id, + patch={ + "status": "success", + "reply_text": "Prior: 4 unmanaged NEs.", + "finished_at": "2026-08-10T06:00:00+00:00", + }, + ) + + captured: dict[str, object] = {} + + def _capture_create(**kwargs: object) -> MagicMock: + captured.update(kwargs) + return MagicMock(id="task-1") + + self.store.oclaw_task_create = _capture_create # type: ignore[method-assign] + + from runtime.scheduler import service as sched_service + from runtime.scheduler.session_resolver import ResolvedSession + + resolved = ResolvedSession( + session_id="sess-sched-1", + tenant_id=self.tenant_id, + user_id=self.user_id, + channel="admin_chat", + account_id="", + external_chat_id="", + external_user_id="", + is_group=False, + ) + original_resolve = sched_service.resolve_scheduled_session + original_worker = sched_service.ensure_worker_started + try: + sched_service.resolve_scheduled_session = MagicMock(return_value=resolved) + sched_service.ensure_worker_started = MagicMock(return_value="worker-1") + out = enqueue_scheduled_job_run(self.store, job=job, mode="scheduled") + finally: + sched_service.resolve_scheduled_session = original_resolve + sched_service.ensure_worker_started = original_worker + + self.assertTrue(out.get("ok"), out) + payload = captured.get("payload") + self.assertIsInstance(payload, dict) + assert isinstance(payload, dict) + self.assertTrue((payload.get("metadata") or {}).get("scheduled_playbook")) + text = str(payload.get("text") or "") + self.assertIn("Scheduled playbook", text) + self.assertIn("Previous run context", text) + self.assertIn("4 unmanaged", text) + + refreshed = self.store.scheduled_job_get(job_id=job.id, tenant_id=self.tenant_id) + assert refreshed is not None + self.assertIn("Query BN EMS", refreshed.recipe_json) + self.assertIn("compiled_from", refreshed.recipe_json) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_scheduler_viewer_username.py b/tests/test_scheduler_viewer_username.py index d46488c2..4b9794c7 100644 --- a/tests/test_scheduler_viewer_username.py +++ b/tests/test_scheduler_viewer_username.py @@ -69,6 +69,9 @@ class ResolveScheduledViewerUsernameTests(unittest.TestCase): job.specialist = "generalist" job.created_by_user_id = str(self.member["id"]) job.schedule_kind = "interval" + job.recipe_json = "{}" + job.source_session_id = "" + job.name = "stretch" run = MagicMock() run.id = "run-1"