本提交为 dsh_mcn_workshop 仓库的首个提交(旧仓 dsh_MCNProject 已停用), 完整固化 C:/Users/Administrator/.dsh 当前磁盘状态。 1) 临时数据清理:移出 sessions/(81M)、logs/、skill-backup-20260906/(14M)、 scripts/__pycache__、storages/workspace.json.bak —— 合计 94M(.dsh 162M → 68M) 2) 凭据脱敏:.credentials.yaml(含明文 DEEPSEEK_API_KEY)移出 git 跟踪并加入 .gitignore;mcn-data-insight/scripts/setup_redfox_key.py 中真实 key 示例改为占位符 3) .gitignore 增补 logs/、skill-backup-*/ 规则 4) 技能结构以当前磁盘状态为准:skills/ = mcn-short-video / impeccable / taste-skill / dsh-multi-user-migration;旧技能目录(mcn-dou-analysis、storyboard-prompt、 short-video-script、browser-harness、mcn-data-insight)与 profiles/web/node_modules/dsh-vision-router 本地已不存在,本提交记为删除
484 lines
17 KiB
JavaScript
484 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 };
|