#!/usr/bin/env node /** * aig-web.mjs — AI 群聊台:一个本地 Web 控制台,用来给群里发号施令。 * * 设计: * · 服务端只做两件事——读总线(直接读 JSONL)+ 调 `aig.mjs` 执行动作(send / ask / wake)。 * 所有业务规则(锁、游标、唤醒通道、深度护栏、模型退档)都留在 aig.mjs 一处,控制台不重复实现。 * · 前端是内嵌的单页(无构建、无依赖):左边群消息,右边成员栏,底部输入框。 * · 「主代理」默认是 dsh:你在输入框里直接说话就是发给它,它有权把活分给 opencode / workbuddy。 * · SSE 长连:新消息、任务进度、任务结束都实时推。 * * 用法: * node aig-web.mjs [--port 3099] [--host 127.0.0.1] [--room main] [--workspace E:\deepseek] */ import http from 'node:http'; import fs from 'node:fs'; import path from 'node:path'; import os from 'node:os'; import { spawn } from 'node:child_process'; import { fileURLToPath } from 'node:url'; const HERE = path.dirname(fileURLToPath(import.meta.url)); const AIG = path.join(HERE, 'aig.mjs'); const UI_FILE = path.join(HERE, 'web', 'ui.html'); const HOME_ROOT = process.env.AIGROUP_HOME || path.join(os.homedir(), '.ai-groups'); const argv = process.argv.slice(2); const flag = (n, d) => { const i = argv.indexOf('--' + n); return i > -1 && argv[i + 1] && !argv[i + 1].startsWith('--') ? argv[i + 1] : d; }; const PORT = Number(flag('port', process.env.AIG_WEB_PORT || 3099)); const HOST = String(flag('host', '127.0.0.1')); const WORKSPACE = path.resolve(String(flag('workspace', process.cwd()))); function currentRoomFromDisk() { try { return fs.readFileSync(path.join(HOME_ROOT, 'current-room'), 'utf8').trim(); } catch { return ''; } } let ROOM = String(flag('room', '')) || currentRoomFromDisk() || 'main'; /* ------------------------------------------------------------ 读总线 */ const roomsRoot = () => path.join(HOME_ROOT, 'rooms'); const roomDir = (room) => path.join(roomsRoot(), room); const readJson = (p, d) => { try { return JSON.parse(fs.readFileSync(p, 'utf8')); } catch { return d; } }; function readMessages(room, since = 0) { let raw = ''; try { raw = fs.readFileSync(path.join(roomDir(room), 'messages.jsonl'), 'utf8'); } catch { return []; } const out = []; for (const line of raw.split('\n')) { const t = line.trim(); if (!t) continue; try { const m = JSON.parse(t); if ((m.seq || 0) > since) out.push(m); } catch { /* 坏行跳过 */ } } return out.sort((a, b) => (a.seq || 0) - (b.seq || 0)); } function readMembers(room) { return readJson(path.join(roomDir(room), 'members.json'), {}); } function readPresence(room, m) { return readJson(path.join(roomDir(room), 'presence', m + '.json'), null); } function readCursor(room, m) { return readJson(path.join(roomDir(room), 'cursors', m + '.json'), { seq: 0 }); } function listRooms() { try { return fs.readdirSync(roomsRoot(), { withFileTypes: true }).filter((e) => e.isDirectory()).map((e) => e.name); } catch { return []; } } function primaryOf(room) { const members = readMembers(room); for (const [name, rec] of Object.entries(members)) if (rec.role === 'primary') return name; return 'dsh'; } /* ------------------------------------------------------------ 调 aig */ function runAig(args, { timeoutMs = 30 * 60 * 1000 } = {}) { return new Promise((resolve) => { const child = spawn(process.execPath, [AIG, ...args], { cwd: WORKSPACE, windowsHide: true, env: { ...process.env, AIGROUP_HOME: HOME_ROOT, AIGROUP_ROOM: ROOM, NO_COLOR: '1' }, }); let out = '', err = ''; const timer = setTimeout(() => { try { child.kill(); } catch {} }, timeoutMs); child.stdout.on('data', (d) => { out += d.toString('utf8'); }); child.stderr.on('data', (d) => { err += d.toString('utf8'); }); child.on('close', (code) => { clearTimeout(timer); resolve({ code, out: out.trim(), err: err.trim() }); }); child.on('error', (e) => { clearTimeout(timer); resolve({ code: -1, out, err: String(e.message) }); }); }); } /* ------------------------------------------------------------ 任务 */ const jobs = new Map(); const runningByMember = new Map(); let jobSeq = 0; const clients = new Set(); /* 自驱模式:默认开。用户发的消息 → 全体参与;agent 之间也可以互相接话。 没有人工放行,所以必须自带刹车:轮次预算(人一说话就重置)+ 每成员冷却 + 正在跑就跳过。 */ const CONFIG_FILE = path.join(HOME_ROOT, 'ui-config.json'); let uiConfig = readJson(CONFIG_FILE, {}); if (typeof uiConfig.auto !== 'boolean') uiConfig.auto = true; // 默认就让大家处于被唤醒状态 if (!Number.isFinite(uiConfig.turnBudget)) uiConfig.turnBudget = 10; // 人一次发言后,自动唤醒的总次数上限 const DEFAULT_TURN_BUDGET = uiConfig.turnBudget; let budgetLeft = uiConfig.turnBudget; let budgetAnnounced = false; const autoHandled = new Set(); const lastWakeAt = new Map(); const COOLDOWN_MS = 2500; function saveUiConfig() { try { fs.mkdirSync(HOME_ROOT, { recursive: true }); fs.writeFileSync(CONFIG_FILE, JSON.stringify(uiConfig, null, 2), 'utf8'); } catch {} } function broadcast(event, data) { const payload = 'event: ' + event + '\ndata: ' + JSON.stringify(data) + '\n\n'; for (const res of clients) { try { res.write(payload); } catch { /* 断了就算了 */ } } } function jobSummary(j) { return { id: j.id, kind: j.kind, member: j.member, note: j.note || null, status: j.status, code: j.code, startedAt: j.startedAt, endedAt: j.endedAt, ms: (j.endedAt || Date.now()) - j.startedAt, tail: j.out.join('').slice(-1200), }; } /** 真的杀干净:Windows 上 aig 自己会再 spawn worker CLI,那是孙进程, * 只 kill 直接子进程会留下孤儿(群里的 opencode 在 #124/#125 指出过这个洞)。 */ function killTree(child) { if (!child || child.killed) return; try { if (process.platform === 'win32' && child.pid) { spawn('taskkill', ['/PID', String(child.pid), '/T', '/F'], { windowsHide: true, stdio: 'ignore' }); } else { child.kill('SIGTERM'); } } catch { /* 忽略 */ } try { child.kill(); } catch { /* 兜底 */ } } function startJob({ kind, member, args, note, deadlineMs }) { const id = 'j' + (++jobSeq); const job = { id, kind, member, args, note, status: 'running', startedAt: Date.now(), endedAt: null, code: null, out: [] }; jobs.set(id, job); if (jobs.size > 40) jobs.delete(jobs.keys().next().value); if (member) runningByMember.set(member, id); broadcast('job', jobSummary(job)); const child = spawn(process.execPath, [AIG, ...args], { cwd: WORKSPACE, windowsHide: true, env: { ...process.env, AIGROUP_HOME: HOME_ROOT, AIGROUP_ROOM: ROOM, NO_COLOR: '1' }, }); job.child = child; // 看门狗:跑太久就掐掉,免得 runningByMember 一直被占(群里 workbuddy 提的那条) const WATCHDOG_MS = Number(deadlineMs || process.env.AIG_WEB_WATCHDOG_MS || 30 * 60 * 1000); job.watchdog = setTimeout(() => { if (job.status !== 'running') return; job.cancelledBy = 'watchdog'; killTree(child); }, WATCHDOG_MS); const push = (buf, stream) => { const text = buf.toString('utf8'); job.out.push(text); if (job.out.length > 2000) job.out.shift(); broadcast('job-chunk', { id, stream, text }); }; child.stdout.on('data', (b) => push(b, 'out')); child.stderr.on('data', (b) => push(b, 'err')); const done = (code) => { if (job.status !== 'running') return; if (job.watchdog) clearTimeout(job.watchdog); job.status = job.cancelledBy ? 'cancelled' : (code === 0 ? 'ok' : 'failed'); job.code = code; job.endedAt = Date.now(); delete job.child; if (job.member && runningByMember.get(job.member) === id) runningByMember.delete(job.member); broadcast('job', jobSummary(job)); pushNewMessages(); // 唤醒产生的回复立刻推给前端 }; child.on('close', done); child.on('error', (e) => { job.out.push(String(e.message)); done(-1); }); return job; } /* ------------------------------------------------------------ 新消息轮询 */ let lastHead = 0; let roomWatcher = null; let pushTimer = null; function headSeq(room) { const msgs = readMessages(room); return msgs.length ? msgs[msgs.length - 1].seq : 0; } /** 合并同一批文件事件,避免一次写入触发多次推送。 */ function schedulePush() { if (pushTimer) return; pushTimer = setTimeout(() => { pushTimer = null; try { pushNewMessages(); } catch { /* 忽略 */ } }, 25); } /** * 用 fs.watch 盯住房间目录:写入 messages.jsonl 立刻推给前端。 * 之前只靠 900ms 轮询,用户自己发的消息要等最多 0.9 秒才出现——这是"点了没反应"的主要来源之一。 */ function watchRoom() { try { roomWatcher?.close(); } catch { /* 忽略 */ } roomWatcher = null; try { roomWatcher = fs.watch(roomDir(ROOM), { persistent: true }, (_ev, file) => { if (!file || String(file).includes('messages')) schedulePush(); }); } catch { /* 房间目录还不存在就等下一次 */ } lastHead = headSeq(ROOM); } function pushNewMessages() { const fresh = readMessages(ROOM, lastHead); if (fresh.length) { lastHead = fresh[fresh.length - 1].seq; broadcast('messages', fresh.map(publicMsg)); maybeAutoWake(fresh); } } /** * 这条消息该叫醒谁(新版语义): * · 显式点名了成员 → 只叫被点名的; * · 广播、或只指名 human → **除了发信人以外的所有成员**(这样 AI 之间才接得上话)。 */ function targetsFor(room, msg) { const all = Object.keys(readMembers(room)).filter((m) => m !== 'human' && m !== 'system' && m !== msg.from); const named = (msg.to || []).filter((x) => all.includes(x)); if (named.length) return named; return all; } /** * 群聊自动接话时的提示词:和"用户点名让主代理干活"是两回事。 * 没有人工监管的自动轮次里,必须禁止自行开新任务/改代码/跑长命令——否则几个人会互相 * 接出一串工程议程,把时间和额度烧穿(第一次自动对话实测就发生了:dsh 单轮跑了 7 分钟)。 */ function chatExtra(room) { return [ '这是群聊里的自动接话(不是用户在派活):用一两句话表态或补充即可。', '不要在新任务上自行开工,不要改任何文件、不要跑长命令、不要贴日志;', '如果你认为某件事该做,就在群里提出建议并等 human 表态,然后结束本轮。', ].join(' '); } /** 自动唤醒的单轮上限:到点就掐,成员锁必须放掉。 */ function autoWakeTimeoutSec() { return Number(uiConfig.autoWakeTimeoutSec || 240); } function startWakeFor(member, why, { fromHuman = false } = {}) { const timeout = autoWakeTimeoutSec(); const args = ['wake', member, '--timeout', String(timeout)]; if (fromHuman && member === primaryOf(ROOM)) args.push('--prompt', primaryExtra(ROOM)); else if (!fromHuman) args.push('--prompt', chatExtra(ROOM)); return startJob({ kind: 'wake', member, args, note: why, deadlineMs: (timeout + 30) * 1000 }); } /** 预算用尽时只播报一次,别刷屏。 */ function announceBudget() { if (budgetAnnounced) return; budgetAnnounced = true; broadcast('budget', { left: 0, exhausted: true }); } function maybeAutoWake(fresh) { if (!uiConfig.auto) return; for (const m of fresh) { if (m.from === 'system' || m.kind === 'dispatch-error') continue; if ((m.tags || []).includes('silent')) continue; // 静默消息不触发自驱 if (m.from === 'human') { // 用户一说话:预算重置,全群重新被唤醒 budgetLeft = uiConfig.turnBudget; budgetAnnounced = false; broadcast('budget', { left: budgetLeft, reset: true }); } for (const t of targetsFor(ROOM, m)) { const key = t + ':' + m.seq; if (autoHandled.has(key)) continue; autoHandled.add(key); if (runningByMember.has(t)) continue; // 正在跑:唤醒本来就按游标取未读,跑完会带上新的 if (Date.now() - (lastWakeAt.get(t) || 0) < COOLDOWN_MS) continue; // 冷却,防抖 if (budgetLeft <= 0) { announceBudget(); continue; } // 预算闸:没有人工监管时的唯一刹车 budgetLeft--; lastWakeAt.set(t, Date.now()); broadcast('budget', { left: budgetLeft }); broadcast('auto', { member: t, seq: m.seq, at: Date.now() }); startWakeFor(t, 'auto:#' + m.seq + ' ' + String(m.text).slice(0, 60), { fromHuman: m.from === 'human' }); } } } function publicMsg(m) { return { seq: m.seq, ts: m.ts, from: m.from, to: m.to || [], kind: m.kind, via: m.via, text: m.text, reply_to: m.reply_to, meta: m.meta || {}, files: m.files || [], }; } // 轮询只作兜底(fs.watch 有时会漏事件);真正的即时性靠上面的文件监听。 setInterval(() => { try { pushNewMessages(); broadcast('tick', { head: lastHead, jobs: jobs.size }); } catch { /* 忽略 */ } }, 500).unref(); /* ------------------------------------------------------------ 主代理提示 */ function primaryExtra(room) { const p = primaryOf(room); return [ '用户是在 Web 群聊台里直接对你(主代理)说的,成员名是 human —— 这条是用户本人的真指令。', '请把这件事办掉:能直接给结论就给结论;需要动手就用你手上的工具去做;需要别人参与就用', '`node "' + AIG + '" ask <成员> "具体问题" --as ' + p + ' --room ' + room + '` 真的把他们叫起来,', '最后用一两句话回群说明:你做了什么、结论是什么、有哪些没做完。', ].join(' '); } /* ------------------------------------------------------------ HTTP */ function sendJson(res, code, obj) { const body = JSON.stringify(obj); res.writeHead(code, { 'content-type': 'application/json; charset=utf-8', 'cache-control': 'no-store' }); res.end(body); } function readBody(req) { return new Promise((resolve) => { let raw = ''; req.on('data', (d) => { raw += d; if (raw.length > 2 * 1024 * 1024) req.destroy(); }); req.on('end', () => { try { resolve(raw ? JSON.parse(raw) : {}); } catch { resolve({}); } }); }); } let uiCache = { mtime: 0, html: '' }; /** 页面主体在 scripts/web/ui.html(方便单独改、单独做语法检查);读不到才用内嵌的兜底页。 */ function pageHtml() { try { const st = fs.statSync(UI_FILE); if (st.mtimeMs !== uiCache.mtime) uiCache = { mtime: st.mtimeMs, html: fs.readFileSync(UI_FILE, 'utf8') }; return uiCache.html; } catch { return PAGE; } } const server = http.createServer(async (req, res) => { const url = new URL(req.url, 'http://localhost'); const p = url.pathname; if (p === '/' || p === '/index.html') { res.writeHead(200, { 'content-type': 'text/html; charset=utf-8', 'cache-control': 'no-store' }); return res.end(pageHtml()); } if (p === '/api/state') { const members = readMembers(ROOM); const primary = primaryOf(ROOM); const now = Date.now(); const rows = Object.entries(members).map(([name, rec]) => { const pr = readPresence(ROOM, name); const age = pr?.lastSeen ? Math.round((now - Date.parse(pr.lastSeen)) / 1000) : null; return { member: name, kind: rec.kind, model: rec.model || null, fallback: rec.fallback_model || null, role: rec.role || null, primary: rec.role === 'primary' || name === primary, note: rec.note || null, cursor: readCursor(ROOM, name).seq || 0, lastSeen: pr?.lastSeen || null, ageSec: age, online: age !== null && age < 300, }; }); const head = headSeq(ROOM); if (head > lastHead && lastHead === 0) lastHead = 0; return sendJson(res, 200, { room: ROOM, rooms: listRooms(), home: HOME_ROOT, roomDir: roomDir(ROOM), workspace: WORKSPACE, primary: primaryOf(ROOM), head, members: rows, auto: uiConfig.auto, budgetLeft, turnBudget: uiConfig.turnBudget, running: [...runningByMember.keys()], jobs: [...jobs.values()].slice(-6).map(jobSummary), }); } if (p === '/api/messages') { const since = Number(url.searchParams.get('since') || 0); const limit = Number(url.searchParams.get('limit') || 300); const all = readMessages(ROOM, since); return sendJson(res, 200, { messages: all.slice(-limit).map(publicMsg), head: headSeq(ROOM) }); } if (p === '/api/stream') { res.writeHead(200, { 'content-type': 'text/event-stream; charset=utf-8', 'cache-control': 'no-cache, no-transform', connection: 'keep-alive', }); res.write('retry: 2000\n\n'); clients.add(res); res.write('event: hello\ndata: ' + JSON.stringify({ room: ROOM, head: lastHead }) + '\n\n'); const ka = setInterval(() => { try { res.write(': keepalive\n\n'); } catch {} }, 15000); req.on('close', () => { clearInterval(ka); clients.delete(res); }); return; } if (p === '/api/send' && req.method === 'POST') { const body = await readBody(req); const text = String(body.text || '').trim(); if (!text) return sendJson(res, 400, { error: 'empty' }); const args = ['send', text, '--as', 'human']; if (Array.isArray(body.to) && body.to.length) args.push('--to', body.to.join(',')); if (body.silent) args.push('--tag', 'silent'); // 静默消息不触发自驱 const r = await runAig(args, { timeoutMs: 30000 }); pushNewMessages(); return sendJson(res, r.code === 0 ? 200 : 500, { result: r }); } if (p === '/api/ask' && req.method === 'POST') { const body = await readBody(req); const text = String(body.text || '').trim(); const target = String(body.target || primaryOf(ROOM)); if (!text) return sendJson(res, 400, { error: 'empty' }); const isPrimary = target === primaryOf(ROOM); const args = ['ask', target, text, '--as', 'human', '--timeout', String(Number(body.timeout || 1800))]; if (isPrimary) args.push('--prompt', primaryExtra(ROOM)); const job = startJob({ kind: 'ask', member: target, args, note: text.slice(0, 80) }); return sendJson(res, 200, { job: jobSummary(job) }); } if (p === '/api/wake' && req.method === 'POST') { const body = await readBody(req); const target = String(body.target || ''); if (!target) return sendJson(res, 400, { error: 'no-target' }); const timeout = Number(body.timeout || autoWakeTimeoutSec()); const args = ['wake', target, '--timeout', String(timeout)]; if (target === primaryOf(ROOM)) args.push('--prompt', primaryExtra(ROOM)); const job = startJob({ kind: 'wake', member: target, args, deadlineMs: (timeout + 30) * 1000 }); return sendJson(res, 200, { job: jobSummary(job) }); } if (p === '/api/primary' && req.method === 'POST') { const body = await readBody(req); const target = String(body.member || ''); const members = readMembers(ROOM); if (!members[target]) return sendJson(res, 400, { error: 'unknown-member' }); for (const [name, rec] of Object.entries(members)) { if (name === target) rec.role = 'primary'; else if (rec.role === 'primary') delete rec.role; } fs.writeFileSync(path.join(roomDir(ROOM), 'members.json'), JSON.stringify(members, null, 2), 'utf8'); broadcast('state', { primary: target }); return sendJson(res, 200, { primary: target }); } if (p === '/api/room' && req.method === 'POST') { const body = await readBody(req); const target = String(body.room || ''); if (!listRooms().includes(target)) return sendJson(res, 400, { error: 'unknown-room' }); ROOM = target; // 切房间后从"当前队尾"开始跟踪,否则整段历史会被当成新消息、把所有人重新唤醒一遍 autoHandled.clear(); watchRoom(); return sendJson(res, 200, { room: ROOM }); } if (p === '/api/jobs') { return sendJson(res, 200, { jobs: [...jobs.values()].map(jobSummary) }); } if (p === '/api/cancel' && req.method === 'POST') { const body = await readBody(req); const id = String(body.id || ''); const target = String(body.member || ''); const list = [...jobs.values()].filter((j) => j.status === 'running' && ((id && j.id === id) || (target && j.member === target))); for (const j of list) { j.cancelledBy = 'user'; killTree(j.child); // 兜底:万一子进程没立刻退,先把状态放掉,别让它一直占着成员 setTimeout(() => { if (j.status === 'running') { j.status = 'cancelled'; j.endedAt = Date.now(); if (j.member) runningByMember.delete(j.member); broadcast('job', jobSummary(j)); } }, 3000); } broadcast('state', { running: [...runningByMember.keys()] }); return sendJson(res, 200, { cancelled: list.map((j) => j.id) }); } if (p === '/api/auto' && req.method === 'POST') { const body = await readBody(req); uiConfig.auto = body.enabled === undefined ? !uiConfig.auto : Boolean(body.enabled); if (Number.isFinite(Number(body.turnBudget)) && Number(body.turnBudget) > 0) { uiConfig.turnBudget = Number(body.turnBudget); budgetLeft = uiConfig.turnBudget; budgetAnnounced = false; } if (uiConfig.auto) { budgetLeft = uiConfig.turnBudget; budgetAnnounced = false; } saveUiConfig(); broadcast('state', { auto: uiConfig.auto, budgetLeft, turnBudget: uiConfig.turnBudget }); return sendJson(res, 200, { auto: uiConfig.auto, budgetLeft, turnBudget: uiConfig.turnBudget }); } res.writeHead(404, { 'content-type': 'text/plain; charset=utf-8' }); res.end('404'); }); /* ------------------------------------------------------------ 前端页面(兜底) */ /* 正常走 scripts/web/ui.html;这份内嵌页只在那个文件被删/读不到时兜底, 是早期版本,点击手感已被前端重构取代,不要用它来改 UI。 */ const PAGE = String.raw`