// 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件のチケットと既知不具合DB ============
const TICKET_TEXT = {
"T-1001": "注文 A-77301 が今月2回引き落とされています。確認して、多く引かれた分を返金してください。",
"T-1002": "レポートページで「CSV エクスポート」を押すと、ボタンが回り続けたまま5分待っても反応がありません。Chrome、社内ネットワークです。",
"T-1003": "有人サポートの電話番号は何番ですか。直接電話で聞きたいです。",
"T-1004": "注文 A-77420 の請求書の宛名が間違っていて、私の個人名で発行されています。会社名義に直してください。",
"T-1005": "スマホアプリでログインしたあと、アバターがずっと表示されません。ウェブ版では正常です。",
"T-1006": "3か月使っていますが、何度問題を報告しても音沙汰がありません。この製品はまだ誰か保守しているのでしょうか。",
};
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. スタブクライアント: 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 = {
// ルーティングノード: 1回の呼び出しで6件を分類する。最終行はわざと正当なラベルにせず、コードにクランプさせる。
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営業日を過ぎても着金が確認できない場合は、このチケットに返信してください。取引IDを添えて決済チャネルに照会します。"
),
],
1124,
168
),
],
// bug ワーカー: 既知不具合DBを照合してから返信し、初版で gate を通過する。
"T-1002#1": [
turn(
"tool_use",
[say("まず既知不具合DBを確認します。"), 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 ワーカー: 初版にその場しのぎの言葉があり、第2版にも同じ言葉が残る——gate のレポートが2回とも同一になる。
"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 ワーカー: 初版はチケットIDが抜けており、第2版で 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 側のアバター URL が旧 CDN ドメインを指したままで、ウェブ版は影響を受けません。修正は今週金曜のバージョン公開に合わせて反映される見込みです。" +
"それまでは一度ログアウトして再度ログインしていただくと、アバターは通常表示に戻ります。"
),
],
1298,
158
),
],
"T-1005#2": [
turn(
"end_turn",
[
say(
"チケット T-1005 の返信: ご報告の事象は既知不具合 KI-91 に該当します。App 側のアバター URL が旧 CDN ドメインを指したままで、ウェブ版は影響を受けません。修正は今週金曜のバージョン公開に合わせて反映される見込みです。" +
"それまでは一度ログアウトして再度ログインしていただくと、アバターは通常表示に戻ります。公開後も表示されない場合は、このチケットにスクリーンショットを添えてください。アカウントを確認します。"
),
],
1466,
172
),
],
};
function makeStubClient(queue) {
let i = 0;
return {
messages: {
async create(req) {
if (!req.model || !req.max_tokens) {
throw new Error("スタブクライアント: 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);
// 1ターン分の 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. 2つのツール ============
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: "作業ディレクトリ配下のテキストファイルを読む。既知不具合DBやチケット原文の確認用。",
input_schema: {
type: "object",
properties: { path: { type: "string", description: "作業ディレクトリからの相対パス" } },
required: ["path"],
},
},
{
name: "lookup_order",
description: "注文IDで請求上の事実を調べる: 金額、課金回数、状態、請求書の宛名。",
input_schema: {
type: "object",
properties: { order_id: { type: "string", description: "A-77301 のような注文ID" } },
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. 3つの委任プロンプト: 目的 / 出力フォーマット / ツール指針 / タスクの境界 ============
const ROUTER_PROMPT = [
"あなたはカスタマーチケットのルーターです。",
"目的: 以下の各チケットを billing(請求、課金、請求書、返金)、bug(機能の不具合)、other(それ以外すべて)に分類すること。",
"出力フォーマット: 1チケット1行、書式は厳密に「チケットID: カテゴリ」、カテゴリは billing / bug / other のいずれか。理由は書かず、他の内容も出力しない。",
"ツール指針: このステップではツールを渡さないので、チケット本文だけで判断すること。システムを調べたと主張しないこと。",
"タスクの境界: 分類のみを行い、返信を書かず、結論を出さず、チケットを統合しない。判断がつかない場合は other にすること。",
].join("\n");
const WORKER_PROMPTS = {
billing: [
"あなたは請求チケットの専門担当で、一度に1件だけ扱います。",
"目的: このチケットの請求上の事実を突き止め、一度で言い切る日本語の返信を作成すること。",
"出力フォーマット: プレーンテキストの一段落。冒頭に「チケット <チケットID> の返信:」と書き、判明した事実、すでに実施した処理、ユーザーが次に何を期待できるかを順に述べる。箇条書きは使わず、社交辞令は書かない。",
"ツール指針: 請求上の事実は必ず lookup_order で調べ、チケット中の注文IDをそのまま渡すこと。見つからなければ見つからないと正直に述べ、チケットの記述から金額や課金回数を推測しないこと。",
"タスクの境界: このチケットの請求部分のみを扱い、注文を変更せず、追加の補償を約束せず、請求と無関係な質問には答えない。「お待ちください」「今しばらくお待ちください」「早急に対応します」のような中身のない言葉は書かない。",
].join("\n"),
bug: [
"あなたは不具合チケットの専門担当で、一度に1件だけ扱います。",
"目的: このチケットが既知の不具合かどうかを判定し、一度で言い切る日本語の返信を作成すること。",
"出力フォーマット: プレーンテキストの一段落。冒頭に「チケット <チケットID> の返信:」と書き、一致した既知不具合の番号と結論、暫定策、修正がいつ来るかを述べる。箇条書きは使わず、社交辞令は書かない。",
"ツール指針: 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(); // 各ノードの完了ごとに1回永続化
log({ node: name, event: "node_end", ms, calls: result.calls ?? 0, tokens: result.tokens ?? 0 });
return result;
}
// ============ 7. ノード1: ルーティング ============
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, [], {});
// 出力の締め付け: 「チケットID: カテゴリ」の行だけを認識し、ホワイトリスト外のカテゴリは 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. ノード2: ファンアウト(セクショニング + 並行プールの上限) ============
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")
: `チケットID ${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. ノード3: 合流(純コード、ペイロードではなく参照を渡す) ============
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. ノード4: レビュー回路(決定的な 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"; // 2ラウンド連続で同じレポート、ループはもう前に進んでいない
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(); // チケット1件を判定するたびに1回永続化
}
return { items, calls, tokens, totalRounds };
}
// ============ 11. ノード5: レポート(純コード) ============
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: review の前で停止、今回の実行に判定はなし");
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 と区別する