Files
mcn-short-video/project/短视频脚本创作/V1.0/mcn-workshop/server.js
T
maogeigei a3ca08eabd fix: 工作台任务补 workspace_scope=workspace,修复会话落「未分组任务」
根因:工作台直接写宿主 automations 表时漏了 workspace_scope 字段。
对比取证:走正规 automation 工具创建的任务该字段为 'workspace'(36 条)→ 会话 is_playground=0 → 正常进工作区分组;
工作台写的为 null(当天 5 条)→ is_playground=1 → 会话落「未分组任务」。
修复:server.js 两处 INSERT INTO automations(createOnceAutomation / /api/run)补该列并填 'workspace'。
实测:新任务 scope=workspace、新会话 is_playground=0;对照组仍为 1。
排除的假因:改 cwd 无效;登记 workspaces 表无效。
2026-10-08 17:09:58 +08:00

1383 lines
76 KiB
JavaScript
Raw 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.
// 短视频工作台 - 本地 Web 服务(零依赖 Node 实现)
// 启动: node server.js [端口] → http://localhost:8900(默认;可传参指定端口,如 node server.js 9000)
// 功能: 技能导航 + 产出内容浏览(默认读桌面 <outputRoot>,根目录名由 config.json 决定)+ dsh 数据功能页(账号/视频/脚本/复盘/周榜,只读)
'use strict';
const http = require('http');
const fs = require('fs');
const path = require('path');
const os = require('os');
const { execFileSync } = require('child_process');
const { URL } = require('url');
const dsh = require('./dsh-data');
const { parseAttachment } = require('./scripts/attachment-parser');
const { loadConfig } = require('./load-config');
// 09-29 新增:本地 CodeBuddy CLI 后端(stdio 直连,替代 Dify;见 cli-backend.js 头部说明)
const { runCli, resolveCli, resolveCredentials } = require('./cli-backend');
// 09-01 默认端口由 8899 改为 8900:本机 8899 被小米 PC 管家 MiPCAudio.exe 系统服务占用(0.0.0.0 监听且自动复活),
// 此时 127.0.0.1:8899 绑定不生效 → 默认直接起在 8900;仍可 `node server.js <端口>` 显式换端口
const ROOT = __dirname;
const PUBLIC_DIR = path.join(ROOT, 'public');
const CONFIG_FILE = path.join(ROOT, 'config.json');
// 09-10 配置收敛:工作台设定(端口/DB名/产出根/榜单目录)统一收进 config.json,代码层一律从 CFG 读,
// 避免调整时到处改(config 缺失/损坏时回退 load-config.cjs 默认值)。argv 端口优先级仍高于 config。
const CFG = loadConfig();
const PORT_BASE = parseInt(process.argv[2], 10) || CFG.port;
// 09-29 AI 后端解析:config.json 的 ai.backend 为准,环境变量 MCN_AI_BACKEND 可临时覆盖
// (覆盖开关用途:不重启服务不动配置文件的前提下,快速在 cli/dify 之间切换做对比与排障)
function aiBackend() {
const e = String(process.env.MCN_AI_BACKEND || '').trim().toLowerCase();
if (e === 'cli' || e === 'dify') return e;
return String((CFG.ai && CFG.ai.backend) || 'dify').toLowerCase();
}
// 09-29 CLI 后端的联网搜索策略(追加到系统提示)。
// 根因:CLI 自带 WebSearch/WebFetch,但默认不会主动用;而每次都搜又会把首字延迟从十几秒拖到几十秒。
// 所以给模型一条明确的「什么时候才搜」的判据,而不是简单开关。
const CLI_SEARCH_HINT =
'你具备联网搜索能力(WebSearch / WebFetch),但搜索很慢,要严格按判据使用。'
+ '【该搜】回答依赖「今天/近期」的外部事实时:当下热点、平台规则与趋势、行业数据、事实核查。'
+ '【不该搜】创作讨论、需求澄清、脚本撰写、文案打磨——这些靠你的知识直接答,一律不搜。'
+ '【搜索预算】中文关键词;最多 1 次 WebSearch;仅当搜索结果指向必须打开的页面时才追加 1 次 WebFetch;'
+ '工具调用总数硬上限 2 次。搜到即答,禁止反复搜或换词重搜。'
+ '搜索失败或超时就直接凭已有信息作答,不要重试。';
// 由 config.json 的 ai.cli 段构造一次 CLI 调用参数(clarify / 自检 / 后续其它入口共用,避免多处硬编码)
function cliBaseOpts(over) {
const c = (CFG.ai && CFG.ai.cli) || {};
const searchOn = c.search !== false;
return Object.assign({
cwd: ROOT,
// MCN_AI_MODEL 可临时覆盖模型(用于不重启、不改配置的 A/B 对比与排障)
model: String(process.env.MCN_AI_MODEL || '').trim() || c.model,
apiKey: c.apiKey,
authToken: c.authToken,
credentialFile: c.credentialFile,
internetEnvironment: c.internetEnvironment,
permissionMode: c.permissionMode,
cliPath: c.path,
timeoutMs: Math.max(30000, parseInt(c.timeoutMs, 10) || 300000),
effort: c.effort || '',
search: searchOn,
tools: c.tools,
strictMcp: c.strictMcp !== false,
maxTurns: c.maxTurns,
appendSystemPrompt: c.appendSystemPrompt || (searchOn ? CLI_SEARCH_HINT : ''),
}, over || {});
}
// 09-10 路径治本(配套「仓库特征判定」):工作台运行时路径一律从 ROOT 动态推导,禁止写死盘符。
// 根因:原写死盘符指向仓库外的固定目录,仓库迁移后该目录不存在 → 浏览器锁失效、会话 cwd 悬空。
const BROWSER_LOCK_FILE = path.join(ROOT, '.browser-lock'); // 浏览器互斥锁文件(与工作台同目录)
// 10-08:默认工作台自身目录;config.json 的 sessionCwd 可指定别的工作区(分组随之改变)
const SESSION_CWD = String(CFG.sessionCwd || '').trim() || ROOT;
// 09-07:prompt 模板文件化(prompts/*.md)——规则迭代只改文件不碰代码;文件缺失/损坏回退 null,调用方自行兜底。
// 模板正文取 ==PROMPT== 与 ==CHANGELOG== 之间(头部说明/文末变更记录不发给模型)
const PROMPTS_DIR = path.join(ROOT, 'prompts');
const _promptCache = new Map();
function loadPrompt(name) {
try {
const p = path.join(PROMPTS_DIR, name.endsWith('.md') ? name : name + '.md');
const st = fs.statSync(p);
const cached = _promptCache.get(p);
if (cached && cached.mtimeMs === st.mtimeMs) return cached.text;
let text = fs.readFileSync(p, 'utf8');
// 定位「独占一行的」==PROMPT== / ==CHANGELOG== 标记(行内出现同名字样一律不算)
// ⚠️ 09-29 加固:旧版用 `(^|\n)\s*TAG\s*(\n|$)` + `m[0].indexOf(tag)` 定位——
// 说明行里若出现该词(如「读取本文件 ==PROMPT== 与 ==CHANGELOG== 之间」),
// `\s*` 可吃下换行,正则会跨行命中说明行里的字样,导致正文被切到只剩几个字
// (实测 clarify.md 被切成 3 字符「 与 」,需求打磨 prompt 整体失效)。
// 改为逐行精确匹配:行 trim 后必须完全等于标记本身,否则跳过。
const findMarker = (s, tag) => {
const lines = s.split(/\r?\n/);
let off = 0;
for (const ln of lines) {
if (ln.trim() === tag) return off;
off += ln.length + 1; // +1 补回被 split 去掉的换行符
}
return -1;
};
const a = findMarker(text, '==PROMPT==');
const b = findMarker(text, '==CHANGELOG==');
if (a >= 0) text = b > a ? text.slice(a + '==PROMPT=='.length, b) : text.slice(a + '==PROMPT=='.length);
text = text.replace(/^\s*\n/, '').replace(/\s+$/, '');
_promptCache.set(p, { mtimeMs: st.mtimeMs, text });
return text;
} catch (e) { return null; }
}
// ---- 默认根目录:桌面/<CFG.outputRoot> ----
const DEFAULT_ROOT = path.join(os.homedir(), 'Desktop', CFG.outputRoot);
// ---- 配置读写(自定义根目录列表;loadConfig 来自 load-config.js,每次读盘取最新值)----
function saveConfig(patch) {
// 合并写盘:保留 config.json 里的其它字段(port/dbName/outputRoot/businessDataDir),只更新传入的键
const merged = Object.assign({}, loadConfig(), patch);
fs.writeFileSync(CONFIG_FILE, JSON.stringify(merged, null, 2), 'utf8');
}
// ---- 允许读取的根目录白名单 ----
function allowedRoots() {
const cfg = loadConfig();
const list = [DEFAULT_ROOT].concat(cfg.extraRoots || []);
// 去重、只保留存在的目录
return list.filter((p, i) => p && fs.existsSync(p) && fs.statSync(p).isDirectory() && list.indexOf(p) === i);
}
// ---- 路径安全:必须落在某个允许根目录内 ----
function resolveSafe(rawPath) {
if (!rawPath) return null;
let p;
try { p = path.resolve(rawPath); } catch (e) { return null; }
for (const root of allowedRoots()) {
const r = path.resolve(root);
if (p === r || p.startsWith(r + path.sep)) {
return { abs: p, root: r };
}
}
return null;
}
// ---- 分析产物入库(09-01 新增):扫描账号产出目录 → 写工作台库 mcn-plugin.db ----
// 背景:账号分析技能产出只落盘文件系统未入库(俊希解析/人设卡全缺),工作台选题等读不到视频内容。
// 实现参照 dsh-plugin-mcn lib/imports.js 的 analyses 入库逻辑(dsh 接口 3080/3081 不可用时直写库)。
function resolveAccountDir(name, root) {
const safeName = String(name || '').replace(/[\\/:*?"<>|]/g, '_').trim();
if (!safeName) return null;
const roots = root ? [root] : allowedRoots();
for (const r of roots) {
const dir = path.join(r, safeName);
try { if (fs.statSync(dir).isDirectory()) return dir; } catch (e) { /* 忽略 */ }
}
return null;
}
function importAccountIntoDb(db, account, dir) {
const stats = { video_analysis: 0, video_source: 0, persona: 0, account_analysis: 0, videos_upserted: 0, unmatched: 0, unmatched_folders: [], skipped: [] };
const now = new Date().toISOString();
const insAnalysis = db.prepare(`INSERT INTO account_video_analysis (video_id, aweme_id, analysis_time, content_json, summary) VALUES (?,?,?,?,?)`);
const insSource = db.prepare(`INSERT INTO account_video_source (video_id, aweme_id, analysis_time, content_json, summary, analysis_json) VALUES (?,?,?,?,?,?)`);
const insPersona = db.prepare(`INSERT INTO account_persona (account_id, analysis_time, content_json, summary) VALUES (?,?,?,?)`);
const insAccAna = db.prepare(`INSERT INTO account_analysis (account_id, analysis_time, content_json, summary) VALUES (?,?,?,?)`);
const acc = db.prepare('SELECT id FROM hot_accounts WHERE account_name=? LIMIT 1').get(account);
const accountId = acc ? acc.id : null;
// 清理历史误插脏数据:09-01 匹配缺陷(标点未规范化)产生 aweme_id=null 且 video_id=null 的无意义记录,重跑入库前先清
db.prepare("DELETE FROM account_video_analysis WHERE aweme_id IS NULL AND video_id IS NULL").run();
db.prepare("DELETE FROM account_video_source WHERE aweme_id IS NULL AND video_id IS NULL").run();
// 规范化标题:去 # 标签 → 删全角/半角标点(xlsx 标题带「!」「“”」而文件夹名无,09-01 匹配失败根因;09-04 补全角问号「?」——folder 生成时删 ? 而 title 未删导致残留干扰)→ 下划线转空格(对齐文件夹 _标签 结构)→ 空格归一
const normTitle = (s) => String(s || '').split('#')[0].replace(/[!!??“”‘’"'`,.…。,、;::;·]/g, '').replace(/[_-]+/g, ' ').replace(/\s+/g, ' ').trim();
// 09-04 修复:文件夹名与 video_title 的规范化不对称——文件夹命名规则「标题去特殊符号」会把词间空格一并吞掉
// ("小主人: 终于"→"小主人终于"),而 normTitle 只删标点保留空格("小主人 终于")→ 两侧精确/包含匹配双失效
// → aweme_id=null 脏入库 → 前端视频详情(按 aweme_id 关联,要求 NOT NULL)查不到解析。
// 治本:比较键统一为「无空格紧凑形态」(normTitle 后再删全部空白),folder 与 title 两侧同构;精确失败后走无空格包含匹配。
const normKey = (s) => normTitle(s).replace(/\s+/g, '');
// 0) aweme_id 权威映射:短视频表格.xlsx(子进程 python,最高优先级)→ account_videos 表(兜底)
const xlsxList = readXlsxMap(dir);
const xlsxByNorm = {};
for (const r of xlsxList) { if (r.aweme_id && r.title) xlsxByNorm[normKey(r.title)] = r.aweme_id; }
const titleRows = db.prepare('SELECT aweme_id, video_title FROM account_videos WHERE video_title IS NOT NULL').all();
const dbByNorm = {};
for (const r of titleRows) { if (r.aweme_id && r.video_title) dbByNorm[normKey(r.video_title)] = String(r.aweme_id); }
// 匹配:无空格紧凑精确优先;包含匹配要求双方 ≥10 字符且按长度降序(最具体优先,防短标题误配)
const matchAweme = (folder) => {
const base = normKey(folder);
if (!base) return null;
if (xlsxByNorm[base]) return xlsxByNorm[base];
if (dbByNorm[base]) return dbByNorm[base];
const cands = [];
for (const [t, aw] of Object.entries(xlsxByNorm)) {
if (base.length >= 10 && t.length >= 10 && (base.includes(t) || t.includes(base))) cands.push([t.length, aw]);
}
if (!cands.length) {
for (const [t, aw] of Object.entries(dbByNorm)) {
if (base.length >= 10 && t.length >= 10 && (base.includes(t) || t.includes(base))) cands.push([t.length, aw]);
}
}
cands.sort((a, b) => b[0] - a[0]);
return cands.length ? cands[0][1] : null;
};
// 0.5) 视频列表入库:短视频表格.xlsx → account_videos(按 account_id+aweme_id 去重,补上工作台视频列表数据)
if (accountId && xlsxList.length) {
const insVideo = db.prepare(`INSERT INTO account_videos (account_id, aweme_id, video_title, video_url, like_count, comment_count, share_count, collect_count, play_count, duration, publish_time, tags, collected_time) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?)`);
for (const v of xlsxList) {
const ex = db.prepare('SELECT id FROM account_videos WHERE account_id=? AND aweme_id=? LIMIT 1').get(accountId, v.aweme_id);
if (!ex) {
insVideo.run(accountId, v.aweme_id, v.title, v.url || null, v.like_count ?? null, v.comment_count ?? null, v.share_count ?? null, v.collect_count ?? null, v.play_count ?? null, v.duration ?? null, v.publish_time || null, v.tags || null, now);
stats.videos_upserted++;
}
}
}
// 1) 视频分析目录(功能三 Step5 产出;功能五产出在 视频对标/ 目录,由 dsh 导入机制处理)
// 字段映射(09-01 修正,对齐 dsh 语义):content.json→account_video_source.content_json(视频脚本);
// analysis.json→account_video_source.analysis_json(视频拆解,MCP 结构化分析);*拆解分析.md→account_video_analysis(视频分析,功能四产物)。
// 原实现把 analysis.json 误写 account_video_analysis 导致「拆解内容进了视频分析」+「视频拆解 tab 空」
const vaDir = path.join(dir, '视频分析');
if (fs.existsSync(vaDir)) {
for (const folder of fs.readdirSync(vaDir)) {
const fdir = path.join(vaDir, folder);
let st; try { st = fs.statSync(fdir); } catch (e) { continue; }
if (!st.isDirectory()) continue;
const awemeId = matchAweme(folder);
// 09-04:awemeId 命中后反查 account_videos.id 一并填 video_id(双键关联,前端 video_id/aweme_id 皆可命中;历史 source 行 video_id 多为 null 靠 aweme_id 兜底)
let videoId = null;
if (awemeId) { const vv = db.prepare('SELECT id FROM account_videos WHERE aweme_id=?').get(awemeId); videoId = vv ? vv.id : null; }
const anPath = path.join(fdir, 'analysis.json');
const coPath = path.join(fdir, 'content.json');
const coExists = fs.existsSync(coPath), anExists = fs.existsSync(anPath);
if (!coExists && !anExists) continue; // 空文件夹(半途写入残留)跳过,不计 unmatched
// 09-04:aweme_id 匹配不到 → 记 unmatched 告警并跳过插入(不再制造双 null 脏记录、避免"入库假成功";补 xlsx/视频列表后重跑 import 即恢复)
if (!awemeId) { stats.unmatched++; stats.unmatched_folders.push(folder); }
if (coExists) {
const content = fs.readFileSync(coPath, 'utf8');
const analysisJson = anExists ? fs.readFileSync(anPath, 'utf8') : null;
if (awemeId) {
db.prepare('DELETE FROM account_video_source WHERE aweme_id=?').run(awemeId);
insSource.run(videoId, awemeId, now, content, '', analysisJson);
stats.video_source++;
}
}
// 拆解分析.md → account_video_analysis(功能四逐条拆解产物,与解析并存于同一文件夹;无则不入)
const decon = fs.readdirSync(fdir).find((f) => f.endsWith('拆解分析.md'));
if (decon) {
const content = fs.readFileSync(path.join(fdir, decon), 'utf8');
if (awemeId) {
db.prepare('DELETE FROM account_video_analysis WHERE aweme_id=?').run(awemeId);
insAnalysis.run(videoId, awemeId, now, content, content.slice(0, 500));
stats.video_analysis++;
}
}
}
}
// 2) 人设卡 + 账号分析(去重,不删旧)
if (accountId) {
const personaPath = path.join(dir, account + '账号设定.md');
if (fs.existsSync(personaPath)) {
const content = fs.readFileSync(personaPath, 'utf8');
const dup = db.prepare('SELECT id FROM account_persona WHERE account_id=? AND content_json=? LIMIT 1').get(accountId, content);
if (!dup) { insPersona.run(accountId, now, content, content.slice(0, 200)); stats.persona++; }
else stats.skipped.push('persona(重复)');
}
const anaPath = path.join(dir, account + '账号数据分析.md');
if (fs.existsSync(anaPath)) {
const content = fs.readFileSync(anaPath, 'utf8');
const dup = db.prepare('SELECT id FROM account_analysis WHERE account_id=? AND content_json=? LIMIT 1').get(accountId, content);
if (!dup) { insAccAna.run(accountId, now, content, content.slice(0, 200)); stats.account_analysis++; }
else stats.skipped.push('account_analysis(重复)');
}
}
return { account, account_id: accountId, dir, xlsx_rows: xlsxList.length, title_map_size: Object.keys(dbByNorm).length, ...stats };
}
// 读短视频表格.xlsx → 视频映射数组(子进程调 python export_video_map.py,openpyxl 读 xlsx;失败返回 [])
function readXlsxMap(dir) {
const xlsxPath = path.join(dir, '短视频表格.xlsx');
if (!fs.existsSync(xlsxPath)) return [];
const py = process.env.PYTHON || 'python';
try {
const out = execFileSync(py, [path.join(ROOT, 'scripts', 'export_video_map.py'), xlsxPath], { encoding: 'utf8', timeout: 20000, windowsHide: true });
const arr = JSON.parse(out);
return Array.isArray(arr) ? arr : [];
} catch (e) {
log('xlsx 视频映射读取失败(用表内匹配兜底): ' + e.message);
return [];
}
}
// ---- 查找账号设定卡 SVG ----
// 产出规范:{产出根}/{达人昵称}/{达人昵称}账号设定卡.svg(功能六产出,与账号设定.md 同目录)
// 查找顺序:规范路径 → 根下直挂 → 账号目录内扫描(防目录名/账号名不一致)
function findAccountCard(name) {
const safeName = String(name || '').replace(/[\\/:*?"<>|]/g, '_').trim();
if (!safeName) return null;
for (const root of allowedRoots()) {
const p1 = path.join(root, safeName, safeName + '账号设定卡.svg');
const p2 = path.join(root, safeName + '账号设定卡.svg');
for (const c of [p1, p2]) {
try { if (fs.existsSync(c) && fs.statSync(c).isFile()) return c; } catch (e) { /* 忽略 */ }
}
}
for (const root of allowedRoots()) {
const dir = path.join(root, safeName);
try {
if (fs.statSync(dir).isDirectory()) {
const hit = fs.readdirSync(dir).find((f) => f.toLowerCase().endsWith('.svg') && f.includes('账号设定卡'));
if (hit) return path.join(dir, hit);
}
} catch (e) { /* 忽略 */ }
}
return null;
}
// ---- 目录树(递归,跳过隐藏项)----
function buildTree(absPath, depth) {
if (depth > 4) return null;
let entries;
try { entries = fs.readdirSync(absPath, { withFileTypes: true }); } catch (e) { return null; }
const dirs = [];
const files = [];
for (const ent of entries) {
if (ent.name.startsWith('.')) continue; // 隐藏文件
if (ent.name === 'node_modules' || ent.name === '.git') continue;
const full = path.join(absPath, ent.name);
try {
if (ent.isDirectory()) {
const child = buildTree(full, depth + 1);
if (child) dirs.push({ name: ent.name, path: full, type: 'dir', children: child });
} else if (ent.isFile()) {
const st = fs.statSync(full);
files.push({ name: ent.name, path: full, type: 'file', ext: path.extname(ent.name).toLowerCase(), size: st.size });
}
} catch (e) { /* 权限/占用跳过 */ }
}
dirs.sort((a, b) => a.name.localeCompare(b.name, 'zh'));
files.sort((a, b) => a.name.localeCompare(b.name, 'zh'));
return dirs.concat(files);
}
// ---- MIME ----
const MIME = {
'.html': 'text/html; charset=utf-8',
'.js': 'text/javascript; charset=utf-8',
'.css': 'text/css; charset=utf-8',
'.json': 'application/json; charset=utf-8',
'.svg': 'image/svg+xml',
'.png': 'image/png',
'.jpg': 'image/jpeg',
'.jpeg': 'image/jpeg',
'.gif': 'image/gif',
'.ico': 'image/x-icon',
'.woff2': 'font/woff2',
};
function mimeOf(p) {
return MIME[path.extname(p).toLowerCase()] || 'application/octet-stream';
}
// ---- 读取文本文件(按扩展名判定预览能力)----
const PREVIEW_EXTS = ['.md', '.markdown', '.txt', '.json', '.csv', '.log'];
function readText(absPath) {
const ext = path.extname(absPath).toLowerCase();
if (!PREVIEW_EXTS.includes(ext)) return null;
const buf = fs.readFileSync(absPath);
// BOM 去除
let text = buf.toString('utf8');
if (text.charCodeAt(0) === 0xFEFF) text = text.slice(1);
return text;
}
// ---- 简易日志 ----
function log(msg) {
console.log('[' + new Date().toLocaleTimeString('zh-CN', { hour12: false }) + '] ' + msg);
}
// ---------- 榜单数据更新任务(09-03 新增:热点数据页「更新榜单」按钮后台任务)----------
// 任务名固定用于防重识别(同名最近一次任务状态 = 是否有更新在跑 / 上次执行结果)
const RANK_UPDATE_NAME = '更新榜单数据';
const RANK_UPDATE_COOLDOWN_MS = 30 * 60 * 1000; // 防重复更新窗口:距上次成功落盘 < 30 分钟拦截(可 force 跳过)
// workbuddy 宿主库(automations/sessions/automation_runs),与 /api/run 同款路径解析
function wbDb() {
const { DatabaseSync } = require('node:sqlite');
return new DatabaseSync(process.env.WORKBUDDY_DB || path.join(os.homedir(), '.workbuddy', 'workbuddy.db'));
}
// 创建一次性后台任务(08-31 实测约定:scheduled_at 与 next_run_at 均取 +5s 未来,调度器按 next_run_at 扫描拾取;
// 模型跟随用户最近活跃手动会话选择,兜底 'auto')——逻辑与 /api/run 完全一致,供更新榜单等新入口复用
// 09-04 connectorIdsArr:任务会话要用的「用户连接器」configId 列表(如 ['custom-mcp:myai-mcp-production'])——
// 宿主会话 register(connectorIds) 只注入列出的连接器,空数组=不注入任何连接器(后台任务「找不到 MCP」根因),
// 依赖 MCP 的任务(视频解析等)必须显式挂载;纯脚本类任务(榜单更新)保持 [] 即可
function createOnceAutomation(name, prompt, skillsArr, connectorIdsArr) {
const db = wbDb();
try {
const now = Date.now();
const id = 'automation-' + now;
const cwd = SESSION_CWD; // 工作台触发的会话归入 mcn-workshop 空间分组(由 ROOT 动态推导)
const d = new Date(now + 5 * 1000);
const pad = (n) => String(n).padStart(2, '0');
const scheduledAt = `${d.getFullYear()}-${pad(d.getMonth() + 1)}-${pad(d.getDate())}T${pad(d.getHours())}:${pad(d.getMinutes())}:${pad(d.getSeconds())}`;
const nextRunAt = now + 5 * 1000;
const uidRow = db.prepare("SELECT user_id FROM sessions WHERE user_id IS NOT NULL AND user_id <> '' ORDER BY created_at DESC LIMIT 1").get();
const ownerUserId = uidRow ? uidRow.user_id : '';
const modelRow = db.prepare(`SELECT model FROM sessions WHERE deleted_at IS NULL AND COALESCE(is_background_automation,0)=0 AND model IS NOT NULL AND model <> '' ORDER BY COALESCE(last_activity_at, updated_at) DESC LIMIT 1`).get();
const modelId = (modelRow && modelRow.model) ? modelRow.model : 'auto';
// 10-08 修复:必须带 workspace_scope='workspace'。走正规 automation 工具创建的任务都带该值(库里 36 条如此),
// 工作台直接写库时漏了它 ⇒ 会话被视为"自动生成的空间"、落进「未分组任务」。
db.prepare(`INSERT INTO automations (id,name,prompt,status,schedule_type,scheduled_at,next_run_at,rrule,cwds,created_at,updated_at,skills_json,connector_ids_json,model_id,permission_mode,owner_user_id,owner_status,workspace_scope)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)`)
.run(id, name, prompt, 'ACTIVE', 'once', scheduledAt, nextRunAt, '', JSON.stringify([cwd]), now, now, JSON.stringify(Array.isArray(skillsArr) ? skillsArr : []), JSON.stringify(Array.isArray(connectorIdsArr) ? connectorIdsArr : []), modelId, 'fullAccess', ownerUserId, 'confirmed', 'workspace');
return { id };
} finally {
db.close();
}
}
// 最近一次「更新榜单数据」任务状态(running/queued/pending = 进行中;done/error/review = 已结束)
function rankLastTaskState() {
try {
const db = wbDb();
try {
const row = db.prepare('SELECT id, status, deleted_at FROM automations WHERE name=? AND deleted_at IS NULL ORDER BY created_at DESC LIMIT 1').get(RANK_UPDATE_NAME);
if (!row) return null;
const st = db.prepare('SELECT running, last_error FROM automation_runtime_state WHERE automation_id=?').get(row.id);
const runs = db.prepare('SELECT status FROM automation_runs WHERE automation_id=? ORDER BY created_at DESC LIMIT 1').get(row.id);
let state = 'queued';
if (st && Number(st.running) === 1) state = 'running';
else if (runs) {
const s = String(runs.status || '').toUpperCase();
if (s === 'IN_PROGRESS') state = 'running';
else if (['DONE', 'COMPLETED', 'SUCCESS', 'FINISHED', 'ACCEPTED'].includes(s)) state = 'done';
else if (['ERROR', 'CANCELLED', 'FAILED', 'INTERRUPTED'].includes(s)) state = 'error';
else if (s === 'PENDING_REVIEW') state = 'review';
else state = 'pending';
}
return { id: row.id, state, error: (st && st.last_error) || null };
} finally { db.close(); }
} catch (e) { return null; }
}
// ---------- 任务进度可见(10-08 新增:用户「我需要看到 才知道执行的是否正常」)----------
// 后台任务只返回 running/done 太粗:用户无法判断"在真跑"还是"卡住了"。
// 解法:automation → 关联它拉起的后台会话 → 读该会话 jsonl 尾部 → 产出
// 已运行时长 / 工具调用统计 / 最近一段动作 / 会话最后写入距今(判断是否停滞)
// 会话 cwd → WorkBuddy projects 目录名(实测:`D:\a\b\V1.0\mcn-workshop` → `d-a-b-V1.0-mcn-workshop`)
function escapeCwdForProjects(cwd) {
const s = String(cwd || '').replace(/\\/g, '/');
const m = s.match(/^([A-Za-z]):\/(.+)$/);
if (!m) return null;
return m[1].toLowerCase() + '-' + m[2].replace(/\//g, '-');
}
const TOOL_NAME_RE = /"name":"(Bash|Read|Write|Edit|Grep|Glob|WebSearch|WebFetch|WebFetch2|Skill|Task|Agent|TodoWrite|ExitPlanMode|mcp__[A-Za-z0-9_]+)"/g;
function readSessionProgress(sessionId, cwd) {
const out = { lastAction: '', toolCalls: {}, updatedAgoMs: null, stalled: false, logBytes: 0 };
try {
const esc = escapeCwdForProjects(cwd);
if (!esc) return out;
const f = path.join(os.homedir(), '.workbuddy', 'projects', esc, sessionId + '.jsonl');
if (!fs.existsSync(f)) return out;
const st = fs.statSync(f);
out.logBytes = st.size;
out.updatedAgoMs = Date.now() - Number(st.mtimeMs);
// 90 秒没再写入即视为停滞(模型正常推进时几乎每秒都在写日志)
out.stalled = out.updatedAgoMs > 90000;
if (st.size <= 0) return out;
const readLen = Math.min(st.size, 256 * 1024);
const fd = fs.openSync(f, 'r');
const buf = Buffer.alloc(readLen);
fs.readSync(fd, buf, 0, readLen, st.size - readLen);
fs.closeSync(fd);
const txt = buf.toString('utf8');
let m;
TOOL_NAME_RE.lastIndex = 0;
while ((m = TOOL_NAME_RE.exec(txt))) out.toolCalls[m[1]] = (out.toolCalls[m[1]] || 0) + 1;
const texts = [];
const re2 = /"text":"((?:[^"\\]|\\.){4,})"/g;
while ((m = re2.exec(txt))) {
let s = m[1];
try { s = JSON.parse('"' + s + '"'); } catch (e) { /* 解析失败保留原串 */ }
if (s && String(s).trim()) texts.push(String(s).trim());
}
if (texts.length) out.lastAction = texts[texts.length - 1].replace(/\s+/g, ' ').slice(0, 160);
} catch (e) { /* 读不到不影响状态判定 */ }
return out;
}
// ---- 静态文件 / API ----
const server = http.createServer((req, res) => {
const u = new URL(req.url, 'http://localhost');
const p = decodeURIComponent(u.pathname);
const q = u.searchParams;
// 统一 JSON 响应
const sendJSON = (code, obj) => {
const body = JSON.stringify(obj);
res.writeHead(code, { 'Content-Type': 'application/json; charset=utf-8', 'Cache-Control': 'no-store' });
res.end(body);
};
const sendErr = (code, msg) => sendJSON(code, { error: msg });
// ---------- API ----------
if (p === '/api/roots') {
return sendJSON(200, {
defaultRoot: DEFAULT_ROOT,
extraRoots: loadConfig().extraRoots,
defaultExists: fs.existsSync(DEFAULT_ROOT),
});
}
if (p === '/api/roots' && req.method === 'POST') {
let body = '';
req.on('data', (c) => { body += c; if (body.length > 1e6) req.destroy(); });
req.on('end', () => {
try {
const { extraRoots } = JSON.parse(body || '{}');
if (!Array.isArray(extraRoots)) return sendErr(400, 'extraRoots 必须为数组');
const cleaned = extraRoots
.map((x) => (typeof x === 'string' ? x.trim() : ''))
.filter((x) => x && fs.existsSync(x) && fs.statSync(x).isDirectory());
saveConfig({ extraRoots: cleaned });
sendJSON(200, { extraRoots: cleaned });
} catch (e) {
sendErr(500, '配置保存失败: ' + e.message);
}
});
return;
}
if (p === '/api/tree') {
const rootParam = q.get('root') || DEFAULT_ROOT;
const safe = resolveSafe(rootParam);
if (!safe) return sendErr(403, '目录不在允许范围,请先在设置中添加');
const children = buildTree(safe.abs, 0) || [];
return sendJSON(200, { root: safe.abs, tree: children });
}
if (p === '/api/read') {
const safe = resolveSafe(q.get('path'));
if (!safe) return sendErr(403, '路径不在允许范围');
try {
const st = fs.statSync(safe.abs);
if (st.size > 5 * 1024 * 1024) return sendErr(413, '文件超过 5MB,不支持在线预览,请下载查看');
const text = readText(safe.abs);
if (text === null) return sendErr(415, '该类型不支持在线预览,请下载后用本地软件打开');
return sendJSON(200, {
name: path.basename(safe.abs),
path: safe.abs,
ext: path.extname(safe.abs).toLowerCase(),
size: st.size,
content: text,
});
} catch (e) {
return sendErr(500, '读取失败: ' + e.message);
}
}
if (p === '/api/download') {
const safe = resolveSafe(q.get('path'));
if (!safe) return sendErr(403, '路径不在允许范围');
try {
const st = fs.statSync(safe.abs);
const fname = encodeURIComponent(path.basename(safe.abs));
res.writeHead(200, {
'Content-Type': mimeOf(safe.abs),
'Content-Length': st.size,
'Content-Disposition': "attachment; filename*=UTF-8''" + fname,
'Cache-Control': 'no-store',
});
fs.createReadStream(safe.abs).pipe(res);
} catch (e) {
return sendErr(500, '下载失败: ' + e.message);
}
return;
}
if (p === '/api/file') {
// 内联输出(无 attachment 头):供 <img>/<object> 直接渲染(如账号设定卡 SVG)
const safe = resolveSafe(q.get('path'));
if (!safe) return sendErr(403, '路径不在允许范围');
try {
const st = fs.statSync(safe.abs);
if (!st.isFile()) return sendErr(404, '文件不存在');
res.writeHead(200, {
'Content-Type': mimeOf(safe.abs),
'Content-Length': st.size,
'Cache-Control': 'no-store',
});
fs.createReadStream(safe.abs).pipe(res);
} catch (e) {
return sendErr(500, '读取失败: ' + e.message);
}
return;
}
// ---------- 数据功能(只读工作台自有数据库副本) ----------
const dshOk = (code, obj) => sendJSON(code, { ok: true, ...obj });
const dshGuard = (res2) => { if (!dsh.dbAvailable()) { sendJSON(200, { ok: false, error: '未检测到工作台数据库(' + CFG.dbName + '),本功能不可用' }); return true; } return false; };
const qInt = (k, def) => { const v = Number(q.get(k)); return Number.isFinite(v) ? v : def; };
if (p === '/api/dsh/meta') {
return sendJSON(200, { ok: true, db: dsh.dbAvailable(), dbPath: dsh.DSH_DB, rankingDir: dsh.RANKING_DIR, ranking: dsh.rankingMeta() });
}
if (p === '/api/dsh/stats') {
if (dshGuard()) return;
return dshOk(200, { stats: dsh.getStats() });
}
if (p === '/api/dsh/recent-scripts') {
if (dshGuard()) return;
return dshOk(200, { items: dsh.listRecentScripts(qInt('limit', 4)) });
}
if (p === '/api/dsh/tracks') {
if (dshGuard()) return;
return dshOk(200, { tracks: dsh.listTracks() });
}
if (p === '/api/dsh/accounts') {
if (dshGuard()) return;
try {
return dshOk(200, dsh.listAccounts({ page: qInt('page', 1), pageSize: qInt('pageSize', 20), search: String(q.get('search') || ''), track: String(q.get('track') || ''), type: String(q.get('type') || 'all') }));
} catch (e) { return sendErr(500, e.message); }
}
if (p === '/api/dsh/account-favorite' && req.method === 'POST') {
if (dshGuard()) return;
let body = '';
req.on('data', (c) => { body += c; if (body.length > 1e6) req.destroy(); });
req.on('end', () => {
try {
const { id, fav } = JSON.parse(body || '{}');
if (!id) return sendErr(400, '缺少 id');
dsh.setFavorite(id, !!fav);
return dshOk(200, { id, favorite: !!fav });
} catch (e) { return sendErr(500, e.message); }
});
return;
}
if (p === '/api/dsh/batch-delete' && req.method === 'POST') {
if (dshGuard()) return;
let body = '';
req.on('data', (c) => { body += c; if (body.length > 1e6) req.destroy(); });
req.on('end', () => {
try {
const { kind, ids } = JSON.parse(body || '{}');
if (!['accounts', 'videos', 'rewrites'].includes(kind)) return sendErr(400, 'kind 不合法');
const r = dsh.batchDelete(kind, ids);
return dshOk(200, { ok: true, deleted: r.deleted });
} catch (e) { return sendErr(500, e.message); }
});
return;
}
if (p === '/api/dsh/script-delete' && req.method === 'POST') {
if (dshGuard()) return;
let body = '';
req.on('data', (c) => { body += c; if (body.length > 1e6) req.destroy(); });
req.on('end', () => {
try {
const { id } = JSON.parse(body || '{}');
if (!(Number(id) > 0)) return sendErr(400, '缺少脚本 id');
const r = dsh.deleteScript(Number(id)); /* R58:连坐删除该脚本的诊断复盘与分镜提示词 */
return dshOk(200, { ok: true, deleted: r.deleted });
} catch (e) { return sendErr(500, e.message); }
});
return;
}
if (p === '/api/dsh/account-videos') {
if (dshGuard()) return;
try {
return dshOk(200, dsh.listAccountVideos(qInt('accountId', 0), qInt('page', 1), qInt('pageSize', 20)));
} catch (e) { return sendErr(500, e.message); }
}
if (p === '/api/dsh/account-persona') {
if (dshGuard()) return;
const r = dsh.getPersona(qInt('accountId', 0));
return dshOk(200, { content: r ? r.content : null, analysisTime: r ? r.analysis_time : null });
}
if (p === '/api/dsh/account-card') {
// 查询账号设定卡 SVG 文件(桌面产出,允许根目录下查找):{账号名}账号设定卡.svg
if (dshGuard()) return;
const hit = findAccountCard(q.get('name'));
return dshOk(200, { exists: !!hit, path: hit || null });
}
if (p === '/api/dsh/account-analysis') {
if (dshGuard()) return;
const r = dsh.getAccountAnalysis(qInt('accountId', 0));
return dshOk(200, { content: r ? r.content : null, analysisTime: r ? r.analysis_time : null });
}
if (p === '/api/dsh/videos') {
if (dshGuard()) return;
return dshOk(200, dsh.listVideos({ page: qInt('page', 1), pageSize: qInt('pageSize', 20), search: String(q.get('search') || ''), parsed: String(q.get('parsed') || ''), days: String(q.get('days') || '') }));
}
if (p === '/api/dsh/video-detail') {
if (dshGuard()) return;
try {
return dshOk(200, dsh.getVideoDetail({ videoId: qInt('videoId', 0), awemeId: q.get('awemeId') }));
} catch (e) { return sendErr(500, e.message); }
}
if (p === '/api/dsh/video-script') {
if (dshGuard()) return;
try {
return dshOk(200, dsh.getVideoScript(qInt('videoId', 0)));
} catch (e) { return sendErr(500, e.message); }
}
if (p === '/api/dsh/rewrites') {
if (dshGuard()) return;
return dshOk(200, dsh.listRewrites(qInt('videoId', 0)));
}
if (p === '/api/dsh/custom-scripts') { /* 09-02 v25:原创选题(无视频载体)脚本列表 */
if (dshGuard()) return;
return dshOk(200, dsh.listCustomScripts({ page: qInt('page', 1), pageSize: qInt('pageSize', 20), search: String(q.get('search') || ''), days: String(q.get('days') || ''), account: String(q.get('account') || '') }));
}
if (p === '/api/dsh/rewrite-stats') { /* 09-02 v25:AI写脚本 tab 徽标计数(参考选题/原创选题),可传 account 按账号统计 */
if (dshGuard()) return;
return dshOk(200, dsh.getRewriteStats(String(q.get('account') || '')));
}
if (p === '/api/dsh/custom-script') { /* 原创选题脚本全文(查看弹窗) */
if (dshGuard()) return;
try { return dshOk(200, { script: dsh.getCustomScript(qInt('id', 0)) }); }
catch (e) { return sendErr(404, e.message); }
}
if (p === '/api/dsh/script-save' && req.method === 'POST') { /* AI创作任务完成回调写库(参考选题传 videoId;原创选题省略→video_id 空;带 topicId 则回填选题关联) */
if (dshGuard()) return;
let body = '';
req.on('data', (c) => { body += c; if (body.length > 2e6) req.destroy(); });
req.on('end', () => {
try {
const { videoId, awemeId, accountId, title, scriptText, topicId } = JSON.parse(body || '{}');
if (!scriptText || !accountId) return sendErr(400, '缺少 scriptText/accountId');
const r = dsh.saveScript({ videoId, awemeId, accountId, title, scriptText, topicId });
return dshOk(200, { ok: true, id: r.id });
} catch (e) { return sendErr(500, e.message); }
});
return;
}
/* ===== 选题列表(09-16 新增):提交创作前先落选题;脚本生成后经 script-save 回填关联 ===== */
if (p === '/api/dsh/topic-save' && req.method === 'POST') {
if (dshGuard()) return;
let body = '';
req.on('data', (c) => { body += c; if (body.length > 2e6) req.destroy(); });
req.on('end', () => {
try {
const { id, accountId, title, req: reqList, dur } = JSON.parse(body || '{}');
if (!accountId) return sendErr(400, '缺少 accountId');
const r = dsh.saveTopic({ id, accountId, title, req: reqList, dur });
return dshOk(200, { ok: true, id: r.id, updated: !!r.updated });
} catch (e) { return sendErr(500, e.message); }
});
return;
}
if (p === '/api/dsh/topics' && req.method === 'GET') {
try { return dshOk(200, { ok: true, items: dsh.listTopics(qInt('accountId', 0)) }); }
catch (e) { return sendErr(500, e.message); }
}
if (p === '/api/dsh/topic-script' && req.method === 'GET') {
try { return dshOk(200, { ok: true, topic: dsh.getTopicScript(qInt('topicId', 0)) }); }
catch (e) { return sendErr(404, e.message); }
}
if (p === '/api/dsh/storyboards') {
if (dshGuard()) return;
return dshOk(200, { items: dsh.listStoryboards(qInt('videoId', 0), qInt('rewriteId', 0)) });
}
if (p === '/api/dsh/review-detail') {
if (dshGuard()) return;
const r = dsh.getReviewDetail(qInt('id', 0));
return r ? dshOk(200, { review: r }) : sendErr(404, '复盘记录不存在');
}
if (p === '/api/dsh/review-full') {
if (dshGuard()) return;
const r = dsh.getReviewFull(qInt('rewriteId', 0));
return r ? dshOk(200, { reviews: r }) : sendErr(404, '改写记录不存在');
}
if (p === '/api/dsh/ranking') {
try {
const r = dsh.getRanking({ board: String(q.get('board') || 'week'), date: String(q.get('date') || ''), category: String(q.get('category') || ''), sortBy: String(q.get('sortBy') || ''), order: String(q.get('order') || ''), page: qInt('page', 1), pageSize: qInt('pageSize', 20) });
return dshOk(200, r);
} catch (e) { return sendErr(500, e.message); }
}
// ---------- 榜单数据更新(09-03 新增:热点数据页「更新榜单」按钮;执行前防重复校验)----------
// GET:返回数据新鲜度(最近成功落盘时刻 = RANKING_DIR 最新 json mtime)+ 最近同名更新任务状态,供前端展示与预拦截
if (p === '/api/dsh/ranking-update' && req.method === 'GET') {
try {
const fr = dsh.rankingFreshness();
return sendJSON(200, {
ok: true,
exists: fr.exists, files: fr.files,
lastUpdateAt: fr.latestMtime || null,
lastUpdateDate: fr.latestDate || '',
lastUpdateAgeMin: fr.latestMtime ? Math.floor((Date.now() - fr.latestMtime) / 60000) : null,
cooldownMin: RANK_UPDATE_COOLDOWN_MS / 60000,
task: rankLastTaskState(),
});
} catch (e) { return sendErr(500, '查询更新状态失败: ' + e.message); }
}
// POST:防重校验(① 更新任务进行中拦截;② 距上次落盘 < 冷却窗口且未 force 拦截)→ 创建后台抓取任务
if (p === '/api/dsh/ranking-update' && req.method === 'POST') {
let body = '';
req.on('data', (c) => { body += c; if (body.length > 1e6) req.destroy(); });
req.on('end', () => {
try {
const { prompt, force } = JSON.parse(body || '{}');
if (!prompt || typeof prompt !== 'string' || !prompt.trim()) return sendErr(400, 'prompt 不能为空');
// 防重①:进行中(running/queued/pending)→ 拦截
const task = rankLastTaskState();
if (task && ['running', 'queued', 'pending'].includes(task.state)) {
return sendJSON(200, { ok: false, error: '已有榜单更新任务在执行中(可在左侧会话栏查看进度),请勿重复提交' });
}
// 防重②:距上次成功落盘 < 冷却窗口且未 force → 拦截(needForce 供前端提示仍可强制更新)
const fr = dsh.rankingFreshness();
if (!force && fr.latestMtime && Date.now() - fr.latestMtime < RANK_UPDATE_COOLDOWN_MS) {
const ageMin = Math.floor((Date.now() - fr.latestMtime) / 60000);
return sendJSON(200, {
ok: false, needForce: true,
error: `榜单数据 ${Math.max(ageMin, 1)} 分钟前刚更新过(${fr.latestDate || ''}),${RANK_UPDATE_COOLDOWN_MS / 60000} 分钟内防重复更新;如确需强制刷新请再次点击并选择「仍要更新」`,
});
}
// 通过 → 创建后台更新任务(skills 挂 mcn-data-insight:三榜抓取脚本与落盘规范所在)
const r = createOnceAutomation(RANK_UPDATE_NAME, prompt.trim(), ['mcn-data-insight']);
log('已提交榜单更新任务: ' + r.id);
return sendJSON(200, { ok: true, id: r.id, message: '榜单更新任务已提交,请在左侧会话栏查看执行' });
} catch (e) { return sendErr(500, '任务提交失败: ' + e.message); }
});
return;
}
// ---------- 提交 AI 任务(Automation 一次性任务,WorkBuddy 自动执行,左侧会话栏可见)----------
if (p === '/api/run' && req.method === 'POST') {
let body = '';
req.on('data', (c) => { body += c; if (body.length > 1e6) req.destroy(); });
req.on('end', () => {
try {
const { prompt, name, skills, connectorIds } = JSON.parse(body || '{}');
if (!prompt || typeof prompt !== 'string' || !prompt.trim()) return sendErr(400, 'prompt 不能为空');
// 09-04 连接器级挂载:connector_ids_json 填「用户连接器」configId 数组(如 custom-mcp:myai-mcp-production)。
// 宿主会话 register(connectorIds) 白名单注入;空数组=不注入任何连接器 → 依赖 MCP 的任务会「找不到 MCP」。
// 仅放行合法 configId 形态(connector:/custom-mcp:/裸名),其余丢弃
const connArr = Array.isArray(connectorIds)
? connectorIds.filter((c) => typeof c === 'string' && /^[A-Za-z0-9:_-]{1,120}$/.test(c.trim())).map((c) => c.trim())
: [];
// 09-01 参数级技能挂载:skills_json 是 automations 表原生字段(default '[]'),
// 客户端调度器创建会话时按该字段挂载技能;前端传技能名数组(如 ['短视频工作台'])即完成任务级技能绑定
// 09-03 模型跟随用户选择:原 model_id 钉死 'deepseek-v4-flash'(绕过用户模型选择,全任务跑 flash)。
// 已改为动态读取最近活跃手动会话的 model(下方 modelRow/modelId,跟随用户在 UI 的 auto/flash/pro 选择);
// 兜底 'auto' = 系统按任务复杂度自动路由。实证:调度器 model_id 有值 → setSessionModel 强制指定会话模型,
// 'auto' 与具体模型 ID 均被接受(auto 会话 model='auto'、任务 ACCEPTED success=1 验证通过)
const skillsArr = Array.isArray(skills) ? skills.filter((s) => typeof s === 'string' && s.trim()).map((s) => s.trim()) : [];
// 09-01 并发方案A:浏览器类任务(账号信息/视频列表/采集/导入)自动注入浏览器锁约束——
// 多个 AI 会话任务并行共享同一 9223 Chrome 实例,必须互斥使用,否则标签页互相导航抢占
const BROWSER_HINT = /账号信息|视频列表|保存并分析|采集|导入账号|网页采集|浏览器/;
const BROWSER_LOCK = BROWSER_LOCK_FILE;
let finalPrompt = prompt.trim();
if (BROWSER_HINT.test(finalPrompt)) {
finalPrompt += `
【浏览器锁约束(必须遵守)】本任务需要操作浏览器(browser-harness / 9223 Chrome)。多个 AI 会话任务可能并行,浏览器是共享单实例,必须互斥使用:
1. 执行任何浏览器操作前,先检查锁文件:${BROWSER_LOCK} 是否存在(bash: ls)
2. 锁存在 → 等待 10 秒后重试,最多重试 18 次(约 3 分钟);若锁文件修改时间已超过 10 分钟视为死锁,可删除后抢占
3. 拿到锁(bash: echo <任务名+时间戳> > ${BROWSER_LOCK})→ 才可操作浏览器
4. 浏览器操作全部完成后(无论成功失败)必须删除锁文件(bash: rm -f ${BROWSER_LOCK})
5. 若确认本任务实际不需要浏览器(如数据已齐),忽略本条约束,直接跳过`;
}
const { DatabaseSync } = require('node:sqlite');
const db = new DatabaseSync(process.env.WORKBUDDY_DB || path.join(os.homedir(), '.workbuddy', 'workbuddy.db'));
const now = Date.now();
const id = 'automation-' + now;
const cwd = SESSION_CWD; // 工作台触发的会话归入 mcn-workshop 空间分组(由 ROOT 动态推导)
const d = new Date(Date.now() + 5 * 1000); // 未来 5 秒(08-31 根因:客户端只对未来 scheduledAt 补算 next_run_at;写 now=过去时间→不补算→调度器扫不到→卡死)
const pad = (n) => String(n).padStart(2, '0');
const scheduledAt = `${d.getFullYear()}-${pad(d.getMonth() + 1)}-${pad(d.getDate())}T${pad(d.getHours())}:${pad(d.getMinutes())}:${pad(d.getSeconds())}`;
const nextRunAt = Date.now() + 5 * 1000; // 08-31 实测:客户端调度器按 next_run_at 扫描;缺此列=null=永不拾取(工作台任务卡死根因)
const uidRow = db.prepare("SELECT user_id FROM sessions WHERE user_id IS NOT NULL AND user_id <> '' ORDER BY created_at DESC LIMIT 1").get();
const ownerUserId = uidRow ? uidRow.user_id : '';
// 09-03 动态跟随用户模型选择:WorkBuddy 模型为会话级(sessions.model,无全局设置),
// 读最近活跃的【手动】会话(非后台自动化)的 model 写入任务 → 用户在 UI 切 auto/flash/pro,任务即跟随;
// 无手动会话记录时兜底 'auto'(自动路由)。实证:调度接受 'auto' 与具体模型 ID('deepseek-v4-pro' 曾成功)
const modelRow = db.prepare(`SELECT model FROM sessions WHERE deleted_at IS NULL AND COALESCE(is_background_automation,0)=0 AND model IS NOT NULL AND model <> '' ORDER BY COALESCE(last_activity_at, updated_at) DESC LIMIT 1`).get();
const modelId = (modelRow && modelRow.model) ? modelRow.model : 'auto';
// 10-08 修复:与 createOnceAutomation 同理,必须带 workspace_scope='workspace',否则会话落「未分组任务」
db.prepare(`INSERT INTO automations (id,name,prompt,status,schedule_type,scheduled_at,next_run_at,rrule,cwds,created_at,updated_at,skills_json,connector_ids_json,model_id,permission_mode,owner_user_id,owner_status,workspace_scope)
VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)`)
.run(id, name || '工作台任务', finalPrompt, 'ACTIVE', 'once', scheduledAt, nextRunAt, '', JSON.stringify([cwd]), now, now, JSON.stringify(skillsArr), JSON.stringify(connArr), modelId, 'fullAccess', ownerUserId, 'confirmed', 'workspace');
db.close();
log('已提交 AI 任务: ' + id + (name ? ' (' + name + ')' : ''));
return sendJSON(200, { ok: true, id, message: '任务已提交,请在左侧会话栏查看执行' });
} catch (e) {
return sendErr(500, '任务提交失败: ' + e.message);
}
});
return;
}
// 查询 AI 任务执行状态(供前端轮询恢复按钮):queued / pending / running / done / error / notfound
if (p === '/api/run/status' && req.method === 'GET') {
const id = q.get('id');
if (!id) return sendErr(400, '缺少 id');
try {
const { DatabaseSync } = require('node:sqlite');
const db = new DatabaseSync(process.env.WORKBUDDY_DB || path.join(os.homedir(), '.workbuddy', 'workbuddy.db'), { readOnly: true });
const auto = db.prepare('SELECT status, deleted_at, last_run_at, created_at FROM automations WHERE id=?').get(id);
if (!auto || auto.deleted_at) { db.close(); return sendJSON(200, { ok: true, id, state: 'notfound' }); }
const st = db.prepare('SELECT running, last_error FROM automation_runtime_state WHERE automation_id=?').get(id);
const runs = db.prepare("SELECT status FROM automation_runs WHERE automation_id=? ORDER BY created_at DESC LIMIT 1").get(id);
// 10-08:关联本任务拉起的后台会话(任务创建后 5 分钟内第一条自动化会话),读其日志产出进度
const t0 = Number(auto.created_at) || 0;
const sess = t0 ? db.prepare(
'SELECT id, title, cwd FROM sessions WHERE is_background_automation=1 AND created_at >= ? AND created_at <= ? ORDER BY created_at ASC LIMIT 1'
).get(t0 - 15000, t0 + 300000) : null;
const prog = sess ? readSessionProgress(sess.id, sess.cwd) : null;
const elapsedMs = t0 ? Date.now() - t0 : null;
db.close();
let state = 'queued';
if (st && Number(st.running) === 1) state = 'running';
else if (runs) {
const s = String(runs.status || '').toUpperCase();
if (s === 'IN_PROGRESS') state = 'running';
else if (s === 'PENDING_REVIEW') state = 'review'; // 任务停在「待用户确认/已中断」,前端应提示而非假装执行中
else if (['DONE', 'COMPLETED', 'SUCCESS', 'FINISHED', 'ACCEPTED'].includes(s)) state = 'done'; // ACCEPTED=客户端调度器终态(成功)
else if (['ERROR', 'CANCELLED', 'FAILED', 'INTERRUPTED'].includes(s)) state = 'error';
else state = 'pending'; // ACCEPTED / QUEUED / PENDING 等 = 排队等待执行
}
return sendJSON(200, {
ok: true, id, state, error: (st && st.last_error) || null,
elapsedMs,
sessionId: (sess && sess.id) || null,
sessionTitle: (sess && sess.title) || null,
lastAction: (prog && prog.lastAction) || '',
toolCalls: (prog && prog.toolCalls) || {},
updatedAgoMs: (prog && prog.updatedAgoMs) != null ? prog.updatedAgoMs : null,
stalled: !!(prog && prog.stalled),
logBytes: (prog && prog.logBytes) || 0,
});
} catch (e) {
return sendErr(500, '状态查询失败: ' + e.message);
}
}
// ---------- AI 生成选题(Dify chat-messages 即时调用,参考技能 S5 选题方法论) ----------
if (p === '/api/ai/topics' && req.method === 'POST') {
let body = '';
req.on('data', (c) => { body += c; if (body.length > 1e6) req.destroy(); });
req.on('end', async () => {
try {
const { accountName } = JSON.parse(body || '{}');
if (!accountName || typeof accountName !== 'string') return sendErr(400, '缺少账号名');
// 1) 查账号定位 + 人设摘要 + 最新视频解析(工作台库)
let contentText = '', personaText = '', videoRefText = '';
try {
const { DatabaseSync } = require('node:sqlite');
const db = new DatabaseSync(path.join(ROOT, CFG.dbName), { readOnly: true });
const acc = db.prepare('SELECT * FROM hot_accounts WHERE account_name=? LIMIT 1').get(accountName);
if (acc) {
contentText = String(acc.content || '').slice(0, 500);
const p = db.prepare('SELECT content_json, summary FROM account_persona WHERE account_id=? ORDER BY id DESC LIMIT 1').get(acc.id);
if (p) personaText = String(p.summary || p.content_json || '').slice(0, 500);
// 09-01:读该账号最新视频解析(account_video_analysis 按 aweme_id 关联),让选题基于真实视频内容
try {
const vrows = db.prepare(
'SELECT va.content_json AS c FROM account_video_analysis va JOIN account_videos v ON v.aweme_id = va.aweme_id WHERE v.account_id=? ORDER BY va.analysis_time DESC LIMIT 3'
).all(acc.id);
if (vrows.length) videoRefText = vrows.map((x) => String(x.c || '')).join('\n---\n').slice(0, 3000);
} catch (e) { /* 视频解析缺失不阻塞 */ }
}
db.close();
} catch (e) { /* 库不可用不阻塞,仍可生成 */ }
// 2) 读 Dify key:环境变量优先,其次技能脚本内置默认
let apiKey = process.env.DIFY_MCN_CYLG_KEY || '';
if (!apiKey) {
const scriptPath = path.join(ROOT, '..', 'scripts', 'MCN_CYLG_API.py');
try {
const src = fs.readFileSync(scriptPath, 'utf8');
const m = src.match(/DEFAULT_CYLG_KEY\s*=\s*"([^"]+)"/);
if (m) apiKey = m[1];
} catch (e) { /* 脚本不存在则走 env */ }
}
if (!apiKey) return sendErr(500, '未配置 DIFY_MCN_CYLG_KEY(环境变量或技能脚本内置 KEY)');
// 3) 构造选题 query(参考 references/创作流程/5_生成短视频选题.md 方法论;09-01 追加视频解析参考)
const query = '你是短视频选题策划师,请为达人「' + accountName + '」结合账号设定与最新视频内容生成 3 个爆款选题方案。\n\n'
+ '【达人信息】\n' + (contentText || '(无账号定位信息)') + '\n'
+ (personaText ? '【人设摘要】\n' + personaText + '\n' : '')
+ (videoRefText ? '【该账号最近视频解析(选题须参考的真实视频内容)】\n' + videoRefText + '\n' : '')
+ '\n要求:\n'
+ '1. 选题要跳出账号核心行为框架,但保留账号结构性符号(固定机制/口头禅等)\n'
+ '2. 每个选题一句话说清主题+切入角度,标注目标体感(如治愈感/爽感/获得感)\n'
+ '3. 只输出 3 个选题,每行一个,格式:1. 主题(切入角度|目标体感)\n'
+ '4. 不要输出其他解释';
// 4) 调 Dify chat-messages(blocking 即时返回)
const difyUrl = String(process.env.DIFY_URL || 'https://mydify.youmanvideo.com/v1').replace(/\/+$/, '');
const ctrl = new AbortController();
const timer = setTimeout(() => ctrl.abort(), 100000);
let resp;
try {
resp = await fetch(difyUrl + '/chat-messages', {
method: 'POST',
headers: { 'Authorization': 'Bearer ' + apiKey, 'Content-Type': 'application/json' },
body: JSON.stringify({ inputs: { is_think: '0', is_online: '0', model: process.env.DIFY_MCN_MODEL || 'gemini' }, query, response_mode: 'blocking', user: 'mcn-workshop' }),
signal: ctrl.signal,
});
} finally { clearTimeout(timer); }
const data = await resp.json().catch(() => ({}));
if (!resp.ok) return sendErr(502, 'AI 接口错误: ' + (data.message || resp.status));
log('AI 选题生成: ' + accountName);
return dshOk(200, { topics: data.answer || '' });
} catch (e) {
return sendErr(500, '选题生成失败: ' + e.message);
}
});
return;
}
// ---------- AI 链路自检(09-29 新增):看当前走哪条后端、CLI 能否定位、凭据是否就位 ----------
// GET /api/ai/status 只报告配置(不发起模型调用,零消耗)
// GET /api/ai/status?probe=1 额外真跑一次极短提示,验证端到端可用(会消耗一次调用)
if (p === '/api/ai/status' && req.method === 'GET') {
(async () => {
const c = (CFG.ai && CFG.ai.cli) || {};
const backend = aiBackend();
const mask = (s) => {
const t = String(s || '');
if (!t) return '';
return t.length <= 8 ? '***' : t.slice(0, 6) + '***' + t.slice(-4);
};
const cred = resolveCredentials({ apiKey: c.apiKey, authToken: c.authToken, credentialFile: c.credentialFile });
const info = {
backend,
cli: {
resolved: null,
pathConfigured: !!c.path,
model: c.model || '',
internetEnvironment: c.internetEnvironment || '',
credential: cred.source + (cred.apiKey || cred.authToken
? ' (' + mask(cred.authToken || cred.apiKey) + ')' : ''),
credentialFile: cred.credentialFile,
credentialFileExists: (() => { try { return fs.existsSync(cred.credentialFile); } catch (e) { return false; } })(),
ready: !!(cred.authToken || cred.apiKey),
},
dify: {
url: String(process.env.DIFY_URL || 'https://mydify.youmanvideo.com/v1'),
keyConfigured: !!(process.env.DIFY_MCN_CYLG_KEY || ''),
},
};
try {
const r = resolveCli(c.path);
info.cli.resolved = r ? { kind: r.kind, cmd: r.cmd, script: r.script || '' } : null;
} catch (e) { info.cli.resolved = null; }
info.cli.found = !!info.cli.resolved;
const url = new URL(req.url, 'http://localhost');
if (url.searchParams.get('probe') === '1' && backend === 'cli') {
const t0 = Date.now();
try {
const r = await runCli(cliBaseOpts({ prompt: '只回复两个字:成功', timeoutMs: 120000 }));
info.probe = {
ok: r.ok, ms: Date.now() - t0,
text: (r.text || '').slice(0, 100),
error: r.error || null, sessionId: r.sessionId || null,
};
} catch (e) {
info.probe = { ok: false, ms: Date.now() - t0, error: e.message };
}
}
return dshOk(200, info);
})().catch((e) => { try { sendErr(500, 'AI 自检失败: ' + e.message); } catch (e2) {} });
return;
}
// ---------- AI 创作需求打磨(09-02 弹窗内对话:阻塞/流式 两种模式;09-07 v27 新增 SSE 流式输出,前端可逐 chunk 渲染 + 思考占位) ----------
if (p === '/api/ai/clarify' && req.method === 'POST') {
let body = '';
req.on('data', (c) => { body += c; if (body.length > 1e6) req.destroy(); });
req.on('end', async () => {
let payload;
try { payload = JSON.parse(body || '{}'); } catch (e) { return sendErr(400, 'JSON 解析失败'); }
const { accountName, turn, reqMd, history, attach, stream, prevOpts } = payload;
if (!accountName || typeof accountName !== 'string') return sendErr(400, '缺少账号名');
if (!turn || typeof turn !== 'string' || !turn.trim()) return sendErr(400, '缺少用户输入');
// 1) 账号设定摘要(工作台库,与选题接口同源)
let personaText = '', contentText = '';
try {
const { DatabaseSync } = require('node:sqlite');
const db = new DatabaseSync(path.join(ROOT, CFG.dbName), { readOnly: true });
const acc = db.prepare('SELECT * FROM hot_accounts WHERE account_name=? LIMIT 1').get(accountName);
if (acc) {
contentText = String(acc.content || '').slice(0, 300);
const p = db.prepare('SELECT content_json, summary FROM account_persona WHERE account_id=? ORDER BY id DESC LIMIT 1').get(acc.id);
if (p) personaText = String(p.summary || p.content_json || '').slice(0, 900);
}
db.close();
} catch (e) { /* 库不可用不阻塞 */ }
// 2) Dify key:环境变量优先,其次技能脚本内置默认
let apiKey = process.env.DIFY_MCN_CYLG_KEY || '';
if (!apiKey) {
const scriptPath = path.join(ROOT, '..', 'scripts', 'MCN_CYLG_API.py');
try {
const src = fs.readFileSync(scriptPath, 'utf8');
const m = src.match(/DEFAULT_CYLG_KEY\s*=\s*"([^"]+)"/);
if (m) apiKey = m[1];
} catch (e) { /* 脚本不存在则走 env */ }
}
if (!apiKey) return sendErr(500, '未配置 DIFY_MCN_CYLG_KEY(环境变量或技能脚本内置 KEY)');
// 3) 构造打磨 query(无状态:携带当前需求 + 最近对话;prompt 模板读自 prompts/clarify.md)
const his = Array.isArray(history) ? history.filter((h) => h && h.u).slice(-2) : [];
const histText = his.map((h) => '你:' + String(h.a || '') + '\n达人:' + String(h.u || '')).join('\n');
const reqMdText = String(reqMd || '').trim();
const attachText = String(attach || '').trim().slice(0, 20000);
// 09-16:历史只带最近 2 轮,且候选文字在入库时已被剥离 → 模型看不到自己给过什么,
// "重新给几个"会原样重复。故由前端把「已给过的候选」一并送来,这里注入提示词。
const prevOptList = Array.isArray(prevOpts)
? [...new Set(prevOpts.filter((s) => typeof s === 'string' && s.trim()).map((s) => s.trim()))].slice(-24)
: [];
const prevOptsText = prevOptList.length
? '【已给过的候选(共 ' + prevOptList.length + ' 条,禁止原样重复;达人要新的时须换角度/换维度)】\n- ' + prevOptList.join('\n- ')
: '';
const tpl = loadPrompt('clarify');
if (!tpl) return sendErr(500, '缺少 prompts/clarify.md 需求打磨模板');
const query = tpl
.replace('{{accountName}}', accountName)
.replace('{{personaText}}', personaText || '(暂无可用账号设定,按通用编导经验澄清)')
.replace('{{contentText}}', contentText ? '【账号定位】\n' + contentText : '')
.replace('{{attachText}}', attachText ? '【附件素材】(达人上传的需求文档/资料,打磨与最终创作都须充分参考;把其中的关键诉求/约束吸收进需求清单)\n' + attachText : '')
.replace('{{reqMdText}}', reqMdText || '(空:等待用户先给出初步想法)')
.replace('{{histText}}', histText ? '【对话历史】\n' + histText : '')
.replace('{{prevOptsText}}', prevOptsText)
.replace('{{turnText}}', turn.trim());
// ★ 09-29 AI 通路切换(config.json → ai.backend)
// cli = 本地 CodeBuddy CLI 进程,stdio 直连(无端口、无 HTTP;需自备 CODEBUDDY_API_KEY)
// dify = 原链路(下方保留,未改动)
// 两条链路对外都产出同一套 SSE 协议:{delta} / {replace} / {error},前端无需感知差异。
if (aiBackend() === 'cli') {
const base = cliBaseOpts({ prompt: query });
if (stream) {
res.writeHead(200, {
'Content-Type': 'text/event-stream; charset=utf-8',
'Cache-Control': 'no-cache, no-transform',
'Connection': 'keep-alive',
'X-Accel-Buffering': 'no',
});
const write = (obj) => { try { res.write('data: ' + JSON.stringify(obj) + '\n\n'); } catch (e) {} };
let acc = '';
let r;
try {
r = await runCli(Object.assign({}, base, {
onDelta: (t) => { acc += t; write({ delta: t }); },
}));
} catch (e) {
r = { ok: false, text: acc, error: e.message };
}
if (!r.ok) {
// 已产出的内容优先保留(模型中途失败 / 超时场景)
if (acc && r.text && r.text !== acc) write({ replace: r.text });
write({ error: r.error || 'CLI 执行失败' });
} else if (r.text && r.text !== acc) {
// 成功但「流过的内容」≠「终态正文」:典型是模型在联网搜索前先吐了一句过程话术,
// 后端已把工具调用前的部分切掉 → 这里用权威正文整体覆盖(前端原生支持 replace)
write({ replace: r.text });
}
log('AI 需求打磨(cli,stream): ' + accountName + ' | ' + acc.length + '字'
+ (r.sessionId ? ' | session ' + r.sessionId : '') + (r.ok ? '' : ' | 失败: ' + r.error));
res.write('data: [DONE]\n\n');
res.end();
return;
}
const r = await runCli(base);
if (!r.ok) return sendErr(502, r.error || 'CLI 执行失败');
log('AI 需求打磨(cli): ' + accountName + ' | ' + (r.text || '').length + '字');
return dshOk(200, { answer: r.text || '' });
}
const difyUrl = String(process.env.DIFY_URL || 'https://mydify.youmanvideo.com/v1').replace(/\/+$/, '');
// 4) 流式 / 阻塞 分支
if (stream) {
res.writeHead(200, {
'Content-Type': 'text/event-stream; charset=utf-8',
'Cache-Control': 'no-cache, no-transform',
'Connection': 'keep-alive',
'X-Accel-Buffering': 'no',
});
const write = (obj) => { try { res.write('data: ' + JSON.stringify(obj) + '\n\n'); } catch (e) {} };
const ctrl = new AbortController();
const timer = setTimeout(() => ctrl.abort(), 100000);
try {
const difyResp = await fetch(difyUrl + '/chat-messages', {
method: 'POST',
headers: { 'Authorization': 'Bearer ' + apiKey, 'Content-Type': 'application/json' },
body: JSON.stringify({ inputs: { is_think: '0', is_online: '0', model: process.env.DIFY_MCN_MODEL || 'gemini' }, query, response_mode: 'streaming', user: 'mcn-workshop' }),
signal: ctrl.signal,
});
if (!difyResp.ok || !difyResp.body) {
write({ error: 'AI 接口错误: HTTP ' + difyResp.status });
res.write('data: [DONE]\n\n');
return res.end();
}
const reader = difyResp.body.getReader();
const dec = new TextDecoder('utf-8');
let buf = '';
let acc = ''; // 已转发出去的全文(用于"增量/累计"双形态判定 + 终态校准)
let auth = ''; // workflow_finished.data.outputs.answer —— 权威全文
let errMsg = '';
while (true) {
const { value, done } = await reader.read();
if (done) break;
buf += dec.decode(value, { stream: true });
const events = buf.split('\n\n');
buf = events.pop() || '';
for (const ev of events) {
const lines = ev.split('\n');
for (const line of lines) {
if (!line.startsWith('data:')) continue;
const pl = line.slice(5).trim();
if (!pl || pl === '[DONE]') continue;
let json = null;
try { json = JSON.parse(pl); } catch (e) { continue; }
if (json.event === 'message' && typeof json.answer === 'string') {
// ⚠️ 事件形态判定(09-16 修复"回复只显示一半"的根因):
// 本 Dify 应用是**工作流型**(16 节点),message.answer 是**下一段增量**(每片约 32~38 字),
// 不是累计值。旧实现按 `answer.slice(lastAnswer.length)` 做前缀差分 → 除第一片外全部算成空串,
// 答案被截断(实测 335 字只出 44 字)。
// 规则:answer 若以"已产出全文"开头 → 视为累计,取增量;否则视为独立增量,整段追加。
// 该规则对"增量型 / 累计型"两种形态都正确(离线用 11 帧实测:与权威全文逐字一致)。
const a = json.answer;
if (!a) continue;
let inc;
if (acc && a.startsWith(acc)) inc = a.slice(acc.length);
else inc = a;
if (inc) { acc += inc; write({ delta: inc }); }
} else if (json.event === 'message_end') {
const md = json.metadata || {};
if (md.error) errMsg = String(md.error);
} else if (json.event === 'workflow_finished') {
const d = json.data || {};
if (d.outputs && typeof d.outputs.answer === 'string') auth = d.outputs.answer;
if (d.status && d.status !== 'succeeded' && !errMsg) {
errMsg = 'workflow ' + d.status + (d.error ? ':' + d.error : '');
}
} else if (json.event === 'error') {
errMsg = errMsg || (json.message || json.code || 'AI 返回错误');
} else if (json.code && !json.event) {
errMsg = errMsg || (json.message || json.code || 'AI 返回错误');
}
}
}
}
// 终态校准:以权威全文为准。若与流式累计不一致(未来 Dify 改分片策略/丢帧),发 replace 整体覆盖,
// 保证前端拿到的永远是完整答案;只在"不完整"时补发,正常情况零额外流量。
if (auth && auth !== acc) write({ replace: auth });
if (errMsg && !auth) write({ error: errMsg });
log('AI 需求打磨(stream): ' + accountName + ' | 流式' + acc.length + '字' + (auth ? ' / 权威' + auth.length + '字' + (auth === acc ? '(一致)' : '(已 replace 校准)') : '') + (errMsg ? ' | 错误: ' + errMsg : ''));
res.write('data: [DONE]\n\n');
res.end();
} catch (e) {
try { write({ error: e.name === 'AbortError' ? 'AI 响应超时(' + Math.round(100000 / 1000) + 's),已保留已生成内容' : e.message }); res.write('data: [DONE]\n\n'); res.end(); } catch (e) {}
} finally { clearTimeout(timer); }
return;
}
// 阻塞模式(兼容)
const ctrl = new AbortController();
const timer = setTimeout(() => ctrl.abort(), 100000);
let difyResp;
try {
difyResp = await fetch(difyUrl + '/chat-messages', {
method: 'POST',
headers: { 'Authorization': 'Bearer ' + apiKey, 'Content-Type': 'application/json' },
body: JSON.stringify({ inputs: { is_think: '0', is_online: '0', model: process.env.DIFY_MCN_MODEL || 'gemini' }, query, response_mode: 'blocking', user: 'mcn-workshop' }),
signal: ctrl.signal,
});
} finally { clearTimeout(timer); }
const data = await difyResp.json().catch(() => ({}));
if (!difyResp.ok) return sendErr(502, 'AI 接口错误: ' + (data.message || difyResp.status));
log('AI 需求打磨: ' + accountName);
return dshOk(200, { answer: data.answer || '' });
});
return;
}
// 需求附件解析(09-02 新增):txt/docx/pdf → 纯文本(供 AI 打磨与创作上下文)
if (p === '/api/req/parse-attachment' && req.method === 'POST') {
let body = '';
req.on('data', (c) => { body += c; if (body.length > 30e6) req.destroy(); }); // 附件 base64 上限约 22MB 原文件
req.on('end', async () => {
try {
const { name, data } = JSON.parse(body || '{}');
if (!name || typeof name !== 'string' || !data || typeof data !== 'string') return sendErr(400, '缺少文件参数');
const buf = Buffer.from(data, 'base64');
if (buf.length > 20e6) return sendErr(413, '文件过大(上限 20MB)');
const r = parseAttachment(name, buf);
if (r.error) return sendErr(400, r.error);
const text = String(r.text || '').slice(0, 30000); // 打磨上下文上限 3 万字符
return dshOk(200, { text, chars: text.length, warn: r.warn || '', sourceChars: r.chars });
} catch (e) {
return sendErr(500, '附件解析失败: ' + e.message);
}
});
return;
}
// 分析产物同步入库(09-01 新增):扫描账号产出目录 → 写工作台库 mcn-plugin.db
if (p === '/api/import/account' && req.method === 'POST') {
let body = '';
req.on('data', (c) => { body += c; if (body.length > 1e6) req.destroy(); });
req.on('end', async () => {
try {
const { account, root } = JSON.parse(body || '{}');
if (!account || typeof account !== 'string') return sendErr(400, '缺少账号名');
const accountDir = resolveAccountDir(account, root);
if (!accountDir) return sendErr(404, '账号产出目录不存在: ' + account);
const { DatabaseSync } = require('node:sqlite');
const db = new DatabaseSync(path.join(ROOT, CFG.dbName));
const res = importAccountIntoDb(db, account, accountDir);
db.close();
log('分析产物入库: ' + account + ' → ' + JSON.stringify(res));
return dshOk(200, res);
} catch (e) {
return sendErr(500, '入库失败: ' + e.message);
}
});
return;
}
// ---------- 静态文件 ----------
let filePath = path.join(PUBLIC_DIR, p === '/' ? 'index.html' : p);
// 防穿越
if (!filePath.startsWith(PUBLIC_DIR)) return sendErr(403, '禁止访问');
fs.stat(filePath, (err, st) => {
if (err || !st.isFile()) return sendErr(404, 'Not Found');
res.writeHead(200, { 'Content-Type': mimeOf(filePath), 'Content-Length': st.size, 'Cache-Control': 'no-store' });
fs.createReadStream(filePath).pipe(res);
});
});
// ---- 会话自动清理:mcn-workshop 空间超 6 小时无活动的已结束会话自动清除 ----
// 09-01 用户需求:工作台触发(cwd = 工作台目录 SESSION_CWD)的 AI 会话任务,超过 6 小时自动清除
// 09-10 路径治本:cwd 由写死盘符改为 ROOT 动态推导(<技能根>/mcn-workshop);
// 清理模式放宽为 '%mcn-work%',同时覆盖历史会话(早期 D:\AgentSkill\mcn-workshop、迁移期 …\mcn-work-shop)与新建会话(…\V1.0\mcn-workshop)
// 判定:cwd 含 mcn-work + status=completed(执行中一律不删)+ 未删过 + 最后活动时间超过 6 小时
// 方式:软删除(置 deleted_at),客户端左侧会话栏按 deleted_at 过滤即不再显示;可逆安全
const SESSION_TTL_MS = 6 * 3600 * 1000; // 6 小时
const SESSION_CWD_PATTERN = '%mcn-work%'; // 覆盖当前 SESSION_CWD(…\V1.0\mcn-workshop)与历史 cwd(早期 D:\AgentSkill\mcn-workshop、迁移期 …\mcn-work-shop)
const SESSION_CLEAN_INTERVAL_MS = 60 * 60 * 1000; // 每小时检查一次
function cleanupExpiredSessions() {
try {
const { DatabaseSync } = require('node:sqlite');
const db = new DatabaseSync(process.env.WORKBUDDY_DB || path.join(os.homedir(), '.workbuddy', 'workbuddy.db'));
const now = Date.now();
const cutoff = now - SESSION_TTL_MS;
const r = db.prepare(`
UPDATE sessions SET deleted_at = ?
WHERE cwd LIKE ?
AND status = 'completed'
AND deleted_at IS NULL
AND COALESCE(last_activity_at, updated_at, created_at, 0) < ?
`).run(now, SESSION_CWD_PATTERN, cutoff);
db.close();
if (Number(r.changes) > 0) log(`会话自动清理:已清除 mcn-workshop 超 6 小时无活动的已完成会话 ${r.changes} 条`);
} catch (e) {
console.error('会话自动清理失败:', e.message);
}
}
cleanupExpiredSessions();
setInterval(cleanupExpiredSessions, SESSION_CLEAN_INTERVAL_MS);
// ---- 端口自动避让 ----
function listen(port) {
server.once('error', (e) => {
if (e.code === 'EADDRINUSE') {
log('端口 ' + port + ' 被占用,尝试 ' + (port + 1));
listen(port + 1);
} else {
console.error('启动失败:', e.message);
process.exit(1);
}
});
server.listen(port, '127.0.0.1', () => {
log('短视频工作台已启动: http://localhost:' + port);
log('默认项目地址: ' + DEFAULT_ROOT + (fs.existsSync(DEFAULT_ROOT) ? '' : '(不存在)'));
log('按 Ctrl+C 停止服务');
});
}
listen(PORT_BASE);