483 lines
17 KiB
JavaScript
483 lines
17 KiB
JavaScript
// 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=<epochMs> 返回该时间之后的新触发事件
|
|||
|
|
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 };
|