// dsh-plugin-mcn-schedule 宿主端:计划任务(定时提醒/周期任务) // - 数据持久化:~/.dsh/mcn-schedule.json // - 注册 /mcn/schedule/* HTTP 路由供前端管理 // - 内置调度器:每 10 秒检查到期任务,触发事件供前端轮询提醒 // - 支持动作:remind(页面提醒)/ agent(唤醒 AI 会话执行任务) import { homedir } from "node:os"; import { join } from "node:path"; import { readFileSync, writeFileSync, existsSync, mkdirSync, readdirSync } from "node:fs"; import { createUserMessage } from "@deepseek-ai/dsh-llm"; const DATA_PATH = join(homedir(), ".dsh", "mcn-schedule.json"); const TICK_MS = 10 * 1000; // 调度检查周期 const MAX_EVENTS = 50; // 事件保留上限 const MAX_HISTORY = 20; // 每个任务触发历史保留上限 const MIN_INTERVAL_MINUTES = 1; // 周期任务最小间隔(分钟) const TASK_AGENT_PREFIX = "dsh_schedule-"; // 任务执行会话前缀 const TASK_WORKSPACE_TITLE = "计划任务"; // 任务会话分组工作区标题 const AGENT_PROVIDER = "deepseek-official"; const AGENT_MODEL = "deepseek-v4-flash"; let taskWorkspacePath = null; // 任务会话工作区(apply 时解析) let taskWorkspace = null; // 「计划任务」工作区对象(含 attachSession) //#region 工作区解析(与 dsh-plugin-mcn 一致:workspace.json 深度优先取根工作区) function resolveWorkspaceRoot() { try { const raw = readFileSync(join(homedir(), ".dsh", "storages", "workspace.json"), "utf8"); const data = JSON.parse(raw); const workspaces = data?.tables?.workspaces; if (workspaces) { const candidates = Object.values(workspaces) .filter((ws) => typeof ws?.path === "string" && existsSync(ws.path)) .sort((a, b) => String(a.path).split(/[\\/]/).length - String(b.path).split(/[\\/]/).length); if (candidates.length > 0) return candidates[0].path; } } catch (e) { /* fall through */ } return "D:\\dshworkspace"; } // 确保「计划任务」工作区存在(目录 + workspaceRegistry 注册),缓存工作区对象 async function ensureTaskWorkspace(ctx) { const root = resolveWorkspaceRoot(); const dir = join(root, TASK_WORKSPACE_TITLE); try { mkdirSync(dir, { recursive: true }); } catch (e) { /* ignore */ } try { const reg = ctx.get("workspaceRegistry"); if (reg) { const existing = reg.list().find((w) => w.path === dir || w.title === TASK_WORKSPACE_TITLE); if (existing) { taskWorkspace = existing; } else { taskWorkspace = await reg.create(dir, TASK_WORKSPACE_TITLE); if (taskWorkspace.title !== TASK_WORKSPACE_TITLE) { try { await taskWorkspace.setTitle(TASK_WORKSPACE_TITLE); } catch (e) { /* 更名失败不阻塞 */ } } console.log(`[dsh-plugin-mcn-schedule] 已创建工作区「${TASK_WORKSPACE_TITLE}」: ${dir}`); } } } catch (e) { console.warn(`[dsh-plugin-mcn-schedule] 工作区注册失败: ${e.message}`); } taskWorkspacePath = dir; return dir; } // 把会话 attach 到「计划任务」分组(内部按 header cwd 校验) async function attachTaskSession(ctx, sessionId) { try { const reg = ctx.get("workspaceRegistry"); const ws = taskWorkspace || (reg && reg.list().find((w) => w.path === taskWorkspacePath || w.title === TASK_WORKSPACE_TITLE)); if (!ws) return; if (!ws.sessionIds.includes(sessionId)) { await ws.attachSession(sessionId); console.log(`[dsh-plugin-mcn-schedule] 会话 ${sessionId} 已归入「计划任务」分组`); } } catch (e) { console.warn(`[dsh-plugin-mcn-schedule] 会话 ${sessionId} 归组失败: ${e.message}`); } } // 启动修复:把磁盘上已有的 dsh_schedule-* 会话重新 attach 到「计划任务」分组 async function repairTaskSessions(ctx) { try { const sessionsRoot = join(homedir(), ".dsh", "sessions"); if (!existsSync(sessionsRoot)) return; const dirs = readdirSync(sessionsRoot, { withFileTypes: true }) .filter((e) => e.isDirectory()) .flatMap((e) => readdirSync(join(sessionsRoot, e.name), { withFileTypes: true }) .filter((s) => s.isDirectory() && s.name.startsWith(TASK_AGENT_PREFIX)) .map((s) => s.name)); for (const sessionId of dirs) { await attachTaskSession(ctx, sessionId); } if (dirs.length > 0) console.log(`[dsh-plugin-mcn-schedule] 计划任务会话分组修复完成(${dirs.length} 个会话)`); } catch (e) { console.warn(`[dsh-plugin-mcn-schedule] 会话分组修复失败: ${e.message}`); } } // 任务会话工作区:优先「计划任务」分组,未就绪时回退根工作区 function resolveTaskWorkspace() { if (taskWorkspacePath) return taskWorkspacePath; return join(resolveWorkspaceRoot(), TASK_WORKSPACE_TITLE); } async function composeDefaultPreset(ctx) { const presets = ctx.get("agentPresets"); if (!presets) return { setup: () => Promise.resolve() }; let resolvedId; // 优先 MCN 模式(含 myai MCP 工具);缺失时回退默认 preset try { resolvedId = (await presets.resolve("mcn")).id; } catch (e) { try { resolvedId = (await presets.resolve(undefined)).id; } catch (e2) { return { setup: () => Promise.resolve() }; } } return { agentPreset: resolvedId, setup: async (agentCtx) => { await presets.mount(agentCtx, resolvedId); }, }; } // 获取(或创建)任务执行会话:默认 dsh_schedule-{taskId},可指定 targetSession async function getOrCreateTaskAgent(ctx, task) { const sessionId = task.targetSession || TASK_AGENT_PREFIX + task.id; const existing = ctx.agents.get(sessionId); if (existing !== undefined) { return { agent: existing, created: false }; } const agentOptions = { provider: AGENT_PROVIDER, model: AGENT_MODEL }; const compose = await composeDefaultPreset(ctx); const cwd = resolveTaskWorkspace(); let handle; try { handle = await ctx.agents.resume({ resumeSessionId: sessionId, agentOptions, setup: compose.setup, }); } catch (e) { handle = await ctx.agents.create({ sessionId, agentOptions, meta: { cwd, agentPreset: compose.agentPreset }, setup: compose.setup, }); } // 归入「计划任务」分组(cwd 与工作区 path 匹配时成功) await attachTaskSession(ctx, sessionId); return { agent: handle.agent }; } //#endregion //#region 存储 let store = { tasks: [], events: [] }; let tickCtx = null; // apply 时注入,供调度器唤醒 AI 会话使用 function loadStore() { try { if (existsSync(DATA_PATH)) { const raw = JSON.parse(readFileSync(DATA_PATH, "utf8")); store = { tasks: Array.isArray(raw.tasks) ? raw.tasks : [], events: Array.isArray(raw.events) ? raw.events : [], }; } } catch (e) { console.warn(`[dsh-plugin-mcn-schedule] 数据读取失败,使用空存储: ${e.message}`); store = { tasks: [], events: [] }; } } function saveStore() { try { writeFileSync(DATA_PATH, JSON.stringify(store, null, 2), "utf8"); } catch (e) { console.warn(`[dsh-plugin-mcn-schedule] 数据写入失败: ${e.message}`); } } function nextId() { return "t" + Date.now().toString(36) + Math.random().toString(36).slice(2, 6); } function computeNextRun(task) { const now = Date.now(); if (task.type === "at") return task.atTs; if (task.type === "every") { const base = task.lastFiredAt || task.createdAt; return base + task.intervalMinutes * 60 * 1000; } return null; } //#endregion //#region 调度器 async function fireTask(ctx, task, now) { const firedAt = new Date(now).toISOString(); let action = task.action || "remind"; let dispatched = null; let error = null; if (action === "agent") { try { const text = task.prompt && task.prompt.trim() ? task.prompt.trim() : `请执行计划任务「${task.name}」。`; const { agent } = await getOrCreateTaskAgent(ctx, task); agent.followup(createUserMessage({ content: [{ type: "text", text }], source: { kind: "user" }, })); dispatched = { sessionId: task.targetSession || TASK_AGENT_PREFIX + task.id }; console.log(`[dsh-plugin-mcn-schedule] 任务「${task.name}」已投递给 AI 会话 ${dispatched.sessionId}`); } catch (e) { error = String(e.message || e).slice(0, 200); console.warn(`[dsh-plugin-mcn-schedule] 任务「${task.name}」AI 投递失败: ${error}`); } } store.events.push({ id: nextId(), taskId: task.id, name: task.name, prompt: task.prompt || "", type: task.type, action, dispatched, error, at: firedAt, }); if (store.events.length > MAX_EVENTS) store.events = store.events.slice(-MAX_EVENTS); task.history = task.history || []; task.history.push({ at: firedAt, type: task.type, action, error: error || undefined }); if (task.history.length > MAX_HISTORY) task.history = task.history.slice(-MAX_HISTORY); if (task.type === "at") { task.firedAt = firedAt; task.enabled = false; // 一次性任务触发后停用 } else { task.lastFiredAt = now; } task.nextRunAt = computeNextRun(task); saveStore(); console.log(`[dsh-plugin-mcn-schedule] 触发任务「${task.name}」 action=${action}${error ? " error=" + error : ""}`); } function tick() { const now = Date.now(); let changed = false; for (const task of store.tasks) { if (!task.enabled) continue; if (typeof task.nextRunAt !== "number") { task.nextRunAt = computeNextRun(task); changed = true; continue; } if (task.nextRunAt <= now) { fireTask(tickCtx, task, now).catch((e) => { console.warn(`[dsh-plugin-mcn-schedule] tick 触发异常: ${e.message}`); }); changed = true; } } if (changed) saveStore(); } //#endregion //#region 路由 function apply(ctx) { tickCtx = ctx; // 启动时确保「计划任务」工作区存在,并修复存量会话分组(不阻塞路由注册) ensureTaskWorkspace(ctx) .then(() => repairTaskSessions(ctx)) .catch((e) => { console.warn(`[dsh-plugin-mcn-schedule] 工作区初始化失败: ${e.message}`); }); const server = ctx.get("webServer"); if (!server) { console.warn("[dsh-plugin-mcn-schedule] webServer 服务不可用,跳过路由注册"); return; } const json = (res, status, data) => { res.writeHead(status, { "Content-Type": "application/json; charset=utf-8" }); res.end(JSON.stringify(data)); }; const readBody = (req) => new Promise((resolve, reject) => { let buf = ""; req.on("data", (c) => { buf += c; if (buf.length > 1e6) req.destroy(); }); req.on("end", () => { try { resolve(buf ? JSON.parse(buf) : {}); } catch (e) { reject(e); } }); req.on("error", reject); }); const parseQuery = (url) => { const q = {}; const i = url.indexOf("?"); if (i >= 0) new URLSearchParams(url.slice(i + 1)).forEach((v, k) => { q[k] = v; }); return q; }; // 任务列表 server.register({ path: "/mcn/schedule/list", exact: true, handler: async (req, res) => { const now = Date.now(); json(res, 200, { ok: true, now, tasks: store.tasks.map((t) => ({ id: t.id, name: t.name, type: t.type, at: t.at || "", intervalMinutes: t.intervalMinutes || 0, prompt: t.prompt || "", action: t.action || "remind", targetSession: t.targetSession || "", enabled: !!t.enabled, createdAt: t.createdAt, firedAt: t.firedAt || null, lastFiredAt: t.lastFiredAt || null, nextRunAt: t.nextRunAt ?? null, history: (t.history || []).slice(-10), })), }); }, }); // 创建任务:{name, type:"at"|"every", at?, intervalMinutes?, prompt?, enabled?} server.register({ path: "/mcn/schedule/create", exact: true, handler: async (req, res) => { try { const b = await readBody(req); const name = String(b.name || "").trim(); if (!name) return json(res, 400, { ok: false, error: "任务名称不能为空" }); const type = b.type === "every" ? "every" : "at"; let atTs = null; let intervalMinutes = 0; if (type === "at") { const atStr = String(b.at || "").trim(); if (!atStr) return json(res, 400, { ok: false, error: "请选择执行时间" }); atTs = new Date(atStr).getTime(); if (!Number.isFinite(atTs)) return json(res, 400, { ok: false, error: "执行时间格式不正确" }); if (atTs <= Date.now()) return json(res, 400, { ok: false, error: "执行时间必须晚于当前时间" }); } else { intervalMinutes = Math.round(Number(b.intervalMinutes) || 0); if (intervalMinutes < MIN_INTERVAL_MINUTES) { return json(res, 400, { ok: false, error: `周期任务间隔不能小于 ${MIN_INTERVAL_MINUTES} 分钟` }); } } const now = Date.now(); const task = { id: nextId(), name, type, at: type === "at" ? String(b.at || "") : "", atTs, intervalMinutes, prompt: String(b.prompt || "").trim(), action: b.action === "agent" ? "agent" : "remind", targetSession: typeof b.targetSession === "string" ? b.targetSession.trim() : "", enabled: b.enabled !== false, createdAt: now, lastFiredAt: null, firedAt: null, history: [], }; task.nextRunAt = computeNextRun(task); store.tasks.push(task); saveStore(); json(res, 200, { ok: true, task: { id: task.id } }); } catch (e) { json(res, 400, { ok: false, error: e.message }); } }, }); // 更新任务:{id, name?, prompt?, enabled?, at?, intervalMinutes?, type?} server.register({ path: "/mcn/schedule/update", exact: true, handler: async (req, res) => { try { const b = await readBody(req); const task = store.tasks.find((t) => t.id === b.id); if (!task) return json(res, 404, { ok: false, error: "任务不存在" }); if (typeof b.name === "string" && b.name.trim()) task.name = b.name.trim(); if (typeof b.prompt === "string") task.prompt = b.prompt.trim(); if (typeof b.enabled === "boolean") task.enabled = b.enabled; if (b.action === "agent" || b.action === "remind") task.action = b.action; if (typeof b.targetSession === "string") task.targetSession = b.targetSession.trim(); if (b.type === "at" || b.type === "every") { task.type = b.type; if (task.type === "at") { const atStr = String(b.at || "").trim(); if (!atStr) return json(res, 400, { ok: false, error: "请选择执行时间" }); const ts = new Date(atStr).getTime(); if (!Number.isFinite(ts)) return json(res, 400, { ok: false, error: "执行时间格式不正确" }); task.at = atStr; task.atTs = ts; task.intervalMinutes = 0; } else { const minutes = Math.round(Number(b.intervalMinutes) || 0); if (minutes < MIN_INTERVAL_MINUTES) { return json(res, 400, { ok: false, error: `周期任务间隔不能小于 ${MIN_INTERVAL_MINUTES} 分钟` }); } task.intervalMinutes = minutes; task.at = ""; task.atTs = null; task.lastFiredAt = null; // 重新锚定周期 } } task.nextRunAt = computeNextRun(task); saveStore(); json(res, 200, { ok: true }); } catch (e) { json(res, 400, { ok: false, error: e.message }); } }, }); // 删除任务:{id} server.register({ path: "/mcn/schedule/delete", exact: true, handler: async (req, res) => { try { const b = await readBody(req); const idx = store.tasks.findIndex((t) => t.id === b.id); if (idx < 0) return json(res, 404, { ok: false, error: "任务不存在" }); store.tasks.splice(idx, 1); saveStore(); json(res, 200, { ok: true }); } catch (e) { json(res, 400, { ok: false, error: e.message }); } }, }); // 立即触发(测试用):{id} server.register({ path: "/mcn/schedule/fire-now", exact: true, handler: async (req, res) => { try { const b = await readBody(req); const task = store.tasks.find((t) => t.id === b.id); if (!task) return json(res, 404, { ok: false, error: "任务不存在" }); await fireTask(ctx, task, Date.now()); json(res, 200, { ok: true }); } catch (e) { json(res, 400, { ok: false, error: e.message }); } }, }); // 事件轮询:?since= 返回该时间之后的新触发事件 server.register({ path: "/mcn/schedule/events", exact: true, handler: async (req, res) => { const q = parseQuery(req.url); const since = Number(q.since) || 0; const fresh = store.events.filter((e) => new Date(e.at).getTime() > since); json(res, 200, { ok: true, events: fresh }); }, }); // 启动调度器 const timer = setInterval(tick, TICK_MS); timer.unref?.(); console.log("[dsh-plugin-mcn-schedule] 计划任务插件已启动"); } const inject = ["webServer", "agents", "workspaceRegistry"]; export { apply, inject };