- 零第三方依赖,Node >=20 原生 ESM - 团队式编排引擎:拆解/路由/执行/审查/合并全真实 LLM - 阶段心跳、单一权威清单守卫、all-keys-failed 如实上报 - H4 会话视图/amend/watchdog 有界重试/产物区 artifacts.json - H5 零依赖三栏控制台
247 lines
7.8 KiB
JavaScript
247 lines
7.8 KiB
JavaScript
/**
|
||
* 中心网关状态存储(H3/H6):任务、节点、审计、死信、SSE 事件总线。
|
||
* 内存态 + 防抖落盘 state/server-state.json;所有写操作串行(mutex)。
|
||
*/
|
||
import { mkdirSync, writeFileSync, existsSync, readFileSync } from "node:fs";
|
||
import { resolve, dirname } from "node:path";
|
||
import { fileURLToPath } from "node:url";
|
||
import { TASK_STATE, shortId } from "./protocol.js";
|
||
|
||
const ROOT = resolve(dirname(fileURLToPath(import.meta.url)), "..");
|
||
|
||
class EventBus {
|
||
constructor() {
|
||
this.subs = new Set();
|
||
this.seq = 0;
|
||
}
|
||
subscribe(fn) {
|
||
this.subs.add(fn);
|
||
return () => this.subs.delete(fn);
|
||
}
|
||
emit(type, data) {
|
||
const evt = { seq: ++this.seq, at: Date.now(), type, data };
|
||
for (const fn of this.subs) {
|
||
try {
|
||
fn(evt);
|
||
} catch {
|
||
/* ignore */
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
export class GatewayStore {
|
||
constructor({ persist = true, stateFile, wal = null } = {}) {
|
||
this.tasks = new Map();
|
||
this.wal = wal; // R2 增强:可选 WAL/journal(崩溃恢复);null 表示关闭
|
||
this.nodes = new Map();
|
||
this.audit = [];
|
||
this.deadLetter = [];
|
||
this.bus = new EventBus();
|
||
this.chain = Promise.resolve();
|
||
this.persist = persist;
|
||
this.stateFile = stateFile || resolve(ROOT, "state", "server-state.json");
|
||
this._saveTimer = null;
|
||
this.stats = { assigned: 0, requeued: 0, dead: 0, completed: 0, failed: 0, duplicateClaims: 0 };
|
||
}
|
||
|
||
/** 串行化所有变更,避免并发互踩。 */
|
||
mutate(fn) {
|
||
const run = this.chain.then(() => fn(this));
|
||
this.chain = run.then(
|
||
() => {},
|
||
() => {},
|
||
);
|
||
return run;
|
||
}
|
||
|
||
load() {
|
||
try {
|
||
if (!existsSync(this.stateFile)) return;
|
||
const d = JSON.parse(readFileSync(this.stateFile, "utf8"));
|
||
for (const [k, v] of d.tasks || []) this.tasks.set(k, v);
|
||
for (const [k, v] of d.nodes || []) {
|
||
v.online = false; // 重启后节点需重新注册
|
||
this.nodes.set(k, v);
|
||
}
|
||
this.audit = d.audit || [];
|
||
this.deadLetter = d.deadLetter || [];
|
||
} catch {
|
||
/* 损坏状态不阻塞启动 */
|
||
}
|
||
}
|
||
|
||
scheduleSave() {
|
||
if (!this.persist) return;
|
||
if (this._saveTimer) return;
|
||
this._saveTimer = setTimeout(() => {
|
||
this._saveTimer = null;
|
||
try {
|
||
mkdirSync(dirname(this.stateFile), { recursive: true });
|
||
writeFileSync(
|
||
this.stateFile,
|
||
JSON.stringify(
|
||
{
|
||
tasks: [...this.tasks.entries()],
|
||
nodes: [...this.nodes.entries()],
|
||
audit: this.audit.slice(-500),
|
||
deadLetter: this.deadLetter,
|
||
},
|
||
null,
|
||
2,
|
||
),
|
||
);
|
||
} catch {
|
||
/* ignore */
|
||
}
|
||
}, 200);
|
||
}
|
||
|
||
/** R2 增强:WAL 钩子(先写 journal 再 apply;幂等键去重,重放安全)。 */
|
||
walAppend(type, { taskId = null, key = null, data = {} } = {}) {
|
||
if (!this.wal) return null;
|
||
return this.wal.append(type, { taskId, key, data });
|
||
}
|
||
|
||
addAudit(entry) {
|
||
const rec = { id: shortId("aud"), at: Date.now(), ...entry };
|
||
this.audit.push(rec);
|
||
if (this.audit.length > 1000) this.audit.shift();
|
||
this.bus.emit("audit", rec);
|
||
return rec;
|
||
}
|
||
|
||
// ---- tasks ----
|
||
createTask(task) {
|
||
const id = task.id || shortId("task");
|
||
const rec = {
|
||
id,
|
||
title: task.title || id,
|
||
prompt: task.prompt || "",
|
||
priority: task.priority ?? 5,
|
||
capabilities: task.capabilities || [],
|
||
dependencies: task.dependencies || [],
|
||
state: task.state || TASK_STATE.QUEUED,
|
||
attempts: 0,
|
||
maxAttempts: task.maxAttempts ?? 3,
|
||
requiresApproval: !!task.requiresApproval,
|
||
callbackUrl: task.callbackUrl || null,
|
||
acceptance: task.acceptance || null,
|
||
requiredRole: task.requiredRole || null,
|
||
tags: task.tags || [],
|
||
parentRun: task.parentRun || null,
|
||
result: null,
|
||
nodeId: null,
|
||
createdAt: Date.now(),
|
||
updatedAt: Date.now(),
|
||
history: [],
|
||
};
|
||
this.tasks.set(id, rec);
|
||
this.walAppend("task.create", { taskId: id, key: "create:" + id, data: { title: rec.title, priority: rec.priority } });
|
||
this.addAudit({ kind: "task.create", taskId: id, priority: rec.priority });
|
||
this.bus.emit("task", { taskId: id, state: rec.state });
|
||
this.scheduleSave();
|
||
return rec;
|
||
}
|
||
|
||
updateTask(id, patch, note) {
|
||
const t = this.tasks.get(id);
|
||
if (!t) return null;
|
||
Object.assign(t, patch, { updatedAt: Date.now() });
|
||
if (note) t.history.push({ at: Date.now(), note });
|
||
this.bus.emit("task", { taskId: id, state: t.state });
|
||
this.scheduleSave();
|
||
return t;
|
||
}
|
||
|
||
// ---- nodes ----
|
||
upsertNode(reg) {
|
||
const prev = this.nodes.get(reg.nodeId);
|
||
const existing = this.nodes.get(reg.nodeId);
|
||
const rec = {
|
||
nodeId: reg.nodeId,
|
||
name: reg.name || reg.meta?.name || reg.nodeId,
|
||
kind: reg.kind,
|
||
capabilities: reg.capabilities,
|
||
maxConcurrency: reg.maxConcurrency || 1,
|
||
model: reg.model || null,
|
||
location: reg.location || "local",
|
||
cost: reg.cost ?? 1,
|
||
online: true,
|
||
lastSeen: Date.now(),
|
||
startedAt: existing?.startedAt || Date.now(),
|
||
completed: existing?.completed || 0,
|
||
failed: existing?.failed || 0,
|
||
// 执行健康熔断状态必须跨重注册保留:坏 CLI(如本机 codex/gemini)会周期性自动重注册,
|
||
// 若在此清零,熔断永远无法生效
|
||
failStreak: existing?.failStreak || 0,
|
||
coolUntil: existing?.coolUntil || 0,
|
||
inFlight: 0, // 重新注册=新会话,在途计数清零
|
||
meta: reg.meta || {},
|
||
};
|
||
if (reg.meta && reg.meta.role) rec.role = reg.meta.role;
|
||
if (prev?.editName) rec.name = prev.editName;
|
||
if (prev?.editCaps?.length) rec.capabilities = prev.editCaps;
|
||
if (prev?.editRole) rec.role = prev.editRole;
|
||
this.nodes.set(reg.nodeId, rec);
|
||
this.addAudit({ kind: "node.register", nodeId: reg.nodeId, kindNode: reg.kind });
|
||
this.bus.emit("node", { nodeId: reg.nodeId, online: true });
|
||
this.scheduleSave();
|
||
return rec;
|
||
}
|
||
|
||
heartbeat(nodeId, load) {
|
||
const n = this.nodes.get(nodeId);
|
||
if (!n) return null;
|
||
n.online = true;
|
||
n.lastSeen = Date.now();
|
||
if (load) {
|
||
n.inFlight = load.inFlight ?? n.inFlight;
|
||
n.completed = load.completed ?? n.completed;
|
||
}
|
||
this.scheduleSave();
|
||
return n;
|
||
}
|
||
|
||
markOffline(nodeId) {
|
||
const n = this.nodes.get(nodeId);
|
||
if (n && n.online) {
|
||
n.online = false;
|
||
n.inFlight = 0;
|
||
// 立即把派给该节点、尚未回报的任务重投队列
|
||
for (const t of this.tasks.values()) {
|
||
if (t.nodeId === nodeId && [TASK_STATE.ASSIGNED, TASK_STATE.RUNNING].includes(t.state)) {
|
||
t.state = TASK_STATE.QUEUED;
|
||
t.nodeId = null;
|
||
t.notBefore = Date.now() + 300;
|
||
t.history.push({ at: Date.now(), note: `node ${nodeId} lost → requeue` });
|
||
}
|
||
}
|
||
this.bus.emit("node", { nodeId, online: false });
|
||
this.addAudit({ kind: "node.offline", nodeId });
|
||
this.scheduleSave();
|
||
}
|
||
}
|
||
|
||
removeNode(nodeId) {
|
||
const n = this.nodes.get(nodeId);
|
||
if (!n) return false;
|
||
if (n.online) this.markOffline(nodeId);
|
||
this.nodes.delete(nodeId);
|
||
this.bus.emit("node", { nodeId, removed: true });
|
||
this.addAudit({ kind: "node.removed", nodeId });
|
||
this.scheduleSave();
|
||
return true;
|
||
}
|
||
|
||
snapshot() {
|
||
return {
|
||
tasks: [...this.tasks.values()],
|
||
nodes: [...this.nodes.values()],
|
||
audit: this.audit.slice(-100),
|
||
deadLetter: this.deadLetter,
|
||
stats: this.stats,
|
||
};
|
||
}
|
||
}
|