Files
admin 856e632807 baseline: dsh 工作区基线快照(2026-09-10,全新历史起点)
本提交为 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 本地已不存在,本提交记为删除
2026-09-10 22:31:45 +08:00

484 lines
17 KiB
JavaScript
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// 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 };