From 47813f6a9412ddd2c3c6ffd73d1ea408c53f0052 Mon Sep 17 00:00:00 2001 From: oliver Date: Fri, 26 Jun 2026 15:56:10 +0800 Subject: [PATCH] feat(scheduler): add scheduled jobs with channel-aware delivery Introduce persisted cron/interval jobs, gateway scheduler loop, worker turns, and admin CRUD/edit UI. Route proactive reminders via the originating chat channel (WhatsApp vs WeChat), harden Weixin outbound with durable queue and PG-compatible polling, and fix chat UI to show each scheduled reminder separately. Co-authored-by: Cursor --- .../versions/002_scheduled_jobs.py | 80 +++ interfaces/admin/routes.py | 188 ++++++ interfaces/admin/static/app.js | 436 ++++++++++++ interfaces/admin/static/chat.html | 2 +- interfaces/admin/static/chat.js | 83 ++- interfaces/admin/static/index.html | 1 + interfaces/gateway/context_builder.py | 2 + interfaces/gateway/dispatcher.py | 2 + interfaces/gateway/server_methods/cron.py | 2 + interfaces/http/fastapi_app.py | 34 + interfaces/http/weixin_ilink_api.py | 65 +- requirements.txt | 1 + runtime/agent_core_attempt.py | 14 +- runtime/agent_core_run.py | 3 +- .../application/gateway/inbound_service.py | 100 ++- runtime/chat/tool_runtime.py | 18 +- runtime/direct_loop.py | 8 + .../weixin_bridge/official_runner.ts | 169 +++++ runtime/scheduler/__init__.py | 3 + runtime/scheduler/channel_delivery.py | 252 +++++++ runtime/scheduler/cron_service.py | 201 ++++++ runtime/scheduler/expressions.py | 77 +++ runtime/scheduler/service.py | 200 ++++++ runtime/scheduler/session_resolver.py | 262 ++++++++ runtime/scheduler/turn_text.py | 50 ++ runtime/scheduler/weixin_delivery.py | 90 +++ runtime/scheduler/worker_turn.py | 184 +++++ runtime/tools/context_inject.py | 70 ++ .../experts/productivity/schedule_tools.py | 334 +++++++++ runtime/tools/tool_validation.py | 13 +- runtime/worker.py | 37 +- svc/persistence/assistant_store_protocol.py | 14 + svc/persistence/scheduled_job_store.py | 634 ++++++++++++++++++ svc/persistence/sqlite_store.py | 251 ++++++- tests/test_admin_scheduled_jobs_api.py | 114 ++++ tests/test_gateway_cron_service.py | 49 ++ tests/test_schedule_duration_parse.py | 20 + tests/test_schedule_tools.py | 59 ++ tests/test_scheduler_ephemeral_context.py | 133 ++++ tests/test_scheduler_expressions.py | 130 ++++ tests/test_scheduler_proactive_delivery.py | 180 +++++ tests/test_scheduler_viewer_username.py | 116 ++++ tests/test_tool_context_inject.py | 93 +++ .../test_tool_runtime_schedule_validation.py | 77 +++ tests/test_weixin_scheduled_delivery.py | 148 ++++ 45 files changed, 4956 insertions(+), 43 deletions(-) create mode 100644 assistant_migrations/versions/002_scheduled_jobs.py create mode 100644 runtime/scheduler/__init__.py create mode 100644 runtime/scheduler/channel_delivery.py create mode 100644 runtime/scheduler/cron_service.py create mode 100644 runtime/scheduler/expressions.py create mode 100644 runtime/scheduler/service.py create mode 100644 runtime/scheduler/session_resolver.py create mode 100644 runtime/scheduler/turn_text.py create mode 100644 runtime/scheduler/weixin_delivery.py create mode 100644 runtime/scheduler/worker_turn.py create mode 100644 runtime/tools/context_inject.py create mode 100644 runtime/tools/experts/productivity/schedule_tools.py create mode 100644 svc/persistence/scheduled_job_store.py create mode 100644 tests/test_admin_scheduled_jobs_api.py create mode 100644 tests/test_gateway_cron_service.py create mode 100644 tests/test_schedule_duration_parse.py create mode 100644 tests/test_schedule_tools.py create mode 100644 tests/test_scheduler_ephemeral_context.py create mode 100644 tests/test_scheduler_expressions.py create mode 100644 tests/test_scheduler_proactive_delivery.py create mode 100644 tests/test_scheduler_viewer_username.py create mode 100644 tests/test_tool_context_inject.py create mode 100644 tests/test_tool_runtime_schedule_validation.py create mode 100644 tests/test_weixin_scheduled_delivery.py diff --git a/assistant_migrations/versions/002_scheduled_jobs.py b/assistant_migrations/versions/002_scheduled_jobs.py new file mode 100644 index 00000000..11826456 --- /dev/null +++ b/assistant_migrations/versions/002_scheduled_jobs.py @@ -0,0 +1,80 @@ +"""Add scheduled_job tables for PostgreSQL deployments.""" + +from __future__ import annotations + +from alembic import op +from sqlalchemy import text + +revision = "002_scheduled_jobs" +down_revision = "001_assistant_pg_initial" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + bind = op.get_bind() + if bind.dialect.name != "postgresql": + return + op.execute( + text( + """ + CREATE TABLE IF NOT EXISTS scheduled_job ( + id TEXT PRIMARY KEY, + tenant_id TEXT NOT NULL, + name TEXT NOT NULL, + description TEXT NOT NULL DEFAULT '', + status TEXT NOT NULL DEFAULT 'active', + schedule_kind TEXT NOT NULL, + schedule_expr TEXT NOT NULL, + timezone TEXT NOT NULL DEFAULT 'Asia/Shanghai', + prompt_text TEXT NOT NULL, + interaction_mode TEXT NOT NULL DEFAULT 'expert', + specialist TEXT NOT NULL DEFAULT 'generalist', + lang TEXT NOT NULL DEFAULT 'zh', + delivery_json TEXT NOT NULL DEFAULT '{}', + source_session_id TEXT, + created_by_user_id TEXT NOT NULL DEFAULT '', + source TEXT NOT NULL DEFAULT 'admin', + next_run_at TEXT, + last_run_at TEXT, + last_run_status TEXT NOT NULL DEFAULT '', + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL + ) + """ + ) + ) + op.execute(text("CREATE INDEX IF NOT EXISTS idx_scheduled_job_due ON scheduled_job(status, next_run_at)")) + op.execute( + text( + """ + CREATE TABLE IF NOT EXISTS scheduled_job_run ( + id TEXT PRIMARY KEY, + job_id TEXT NOT NULL REFERENCES scheduled_job(id) ON DELETE CASCADE, + tenant_id TEXT NOT NULL, + status TEXT NOT NULL, + scheduled_at TEXT NOT NULL, + started_at TEXT, + finished_at TEXT, + session_id TEXT, + oclaw_task_id TEXT, + run_id TEXT, + reply_text TEXT NOT NULL DEFAULT '', + delivery_status_json TEXT NOT NULL DEFAULT '{}', + error TEXT NOT NULL DEFAULT '', + created_at TEXT NOT NULL + ) + """ + ) + ) + op.execute( + text("CREATE INDEX IF NOT EXISTS idx_scheduled_job_run_job ON scheduled_job_run(job_id, created_at DESC)") + ) + + +def downgrade() -> None: + bind = op.get_bind() + if bind.dialect.name != "postgresql": + return + op.execute(text("DROP TABLE IF EXISTS scheduled_job_run")) + op.execute(text("DROP TABLE IF EXISTS scheduled_job")) diff --git a/interfaces/admin/routes.py b/interfaces/admin/routes.py index 80cdc25b..9f8a0753 100644 --- a/interfaces/admin/routes.py +++ b/interfaces/admin/routes.py @@ -1986,6 +1986,194 @@ def build_admin_router() -> APIRouter: }, } + @router.get("/admin/api/scheduled-jobs/meta/targets") + def api_scheduled_jobs_meta_targets( + tenant_id: str = Query(default=""), + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + from runtime.scheduler.session_resolver import resolve_weixin_binding + + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:read") + tid = str(tenant_id or ctx.get("tenant_id") or "default").strip() + _require_tenant_scope(ctx, tid) + aid = str(os.getenv("AIA_WHATSAPP_ACCOUNT_ID") or "wa-default").strip() + known = store.list_whatsapp_known_groups(tenant_id=tid, account_id=aid) + weixin = resolve_weixin_binding(store, tenant_id=tid) + return {"ok": True, "whatsapp_groups": known, "weixin_binding": weixin, "whatsapp_account_id": aid} + + @router.get("/admin/api/scheduled-jobs") + def api_scheduled_jobs_list( + status: str | None = Query(default=None), + limit: int = Query(default=100), + offset: int = Query(default=0), + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:read") + tenant_id = str(ctx.get("tenant_id") or "") + rows = store.scheduled_job_list( + tenant_id=tenant_id, + status=str(status or "").strip() or None, + limit=max(1, min(int(limit or 100), 300)), + offset=max(0, int(offset or 0)), + ) + return {"ok": True, "items": [store.scheduled_job_to_dict(r) for r in rows]} + + @router.get("/admin/api/scheduled-jobs/{job_id}") + def api_scheduled_jobs_get( + job_id: str, + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:read") + tenant_id = str(ctx.get("tenant_id") or "") + row = store.scheduled_job_get(job_id=str(job_id), tenant_id=tenant_id) + return {"ok": True, "job": None if not row else store.scheduled_job_to_dict(row)} + + @router.post("/admin/api/scheduled-jobs") + def api_scheduled_jobs_create( + payload: dict[str, Any], + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + from runtime.scheduler.cron_service import build_default_delivery + from runtime.scheduler.expressions import normalize_schedule_kind + + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:runtime:write") + tenant_id = str(ctx.get("tenant_id") or "") + name = str(payload.get("name") or "").strip() + prompt_text = str(payload.get("prompt_text") or "").strip() + schedule_kind = normalize_schedule_kind(payload.get("schedule_kind") or "cron") + schedule_expr = str(payload.get("schedule_expr") or "").strip() + if not name or not prompt_text or not schedule_expr: + return {"ok": False, "error": "name, prompt_text, schedule_expr are required"} + delivery = payload.get("delivery") if isinstance(payload.get("delivery"), dict) else None + if delivery is None: + wa_chat = str((payload.get("whatsapp") or {}).get("chat_id") if isinstance(payload.get("whatsapp"), dict) else payload.get("whatsapp_chat_id") or "") + delivery = build_default_delivery(store=store, tenant_id=tenant_id, whatsapp_chat_id=wa_chat) + row = store.scheduled_job_create( + tenant_id=tenant_id, + name=name, + prompt_text=prompt_text, + schedule_kind=schedule_kind, + schedule_expr=schedule_expr, + timezone_name=str(payload.get("timezone") or "Asia/Shanghai"), + description=str(payload.get("description") or ""), + interaction_mode=normalize_interaction_mode(payload.get("interaction_mode") or "expert"), + specialist=normalize_requested_specialist(payload.get("specialist") or "generalist"), + lang=str(payload.get("lang") or "zh"), + delivery=delivery, + source_session_id=str(payload.get("source_session_id") or "").strip() or None, + created_by_user_id=str(ctx.get("user_id") or ""), + source="admin", + ) + store.add_admin_audit_log( + actor_tenant_id=tenant_id, + actor_user_id=str(ctx.get("user_id") or ""), + action="scheduled_job_create", + target_type="scheduled_job", + target_id=row.id, + status="ok", + ) + return {"ok": True, "job": store.scheduled_job_to_dict(row)} + + @router.patch("/admin/api/scheduled-jobs/{job_id}") + def api_scheduled_jobs_patch( + job_id: str, + payload: dict[str, Any], + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:runtime:write") + tenant_id = str(ctx.get("tenant_id") or "") + row = store.scheduled_job_update(tenant_id=tenant_id, job_id=str(job_id), patch=dict(payload or {})) + if not row: + return {"ok": False, "error": "job_not_found"} + store.add_admin_audit_log( + actor_tenant_id=tenant_id, + actor_user_id=str(ctx.get("user_id") or ""), + action="scheduled_job_update", + target_type="scheduled_job", + target_id=str(job_id), + status="ok", + ) + return {"ok": True, "job": store.scheduled_job_to_dict(row)} + + @router.post("/admin/api/scheduled-jobs/{job_id}/pause") + def api_scheduled_jobs_pause( + job_id: str, + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:runtime:write") + tenant_id = str(ctx.get("tenant_id") or "") + ok = store.scheduled_job_set_status(tenant_id=tenant_id, job_id=str(job_id), status="paused") + return {"ok": bool(ok)} + + @router.post("/admin/api/scheduled-jobs/{job_id}/resume") + def api_scheduled_jobs_resume( + job_id: str, + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:runtime:write") + tenant_id = str(ctx.get("tenant_id") or "") + ok = store.scheduled_job_set_status(tenant_id=tenant_id, job_id=str(job_id), status="active") + return {"ok": bool(ok)} + + @router.delete("/admin/api/scheduled-jobs/{job_id}") + def api_scheduled_jobs_delete( + job_id: str, + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:runtime:write") + tenant_id = str(ctx.get("tenant_id") or "") + ok = store.scheduled_job_delete(tenant_id=tenant_id, job_id=str(job_id)) + return {"ok": bool(ok)} + + @router.post("/admin/api/scheduled-jobs/{job_id}/run-now") + def api_scheduled_jobs_run_now( + job_id: str, + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + from runtime.scheduler.service import run_scheduled_job_now + + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:runtime:write") + tenant_id = str(ctx.get("tenant_id") or "") + out = run_scheduled_job_now(store, tenant_id=tenant_id, job_id=str(job_id)) + return out + + @router.get("/admin/api/scheduled-jobs/{job_id}/runs") + def api_scheduled_jobs_runs( + job_id: str, + limit: int = Query(default=50), + offset: int = Query(default=0), + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + store = get_assistant_store() + ctx = _resolve_auth(store, authorization) + _require_permission(ctx, "admin:read") + tenant_id = str(ctx.get("tenant_id") or "") + rows = store.scheduled_job_run_list( + job_id=str(job_id), + tenant_id=tenant_id, + limit=max(1, min(int(limit or 50), 200)), + offset=max(0, int(offset or 0)), + ) + return {"ok": True, "items": [store.scheduled_job_run_to_dict(r) for r in rows]} + @router.get("/admin/api/replay/turn") def api_replay_turn( session_id: str, diff --git a/interfaces/admin/static/app.js b/interfaces/admin/static/app.js index d68694ec..2c925119 100644 --- a/interfaces/admin/static/app.js +++ b/interfaces/admin/static/app.js @@ -5,6 +5,7 @@ const I18N = { "nav.models": "模型管理", "nav.apiGrants": "API 使用授权", "nav.stack": "运行时", + "nav.scheduledJobs": "定时任务", "nav.users": "用户管理", "nav.workspacePaths": "工作区路径", "nav.memory": "记忆", @@ -18,6 +19,41 @@ const I18N = { "notice.noLogin": "v2 已启用登录鉴权:请仅在内网访问", "action.refresh": "刷新", "title.stack": "运行时", + "title.scheduledJobs": "定时任务", + "scheduledJobs.all": "全部", + "scheduledJobs.statusActive": "运行中", + "scheduledJobs.statusPaused": "已暂停", + "scheduledJobs.loading": "Loading…", + "scheduledJobs.count": "{count} 个任务", + "scheduledJobs.colName": "名称", + "scheduledJobs.colSchedule": "计划", + "scheduledJobs.colStatus": "状态", + "scheduledJobs.colNextRun": "下次运行", + "scheduledJobs.colLastRun": "上次运行", + "scheduledJobs.colSpecialist": "专家", + "scheduledJobs.colDelivery": "投递", + "scheduledJobs.colActions": "操作", + "scheduledJobs.menuTitle": "任务操作", + "scheduledJobs.viewRuns": "查看运行记录", + "scheduledJobs.pause": "暂停", + "scheduledJobs.resume": "恢复", + "scheduledJobs.runNow": "立即运行", + "scheduledJobs.delete": "删除", + "scheduledJobs.deleteConfirm": "确定删除该定时任务?", + "scheduledJobs.triggered": "已触发", + "scheduledJobs.createTitle": "新建任务", + "scheduledJobs.create": "创建", + "scheduledJobs.edit": "编辑", + "scheduledJobs.editTitle": "编辑任务", + "scheduledJobs.updated": "已保存", + "scheduledJobs.cancel": "取消", + "scheduledJobs.runHistory": "运行记录", + "scheduledJobs.runHistoryHint": "最近运行(自动刷新)", + "scheduledJobs.noRuns": "暂无运行记录", + "scheduledJobs.weixinFixed": "微信固定投递:{id}", + "scheduledJobs.weixinMissing": "未找到微信绑定(需管理员渠道身份)", + "scheduledJobs.deliveryWeixin": "微信", + "scheduledJobs.deliveryWhatsapp": "WhatsApp", "title.users": "用户管理", "title.workspacePaths": "工作区路径(按用户)", "title.memory": "记忆", @@ -447,6 +483,7 @@ const I18N = { "nav.models": "Models", "nav.apiGrants": "API access", "nav.stack": "Runtime", + "nav.scheduledJobs": "Scheduled Jobs", "nav.users": "Users", "nav.workspacePaths": "Workspace paths", "nav.memory": "Memory", @@ -460,6 +497,41 @@ const I18N = { "notice.noLogin": "v2 login enabled: internal network only", "action.refresh": "Refresh", "title.stack": "Runtime", + "title.scheduledJobs": "Scheduled Jobs", + "scheduledJobs.all": "All", + "scheduledJobs.statusActive": "active", + "scheduledJobs.statusPaused": "paused", + "scheduledJobs.loading": "Loading…", + "scheduledJobs.count": "{count} job(s)", + "scheduledJobs.colName": "name", + "scheduledJobs.colSchedule": "schedule", + "scheduledJobs.colStatus": "status", + "scheduledJobs.colNextRun": "next_run", + "scheduledJobs.colLastRun": "last_run", + "scheduledJobs.colSpecialist": "specialist", + "scheduledJobs.colDelivery": "delivery", + "scheduledJobs.colActions": "actions", + "scheduledJobs.menuTitle": "Job actions", + "scheduledJobs.viewRuns": "View runs", + "scheduledJobs.pause": "Pause", + "scheduledJobs.resume": "Resume", + "scheduledJobs.runNow": "Run now", + "scheduledJobs.delete": "Delete", + "scheduledJobs.deleteConfirm": "Delete this scheduled job?", + "scheduledJobs.triggered": "Triggered.", + "scheduledJobs.createTitle": "Create job", + "scheduledJobs.create": "Create", + "scheduledJobs.edit": "Edit", + "scheduledJobs.editTitle": "Edit job", + "scheduledJobs.updated": "Saved.", + "scheduledJobs.cancel": "Cancel", + "scheduledJobs.runHistory": "Run history", + "scheduledJobs.runHistoryHint": "Recent runs (auto refresh)", + "scheduledJobs.noRuns": "No runs yet", + "scheduledJobs.weixinFixed": "Weixin fixed target: {id}", + "scheduledJobs.weixinMissing": "Weixin binding not found (administrator channel identity required)", + "scheduledJobs.deliveryWeixin": "WeChat", + "scheduledJobs.deliveryWhatsapp": "WhatsApp", "title.users": "Users", "title.workspacePaths": "Workspace paths (per user)", "title.memory": "Memory", @@ -9221,6 +9293,366 @@ async function renderSkills() { ]); } +function formatScheduledJobDelivery(job) { + const d = (job && job.delivery && typeof job.delivery === "object") ? job.delivery : {}; + const parts = []; + if (d.weixin && d.weixin.enabled) parts.push(t("scheduledJobs.deliveryWeixin")); + if (d.whatsapp && d.whatsapp.enabled) parts.push(t("scheduledJobs.deliveryWhatsapp")); + return parts.length ? parts.join(" + ") : "—"; +} + +function buildScheduledJobDeliveryPayload(job, waChatId) { + const existing = (job && job.delivery && typeof job.delivery === "object") ? job.delivery : {}; + const wa = existing.whatsapp && typeof existing.whatsapp === "object" ? { ...existing.whatsapp } : {}; + const wx = + existing.weixin && typeof existing.weixin === "object" + ? { ...existing.weixin } + : { enabled: true, fixed: true }; + const chat = String(waChatId || "").trim(); + wa.enabled = Boolean(chat); + wa.chat_id = chat; + wa.target_type = chat.includes("@g.us") ? "group" : "direct"; + if (wx.enabled === undefined) wx.enabled = true; + return { whatsapp: wa, weixin: wx }; +} + +async function renderScheduledJobs() { + let scheduledJobMenuEl = null; + let editingJob = null; + const canWrite = hasPermission("admin:runtime:write"); + const closeScheduledJobMenu = () => { + if (scheduledJobMenuEl && scheduledJobMenuEl.parentNode) { + scheduledJobMenuEl.remove(); + } + scheduledJobMenuEl = null; + }; + + const status = el("select", {}, [ + el("option", { value: "", text: t("scheduledJobs.all") }), + el("option", { value: "active", text: t("scheduledJobs.statusActive") }), + el("option", { value: "paused", text: t("scheduledJobs.statusPaused") }), + ]); + const tbody = el("tbody", {}); + const runsBox = el("pre", { class: "code", style: "max-height:220px;overflow:auto;white-space:pre-wrap;" }); + const msg = el("div", { class: "muted", text: "" }); + + const editModal = el("div", { class: "session-monitor-modal", style: "display:none;" }); + const editTitle = el("div", { class: "card__title", text: t("scheduledJobs.editTitle") }); + const editNameInput = el("input", { class: "input", placeholder: "name" }); + const editKindInput = el("select", {}, [ + el("option", { value: "cron", text: "cron" }), + el("option", { value: "once", text: "once" }), + el("option", { value: "interval", text: "interval" }), + ]); + const editExprInput = el("input", { class: "input", placeholder: "schedule_expr (cron / ISO / seconds)" }); + const editPromptInput = el("textarea", { class: "input", rows: "4", placeholder: "prompt_text" }); + const editSpecialistInput = el("input", { class: "input", placeholder: "specialist" }); + const editWaChatInput = el("input", { class: "input", placeholder: "whatsapp chat_id (optional)" }); + const editWeixinInfo = el("div", { class: "muted", text: "" }); + const closeEditModal = () => { + editingJob = null; + editModal.style.display = "none"; + }; + const openEditModal = (job) => { + if (!job || !job.id) return; + editingJob = job; + editTitle.textContent = `${t("scheduledJobs.editTitle")}: ${String(job.name || job.id)}`; + editNameInput.value = String(job.name || ""); + editKindInput.value = String(job.schedule_kind || "cron"); + editExprInput.value = String(job.schedule_expr || ""); + editPromptInput.value = String(job.prompt_text || ""); + editSpecialistInput.value = String(job.specialist || "generalist"); + const wa = job.delivery && job.delivery.whatsapp ? job.delivery.whatsapp : {}; + editWaChatInput.value = String(wa.chat_id || ""); + const wx = job.delivery && job.delivery.weixin ? job.delivery.weixin : {}; + editWeixinInfo.textContent = wx.external_user_id + ? tf("scheduledJobs.weixinFixed", { id: String(wx.external_user_id) }) + : weixinInfo.textContent || t("scheduledJobs.weixinMissing"); + editModal.style.display = "flex"; + setTimeout(() => editNameInput.focus(), 0); + }; + editModal.addEventListener("click", (e) => { + if (e.target === editModal) closeEditModal(); + }); + const editSaveBtn = el("button", { + class: "btn btn--primary", + text: t("action.save"), + onclick: async () => { + if (!editingJob || !editingJob.id) return; + const name = String(editNameInput.value || "").trim(); + const promptText = String(editPromptInput.value || "").trim(); + const scheduleExpr = String(editExprInput.value || "").trim(); + if (!name || !promptText || !scheduleExpr) { + msg.textContent = "name, prompt_text, schedule_expr are required"; + return; + } + try { + await apiRequest("PATCH", `/admin/api/scheduled-jobs/${encodeURIComponent(String(editingJob.id))}`, { + name, + schedule_kind: editKindInput.value, + schedule_expr: scheduleExpr, + prompt_text: promptText, + specialist: String(editSpecialistInput.value || "generalist").trim() || "generalist", + delivery: buildScheduledJobDeliveryPayload(editingJob, editWaChatInput.value), + }); + closeEditModal(); + msg.textContent = t("scheduledJobs.updated"); + await loadJobs(); + } catch (e) { + msg.textContent = String(e && e.message ? e.message : e); + } + }, + }); + const editModalCard = el("div", { class: "card session-monitor-modal__card", style: "width:min(640px,96vw);" }, [ + editTitle, + editWeixinInfo, + editNameInput, + el("div", { class: "row", style: "gap:8px;" }, [editKindInput, editExprInput]), + editPromptInput, + editSpecialistInput, + editWaChatInput, + el("div", { class: "row", style: "gap:8px;justify-content:flex-end;margin-top:10px;" }, [ + el("button", { class: "btn", text: t("scheduledJobs.cancel"), onclick: closeEditModal }), + editSaveBtn, + ]), + ]); + editModal.appendChild(editModalCard); + + async function loadLatestRuns(items) { + const jobs = Array.isArray(items) ? items : []; + if (!jobs.length) { + runsBox.textContent = t("scheduledJobs.noRuns"); + return; + } + const withLast = jobs.filter((j) => String(j.last_run_at || "").trim()); + const target = withLast.length ? withLast[0] : jobs[0]; + if (!target || !target.id) { + runsBox.textContent = t("scheduledJobs.noRuns"); + return; + } + const runs = await apiGet(`/admin/api/scheduled-jobs/${encodeURIComponent(String(target.id))}/runs?limit=10`); + const runItems = runs.items || []; + runsBox.textContent = runItems.length + ? JSON.stringify(runItems, null, 2) + : t("scheduledJobs.noRuns"); + } + + async function loadJobs() { + closeScheduledJobMenu(); + msg.textContent = t("scheduledJobs.loading"); + const st = String(status.value || "").trim(); + const q = st ? `?status=${encodeURIComponent(st)}` : ""; + const resp = await apiGet(`/admin/api/scheduled-jobs${q}`); + tbody.replaceChildren(); + for (const job of resp.items || []) { + const jobId = String(job.id || ""); + const btnMore = el("button", { + class: "chat-sess-more", + text: "⋯", + title: t("scheduledJobs.menuTitle"), + onclick: (ev) => { + ev.stopPropagation(); + closeScheduledJobMenu(); + const menu = el("div", { class: "chat-sess-menu-pop", style: "position:fixed;z-index:250;" }, [ + el("button", { + class: "chat-sess-menu-item", + text: t("scheduledJobs.viewRuns"), + onclick: async () => { + closeScheduledJobMenu(); + const runs = await apiGet(`/admin/api/scheduled-jobs/${encodeURIComponent(jobId)}/runs`); + runsBox.textContent = JSON.stringify(runs.items || [], null, 2); + }, + }), + ...(canWrite + ? [ + el("button", { + class: "chat-sess-menu-item", + text: t("scheduledJobs.edit"), + onclick: () => { + closeScheduledJobMenu(); + openEditModal(job); + }, + }), + ] + : []), + el("button", { + class: "chat-sess-menu-item", + text: job.status === "active" ? t("scheduledJobs.pause") : t("scheduledJobs.resume"), + onclick: async () => { + closeScheduledJobMenu(); + const path = + job.status === "active" + ? `/admin/api/scheduled-jobs/${encodeURIComponent(jobId)}/pause` + : `/admin/api/scheduled-jobs/${encodeURIComponent(jobId)}/resume`; + await apiPost(path, {}); + await loadJobs(); + }, + }), + el("button", { + class: "chat-sess-menu-item", + text: t("scheduledJobs.runNow"), + onclick: async () => { + closeScheduledJobMenu(); + await apiPost(`/admin/api/scheduled-jobs/${encodeURIComponent(jobId)}/run-now`, {}); + msg.textContent = t("scheduledJobs.triggered"); + }, + }), + ...(canWrite + ? [ + el("button", { + class: "chat-sess-menu-item", + text: t("scheduledJobs.delete"), + onclick: async () => { + closeScheduledJobMenu(); + if (!window.confirm(t("scheduledJobs.deleteConfirm"))) return; + await fetch(resolveAdminApiUrl(`/admin/api/scheduled-jobs/${encodeURIComponent(jobId)}`), { + method: "DELETE", + headers: { authorization: `Bearer ${getStoredAuthToken()}`, accept: "application/json" }, + }); + await loadJobs(); + }, + }), + ] + : []), + ]); + const rect = ev.currentTarget.getBoundingClientRect(); + document.body.appendChild(menu); + const mrect = menu.getBoundingClientRect(); + const pad = 8; + let left = rect.left - 120; + let top = rect.bottom + 4; + if (top + mrect.height > window.innerHeight - pad) { + top = rect.top - 4 - mrect.height; + } + left = Math.max(pad, Math.min(left, window.innerWidth - pad - mrect.width)); + top = Math.max(pad, Math.min(top, window.innerHeight - pad - mrect.height)); + menu.style.left = `${left}px`; + menu.style.top = `${top}px`; + scheduledJobMenuEl = menu; + const close = (e) => { + if (!menu.contains(e.target)) { + closeScheduledJobMenu(); + document.removeEventListener("click", close); + } + }; + setTimeout(() => document.addEventListener("click", close), 0); + }, + }); + const lastRun = job.last_run_at + ? `${formatSystemLocalDateTime(String(job.last_run_at))}${job.last_run_status ? ` (${job.last_run_status})` : ""}` + : "—"; + const tr = el("tr", {}, [ + tdCell(job.name || "", 24), + tdCell(`${job.schedule_kind}:${job.schedule_expr}`, 28), + tdCell(job.status || "", 10), + tdCell(job.next_run_at || "—", 20), + tdCell(lastRun, 24), + tdCell(job.specialist || "", 14), + tdCell(formatScheduledJobDelivery(job), 12), + el("td", { class: "table__cell--actions" }, [ + el("div", { class: "table__cell-actions" }, [btnMore]), + ]), + ]); + tbody.appendChild(tr); + } + msg.textContent = tf("scheduledJobs.count", { count: String((resp.items || []).length) }); + await loadLatestRuns(resp.items || []); + } + + const nameInput = el("input", { class: "input", placeholder: "name" }); + const kindInput = el("select", {}, [ + el("option", { value: "cron", text: "cron" }), + el("option", { value: "once", text: "once" }), + el("option", { value: "interval", text: "interval" }), + ]); + const exprInput = el("input", { class: "input", placeholder: "schedule_expr (cron / ISO / seconds)" }); + const promptInput = el("textarea", { class: "input", rows: "3", placeholder: "prompt_text" }); + const specialistInput = el("input", { class: "input", placeholder: "specialist", value: "generalist" }); + const waChatInput = el("input", { class: "input", placeholder: "whatsapp chat_id (optional)" }); + const weixinInfo = el("div", { class: "muted", text: "" }); + + async function loadMeta() { + const meta = await apiGet("/admin/api/scheduled-jobs/meta/targets"); + const wx = meta.weixin_binding || {}; + weixinInfo.textContent = wx.external_user_id + ? tf("scheduledJobs.weixinFixed", { id: String(wx.external_user_id) }) + : t("scheduledJobs.weixinMissing"); + } + + await loadMeta(); + await loadJobs(); + + const createCardChildren = [ + el("div", { class: "card__title", text: t("scheduledJobs.createTitle") }), + weixinInfo, + nameInput, + el("div", { class: "row", style: "gap:8px;" }, [kindInput, exprInput]), + promptInput, + specialistInput, + waChatInput, + el("button", { + class: "btn btn--primary", + text: t("scheduledJobs.create"), + onclick: async () => { + await apiPost("/admin/api/scheduled-jobs", { + name: nameInput.value, + schedule_kind: kindInput.value, + schedule_expr: exprInput.value, + prompt_text: promptInput.value, + specialist: specialistInput.value, + whatsapp_chat_id: waChatInput.value, + delivery: buildScheduledJobDeliveryPayload(null, waChatInput.value), + }); + await loadJobs(); + }, + }), + ]; + + return el("div", {}, [ + editModal, + el("div", { class: "card" }, [ + el("div", { class: "card__title", text: t("title.scheduledJobs") }), + el("div", { class: "row", style: "gap:8px;align-items:center;" }, [ + status, + el("button", { class: "btn", text: t("action.refresh"), onclick: () => loadJobs() }), + msg, + ]), + el("div", { class: "table-wrap" }, [ + el("table", { class: "table table--compact" }, [ + el("colgroup", {}, [ + el("col", { style: "width:16%" }), + el("col", { style: "width:18%" }), + el("col", { style: "width:8%" }), + el("col", { style: "width:14%" }), + el("col", { style: "width:14%" }), + el("col", { style: "width:10%" }), + el("col", { style: "width:10%" }), + el("col", { style: "width:10%" }), + ]), + el("thead", {}, [ + el("tr", {}, [ + el("th", { text: t("scheduledJobs.colName") }), + el("th", { text: t("scheduledJobs.colSchedule") }), + el("th", { text: t("scheduledJobs.colStatus") }), + el("th", { text: t("scheduledJobs.colNextRun") }), + el("th", { text: t("scheduledJobs.colLastRun") }), + el("th", { text: t("scheduledJobs.colSpecialist") }), + el("th", { text: t("scheduledJobs.colDelivery") }), + el("th", { text: t("scheduledJobs.colActions") }), + ]), + ]), + tbody, + ]), + ]), + ]), + ...(canWrite ? [el("div", { class: "card" }, createCardChildren)] : []), + el("div", { class: "card" }, [ + el("div", { class: "card__title", text: `${t("scheduledJobs.runHistory")} · ${t("scheduledJobs.runHistoryHint")}` }), + runsBox, + ]), + ]); +} + async function router() { const route = getRoute(); const page = route.page; @@ -9243,6 +9675,7 @@ async function router() { document.querySelectorAll(".nav__item").forEach((a) => { const p = String(a.dataset.page || ""); if (p === "stack" && !hasPermission("admin:runtime:write")) a.style.display = "none"; + else if (p === "scheduled-jobs" && !hasPermission("admin:read")) a.style.display = "none"; else if (p === "users" && !hasPermission("admin:user:read")) a.style.display = "none"; else if (p === "session-monitor" && !isAdministratorUsername()) a.style.display = "none"; else if (p === "admin-audit" && !hasPermission("admin:user:write")) a.style.display = "none"; @@ -9271,6 +9704,9 @@ async function router() { ); view = await renderStack(); } + else if (page === "scheduled-jobs") { + view = hasPermission("admin:read") ? await renderScheduledJobs() : forbiddenCard(); + } else if (page === "users") { view = hasPermission("admin:user:read") ? await renderUserManagement() : forbiddenCard(); } else if (page === "memory") view = await renderMemory(); diff --git a/interfaces/admin/static/chat.html b/interfaces/admin/static/chat.html index e682796c..f28f22a9 100644 --- a/interfaces/admin/static/chat.html +++ b/interfaces/admin/static/chat.html @@ -574,7 +574,7 @@ diff --git a/interfaces/admin/static/chat.js b/interfaces/admin/static/chat.js index a308035c..1ddf522f 100644 --- a/interfaces/admin/static/chat.js +++ b/interfaces/admin/static/chat.js @@ -408,6 +408,58 @@ function _normalizeEventType(v) { return String(v || "").trim().toLowerCase(); } +function _isAssistantBodyEventType(eventType) { + const et = _normalizeEventType(eventType); + return !et || et === "assistant_text" || et === "assistant" || et === "scheduled_reminder"; +} + +function _parseEventPayload(raw) { + let ep = raw; + if (typeof ep === "string" && String(ep).trim()) { + try { + ep = JSON.parse(ep); + } catch (_) { + ep = null; + } + } + return ep && typeof ep === "object" && !Array.isArray(ep) ? ep : null; +} + +function _isScheduledProactiveMessage(m) { + const eventType = _normalizeEventType(m && m.event_type); + if (eventType === "scheduled_reminder") return true; + const ep = _parseEventPayload(m && m.event_payload); + if (ep && ep.scheduled_proactive) return true; + const rc = String((ep && ep.reasoning_content) || "").trim(); + if (!rc) return false; + if (!/定时主动提醒|定时任务模式|scheduled reminder|proactive reminder/i.test(rc)) return false; + return eventType === "assistant_text" || eventType === "assistant" || !eventType; +} + +function _messageTurnUuid(m) { + return String((m && m.turn_uuid) || "").trim(); +} + +function _pushScheduledAssistantRow(rows, m, content, eventType) { + const body = String(content || "").trim(); + const attsParsed = parseAttachments(m && m.attachments); + if (!body && !attsParsed.length) return; + const piece = { + kind: "assistant_text", + text: content, + assistantEventType: eventType || "assistant_text", + }; + if (attsParsed.length) piece.attachments = attsParsed; + rows.push({ + role: "assistant", + content: body, + timestamp: (m && m.timestamp) != null ? m.timestamp : "", + attachments: m && m.attachments ? m.attachments : null, + _items: [piece], + _message_ids: m && m.id != null ? [m.id] : [], + }); +} + function _collapsedBlockNode(title, text) { const raw = String(text || ""); const clipped = raw.length > REASONING_BLOCK_MAX_CHARS ? raw.slice(0, REASONING_BLOCK_MAX_CHARS) : raw; @@ -610,6 +662,15 @@ function _buildRenderRows(msgs) { continue; } if (role === "assistant" || role === "tool" || role === "function") { + const turnUuid = _messageTurnUuid(m); + if (agg && turnUuid && agg._turn_uuid && turnUuid !== agg._turn_uuid) { + flush(); + } + if (role === "assistant" && _isScheduledProactiveMessage(m)) { + flush(); + _pushScheduledAssistantRow(rows, m, content, eventType); + continue; + } if (!agg) { agg = { role: "assistant", @@ -618,20 +679,16 @@ function _buildRenderRows(msgs) { attachments: null, _items: [], _message_ids: [], + _turn_uuid: turnUuid, }; + } else if (turnUuid && !agg._turn_uuid) { + agg._turn_uuid = turnUuid; } if (m && m.id != null) agg._message_ids.push(m.id); if (role === "assistant") { // thinking_mode_enabled: reasoning lives in event_payload.reasoning_content (no separate reasoning rows). - let ep = m && m.event_payload; - if (typeof ep === "string" && String(ep).trim()) { - try { - ep = JSON.parse(ep); - } catch (_) { - ep = null; - } - } - if (ep && typeof ep === "object" && !Array.isArray(ep)) { + const ep = _parseEventPayload(m && m.event_payload); + if (ep) { const rc = String(ep.reasoning_content || "").trim(); if (rc) { agg._items.push({ kind: "reasoning", text: rc }); @@ -649,7 +706,7 @@ function _buildRenderRows(msgs) { if (attsParsed.length) piece.attachments = attsParsed; agg._items.push(piece); } - } else if (eventType === "assistant_text" || eventType === "assistant" || !eventType) { + } else if (_isAssistantBodyEventType(eventType)) { const hasText = !!String(content || "").trim(); const attsParsed = parseAttachments(m && m.attachments); if (hasText || attsParsed.length) { @@ -4308,7 +4365,7 @@ async function renderChatUi() { ).trim(); const hasBodyText = out.some((row) => { const et = _normalizeEventType(row.event_type); - return et === "assistant_text" || et === "tool_call"; + return _isAssistantBodyEventType(et) || et === "tool_call"; }); if (textFallback && !hasBodyText) { out.push({ ...base, content: textFallback, event_type: fallbackEventType }); @@ -4327,7 +4384,7 @@ async function renderChatUi() { if (out.length) { const _isBodyRow = (row) => { const et = _normalizeEventType(row && row.event_type); - return et === "assistant_text" || et === "tool_call"; + return _isAssistantBodyEventType(et) || et === "tool_call"; }; let attachIdx = out.length - 1; for (let i = out.length - 1; i >= 0; i -= 1) { @@ -4408,7 +4465,7 @@ async function renderChatUi() { // Recovery should only accept visible assistant body, not intermediate // reasoning/tool-call events; otherwise we may terminate on a partial line. const et = String((m && m.event_type) || "").trim().toLowerCase(); - if (et && et !== "assistant_text" && et !== "assistant") continue; + if (et && !_isAssistantBodyEventType(et)) continue; if (!String((m && m.content) || "").trim()) continue; { last = m; diff --git a/interfaces/admin/static/index.html b/interfaces/admin/static/index.html index 32621ceb..6ca0b2d1 100644 --- a/interfaces/admin/static/index.html +++ b/interfaces/admin/static/index.html @@ -65,6 +65,7 @@