diff --git a/interfaces/admin/chat_api.py b/interfaces/admin/chat_api.py index 8b50dff0..760f785f 100644 --- a/interfaces/admin/chat_api.py +++ b/interfaces/admin/chat_api.py @@ -2320,6 +2320,72 @@ def include_chat_routes(router: APIRouter, *, resolve_auth: Callable[[SqliteStor }, ) + @chat.get("/jobs") + def api_chat_list_jobs( + limit: int = Query(default=50, ge=1, le=100), + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + store = get_assistant_store() + resolve_auth(store, authorization) + from svc.jobs.background_jobs import get_job_store + + return get_job_store().list_jobs(limit=limit) + + @chat.get("/jobs/{job_id}") + def api_chat_get_job( + job_id: str, + log_tail_chars: int = Query(default=4000, ge=200, le=50000), + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + store = get_assistant_store() + resolve_auth(store, authorization) + from svc.jobs.background_jobs import get_job_store + + out = get_job_store().get(job_id, log_tail_chars=log_tail_chars) + if not out.get("ok") and out.get("error") == "job_not_found": + raise HTTPException(status_code=404, detail="job_not_found") + return out + + @chat.post("/jobs/{job_id}/cancel") + def api_chat_cancel_job( + job_id: str, + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + store = get_assistant_store() + resolve_auth(store, authorization) + from svc.jobs.background_jobs import get_job_store + + out = get_job_store().cancel(job_id) + if not out.get("ok") and out.get("error") == "job_not_found": + raise HTTPException(status_code=404, detail="job_not_found") + return out + + @chat.post("/jobs/cancel-running") + def api_chat_cancel_all_running_jobs( + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + store = get_assistant_store() + resolve_auth(store, authorization) + from svc.jobs.background_jobs import get_job_store + + return get_job_store().cancel_all_running() + + @chat.delete("/jobs/{job_id}") + def api_chat_purge_job( + job_id: str, + authorization: str | None = Header(default=None), + ) -> dict[str, Any]: + store = get_assistant_store() + resolve_auth(store, authorization) + from svc.jobs.background_jobs import get_job_store + + out = get_job_store().purge(job_id) + if not out.get("ok") and out.get("error") == "job_not_found": + raise HTTPException(status_code=404, detail="job_not_found") + if not out.get("ok") and out.get("error") == "job_still_running": + raise HTTPException(status_code=409, detail="job_still_running") + return out + router.include_router(chat) diff --git a/interfaces/admin/static/chat.html b/interfaces/admin/static/chat.html index 0c3aec53..804cc607 100644 --- a/interfaces/admin/static/chat.html +++ b/interfaces/admin/static/chat.html @@ -516,6 +516,114 @@ line-height: 1.45; color: var(--ds-text, #e8e8e8); } + .chat-nav__toolbar { + display: flex; + flex-direction: column; + gap: 8px; + } + .chat-nav__jobs { + display: flex; + align-items: center; + justify-content: space-between; + gap: 8px; + width: 100%; + padding: 8px 10px; + border-radius: 8px; + border: 1px solid var(--ds-border, rgba(255, 255, 255, 0.12)); + background: transparent; + color: inherit; + cursor: pointer; + font-size: 13px; + } + .chat-nav__jobs:hover { + background: var(--chat-sess-hover-bg, rgba(255, 255, 255, 0.05)); + } + .chat-nav__jobs--active { + border-color: rgba(245, 158, 11, 0.55); + background: rgba(245, 158, 11, 0.12); + } + .chat-jobs-badge { + min-width: 1.25rem; + height: 1.25rem; + padding: 0 6px; + border-radius: 999px; + color: #fff; + font-size: 11px; + line-height: 1.25rem; + text-align: center; + font-weight: 600; + flex: 0 0 auto; + } + .chat-jobs-badge--idle { + background: rgba(148, 163, 184, 0.45); + } + .chat-jobs-badge--hot { + background: rgba(239, 68, 68, 0.9); + box-shadow: 0 0 0 2px rgba(239, 68, 68, 0.25); + } + .chat-jobs-row { + border: 1px solid var(--ds-border, rgba(255, 255, 255, 0.12)); + border-radius: 10px; + padding: 10px; + background: rgba(255, 255, 255, 0.03); + } + .chat-jobs-row--running { + border-color: rgba(245, 158, 11, 0.45); + } + .chat-jobs-row__top { + display: flex; + align-items: center; + justify-content: space-between; + gap: 8px; + margin-bottom: 6px; + } + .chat-jobs-row__name { + font-weight: 600; + font-size: 14px; + } + .chat-jobs-row__meta { + font-size: 12px; + line-height: 1.4; + margin-bottom: 8px; + } + .chat-jobs-row__cmd { + white-space: nowrap; + overflow: hidden; + text-overflow: ellipsis; + max-width: 100%; + } + .chat-jobs-pill { + font-size: 11px; + padding: 2px 8px; + border-radius: 999px; + border: 1px solid rgba(255, 255, 255, 0.15); + } + .chat-jobs-pill--running { + background: rgba(245, 158, 11, 0.2); + border-color: rgba(245, 158, 11, 0.4); + } + .chat-jobs-pill--succeeded { + background: rgba(34, 197, 94, 0.18); + border-color: rgba(34, 197, 94, 0.35); + } + .chat-jobs-pill--failed, + .chat-jobs-pill--timeout, + .chat-jobs-pill--cancelled { + background: rgba(239, 68, 68, 0.15); + border-color: rgba(239, 68, 68, 0.35); + } + .chat-jobs-log { + margin: 8px 0 0; + padding: 8px; + max-height: 220px; + overflow: auto; + font-size: 11px; + line-height: 1.35; + white-space: pre-wrap; + word-break: break-word; + background: rgba(0, 0, 0, 0.25); + border-radius: 8px; + } /* First paint before chat.js runs (desktop / slow network): avoid empty #app “black hole”. */ .chat-boot-splash { flex: 1; diff --git a/interfaces/admin/static/chat.js b/interfaces/admin/static/chat.js index 26bb3356..9510e37d 100644 --- a/interfaces/admin/static/chat.js +++ b/interfaces/admin/static/chat.js @@ -34,6 +34,24 @@ const I18N = { "chat.fork": "复制会话", "chat.audit": "审计与追踪", "chat.myProfile": "设置", + "chat.jobs": "后台任务", + "chat.jobsTitle": "后台任务", + "chat.jobsEmpty": "暂无后台任务", + "chat.jobsRefresh": "刷新", + "chat.jobsKill": "终止", + "chat.jobsKillConfirm": "确认终止该任务?进程将被强制结束。", + "chat.jobsKillAll": "全部终止", + "chat.jobsKillAllConfirm": "确认终止全部运行中的任务?", + "chat.jobsPurge": "清除记录", + "chat.jobsPurgeConfirm": "清除该已结束任务的本地记录与日志?", + "chat.jobsClose": "关闭", + "chat.jobsRunning": "运行中 {n}", + "chat.jobsHint": "Agent 启动的长任务会列在这里;对话结束后进程仍可能继续。可手动终止以防风险。", + "chat.jobsStatus.running": "运行中", + "chat.jobsStatus.succeeded": "成功", + "chat.jobsStatus.failed": "失败", + "chat.jobsStatus.timeout": "超时", + "chat.jobsStatus.cancelled": "已取消", "chat.attach": "附件", "chat.tools": "推理", "chat.tools.hidden": "推理已隐藏", @@ -214,6 +232,24 @@ const I18N = { "chat.fork": "Duplicate chat", "chat.audit": "Audit & trace", "chat.myProfile": "Settings", + "chat.jobs": "Background jobs", + "chat.jobsTitle": "Background jobs", + "chat.jobsEmpty": "No background jobs", + "chat.jobsRefresh": "Refresh", + "chat.jobsKill": "Kill", + "chat.jobsKillConfirm": "Kill this job? The process will be force-terminated.", + "chat.jobsKillAll": "Kill all running", + "chat.jobsKillAllConfirm": "Kill all running jobs?", + "chat.jobsPurge": "Purge", + "chat.jobsPurgeConfirm": "Remove local records/logs for this finished job?", + "chat.jobsClose": "Close", + "chat.jobsRunning": "{n} running", + "chat.jobsHint": "Long agent jobs appear here and may keep running after the chat turn ends. Kill them manually if needed.", + "chat.jobsStatus.running": "running", + "chat.jobsStatus.succeeded": "succeeded", + "chat.jobsStatus.failed": "failed", + "chat.jobsStatus.timeout": "timeout", + "chat.jobsStatus.cancelled": "cancelled", "chat.attach": "Attach", "chat.tools": "Reasoning", "chat.tools.hidden": "Reasoning hidden", @@ -892,6 +928,257 @@ function showToast(text, { kind = "info", ttlMs = 4200 } = {}) { return node; } +let _jobsPanelPollTimer = null; +let _jobsBadgePollTimer = null; +let _jobsBadgeEl = null; +let _jobsBtnLabelEl = null; + +function _jobStatusLabel(status) { + const s = String(status || "").trim().toLowerCase(); + const key = `chat.jobsStatus.${s}`; + const labeled = t(key); + return labeled === key ? s || "-" : labeled; +} + +function _fmtJobTs(ts) { + const n = Number(ts || 0); + if (!n) return "-"; + try { + return new Date(n * 1000).toLocaleString(currentLang === "zh" ? "zh-CN" : "en-US"); + } catch (_) { + return String(n); + } +} + +function updateJobsBadge(runningCount) { + const n = Math.max(0, parseInt(String(runningCount || 0), 10) || 0); + if (_jobsBtnLabelEl) { + _jobsBtnLabelEl.textContent = t("chat.jobs"); + } + if (!_jobsBadgeEl) return; + _jobsBadgeEl.textContent = String(n); + _jobsBadgeEl.hidden = false; + _jobsBadgeEl.setAttribute("aria-label", t("chat.jobsRunning", { n })); + if (n > 0) { + _jobsBadgeEl.classList.add("chat-jobs-badge--hot"); + _jobsBadgeEl.classList.remove("chat-jobs-badge--idle"); + } else { + _jobsBadgeEl.classList.remove("chat-jobs-badge--hot"); + _jobsBadgeEl.classList.add("chat-jobs-badge--idle"); + } + const btn = _jobsBadgeEl.closest && _jobsBadgeEl.closest(".chat-nav__jobs"); + if (btn) { + if (n > 0) btn.classList.add("chat-nav__jobs--active"); + else btn.classList.remove("chat-nav__jobs--active"); + btn.title = t("chat.jobsRunning", { n }); + } +} + +async function refreshJobsBadge() { + try { + const r = await apiGet("/admin/api/chat/jobs?limit=50"); + updateJobsBadge(r && r.running_count); + } catch (_) {} +} + +function startJobsBadgePoller() { + stopJobsBadgePoller(); + refreshJobsBadge(); + // Keep sidebar count fresh without opening the panel. + _jobsBadgePollTimer = setInterval(refreshJobsBadge, 4000); +} + +function stopJobsBadgePoller() { + if (_jobsBadgePollTimer) { + clearInterval(_jobsBadgePollTimer); + _jobsBadgePollTimer = null; + } +} + +async function openBackgroundJobsPanel() { + document.querySelectorAll(".chat-jobs-backdrop").forEach((n) => n.remove()); + if (_jobsPanelPollTimer) { + clearInterval(_jobsPanelPollTimer); + _jobsPanelPollTimer = null; + } + + const backdrop = el("div", { + class: "chat-confirm-backdrop chat-jobs-backdrop", + style: "z-index:9998;", + }); + const card = el("div", { + class: "chat-confirm-card chat-jobs-card", + style: "max-width:920px;width:94vw;max-height:82vh;display:flex;flex-direction:column;gap:10px;", + }); + const title = el("div", { class: "card__title", text: t("chat.jobsTitle") }); + const closeBtn = el("button", { type: "button", class: "btn", text: t("chat.jobsClose") }); + const refreshBtn = el("button", { type: "button", class: "btn", text: t("chat.jobsRefresh") }); + const killAllBtn = el("button", { + type: "button", + class: "btn btn--danger", + text: t("chat.jobsKillAll"), + }); + const head = el("div", { class: "row", style: "gap:8px;justify-content:space-between;align-items:center;flex-wrap:wrap;" }, [ + title, + el("div", { class: "row", style: "gap:8px;" }, [refreshBtn, killAllBtn, closeBtn]), + ]); + const hint = el("div", { class: "muted", style: "font-size:12px;line-height:1.45;", text: t("chat.jobsHint") }); + const summary = el("div", { class: "muted", style: "font-size:12px;" }); + const list = el("div", { + class: "chat-jobs-list", + style: "overflow:auto;flex:1;min-height:220px;max-height:58vh;display:flex;flex-direction:column;gap:8px;", + }); + + const close = () => { + if (_jobsPanelPollTimer) { + clearInterval(_jobsPanelPollTimer); + _jobsPanelPollTimer = null; + } + try { + backdrop.remove(); + } catch (_) {} + refreshJobsBadge(); + }; + closeBtn.addEventListener("click", close); + backdrop.addEventListener("click", (ev) => { + if (ev.target === backdrop) close(); + }); + + const paint = async () => { + list.innerHTML = ""; + list.appendChild(el("div", { class: "muted", text: t("chat.loading") })); + try { + const r = await apiGet("/admin/api/chat/jobs?limit=50"); + const jobs = Array.isArray(r && r.jobs) ? r.jobs : []; + const running = Number((r && r.running_count) || 0); + updateJobsBadge(running); + summary.textContent = t("chat.jobsRunning", { n: running }) + ` · total ${Number((r && r.total) || jobs.length)}`; + list.innerHTML = ""; + if (!jobs.length) { + list.appendChild(el("div", { class: "muted", text: t("chat.jobsEmpty") })); + return; + } + for (const job of jobs) { + const status = String(job.status || ""); + const runningJob = status === "running"; + const row = el("div", { class: `chat-jobs-row chat-jobs-row--${status || "unknown"}` }); + const top = el("div", { class: "chat-jobs-row__top" }, [ + el("span", { class: "chat-jobs-row__name", text: String(job.name || job.job_id || "-") }), + el("span", { class: `chat-jobs-pill chat-jobs-pill--${status}`, text: _jobStatusLabel(status) }), + ]); + const meta = el( + "div", + { class: "chat-jobs-row__meta muted" }, + [ + el("div", { text: `id: ${job.job_id || "-"}` }), + el("div", { text: `pid: ${job.pid != null ? job.pid : "-"} · timeout: ${job.timeout_s || "-"}s` }), + el("div", { text: `created: ${_fmtJobTs(job.created_at)} · finished: ${_fmtJobTs(job.finished_at)}` }), + el("div", { + class: "chat-jobs-row__cmd", + text: String(job.command || ""), + title: String(job.command || ""), + }), + ], + ); + const actions = el("div", { class: "chat-jobs-row__actions row", style: "gap:8px;" }); + if (runningJob) { + const killBtn = el("button", { + type: "button", + class: "btn btn--danger", + text: t("chat.jobsKill"), + onclick: async () => { + if (!window.confirm(t("chat.jobsKillConfirm"))) return; + killBtn.disabled = true; + try { + await apiPost(`/admin/api/chat/jobs/${encodeURIComponent(job.job_id)}/cancel`, {}); + await paint(); + } catch (e) { + showToast(`${t("chat.error")}: ${String(e)}`, { kind: "error" }); + } finally { + killBtn.disabled = false; + } + }, + }); + actions.appendChild(killBtn); + } else { + const purgeBtn = el("button", { + type: "button", + class: "btn", + text: t("chat.jobsPurge"), + onclick: async () => { + if (!window.confirm(t("chat.jobsPurgeConfirm"))) return; + purgeBtn.disabled = true; + try { + await apiDelete(`/admin/api/chat/jobs/${encodeURIComponent(job.job_id)}`); + await paint(); + } catch (e) { + showToast(`${t("chat.error")}: ${String(e)}`, { kind: "error" }); + } finally { + purgeBtn.disabled = false; + } + }, + }); + actions.appendChild(purgeBtn); + } + const detailBtn = el("button", { + type: "button", + class: "btn", + text: currentLang === "zh" ? "日志" : "Logs", + onclick: async () => { + detailBtn.disabled = true; + try { + const d = await apiGet(`/admin/api/chat/jobs/${encodeURIComponent(job.job_id)}?log_tail_chars=6000`); + const pre = el("pre", { + class: "chat-jobs-log", + text: + `status=${d.status} exit=${d.exit_code}\n\n--- stdout ---\n${d.stdout_tail || ""}\n\n--- stderr ---\n${d.stderr_tail || ""}`, + }); + const existing = row.querySelector(".chat-jobs-log"); + if (existing) existing.remove(); + row.appendChild(pre); + } catch (e) { + showToast(`${t("chat.error")}: ${String(e)}`, { kind: "error" }); + } finally { + detailBtn.disabled = false; + } + }, + }); + actions.appendChild(detailBtn); + row.appendChild(top); + row.appendChild(meta); + row.appendChild(actions); + list.appendChild(row); + } + } catch (e) { + list.innerHTML = ""; + list.appendChild(el("div", { class: "muted", text: `${t("chat.error")}: ${String(e)}` })); + } + }; + + refreshBtn.addEventListener("click", () => paint()); + killAllBtn.addEventListener("click", async () => { + if (!window.confirm(t("chat.jobsKillAllConfirm"))) return; + killAllBtn.disabled = true; + try { + await apiPost("/admin/api/chat/jobs/cancel-running", {}); + await paint(); + } catch (e) { + showToast(`${t("chat.error")}: ${String(e)}`, { kind: "error" }); + } finally { + killAllBtn.disabled = false; + } + }); + + card.appendChild(head); + card.appendChild(hint); + card.appendChild(summary); + card.appendChild(list); + backdrop.appendChild(card); + document.body.appendChild(backdrop); + await paint(); + _jobsPanelPollTimer = setInterval(paint, 4000); +} + async function openWikiPreviewModal({ sessionId, path }) { const sid = String(sessionId || ""); const p = String(path || "").replace(/\\/g, "/").replace(/^\//, ""); @@ -2737,6 +3024,12 @@ function syncAuthUserLabel() { "data-menu-action": "profile", text: t("chat.myProfile"), }), + el("button", { + type: "button", + class: "chat-sess-menu-item", + "data-menu-action": "jobs", + text: t("chat.jobs"), + }), ]; items.push(el("div", { class: "chat-sess-menu-sep" })); items.push( @@ -4041,6 +4334,25 @@ async function renderChatUi() { }), ); + const jobsBadge = el("span", { class: "chat-jobs-badge chat-jobs-badge--idle", text: "0" }); + _jobsBadgeEl = jobsBadge; + const jobsLabel = el("span", { class: "chat-nav__jobsLabel", text: t("chat.jobs") }); + _jobsBtnLabelEl = jobsLabel; + const btnJobs = el("button", { + type: "button", + class: "chat-nav__jobs", + title: t("chat.jobsRunning", { n: 0 }), + onclick: () => { + openBackgroundJobsPanel().catch((e) => showToast(`${t("chat.error")}: ${String(e)}`, { kind: "error" })); + }, + }); + btnJobs.appendChild(jobsLabel); + btnJobs.appendChild(jobsBadge); + startJobsBadgePoller(); + // Immediate paint so count is visible before first interval tick. + updateJobsBadge(0); + refreshJobsBadge(); + class OclawWsChatTransport { constructor({ tokenProvider }) { this.tokenProvider = tokenProvider; @@ -5490,7 +5802,7 @@ ${autoLimit ? `
auto-added claus el("div", { class: "chat-nav__top" }, [ buildChatBrandLogoNode(), ]), - el("div", { class: "chat-nav__toolbar" }, [btnNew]), + el("div", { class: "chat-nav__toolbar" }, [btnNew, btnJobs]), el("div", { class: "chat-nav__scroll" }, [sessionsListEl, loadMoreWrap]), navFooter, ]), @@ -5588,6 +5900,10 @@ document.body.addEventListener("click", async (e) => { openAdminFromChat("profile"); return; } + if (action === "jobs") { + openBackgroundJobsPanel().catch((err) => showToast(`${t("chat.error")}: ${String(err)}`, { kind: "error" })); + return; + } if (action === "logout") { try { await apiPost("/admin/api/auth/logout", {}); diff --git a/runtime/tools/public/job_tools.py b/runtime/tools/public/job_tools.py new file mode 100644 index 00000000..a510ceb9 --- /dev/null +++ b/runtime/tools/public/job_tools.py @@ -0,0 +1,173 @@ +from __future__ import annotations + +from typing import Any + +from runtime.tools.base import ToolSpec +from runtime.tools.path_guard import resolve_workspace_path +from svc.jobs.background_jobs import ( + DEFAULT_TIMEOUT_S, + MAX_TIMEOUT_S, + get_job_store, + is_shell_exec_enabled, +) + + +def start_job_tool() -> ToolSpec: + def _handler(args: dict[str, Any]) -> dict[str, Any]: + enabled, hint = is_shell_exec_enabled() + if not enabled: + return {"ok": False, "error": "disabled", "hint": hint} + command = str(args.get("command") or "").strip() + if not command: + return {"ok": False, "error": "command_required"} + cwd_raw = str(args.get("cwd") or "").strip() or "." + try: + cwd = str(resolve_workspace_path(cwd_raw)) + except ValueError as exc: + return {"ok": False, "error": str(exc)} + timeout_s = args.get("timeout_s") + name = str(args.get("name") or "").strip() + notify = args.get("notify") if isinstance(args.get("notify"), dict) else None + return get_job_store().start( + command=command, + cwd=cwd, + timeout_s=int(timeout_s) if timeout_s is not None else DEFAULT_TIMEOUT_S, + name=name, + notify=notify, + ) + + return ToolSpec( + name="start_job", + description=( + "Start a long-running shell command in the background and return job_id immediately " + f"(default timeout {DEFAULT_TIMEOUT_S}s / 2h, max {MAX_TIMEOUT_S}s / 3h). " + "The process keeps running after this agent turn ends or the chat disconnects. " + "Tell the user the job_id and end the turn — do NOT sleep for hours. " + "Resume later with get_job/list_jobs. Optional notify={channel,chat_id,...} pings the " + "channel when the job finishes. Same enable gate as run_command (AIA_ENABLE_RUN_COMMAND)." + ), + parameters={ + "type": "object", + "properties": { + "command": {"type": "string", "description": "Shell command to run in background."}, + "cwd": {"type": "string", "description": "Working directory (workspace-relative or allowed path)."}, + "timeout_s": { + "type": "integer", + "description": f"Kill after N seconds (default {DEFAULT_TIMEOUT_S}, max {MAX_TIMEOUT_S}).", + "default": DEFAULT_TIMEOUT_S, + }, + "name": {"type": "string", "description": "Optional human label for the job."}, + "notify": { + "type": "object", + "description": ( + "Optional channel ping on completion. " + "Fields: channel (whatsapp/weixin), chat_id, account_id?, tenant_id?, " + "context_token? (weixin), message? (custom text)." + ), + "properties": { + "channel": {"type": "string"}, + "chat_id": {"type": "string"}, + "account_id": {"type": "string"}, + "tenant_id": {"type": "string"}, + "context_token": {"type": "string"}, + "message": {"type": "string"}, + }, + "additionalProperties": False, + }, + }, + "required": ["command"], + "additionalProperties": False, + }, + handler=_handler, + tags=frozenset({"public", "exec", "job"}), + risk_level="high", + read_only=False, + timeout_s=30.0, + ) + + +def get_job_tool() -> ToolSpec: + def _handler(args: dict[str, Any]) -> dict[str, Any]: + job_id = str(args.get("job_id") or "").strip() + log_tail = int(args.get("log_tail_chars") or 4000) + return get_job_store().get(job_id, log_tail_chars=log_tail) + + return ToolSpec( + name="get_job", + description=( + "Get background job status by job_id (running/succeeded/failed/timeout/cancelled), " + "exit_code, and stdout/stderr tails. Poll until done=true." + ), + parameters={ + "type": "object", + "properties": { + "job_id": {"type": "string", "description": "Job id returned by start_job."}, + "log_tail_chars": { + "type": "integer", + "description": "Max characters of each log tail (default 4000).", + "default": 4000, + }, + }, + "required": ["job_id"], + "additionalProperties": False, + }, + handler=_handler, + tags=frozenset({"public", "job", "read"}), + risk_level="low", + read_only=True, + timeout_s=10.0, + ) + + +def cancel_job_tool() -> ToolSpec: + def _handler(args: dict[str, Any]) -> dict[str, Any]: + enabled, hint = is_shell_exec_enabled() + if not enabled: + return {"ok": False, "error": "disabled", "hint": hint} + job_id = str(args.get("job_id") or "").strip() + return get_job_store().cancel(job_id) + + return ToolSpec( + name="cancel_job", + description="Cancel a running background job (best-effort process tree kill).", + parameters={ + "type": "object", + "properties": { + "job_id": {"type": "string", "description": "Job id returned by start_job."}, + }, + "required": ["job_id"], + "additionalProperties": False, + }, + handler=_handler, + tags=frozenset({"public", "exec", "job"}), + risk_level="high", + read_only=False, + timeout_s=20.0, + ) + + +def list_jobs_tool() -> ToolSpec: + def _handler(args: dict[str, Any]) -> dict[str, Any]: + limit = int(args.get("limit") or 20) + return get_job_store().list_jobs(limit=limit) + + return ToolSpec( + name="list_jobs", + description="List recent background jobs (id, name, status, timestamps).", + parameters={ + "type": "object", + "properties": { + "limit": {"type": "integer", "description": "Max jobs to return (default 20).", "default": 20}, + }, + "required": [], + "additionalProperties": False, + }, + handler=_handler, + tags=frozenset({"public", "job", "read"}), + risk_level="low", + read_only=True, + timeout_s=10.0, + ) + + +__all__ = ["start_job_tool", "get_job_tool", "cancel_job_tool", "list_jobs_tool"] diff --git a/runtime/tools/public/sleep_tool.py b/runtime/tools/public/sleep_tool.py new file mode 100644 index 00000000..b97d6642 --- /dev/null +++ b/runtime/tools/public/sleep_tool.py @@ -0,0 +1,63 @@ +from __future__ import annotations + +import time +from typing import Any + +from runtime.tools.base import ToolSpec + +_MAX_SLEEP_S = 120 + + +def sleep_tool() -> ToolSpec: + def _handler(args: dict[str, Any]) -> dict[str, Any]: + try: + seconds = float(args.get("seconds") or 0) + except Exception: + return {"ok": False, "error": "invalid_seconds"} + if seconds <= 0: + return {"ok": False, "error": "seconds_must_be_positive"} + if seconds > _MAX_SLEEP_S: + return { + "ok": False, + "error": "seconds_too_large", + "max_seconds": _MAX_SLEEP_S, + "hint": ( + "For multi-hour work use start_job + get_job, or schedule_create. " + f"sleep is only for short poll gaps (max {_MAX_SLEEP_S}s)." + ), + } + started = time.time() + time.sleep(seconds) + return { + "ok": True, + "slept_seconds": round(time.time() - started, 3), + "requested_seconds": seconds, + } + + return ToolSpec( + name="sleep", + description=( + f"Block the current turn for N seconds (max {_MAX_SLEEP_S}). " + "Use between get_job polls or after a change that needs a short settle time. " + "Do NOT use for multi-hour waits — use start_job or schedule_create instead." + ), + parameters={ + "type": "object", + "properties": { + "seconds": { + "type": "number", + "description": f"Seconds to wait (1–{_MAX_SLEEP_S}).", + }, + }, + "required": ["seconds"], + "additionalProperties": False, + }, + handler=_handler, + tags=frozenset({"public", "utility"}), + risk_level="low", + read_only=True, + timeout_s=float(_MAX_SLEEP_S + 5), + ) + + +__all__ = ["sleep_tool"] diff --git a/skills/_workspace/public/background-jobs/SKILL.md b/skills/_workspace/public/background-jobs/SKILL.md new file mode 100644 index 00000000..192bcc4f --- /dev/null +++ b/skills/_workspace/public/background-jobs/SKILL.md @@ -0,0 +1,74 @@ +--- +name: background-jobs +description: "长时间任务(数分钟到数小时)用 start_job/get_job;同轮不要干等;断线后任务可继续,靠 job_id 恢复或 notify 通知。" +--- + +# 后台长任务 — Background Jobs + +## 核心事实 + +1. **同轮等不了两小时**(也不该等)。`sleep` 最多 120 秒。 +2. `start_job` 立刻返回 `job_id`;**agent 断线 / 本轮结束,进程仍可继续跑**(gateway 进程在的前提下)。 +3. 正确姿势是 **交作业 → 告知用户 job_id → 结束本轮**;以后再查,或靠完成通知。 + +## 推荐流程(1–2 小时任务) + +```text +1. write_file 写自包含脚本(输入/输出路径写死) +2. start_job(command=..., timeout_s=7200, name=..., notify={...可选}) +3. 回复用户:已提交后台任务,job_id=...,跑完会通知 / 请稍后问「查任务 job_xxx」 +4. 结束本轮(不要 sleep 循环两小时) +5. 之后任意一轮:get_job(job_id) 或 list_jobs → 读产出 → 需要时 save_deliverable_attachment +``` + +## 断线后会怎样 + +| 情况 | 行为 | +|------|------| +| 本轮对话结束 / 用户离开 | 任务继续;磁盘有 `data/jobs//` | +| 用户过一会再问 | `get_job` / `list_jobs` 恢复状态 | +| Gateway 重启但子进程还在 | 后台 reaper(约 15s)按超时/pid 对账回收;Chat 徽章也会刷新 | +| Gateway 整机挂掉 | 子进程可能被系统清掉;下次 reaper/`get_job` 标失败/超时 | + +## notify(可选,推荐渠道场景) + +`start_job` 可带: + +```json +"notify": { + "channel": "whatsapp", + "chat_id": "<群或会话 id>", + "account_id": "wa-default", + "message": "可选自定义文案" +} +``` + +任务结束(成功/失败/超时/取消)后会往该渠道塞一条出站提醒(含 `job_id`)。 +微信需额外 `context_token`(与现有 weixin 出站一致)。 + +**WebChat**:没有渠道 notify 时,靠用户稍后再问 + `list_jobs`。 + +## 工具 + +| 工具 | 作用 | +|------|------| +| `start_job` | 后台启动,返回 `job_id`(默认 2h,最大 3h) | +| `get_job` | 查状态与日志尾;`done=true` 表示结束 | +| `list_jobs` | 近期任务 | +| `cancel_job` | 取消 | +| `sleep` | 仅短轮询间隙(≤120s),不是长等 | + +与 `run_command` 同一开关:`AIA_ENABLE_RUN_COMMAND`。 + +## Chat 可视化 + +Chat 左侧导航 **「后台任务」** 右侧数字徽章持续显示运行中数量(约每 4 秒刷新,无需点开面板);用户菜单也有入口: +- 有任务时按钮高亮,右侧数字徽章变红 +- 点开后可:列表 / **终止** / **全部终止** / **清除记录** / 日志 +- 面板打开时约每 4 秒自动刷新列表 + +## 不要做 + +- 同轮 `sleep`/`get_job` 死循环拖两小时 +- 用 `run_command` 跑多小时脚本 +- 不告诉用户 `job_id` 就结束(断线后难找回,除非用户记得或去 Chat「后台任务」/ `list_jobs`) diff --git a/svc/jobs/__init__.py b/svc/jobs/__init__.py new file mode 100644 index 00000000..4eb71294 --- /dev/null +++ b/svc/jobs/__init__.py @@ -0,0 +1 @@ +"""Background job package.""" diff --git a/svc/jobs/background_jobs.py b/svc/jobs/background_jobs.py new file mode 100644 index 00000000..f26caa40 --- /dev/null +++ b/svc/jobs/background_jobs.py @@ -0,0 +1,713 @@ +from __future__ import annotations + +import json +import os +import signal +import subprocess +import threading +import time +import uuid +from dataclasses import asdict, dataclass +from pathlib import Path +from typing import Any, Optional + +from svc.config.paths import PROJECT_ROOT + +DEFAULT_TIMEOUT_S = 7200 # 2 hours +MAX_TIMEOUT_S = 10800 # 3 hours +MAX_CONCURRENT_RUNNING = 4 +_REAPER_INTERVAL_S = 15 +_META_NAME = "meta.json" +_STDOUT_NAME = "stdout.log" +_STDERR_NAME = "stderr.log" +_STATUS_RUNNING = "running" +_STATUS_SUCCEEDED = "succeeded" +_STATUS_FAILED = "failed" +_STATUS_TIMEOUT = "timeout" +_STATUS_CANCELLED = "cancelled" +_TERMINAL = {_STATUS_SUCCEEDED, _STATUS_FAILED, _STATUS_TIMEOUT, _STATUS_CANCELLED} + + +def jobs_dir() -> Path: + override = str(os.getenv("AIA_JOBS_DIR") or os.getenv("OPS_JOBS_DIR") or "").strip() + if override: + p = Path(override).expanduser().resolve() + else: + # Prefer data/jobs under project data root. + data = (Path(PROJECT_ROOT) / "data").resolve() + nested = (Path(PROJECT_ROOT) / "oclaw" / "data").resolve() + root = nested if nested.exists() else data + p = (root / "jobs").resolve() + p.mkdir(parents=True, exist_ok=True) + return p + + +def _utc_ts() -> int: + return int(time.time()) + + +def _clamp_timeout(raw: Any) -> int: + try: + val = int(raw) + except Exception: + val = DEFAULT_TIMEOUT_S + return max(1, min(val, MAX_TIMEOUT_S)) + + +def is_shell_exec_enabled() -> tuple[bool, str]: + """Same gate as run_command (Admin tool policy / AIA_ENABLE_RUN_COMMAND).""" + + def _truthy(v: str | None) -> bool: + return str(v or "").strip().lower() in {"1", "true", "yes", "on"} + + try: + enabled: bool | None = None + dbp = str(os.getenv("OPS_ASSISTANT_DB_PATH") or "").strip() + try: + if not dbp: + from svc.config.paths import db_path + + dbp = str(db_path() or "").strip() + except Exception: + dbp = "" + if dbp: + try: + from svc.persistence.sqlite_store import SqliteStore + + raw_db = str(SqliteStore(dbp).get_setting("AIA_ENABLE_RUN_COMMAND") or "").strip() + if raw_db: + enabled = _truthy(raw_db) + except Exception: + enabled = None + if enabled is None: + raw_env = str(os.getenv("AIA_ENABLE_RUN_COMMAND") or "").strip() + enabled = _truthy(raw_env) if raw_env else False + if enabled: + return True, "" + return ( + False, + "run_command/start_job is off: enable AIA_ENABLE_RUN_COMMAND in Admin → Tool policy " + "or set env AIA_ENABLE_RUN_COMMAND=1.", + ) + except Exception: + return False, "shell_exec_gate_failed" + + +@dataclass +class JobMeta: + job_id: str + command: str + cwd: str + status: str + created_at: int + timeout_s: int + name: str = "" + pid: int | None = None + started_at: int | None = None + finished_at: int | None = None + exit_code: int | None = None + error: str | None = None + # Optional channel ping when job ends (agent turn may already be gone). + notify: dict[str, Any] | None = None + notify_result: dict[str, Any] | None = None + + def to_dict(self) -> dict[str, Any]: + return asdict(self) + + @staticmethod + def from_dict(d: dict[str, Any]) -> "JobMeta": + notify = d.get("notify") if isinstance(d.get("notify"), dict) else None + notify_result = d.get("notify_result") if isinstance(d.get("notify_result"), dict) else None + return JobMeta( + job_id=str(d.get("job_id") or ""), + command=str(d.get("command") or ""), + cwd=str(d.get("cwd") or ""), + status=str(d.get("status") or ""), + created_at=int(d.get("created_at") or 0), + timeout_s=int(d.get("timeout_s") or DEFAULT_TIMEOUT_S), + name=str(d.get("name") or ""), + pid=int(d["pid"]) if d.get("pid") is not None else None, + started_at=int(d["started_at"]) if d.get("started_at") is not None else None, + finished_at=int(d["finished_at"]) if d.get("finished_at") is not None else None, + exit_code=int(d["exit_code"]) if d.get("exit_code") is not None else None, + error=str(d["error"]) if d.get("error") is not None else None, + notify=notify, + notify_result=notify_result, + ) + + +class BackgroundJobStore: + """Disk-backed background shell jobs for long-running agent scripts.""" + + def __init__(self, root_dir: str | Path | None = None, *, enable_reaper: bool = False): + self.root = Path(root_dir) if root_dir is not None else jobs_dir() + self.root.mkdir(parents=True, exist_ok=True) + self._lock = threading.RLock() + self._watchers: dict[str, threading.Thread] = {} + self._procs: dict[str, subprocess.Popen[Any]] = {} + self._log_handles: dict[str, tuple[Any, Any]] = {} + self._reaper_stop = threading.Event() + self._reaper_thread: threading.Thread | None = None + if enable_reaper: + self._ensure_reaper() + + def _ensure_reaper(self) -> None: + if self._reaper_thread is not None and self._reaper_thread.is_alive(): + return + self._reaper_stop.clear() + t = threading.Thread(target=self._reaper_loop, name="oclaw-job-reaper", daemon=True) + self._reaper_thread = t + t.start() + + def _reaper_loop(self) -> None: + # First pass soon after boot so dirty metas from a prior crash are cleared. + try: + self.reap() + except Exception: + pass + while not self._reaper_stop.wait(_REAPER_INTERVAL_S): + try: + self.reap() + except Exception: + pass + + def reap(self) -> dict[str, Any]: + """Reconcile all running jobs (timeout / dead pid). Safe to call anytime.""" + scanned = 0 + changed = 0 + if not self.root.exists(): + return {"ok": True, "scanned": 0, "changed": 0} + for p in list(self.root.iterdir()): + if not p.is_dir(): + continue + meta = self._read_meta(p.name) + if meta is None or meta.status != _STATUS_RUNNING: + continue + scanned += 1 + before = meta.status + after = self._reconcile(meta) + if after.status != before: + changed += 1 + return {"ok": True, "scanned": scanned, "changed": changed} + + def _job_dir(self, job_id: str) -> Path: + return self.root / job_id + + def _meta_path(self, job_id: str) -> Path: + return self._job_dir(job_id) / _META_NAME + + def _read_meta(self, job_id: str) -> Optional[JobMeta]: + mp = self._meta_path(job_id) + if not mp.exists(): + return None + try: + obj = json.loads(mp.read_text(encoding="utf-8")) + if not isinstance(obj, dict): + return None + return JobMeta.from_dict(obj) + except Exception: + return None + + def _write_meta(self, meta: JobMeta) -> None: + d = self._job_dir(meta.job_id) + d.mkdir(parents=True, exist_ok=True) + tmp = self._meta_path(meta.job_id).with_suffix(".tmp") + tmp.write_text(json.dumps(meta.to_dict(), ensure_ascii=False, indent=2), encoding="utf-8") + os.replace(tmp, self._meta_path(meta.job_id)) + + def _count_running(self) -> int: + """Count live running jobs after reconciling stale metas.""" + n = 0 + if not self.root.exists(): + return 0 + for p in list(self.root.iterdir()): + if not p.is_dir(): + continue + meta = self._read_meta(p.name) + if meta is None: + continue + if meta.status == _STATUS_RUNNING: + meta = self._reconcile(meta) + if meta.status == _STATUS_RUNNING: + n += 1 + return n + + def _pid_alive(self, pid: int | None) -> bool: + if not pid or pid <= 0: + return False + try: + if os.name == "nt": + out = subprocess.run( + ["tasklist", "/FI", f"PID eq {pid}"], + capture_output=True, + text=True, + timeout=5, + creationflags=getattr(subprocess, "CREATE_NO_WINDOW", 0), + ) + return str(pid) in (out.stdout or "") + os.kill(pid, 0) + return True + except Exception: + return False + + def _kill_tree(self, pid: int | None) -> None: + if not pid or pid <= 0: + return + try: + if os.name == "nt": + subprocess.run( + ["taskkill", "/PID", str(pid), "/T", "/F"], + capture_output=True, + text=True, + timeout=15, + creationflags=getattr(subprocess, "CREATE_NO_WINDOW", 0), + ) + else: + try: + os.killpg(pid, signal.SIGKILL) + except Exception: + os.kill(pid, signal.SIGKILL) + except Exception: + pass + + def _close_logs(self, job_id: str) -> None: + handles = self._log_handles.pop(job_id, None) + if not handles: + return + for h in handles: + try: + h.close() + except Exception: + pass + + def _finalize( + self, + job_id: str, + *, + status: str, + exit_code: int | None, + error: str | None = None, + force: bool = False, + ) -> JobMeta | None: + with self._lock: + meta = self._read_meta(job_id) + if meta is None: + self._close_logs(job_id) + self._procs.pop(job_id, None) + return None + if meta.status in _TERMINAL: + # Cancel may race the watcher (kill → non-zero exit → failed). + # Allow an intentional cancel to override a just-recorded failed/timeout. + if not ( + force + and status == _STATUS_CANCELLED + and meta.status in {_STATUS_FAILED, _STATUS_TIMEOUT} + ): + self._close_logs(job_id) + self._procs.pop(job_id, None) + return meta + meta.status = status + meta.exit_code = exit_code + meta.finished_at = _utc_ts() + if error: + meta.error = error + elif force and status == _STATUS_CANCELLED: + meta.error = "cancelled_by_user" + self._write_meta(meta) + self._close_logs(job_id) + self._procs.pop(job_id, None) + self._watchers.pop(job_id, None) + # Notify outside the lock (may hit SQLite / network). + self._notify_complete(job_id) + return self._read_meta(job_id) + + def _notify_complete(self, job_id: str) -> None: + meta = self._read_meta(job_id) + if meta is None or not meta.notify or meta.notify_result: + return + notify = dict(meta.notify) + channel = str(notify.get("channel") or "").strip().lower() + chat_id = str(notify.get("chat_id") or "").strip() + if not channel or not chat_id: + meta.notify_result = {"ok": False, "error": "notify_channel_or_chat_id_missing"} + self._write_meta(meta) + return + if channel in {"wechat", "weixin"}: + channel = "weixin" + custom = str(notify.get("message") or "").strip() + text = custom or ( + f"后台任务已结束\n" + f"job_id={meta.job_id}\n" + f"name={meta.name}\n" + f"status={meta.status}\n" + f"exit_code={meta.exit_code}\n" + f"可在对话中发送:查任务 {meta.job_id}" + ) + account_id = str(notify.get("account_id") or "").strip() + if not account_id and channel == "whatsapp": + account_id = str(os.getenv("AIA_WHATSAPP_ACCOUNT_ID") or "wa-default").strip() + tenant_id = str(notify.get("tenant_id") or "").strip() + try: + from svc.config.paths import db_path + from svc.persistence.sqlite_store import SqliteStore + + store = SqliteStore(db_path()) + source_payload = { + "kind": "background_job_complete", + "job_id": meta.job_id, + "status": meta.status, + } + ctx = str(notify.get("context_token") or "").strip() + if ctx: + source_payload["context_token"] = ctx + msg_id = store.enqueue_channel_outbound_message( + channel=channel, + chat_id=chat_id, + text=text, + tenant_id=tenant_id, + account_id=account_id, + source=json.dumps(source_payload, ensure_ascii=False), + ) + meta.notify_result = {"ok": True, "message_id": msg_id, "channel": channel, "chat_id": chat_id} + except Exception as exc: + meta.notify_result = {"ok": False, "error": f"{type(exc).__name__}: {exc}"} + self._write_meta(meta) + + def _watch(self, job_id: str, proc: subprocess.Popen[Any], timeout_s: int) -> None: + try: + try: + code = proc.wait(timeout=timeout_s) + except subprocess.TimeoutExpired: + self._kill_tree(proc.pid) + try: + proc.wait(timeout=10) + except Exception: + pass + self._finalize(job_id, status=_STATUS_TIMEOUT, exit_code=None, error="job_timeout") + return + status = _STATUS_SUCCEEDED if int(code) == 0 else _STATUS_FAILED + self._finalize(job_id, status=status, exit_code=int(code)) + except Exception as exc: + self._finalize(job_id, status=_STATUS_FAILED, exit_code=None, error=f"{type(exc).__name__}: {exc}") + + def _reconcile(self, meta: JobMeta) -> JobMeta: + if meta.status != _STATUS_RUNNING: + return meta + started = int(meta.started_at or meta.created_at or 0) + if started and (_utc_ts() - started) > int(meta.timeout_s or DEFAULT_TIMEOUT_S): + self._kill_tree(meta.pid) + proc = self._procs.get(meta.job_id) + if proc is not None: + try: + proc.wait(timeout=5) + except Exception: + pass + self._finalize(meta.job_id, status=_STATUS_TIMEOUT, exit_code=None, error="job_timeout_reconcile") + return self._read_meta(meta.job_id) or meta + + proc = self._procs.get(meta.job_id) + if proc is not None: + rc = proc.poll() + if rc is not None: + status = _STATUS_SUCCEEDED if int(rc) == 0 else _STATUS_FAILED + self._finalize(meta.job_id, status=status, exit_code=int(rc)) + return self._read_meta(meta.job_id) or meta + return meta + # Process not tracked in this process (gateway/agent restart): infer from pid. + if meta.pid and self._pid_alive(meta.pid): + # Still running without local watcher — timeout enforced above on later polls. + return meta + # Pid gone — mark failed/unknown unless already terminal on disk. + self._finalize( + meta.job_id, + status=_STATUS_FAILED, + exit_code=None, + error="process_exited_untracked", + ) + return self._read_meta(meta.job_id) or meta + + def start( + self, + *, + command: str, + cwd: str, + timeout_s: int | None = None, + name: str = "", + notify: dict[str, Any] | None = None, + ) -> dict[str, Any]: + cmd = str(command or "").strip() + if not cmd: + return {"ok": False, "error": "command_required"} + notify_clean: dict[str, Any] | None = None + if isinstance(notify, dict) and notify: + channel = str(notify.get("channel") or "").strip().lower() + chat_id = str(notify.get("chat_id") or "").strip() + if channel and chat_id: + notify_clean = { + "channel": channel, + "chat_id": chat_id, + "account_id": str(notify.get("account_id") or "").strip(), + "tenant_id": str(notify.get("tenant_id") or "").strip(), + "context_token": str(notify.get("context_token") or "").strip(), + "message": str(notify.get("message") or "").strip(), + } + with self._lock: + if self._count_running() >= MAX_CONCURRENT_RUNNING: + return { + "ok": False, + "error": "too_many_running_jobs", + "max_concurrent": MAX_CONCURRENT_RUNNING, + } + job_id = "job_" + uuid.uuid4().hex[:12] + timeout = _clamp_timeout(timeout_s if timeout_s is not None else DEFAULT_TIMEOUT_S) + workdir = str(Path(cwd).resolve()) + Path(workdir).mkdir(parents=True, exist_ok=True) + jdir = self._job_dir(job_id) + jdir.mkdir(parents=True, exist_ok=True) + stdout_path = jdir / _STDOUT_NAME + stderr_path = jdir / _STDERR_NAME + stdout_f = open(stdout_path, "w", encoding="utf-8", errors="replace") + stderr_f = open(stderr_path, "w", encoding="utf-8", errors="replace") + now = _utc_ts() + meta = JobMeta( + job_id=job_id, + command=cmd, + cwd=workdir, + status=_STATUS_RUNNING, + created_at=now, + started_at=now, + timeout_s=timeout, + name=str(name or "").strip() or job_id, + notify=notify_clean, + ) + run_kwargs: dict[str, Any] = { + "cwd": workdir, + "shell": True, + "stdout": stdout_f, + "stderr": stderr_f, + "text": True, + } + if os.name == "nt": + startupinfo = subprocess.STARTUPINFO() + startupinfo.dwFlags |= subprocess.STARTF_USESHOWWINDOW + startupinfo.wShowWindow = 0 + run_kwargs["startupinfo"] = startupinfo + run_kwargs["creationflags"] = subprocess.CREATE_NO_WINDOW + else: + run_kwargs["start_new_session"] = True + try: + proc = subprocess.Popen(cmd, **run_kwargs) + except Exception as exc: + stdout_f.close() + stderr_f.close() + meta.status = _STATUS_FAILED + meta.finished_at = _utc_ts() + meta.error = f"{type(exc).__name__}: {exc}" + self._write_meta(meta) + return {"ok": False, "error": "spawn_failed", "detail": str(exc), "job": meta.to_dict()} + meta.pid = int(proc.pid) if proc.pid else None + self._write_meta(meta) + self._procs[job_id] = proc + self._log_handles[job_id] = (stdout_f, stderr_f) + t = threading.Thread(target=self._watch, args=(job_id, proc, timeout), daemon=True, name=f"job-{job_id}") + self._watchers[job_id] = t + t.start() + return { + "ok": True, + "job_id": job_id, + "status": meta.status, + "pid": meta.pid, + "timeout_s": timeout, + "cwd": workdir, + "name": meta.name, + "log_dir": str(jdir), + "notify": bool(notify_clean), + "hint": ( + "Job started in background and will keep running if this agent turn ends. " + "Tell the user the job_id, then end the turn (do NOT sleep for hours). " + "Later: get_job/list_jobs to resume. Optional notify pings the channel on completion." + ), + } + + def get(self, job_id: str, *, log_tail_chars: int = 4000) -> dict[str, Any]: + aid = str(job_id or "").strip() + if not aid: + return {"ok": False, "error": "job_id_required"} + meta = self._read_meta(aid) + if meta is None: + return {"ok": False, "error": "job_not_found", "job_id": aid} + meta = self._reconcile(meta) + tail = max(200, min(int(log_tail_chars or 4000), 50000)) + jdir = self._job_dir(aid) + + def _tail(path: Path) -> str: + if not path.exists(): + return "" + try: + data = path.read_text(encoding="utf-8", errors="replace") + except Exception: + return "" + if len(data) <= tail: + return data + return data[-tail:] + + out = { + "ok": True, + "job_id": meta.job_id, + "name": meta.name, + "status": meta.status, + "command": meta.command, + "cwd": meta.cwd, + "pid": meta.pid, + "timeout_s": meta.timeout_s, + "created_at": meta.created_at, + "started_at": meta.started_at, + "finished_at": meta.finished_at, + "exit_code": meta.exit_code, + "error": meta.error, + "log_dir": str(jdir), + "stdout_tail": _tail(jdir / _STDOUT_NAME), + "stderr_tail": _tail(jdir / _STDERR_NAME), + "done": meta.status in _TERMINAL, + "notify_result": meta.notify_result, + } + if not out["done"]: + out["hint"] = ( + "still running; prefer ending this turn and checking later with get_job. " + "Short sleep+poll only for near-term completion (minutes), not multi-hour waits." + ) + return out + + def cancel(self, job_id: str) -> dict[str, Any]: + aid = str(job_id or "").strip() + if not aid: + return {"ok": False, "error": "job_id_required"} + with self._lock: + meta = self._read_meta(aid) + if meta is None: + return {"ok": False, "error": "job_not_found", "job_id": aid} + if meta.status in _TERMINAL: + return {"ok": True, "job_id": aid, "status": meta.status, "already_done": True} + proc = self._procs.get(aid) + pid = meta.pid or (proc.pid if proc else None) + self._kill_tree(pid) + if proc is not None: + try: + proc.wait(timeout=10) + except Exception: + pass + # Finalize + notify outside the lock so other start/get/list calls are not blocked. + self._finalize( + aid, + status=_STATUS_CANCELLED, + exit_code=None, + error="cancelled_by_user", + force=True, + ) + latest = self._read_meta(aid) + return {"ok": True, "job_id": aid, "status": (latest.status if latest else _STATUS_CANCELLED)} + + def cancel_all_running(self) -> dict[str, Any]: + listed = self.list_jobs(limit=100) + cancelled: list[str] = [] + for row in listed.get("jobs") or []: + if str(row.get("status") or "") != _STATUS_RUNNING: + continue + jid = str(row.get("job_id") or "") + if not jid: + continue + out = self.cancel(jid) + if out.get("ok"): + cancelled.append(jid) + return {"ok": True, "cancelled": cancelled, "count": len(cancelled)} + + def purge(self, job_id: str) -> dict[str, Any]: + """Remove finished job files from disk. Running jobs must be cancelled first.""" + aid = str(job_id or "").strip() + if not aid: + return {"ok": False, "error": "job_id_required"} + with self._lock: + meta = self._read_meta(aid) + if meta is None: + return {"ok": False, "error": "job_not_found", "job_id": aid} + if meta.status == _STATUS_RUNNING: + return {"ok": False, "error": "job_still_running", "hint": "cancel first"} + self._close_logs(aid) + self._procs.pop(aid, None) + self._watchers.pop(aid, None) + jdir = self._job_dir(aid) + try: + import shutil + + shutil.rmtree(jdir, ignore_errors=True) + except Exception as exc: + return {"ok": False, "error": "purge_failed", "detail": str(exc)} + return {"ok": True, "job_id": aid, "purged": True} + + def list_jobs(self, *, limit: int = 20) -> dict[str, Any]: + lim = max(1, min(int(limit or 20), 100)) + items: list[JobMeta] = [] + for p in self.root.iterdir(): + if not p.is_dir(): + continue + meta = self._read_meta(p.name) + if meta is None: + continue + if meta.status == _STATUS_RUNNING: + meta = self._reconcile(meta) + items.append(meta) + items.sort(key=lambda m: int(m.created_at or 0), reverse=True) + jobs = [] + for m in items[:lim]: + cmd = str(m.command or "") + jobs.append( + { + "job_id": m.job_id, + "name": m.name, + "status": m.status, + "pid": m.pid, + "command": cmd if len(cmd) <= 240 else cmd[:237] + "...", + "cwd": m.cwd, + "created_at": m.created_at, + "started_at": m.started_at, + "finished_at": m.finished_at, + "exit_code": m.exit_code, + "timeout_s": m.timeout_s, + "error": m.error, + } + ) + # Count all running, not only truncated page. + running_all = sum(1 for m in items if m.status == _STATUS_RUNNING) + return { + "ok": True, + "jobs": jobs, + "running_count": running_all, + "total": len(items), + } + + +_STORE: BackgroundJobStore | None = None +_STORE_LOCK = threading.Lock() + + +def get_job_store(root_dir: str | Path | None = None) -> BackgroundJobStore: + global _STORE + if root_dir is not None: + return BackgroundJobStore(root_dir=root_dir, enable_reaper=False) + with _STORE_LOCK: + if _STORE is None: + _STORE = BackgroundJobStore(enable_reaper=True) + else: + _STORE._ensure_reaper() + return _STORE + + +__all__ = [ + "BackgroundJobStore", + "DEFAULT_TIMEOUT_S", + "MAX_TIMEOUT_S", + "JobMeta", + "get_job_store", + "is_shell_exec_enabled", + "jobs_dir", +] diff --git a/tests/test_background_jobs.py b/tests/test_background_jobs.py new file mode 100644 index 00000000..9d733086 --- /dev/null +++ b/tests/test_background_jobs.py @@ -0,0 +1,245 @@ +from __future__ import annotations + +import json +import sys +import time + +from runtime.tools.public.job_tools import cancel_job_tool, get_job_tool, list_jobs_tool, start_job_tool +from runtime.tools.public.sleep_tool import sleep_tool +from svc.jobs.background_jobs import BackgroundJobStore, get_job_store + + +def test_sleep_bounds() -> None: + spec = sleep_tool() + assert spec.handler({"seconds": 0.2}).get("ok") is True + bad = spec.handler({"seconds": 999}) + assert bad.get("ok") is False + assert bad.get("error") == "seconds_too_large" + + +def test_start_and_get_job_succeeds(tmp_path, monkeypatch) -> None: + store = BackgroundJobStore(root_dir=tmp_path / "jobs") + monkeypatch.setattr("runtime.tools.public.job_tools.get_job_store", lambda: store) + monkeypatch.setattr("runtime.tools.public.job_tools.is_shell_exec_enabled", lambda: (True, "")) + monkeypatch.setattr( + "runtime.tools.public.job_tools.resolve_workspace_path", + lambda raw: tmp_path / "ws", + ) + (tmp_path / "ws").mkdir(parents=True, exist_ok=True) + + py = sys.executable + started = start_job_tool().handler( + { + "command": f'"{py}" -c "print(123); import time; time.sleep(0.3)"', + "cwd": str(tmp_path / "ws"), + "timeout_s": 30, + "name": "quick", + } + ) + assert started.get("ok") is True + job_id = started["job_id"] + + deadline = time.time() + 10 + last = {} + while time.time() < deadline: + last = get_job_tool().handler({"job_id": job_id}) + if last.get("done"): + break + time.sleep(0.1) + assert last.get("ok") is True + assert last.get("status") == "succeeded" + assert last.get("exit_code") == 0 + assert "123" in str(last.get("stdout_tail") or "") + + +def test_get_job_timeout(tmp_path, monkeypatch) -> None: + store = BackgroundJobStore(root_dir=tmp_path / "jobs") + monkeypatch.setattr("runtime.tools.public.job_tools.get_job_store", lambda: store) + monkeypatch.setattr("runtime.tools.public.job_tools.is_shell_exec_enabled", lambda: (True, "")) + monkeypatch.setattr( + "runtime.tools.public.job_tools.resolve_workspace_path", + lambda raw: tmp_path / "ws", + ) + (tmp_path / "ws").mkdir(parents=True, exist_ok=True) + + py = sys.executable + started = start_job_tool().handler( + { + "command": f'"{py}" -c "import time; time.sleep(30)"', + "timeout_s": 1, + } + ) + assert started.get("ok") is True + job_id = started["job_id"] + deadline = time.time() + 15 + last = {} + while time.time() < deadline: + last = get_job_tool().handler({"job_id": job_id}) + if last.get("done"): + break + time.sleep(0.2) + assert last.get("status") == "timeout" + assert last.get("done") is True + + +def test_cancel_job(tmp_path, monkeypatch) -> None: + store = BackgroundJobStore(root_dir=tmp_path / "jobs") + monkeypatch.setattr("runtime.tools.public.job_tools.get_job_store", lambda: store) + monkeypatch.setattr("runtime.tools.public.job_tools.is_shell_exec_enabled", lambda: (True, "")) + monkeypatch.setattr( + "runtime.tools.public.job_tools.resolve_workspace_path", + lambda raw: tmp_path / "ws", + ) + (tmp_path / "ws").mkdir(parents=True, exist_ok=True) + + py = sys.executable + started = start_job_tool().handler( + {"command": f'"{py}" -c "import time; time.sleep(60)"', "timeout_s": 120} + ) + assert started.get("ok") is True + job_id = started["job_id"] + cancelled = cancel_job_tool().handler({"job_id": job_id}) + assert cancelled.get("ok") is True + got = get_job_tool().handler({"job_id": job_id}) + assert got.get("status") == "cancelled" + assert got.get("done") is True + + +def test_start_job_requires_enable_gate(monkeypatch, tmp_path) -> None: + monkeypatch.setattr("runtime.tools.public.job_tools.is_shell_exec_enabled", lambda: (False, "off")) + monkeypatch.setattr( + "runtime.tools.public.job_tools.resolve_workspace_path", + lambda raw: tmp_path, + ) + out = start_job_tool().handler({"command": "echo hi"}) + assert out.get("ok") is False + assert out.get("error") == "disabled" + + +def test_list_jobs(tmp_path, monkeypatch) -> None: + store = BackgroundJobStore(root_dir=tmp_path / "jobs") + monkeypatch.setattr("runtime.tools.public.job_tools.get_job_store", lambda: store) + monkeypatch.setattr("runtime.tools.public.job_tools.is_shell_exec_enabled", lambda: (True, "")) + monkeypatch.setattr( + "runtime.tools.public.job_tools.resolve_workspace_path", + lambda raw: tmp_path / "ws", + ) + (tmp_path / "ws").mkdir(parents=True, exist_ok=True) + py = sys.executable + start_job_tool().handler({"command": f'"{py}" -c "print(1)"', "name": "a"}) + listed = list_jobs_tool().handler({"limit": 5}) + assert listed.get("ok") is True + assert len(listed.get("jobs") or []) >= 1 + + +def test_reconcile_timeout_for_orphan_meta(tmp_path) -> None: + store = BackgroundJobStore(root_dir=tmp_path / "jobs") + job_id = "job_orphan1" + jdir = tmp_path / "jobs" / job_id + jdir.mkdir(parents=True) + meta = { + "job_id": job_id, + "command": "sleep 999", + "cwd": str(tmp_path), + "status": "running", + "created_at": int(time.time()) - 100, + "started_at": int(time.time()) - 100, + "timeout_s": 1, + "name": "orphan", + "pid": 99999999, + } + (jdir / "meta.json").write_text(json.dumps(meta), encoding="utf-8") + (jdir / "stdout.log").write_text("", encoding="utf-8") + (jdir / "stderr.log").write_text("", encoding="utf-8") + got = store.get(job_id) + assert got.get("done") is True + assert got.get("status") == "timeout" + + +def test_count_running_reconciles_stale_meta(tmp_path) -> None: + store = BackgroundJobStore(root_dir=tmp_path / "jobs") + # Four stale running metas with dead pids would previously block new starts. + for i in range(4): + job_id = f"job_stale{i}" + jdir = tmp_path / "jobs" / job_id + jdir.mkdir(parents=True) + meta = { + "job_id": job_id, + "command": "sleep 999", + "cwd": str(tmp_path), + "status": "running", + "created_at": int(time.time()) - 10, + "started_at": int(time.time()) - 10, + "timeout_s": 7200, + "name": job_id, + "pid": 90000000 + i, + } + (jdir / "meta.json").write_text(json.dumps(meta), encoding="utf-8") + (jdir / "stdout.log").write_text("", encoding="utf-8") + (jdir / "stderr.log").write_text("", encoding="utf-8") + assert store._count_running() == 0 + py = sys.executable + started = store.start(command=f'"{py}" -c "print(1)"', cwd=str(tmp_path), timeout_s=30) + assert started.get("ok") is True + + +def test_reap_timeouts_without_get(tmp_path) -> None: + store = BackgroundJobStore(root_dir=tmp_path / "jobs") + job_id = "job_reap1" + jdir = tmp_path / "jobs" / job_id + jdir.mkdir(parents=True) + meta = { + "job_id": job_id, + "command": "sleep 999", + "cwd": str(tmp_path), + "status": "running", + "created_at": int(time.time()) - 100, + "started_at": int(time.time()) - 100, + "timeout_s": 1, + "name": "reap", + "pid": 99999998, + } + (jdir / "meta.json").write_text(json.dumps(meta), encoding="utf-8") + (jdir / "stdout.log").write_text("", encoding="utf-8") + (jdir / "stderr.log").write_text("", encoding="utf-8") + out = store.reap() + assert out.get("ok") is True + assert out.get("changed", 0) >= 1 + got = store._read_meta(job_id) + assert got is not None + assert got.status == "timeout" + + +def test_notify_on_complete(tmp_path, monkeypatch) -> None: + store = BackgroundJobStore(root_dir=tmp_path / "jobs") + calls: list[dict] = [] + + class _Store: + def enqueue_channel_outbound_message(self, **kwargs): + calls.append(kwargs) + return "msg-1" + + monkeypatch.setattr("svc.config.paths.db_path", lambda: str(tmp_path / "x.sqlite")) + monkeypatch.setattr("svc.persistence.sqlite_store.SqliteStore", lambda *_a, **_k: _Store()) + + py = sys.executable + started = store.start( + command=f'"{py}" -c "print(\'ok\')"', + cwd=str(tmp_path), + timeout_s=30, + name="n1", + notify={"channel": "whatsapp", "chat_id": "628@g.us", "message": "done-test"}, + ) + assert started.get("ok") is True + job_id = started["job_id"] + deadline = time.time() + 10 + last = {} + while time.time() < deadline: + last = store.get(job_id) + if last.get("done") and last.get("notify_result"): + break + time.sleep(0.1) + assert last.get("status") == "succeeded" + assert calls and calls[0]["chat_id"] == "628@g.us" + assert "done-test" in calls[0]["text"] + assert (last.get("notify_result") or {}).get("ok") is True diff --git a/tests/test_chat_jobs_api.py b/tests/test_chat_jobs_api.py new file mode 100644 index 00000000..90682d9f --- /dev/null +++ b/tests/test_chat_jobs_api.py @@ -0,0 +1,57 @@ +from __future__ import annotations + +from fastapi import FastAPI +from fastapi.testclient import TestClient + +from interfaces.admin.chat_api import include_chat_routes +from svc.jobs.background_jobs import BackgroundJobStore + + +def test_chat_jobs_api_list_cancel_purge(tmp_path, monkeypatch) -> None: + store = BackgroundJobStore(root_dir=tmp_path / "jobs") + monkeypatch.setattr("svc.jobs.background_jobs.get_job_store", lambda root_dir=None: store) + # Also patch the import site used inside route handlers + import svc.jobs.background_jobs as bj + + monkeypatch.setattr(bj, "get_job_store", lambda root_dir=None: store) + + class _FakeSqlite: + pass + + def _resolve_auth(_store, _auth): + return {"tenant_id": "t1", "user_id": "u1", "username": "u1"} + + monkeypatch.setattr("interfaces.admin.chat_api.get_assistant_store", lambda: _FakeSqlite()) + + app = FastAPI() + from fastapi import APIRouter + + router = APIRouter() + include_chat_routes(router, resolve_auth=_resolve_auth) + app.include_router(router) + client = TestClient(app) + + started = store.start(command="echo hi", cwd=str(tmp_path), timeout_s=5, name="ui-job") + assert started.get("ok") is True + job_id = started["job_id"] + + r = client.get("/admin/api/chat/jobs", headers={"Authorization": "Bearer x"}) + assert r.status_code == 200 + body = r.json() + assert body.get("ok") is True + assert any(j.get("job_id") == job_id for j in body.get("jobs") or []) + + # wait finish then purge + import time + + deadline = time.time() + 8 + while time.time() < deadline: + g = store.get(job_id) + if g.get("done"): + break + time.sleep(0.05) + assert store.get(job_id).get("done") is True + + d = client.delete(f"/admin/api/chat/jobs/{job_id}", headers={"Authorization": "Bearer x"}) + assert d.status_code == 200 + assert d.json().get("purged") is True