- 零第三方依赖,Node >=20 原生 ESM - 团队式编排引擎:拆解/路由/执行/审查/合并全真实 LLM - 阶段心跳、单一权威清单守卫、all-keys-failed 如实上报 - H4 会话视图/amend/watchdog 有界重试/产物区 artifacts.json - H5 零依赖三栏控制台
151 lines
6.2 KiB
JavaScript
151 lines
6.2 KiB
JavaScript
/**
|
||
* DSH 桌面端插件(cordis plugin):scheduler-gateway-merged
|
||
* 与 standalone 共享同一中心网关 HTTP API(默认 127.0.0.1:4180)。
|
||
* 注册四个工具:
|
||
* gateway_status 总览(队列/在途/节点/死信)
|
||
* gateway_run 建任务 / 触发真实 LLM 流水线 / 后台起 server
|
||
* gateway_board 任务清单
|
||
* gateway_nodes 节点健康与能力
|
||
* 零运行时依赖(defineTool 由 DSH 宿主 @deepseek-ai/dsh-tools 提供)。
|
||
*
|
||
* 注意:cordis 加载器要求插件入口导出「函数或带 apply 的对象」(unwrapExports 取 default),
|
||
* 所以本文件 default 导出必须是 { name, inject, apply } 形态,工具在 apply(ctx) 里经
|
||
* ctx.tools.register 注册——不能在模块顶层直接注册。
|
||
*/
|
||
import { defineTool } from "@deepseek-ai/dsh-tools";
|
||
import { spawn } from "node:child_process";
|
||
import { fileURLToPath } from "node:url";
|
||
import { dirname, resolve } from "node:path";
|
||
|
||
const HERE = dirname(fileURLToPath(import.meta.url));
|
||
const ROOT = resolve(HERE, "..");
|
||
const GW = (typeof process !== "undefined" && process.env?.GATEWAY_URL) || "http://127.0.0.1:4180";
|
||
|
||
async function api(path, opts) {
|
||
const r = await fetch(GW + path, opts);
|
||
const j = await r.json().catch(() => ({}));
|
||
if (!r.ok) throw new Error(`网关 HTTP ${r.status}: ${JSON.stringify(j).slice(0, 200)}`);
|
||
return j;
|
||
}
|
||
const post = (path, body) =>
|
||
api(path, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify(body || {}) });
|
||
|
||
function textOut(v) {
|
||
return {
|
||
schema: { type: "object", additionalProperties: true, properties: { text: { type: "string" } } },
|
||
render: (_a, v) => [{ type: "text", text: v.text }],
|
||
};
|
||
}
|
||
|
||
const gatewayStatus = defineTool({
|
||
name: "gateway_status",
|
||
description: "查看调度网关总览:队列长度、在途、在线节点与负载、重试/死信计数、任务状态分布。",
|
||
parameters: {},
|
||
output: textOut(),
|
||
async execute() {
|
||
const s = await api("/api/status");
|
||
const text =
|
||
`队列 ${s.queue}|在途 ${s.inFlight}|节点 ${s.nodesOnline}/${s.nodes}|重试 ${s.retry}|死信 ${s.deadLetter}\n` +
|
||
s.nodeHealth
|
||
.map(
|
||
(n) =>
|
||
`- ${n.name || n.nodeId}(${n.nodeId})[${n.kind}] ${n.online ? "🟢 在线" : "⚪ 离线"} 负载${n.load} caps=${n.caps.join("/")}`,
|
||
)
|
||
.join("\n");
|
||
return { ...s, text };
|
||
},
|
||
});
|
||
|
||
const gatewayRun = defineTool({
|
||
name: "gateway_run",
|
||
description:
|
||
"启动网关后台服务(serve),或提交任务(task),或触发真实大模型多智能体流水线(pipeline:planner→跨节点workers→reviewer→merger)。",
|
||
parameters: {
|
||
action: { type: "string", required: true, description: "serve | task | pipeline" },
|
||
goal: { type: "string", description: "task/pipeline 的任务内容" },
|
||
priority: { type: "number", description: "优先级 1-10,默认 5" },
|
||
capabilities: { type: "string", description: "能力标签逗号分隔,如 worker,shell" },
|
||
},
|
||
output: textOut(),
|
||
async execute(args) {
|
||
if (args.action === "serve") {
|
||
const child = spawn(process.execPath, [resolve(ROOT, "src", "cli-live.js"), "server"], {
|
||
detached: true,
|
||
stdio: "ignore",
|
||
env: process.env,
|
||
});
|
||
child.unref();
|
||
return { pid: child.pid, text: `中心网关已后台启动 pid=${child.pid},面板 ${GW}/panel` };
|
||
}
|
||
if (args.action === "pipeline") {
|
||
const r = await post("/api/pipeline", { goal: args.goal });
|
||
return { text: `流水线已提交:${r.run?.runId}\n约 1-2 分钟后可用 gateway_status / gateway_board 查看。` };
|
||
}
|
||
const r = await post("/api/tasks", {
|
||
title: (args.goal || "").slice(0, 40),
|
||
prompt: args.goal,
|
||
priority: args.priority || 5,
|
||
capabilities: String(args.capabilities || "worker").split(",").map((s) => s.trim()).filter(Boolean),
|
||
});
|
||
return { taskId: r.task.id, text: `任务已入队 ${r.task.id}(优先级 ${r.task.priority})` };
|
||
},
|
||
});
|
||
|
||
const gatewayBoard = defineTool({
|
||
name: "gateway_board",
|
||
description: "列出网关上的任务(可按状态过滤),返回编号/标题/状态/节点/尝试/评分。",
|
||
parameters: {
|
||
status: { type: "string", description: "可选:queued/assigned/running/done/dead" },
|
||
},
|
||
output: textOut(),
|
||
async execute(args) {
|
||
const { tasks } = await api("/api/tasks");
|
||
const list = (args.status ? tasks.filter((t) => t.state === args.status) : tasks).slice(-50);
|
||
const lines = list
|
||
.slice()
|
||
.reverse()
|
||
.map((t) => `- ${t.id.slice(-8)} [${t.state}] ${(t.title || "").slice(0, 50)} →${t.nodeId || "-"} 试${t.attempts} 分${t.result?.score ?? "-"}`);
|
||
return { count: list.length, text: `共 ${list.length} 个任务:\n${lines.join("\n") || "(无)"}` };
|
||
},
|
||
});
|
||
|
||
const gatewayNodes = defineTool({
|
||
name: "gateway_nodes",
|
||
description: "查看注册到中心网关的节点:native 直连 / adapter 兼容转化层,及其能力、负载、健康。",
|
||
parameters: {},
|
||
output: textOut(),
|
||
async execute() {
|
||
const { nodes } = await api("/api/nodes");
|
||
const text = nodes
|
||
.map((n) => {
|
||
const nm = n.name || n.meta?.name || n.nodeId; // agent 显示名称,缺省回退 nodeId
|
||
const st = n.online ? "🟢 在线" : "⚪ 离线";
|
||
const mdl = n.model ? ` · ${n.model}` : "";
|
||
return `- ${nm}(${n.nodeId})[${n.kind}] ${st}${mdl} 能力=${n.capabilities.join("/")} 负载${n.inFlight}/${n.maxConcurrency} 完成${n.completed}/失败${n.failed}`;
|
||
})
|
||
.join("\n");
|
||
return { count: nodes.length, text: text || "(暂无节点注册)" };
|
||
},
|
||
});
|
||
|
||
const tools = [gatewayStatus, gatewayRun, gatewayBoard, gatewayNodes];
|
||
|
||
/** 稳定 cordis 插件名。 */
|
||
const name = "scheduler-gateway-merged";
|
||
/** 挂载前需就绪的宿主服务。 */
|
||
const inject = ["tools"];
|
||
|
||
/** 在宿主上下文中注册四个网关工具。 */
|
||
function apply(ctx) {
|
||
ctx.effect(() => {
|
||
const disposers = tools.map((tool) => ctx.tools.register(tool));
|
||
return () => {
|
||
for (const dispose of disposers) dispose();
|
||
};
|
||
}, "scheduler-gateway-merged: tools");
|
||
}
|
||
|
||
const plugin = { name, inject, apply };
|
||
|
||
export { plugin as default, name, inject, apply, tools, gatewayStatus, gatewayRun, gatewayBoard, gatewayNodes };
|