// orchestrate.mjs —— 工单批处理小图:路由 → 扇出 → 汇合 → 评审回路 → 汇报
// 零依赖,node orchestrate.mjs 直接跑。模型客户端是按固定队列出牌的桩。
import fs from "node:fs";
import path from "node:path";
// ============ 0. 常量与目录 ============
const MODEL = "claude-sonnet-5";
const MAX_TURNS = 6; // 单个节点内部的循环上界(本系列第 7 门课 阀1)
const POOL_SIZE = Math.max(1, Number(process.env.POOL_SIZE) || 2); // 扇出并发上界(第 3 课);0/非法值兜底为 1
const STUB_LATENCY_MS = 60; // 桩的固定延迟,替代真实网络往返,好让耗时列有东西可看
const MAX_REVIEW_ROUNDS = 3; // 评审回路的最大重写轮数(第 5 课)
const FILLER_WORDS = ["稍等", "请耐心等待", "尽快处理"];
const CATEGORIES = ["billing", "bug", "other"];
const ROOT = process.cwd();
const INBOX = path.join(ROOT, "inbox");
const OUT = path.join(ROOT, "out");
const KB = path.join(ROOT, "kb");
const STATE_PATH = path.join(ROOT, "run-state.json");
const LOG_PATH = path.join(ROOT, "run.jsonl");
// ============ 1. 输入:inbox/ 里的 6 条工单与一份已知问题库 ============
const TICKET_TEXT = {
"T-1001": "订单 A-77301 这个月被扣了两次款,麻烦查一下,多扣的那笔退给我。",
"T-1002": "在报表页点「导出 CSV」,按钮一直转圈,等了五分钟也没反应。Chrome,公司网络。",
"T-1003": "你们的人工客服电话是多少?我想直接打电话问。",
"T-1004": "订单 A-77420 的发票抬头开错了,开成了我的个人名字,要改成公司抬头。",
"T-1005": "手机 App 上登录之后头像一直不显示,网页端是正常的。",
"T-1006": "用了三个月,问题提了好几次都没下文,这产品到底还有没有人维护?",
};
const KNOWN_ISSUES = [
"## KI-88 报表页导出 CSV 无响应",
"影响:点击导出后按钮持续转圈,后台导出队列积压。状态:已在 3.4.2 修复,等待发版。",
"临时办法:改用同一页的「导出 XLSX」,数据列完全一致。",
"",
"## KI-91 移动端头像不显示",
"影响:App 端头像 URL 仍指向旧 CDN 域名,网页端不受影响。状态:修复中,预计本周五随版本发布。",
"临时办法:退出登录后重新登录一次,头像通常会恢复显示。",
].join("\n");
function seedWorkspace() {
fs.mkdirSync(INBOX, { recursive: true });
fs.mkdirSync(OUT, { recursive: true });
fs.mkdirSync(KB, { recursive: true });
for (const [id, text] of Object.entries(TICKET_TEXT)) {
fs.writeFileSync(path.join(INBOX, `${id}.txt`), `${text}\n`);
}
fs.writeFileSync(path.join(KB, "known-issues.md"), `${KNOWN_ISSUES}\n`);
}
function loadInbox() {
return fs
.readdirSync(INBOX)
.filter((f) => f.endsWith(".txt"))
.sort()
.map((f) => ({
id: path.basename(f, ".txt"),
text: fs.readFileSync(path.join(INBOX, f), "utf8").trim(),
}));
}
// ============ 2. 桩 client:SCRIPTS 按工单 id 钉死每次回复 ============
const say = (text) => ({ type: "text", text });
const call = (id, name, input) => ({ type: "tool_use", id, name, input });
const turn = (stop_reason, content, inTok, outTok) => ({
stop_reason,
content,
usage: { input_tokens: inTok, output_tokens: outTok },
});
const SCRIPTS = {
// 路由节点:一次调用分完六条。注意最后一行故意不是合法标签,由代码收紧。
router: [
turn(
"end_turn",
[
say(
[
"T-1001: billing",
"T-1002: bug",
"T-1003: other",
"T-1004: billing",
"T-1005: bug",
"T-1006: 投诉",
].join("\n")
),
],
624,
96
),
],
// billing 工人:查订单后一次写清,第一版就过 gate。
"T-1001#1": [
turn(
"tool_use",
[say("先查订单流水。"), call("toolu_a1", "lookup_order", { order_id: "A-77301" })],
781,
92
),
turn(
"end_turn",
[
say(
"工单 T-1001 回复:订单 A-77301 本月确实扣款 2 次,每笔 399.00 元,系统标记为重复扣款。" +
"多扣的那一笔已提交退款,按原路退回,到账时间取决于发卡行,一般 3-5 个工作日。" +
"如果第 5 个工作日仍未到账,请在这条工单里回复,我们会带流水号找支付渠道核查。"
),
],
1124,
168
),
],
// bug 工人:查已知问题库后回复,第一版就过 gate。
"T-1002#1": [
turn(
"tool_use",
[say("先查已知问题库。"), call("toolu_b1", "read_file", { path: "kb/known-issues.md" })],
763,
88
),
turn(
"end_turn",
[
say(
"工单 T-1002 回复:你遇到的是已知问题 KI-88——报表页导出 CSV 时后台队列积压,按钮会一直转圈。" +
"修复已合入 3.4.2,随下个版本发布后自动生效。在那之前可以改用同一页的「导出 XLSX」,数据列完全一致。" +
"版本发布当天我们会在这条工单里通知你。"
),
],
1312,
176
),
],
// billing 工人:第一版带敷衍词,第二版还是同一句敷衍词 —— 两轮 gate 报告完全相同。
"T-1004#1": [
turn(
"tool_use",
[say("先确认这张发票的抬头。"), call("toolu_c1", "lookup_order", { order_id: "A-77420" })],
786,
90
),
turn(
"end_turn",
[
say(
"工单 T-1004:发票抬头变更需要财务复核,订单 A-77420 的申请我已经提上去了,请您稍等。"
),
],
1133,
96
),
],
"T-1004#2": [
turn(
"end_turn",
[
say(
"工单 T-1004:抬头从个人改成公司,需要财务在系统里作废重开一次发票,申请还排在复核队列里,请您稍等,复核完成后我会把新发票发到这条工单。"
),
],
1291,
104
),
],
// bug 工人:第一版漏了工单号,第二版按 gate 报告补上。
"T-1005#1": [
turn(
"tool_use",
[say("查一下这是不是已知问题。"), call("toolu_d1", "read_file", { path: "kb/known-issues.md" })],
752,
86
),
turn(
"end_turn",
[
say(
"这是已知问题 KI-91——App 端头像走的还是旧 CDN 域名,网页端不受影响,修复预计本周五随版本发布。" +
"在那之前退出登录再重新登录一次,头像通常会恢复显示。"
),
],
1298,
158
),
],
"T-1005#2": [
turn(
"end_turn",
[
say(
"工单 T-1005 回复:这是已知问题 KI-91——App 端头像走的还是旧 CDN 域名,网页端不受影响,修复预计本周五随版本发布。" +
"在那之前退出登录再重新登录一次,头像通常会恢复显示。发版后如果仍是空白,请在这条工单里补一张截图,我们再查你的账号。"
),
],
1466,
172
),
],
};
function makeStubClient(queue) {
let i = 0;
return {
messages: {
async create(req) {
if (!req.model || !req.max_tokens) {
throw new Error("桩 client:create 必须带 model 与 max_tokens");
}
if (i >= queue.length) {
throw new Error(`桩队列耗尽:第 ${i + 1} 次请求没有预设响应`);
}
await new Promise((r) => setTimeout(r, STUB_LATENCY_MS));
return queue[i++];
},
},
};
}
// 计量包在客户端外面,循环内部一行不动。
function metered(client) {
const meter = { calls: 0, tokens: 0 };
const wrapped = {
messages: {
async create(req) {
const res = await client.messages.create(req);
meter.calls += 1;
meter.tokens += res.usage.input_tokens + res.usage.output_tokens;
return res;
},
},
};
return { client: wrapped, meter };
}
// ============ 3. 节点内部:本系列第 7 门课的那个循环,原样搬来 ============
async function runAgent(client, system, userInput, tools, toolImpls) {
const messages = [{ role: "user", content: userInput }];
let turns = 0;
let response = await client.messages.create({
model: MODEL,
max_tokens: 1024,
system,
tools,
messages,
});
while (response.stop_reason === "tool_use") {
// —— 阀1:最大轮次。放在循环体最前、turns++ 之前 ——
if (turns >= MAX_TURNS) {
return `已达最大轮次 ${MAX_TURNS},主动收手(可能任务过难或模型卡住)`;
}
turns++;
// 把模型这一轮的完整响应(assistant 角色)追加进历史
messages.push({ role: "assistant", content: response.content });
// 执行这一轮所有 tool_use 块,各自打包成 tool_result
const toolResults = await runToolUses(response.content, toolImpls);
// 一轮里的所有 tool_result 放进紧随其后的同一条 user 消息
messages.push({ role: "user", content: toolResults });
// 带着变长后的历史再发一次,循环回到 while 判断
response = await client.messages.create({
model: MODEL,
max_tokens: 1024,
system,
tools,
messages,
});
}
// stop_reason 不再是 tool_use,取出最终文字返回
return response.content.find((b) => b.type === "text")?.text ?? "";
}
async function runToolUses(content, toolImpls) {
const toolUseBlocks = content.filter((b) => b.type === "tool_use");
return Promise.all(
toolUseBlocks.map(async (block) => {
const impl = toolImpls[block.name];
try {
const output = await impl(block.input);
return {
type: "tool_result",
tool_use_id: block.id,
content: output,
};
} catch (err) {
return {
type: "tool_result",
tool_use_id: block.id,
content: `工具执行出错: ${err.message}`,
is_error: true,
};
}
})
);
}
// ============ 4. 两个工具 ============
const ORDERS = {
"A-77301": { order_id: "A-77301", amount_cents: 39900, charged_times: 2, status: "duplicate_charge", invoice_title: "李明(个人)" },
"A-77420": { order_id: "A-77420", amount_cents: 128000, charged_times: 1, status: "paid", invoice_title: "李明(个人)" },
};
const TOOLS = [
{
name: "read_file",
description: "读取工作目录下的一个文本文件,用于查已知问题库或工单原文。",
input_schema: {
type: "object",
properties: { path: { type: "string", description: "相对工作目录的路径" } },
required: ["path"],
},
},
{
name: "lookup_order",
description: "按订单号查账单事实:金额、扣款次数、状态、发票抬头。",
input_schema: {
type: "object",
properties: { order_id: { type: "string", description: "形如 A-77301 的订单号" } },
required: ["order_id"],
},
},
];
const toolImpls = {
read_file({ path: rel }) {
const full = path.resolve(ROOT, rel);
if (!full.startsWith(ROOT)) throw new Error("越界路径");
return fs.readFileSync(full, "utf8");
},
lookup_order({ order_id }) {
const row = ORDERS[order_id];
if (!row) throw new Error(`查无此订单:${order_id}`);
return JSON.stringify(row);
},
};
// ============ 5. 三份派工提示词:目标 / 输出格式 / 工具指引 / 任务边界 ============
const ROUTER_PROMPT = [
"你是客服工单分拣员。",
"目标:把下面每一条工单归到 billing(账单、扣款、发票、退款)、bug(功能故障)、other(其余)三类之一。",
"输出格式:每行一条,格式严格为「工单号: 类别」,类别只能是 billing / bug / other,不要写理由,不要输出别的内容。",
"工具指引:这一步不给你任何工具,只看工单文本判断,不要声称查过系统。",
"任务边界:只分类,不写回复、不下结论、不合并工单;拿不准就归 other。",
].join("\n");
const WORKER_PROMPTS = {
billing: [
"你是账单工单专员,一次只处理一条工单。",
"目标:查清这条工单的账单事实,给出一次把话说完的中文回复。",
"输出格式:一段纯文本,开头写「工单 <工单号> 回复:」,依次交代查到的事实、已经做的处理、用户下一步能预期什么;不用列表,不写寒暄。",
"工具指引:账单事实一律用 lookup_order 查,工单里的订单号原样传入;查不到就如实说查不到,不要从工单描述里推断金额或扣款次数。",
"任务边界:只处理这条工单的账单部分,不修改订单、不承诺额外补偿、不回答与账单无关的问题;不要写「稍等」「请耐心等待」「尽快处理」这类没有信息量的话。",
].join("\n"),
bug: [
"你是故障工单专员,一次只处理一条工单。",
"目标:判断这条工单是不是已知问题,给出一次把话说完的中文回复。",
"输出格式:一段纯文本,开头写「工单 <工单号> 回复:」,依次交代命中的已知问题编号与结论、临时办法、修复什么时候到;不用列表,不写寒暄。",
"工具指引:用 read_file 读 kb/known-issues.md 对照,命中就引用其中的编号;没命中就说没命中,不要自己编一个问题编号。",
"任务边界:只做问题定位与回复,不下线功能、不承诺具体到分钟的修复时间、不索要账号密码;不要写「稍等」「请耐心等待」「尽快处理」这类没有信息量的话。",
].join("\n"),
};
// other 类不进模型:一段纯代码模板。不是每个节点都得是模型。
const otherTemplate = (id) =>
`工单 ${id} 已收到。这条工单不涉及账单,也不是功能故障,已经转给客服组人工跟进:` +
`工作日 9:00-18:00 可拨打 400-000-1234 直接沟通,也可以在这条工单里补充信息,回复都会记在这条工单下。`;
// ============ 6. 观测:JSONL 结构化日志 + run-state.json 逐步留痕 ============
const RUN_ID = `run-${Date.now().toString(36)}`;
function initLog() {
fs.writeFileSync(LOG_PATH, "");
}
function log(fields) {
const line = { ts: new Date().toISOString(), run_id: RUN_ID, ...fields };
fs.appendFileSync(LOG_PATH, `${JSON.stringify(line)}\n`);
}
const state = {
version: 1,
run_id: RUN_ID,
started_at: new Date().toISOString(),
updated_at: null,
nodes: {},
tickets: {},
};
// 原子写:先写 .tmp 再 rename(本系列第 9 门课的家法)
function saveState() {
state.updated_at = new Date().toISOString();
const tmp = `${STATE_PATH}.tmp`;
fs.writeFileSync(tmp, JSON.stringify(state, null, 2));
fs.renameSync(tmp, STATE_PATH);
}
async function node(name, fn) {
const t0 = Date.now();
log({ node: name, event: "node_start" });
const result = await fn();
const ms = Date.now() - t0;
state.nodes[name] = {
ms,
calls: result.calls ?? 0,
tokens: result.tokens ?? 0,
status: result.status ?? "ok",
};
saveState(); // 每个节点跑完留一次痕
log({ node: name, event: "node_end", ms, calls: result.calls ?? 0, tokens: result.tokens ?? 0 });
return result;
}
// ============ 7. 节点一:路由 ============
async function routeNode(tickets) {
const { client, meter } = metered(makeStubClient(SCRIPTS.router));
const input = tickets.map((t) => `${t.id}: ${t.text}`).join("\n");
const text = await runAgent(client, ROUTER_PROMPT, input, [], {});
// 输出收紧:只认「工单号: 类别」这一种行,类别不在白名单里的一律落到 other
const parsed = new Map();
for (const line of text.split("\n")) {
const m = line.match(/^\s*(T-\d+)\s*:\s*(\S+)\s*$/);
if (!m) continue;
const [, id, raw] = m;
const category = CATEGORIES.includes(raw) ? raw : "other";
if (category !== raw) log({ node: "route", event: "clamped", ticket: id, raw, category });
parsed.set(id, category);
}
const routed = tickets.map((t) => ({ ...t, category: parsed.get(t.id) ?? "other" }));
for (const t of routed) log({ node: "route", event: "routed", ticket: t.id, category: t.category });
return { routed, calls: meter.calls, tokens: meter.tokens };
}
// ============ 8. 节点二:扇出(分段 + 并发池上界)============
async function runPool(items, limit, worker) {
const results = new Array(items.length);
let next = 0;
const runners = Array.from({ length: Math.min(limit, items.length) }, async () => {
while (next < items.length) {
const i = next;
next += 1;
results[i] = await worker(items[i]);
}
});
await Promise.all(runners);
return results;
}
async function callWorker(ticket, round, extra) {
const key = `${ticket.id}#${round}`;
const queue = SCRIPTS[key];
if (!queue) throw new Error(`桩脚本缺失:${key}`);
const { client, meter } = metered(makeStubClient(queue));
const input = extra
? [
`下面是你上一版对工单 ${ticket.id} 的回复:`,
"---",
extra.prev,
"---",
`确定性检查没有通过,报告:${extra.report}`,
"只修报告里点名的问题,重写一版完整回复。",
].join("\n")
: `工单号 ${ticket.id}\n用户原文:${ticket.text}`;
const text = await runAgent(client, WORKER_PROMPTS[ticket.category], input, TOOLS, toolImpls);
log({ node: extra ? "review" : "fanout", event: "worker_done", ticket: ticket.id, round, calls: meter.calls, tokens: meter.tokens });
return { text, calls: meter.calls, tokens: meter.tokens };
}
async function fanoutNode(routed) {
let calls = 0;
let tokens = 0;
const drafts = await runPool(routed, POOL_SIZE, async (ticket) => {
if (ticket.category === "other") {
const text = otherTemplate(ticket.id);
log({ node: "fanout", event: "template_done", ticket: ticket.id });
return { ticket, handler: "template", text };
}
const r = await callWorker(ticket, 1);
calls += r.calls;
tokens += r.tokens;
return { ticket, handler: `worker:${ticket.category}`, text: r.text };
});
return { drafts, calls, tokens };
}
// ============ 9. 节点三:汇合(纯代码,传引用不传载荷)============
function oneLineOf(text) {
const head = text.split("。")[0];
return head.length > 22 ? `${head.slice(0, 22)}…` : head;
}
function mergeNode(drafts) {
const items = drafts.map((d) => {
const rel = path.join("out", `${d.ticket.id}.txt`);
fs.writeFileSync(path.join(ROOT, rel), `${d.text}\n`);
const item = {
id: d.ticket.id,
category: d.ticket.category,
handler: d.handler,
file: rel,
oneLine: oneLineOf(d.text),
};
state.tickets[item.id] = {
category: item.category,
handler: item.handler,
file: item.file,
one_line: item.oneLine,
gate_rounds: 0,
gate_reports: [],
stop: null,
status: "drafted",
};
log({ node: "merge", event: "collected", ticket: item.id, file: item.file, chars: d.text.length });
return item;
});
return { items };
}
// ============ 10. 节点四:评审回路(确定性 gate 优先,查-修-再查)============
function gateCheck(ticketId, reply) {
const problems = [];
if (!reply.includes(ticketId)) problems.push("missing_ticket_id");
for (const w of FILLER_WORDS) {
if (reply.includes(w)) problems.push(`filler_word:${w}`);
}
return { pass: problems.length === 0, report: problems.join(" | ") };
}
async function reviewNode(items, byId) {
let calls = 0;
let tokens = 0;
let totalRounds = 0;
for (const item of items) {
const full = path.join(ROOT, item.file);
let reply = fs.readFileSync(full, "utf8").trim(); // 载荷从文件读,不从上一节点带过来
let rounds = 0;
let lastReport = null;
const reports = [];
let verdict = null;
let gate = gateCheck(item.id, reply);
log({ node: "review", event: "gate", ticket: item.id, round: 0, pass: gate.pass, report: gate.report });
while (!gate.pass) {
reports.push(gate.report);
if (rounds >= MAX_REVIEW_ROUNDS) {
verdict = "max_rounds";
break;
}
if (gate.report === lastReport) {
verdict = "no_progress"; // 连续两轮报告一模一样,回路不再往前走
break;
}
if (item.handler === "template") {
verdict = "no_rewriter"; // 纯代码模板没有可回炉的工人,直接交人
break;
}
lastReport = gate.report;
rounds += 1;
totalRounds += 1;
const r = await callWorker(byId.get(item.id), rounds + 1, { prev: reply, report: gate.report });
calls += r.calls;
tokens += r.tokens;
reply = r.text;
fs.writeFileSync(full, `${reply}\n`);
gate = gateCheck(item.id, reply);
log({ node: "review", event: "gate", ticket: item.id, round: rounds, pass: gate.pass, report: gate.report });
}
const rec = state.tickets[item.id];
rec.gate_rounds = rounds;
rec.gate_reports = reports;
rec.stop = gate.pass ? "gate_pass" : verdict;
rec.status = gate.pass ? "pass" : "needs_human";
rec.one_line = oneLineOf(reply);
item.oneLine = rec.one_line;
item.status = rec.status;
item.rounds = rounds;
item.stop = rec.stop;
saveState(); // 一条工单判完留一次痕
}
return { items, calls, tokens, totalRounds };
}
// ============ 11. 节点五:汇报(纯代码)============
const pad = (s, n) => {
const w = [...String(s)].reduce((a, c) => a + (c.charCodeAt(0) > 127 ? 2 : 1), 0);
return String(s) + " ".repeat(Math.max(1, n - w));
};
function reportNode(items, totalRounds) {
console.log("\n=== 全图执行汇总 ===");
console.log(pad("节点", 10) + pad("耗时", 8) + pad("模型调用", 12) + pad("token", 9) + pad("gate 轮数", 12) + "状态");
// 汇报节点自己不进这张表:它就是这张表;它的耗时由外层 node() 记进 run-state.json
const order = ["route", "fanout", "merge", "review"];
for (const name of order) {
const n = state.nodes[name];
if (!n) continue;
const rounds = name === "review" ? String(totalRounds) : "-";
console.log(pad(name, 10) + pad(`${n.ms}ms`, 8) + pad(n.calls, 12) + pad(n.tokens, 9) + pad(rounds, 12) + n.status);
}
console.log("\n=== 逐条工单 ===");
console.log(pad("工单", 9) + pad("类别", 10) + pad("处理者", 18) + pad("gate 轮数", 12) + pad("停止原因", 16) + "状态");
for (const it of items) {
console.log(
pad(it.id, 9) + pad(it.category, 10) + pad(it.handler, 18) + pad(it.rounds, 12) + pad(it.stop, 16) + it.status
);
}
const needsHuman = items.filter((it) => it.status === "needs_human");
console.log(`\n产出目录 out/:${items.length} 份回复;需要人工接手:${needsHuman.length} 条`);
for (const it of needsHuman) {
console.log(` - ${it.id}(${it.stop}):${it.oneLine}`);
}
console.log(`留痕:run-state.json / run.jsonl(run_id=${RUN_ID})`);
return { needsHuman: needsHuman.length };
}
// ============ 12. 主流程:计划就是下面这十几行 ============
async function main() {
seedWorkspace();
initLog();
saveState();
const tickets = loadInbox();
console.log(`inbox/ 收到 ${tickets.length} 条工单:${tickets.map((t) => t.id).join(", ")}`);
const { routed } = await node("route", () => routeNode(tickets));
console.log(`[route] ${routed.map((t) => `${t.id}=${t.category}`).join(" ")}`);
const { drafts } = await node("fanout", () => fanoutNode(routed));
console.log(`[fanout] 并发上界 ${POOL_SIZE},产出 ${drafts.length} 份初稿`);
const { items } = await node("merge", async () => mergeNode(drafts));
console.log(`[merge] 写入 out/ ${items.length} 份,向下只传引用与一行摘要`);
if (process.env.STOP_AFTER === "merge") {
console.log("[stop] STOP_AFTER=merge:在评审之前停下,本次不做判定");
process.exit(2);
}
const byId = new Map(routed.map((t) => [t.id, t]));
const reviewed = await node("review", () => reviewNode(items, byId));
console.log(`[review] gate 重写共 ${reviewed.totalRounds} 轮`);
const { needsHuman } = await node("report", async () => reportNode(items, reviewed.totalRounds));
process.exit(needsHuman > 0 ? 1 : 0);
}
main().catch((e) => { console.error(e); process.exit(3); }); // 崩溃退 3,与 needs_human 的 1 区分开