scheduler-gateway/lib/index.js
Liuxinyu176 bbb364cd69 feat: DSH 调度网关(互联网关)v6.0.0 — 四模型合并版 R1+R2+R3
- 零第三方依赖,Node >=20 原生 ESM
- 团队式编排引擎:拆解/路由/执行/审查/合并全真实 LLM
- 阶段心跳、单一权威清单守卫、all-keys-failed 如实上报
- H4 会话视图/amend/watchdog 有界重试/产物区 artifacts.json
- H5 零依赖三栏控制台
2026-10-09 23:20:27 +08:00

151 lines
6.2 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 桌面端插件(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 };