check-examples-routing-tasks.ts
Source: scripts/check-examples-routing-tasks.ts · Download source · Example guide
text
/** Task acceptance: physical inputs, actual tools, HTTP inference, durable effects and reopened artifacts. */
import assert from "node:assert/strict";
import { execFile } from "node:child_process";
import { createHash, randomBytes, randomInt } from "node:crypto";
import { mkdir, mkdtemp, readFile, writeFile } from "node:fs/promises";
import { join, resolve } from "node:path";
import { fileURLToPath } from "node:url";
import { parseArgs, promisify } from "node:util";
import { createDitto, loadRuntimeConfigFile } from "@codesoul-co/ditto/runtime";
import { createInferWorker, type NodeResult, type SampleOutput } from "@codesoul-co/ditto/worker/infer";
import { createInteractionWorker } from "@codesoul-co/ditto/worker/interaction";
import { createFileTools } from "../examples/_shared/tools/file-ingestion/index.ts";
import { PickupTaskStore } from "../examples/_shared/tools/pickup-task-store.ts";
import { runStateRouting } from "../examples/control-flow/routing/state-routing.ts";
import { runConditional } from "../examples/control-flow/routing/conditional.ts";
import { runBranchMerge } from "../examples/control-flow/routing/branch-merge.ts";
import { runFileTask } from "../examples/control-flow/routing/file-type.ts";
import { runRisk, resumeRisk } from "../examples/control-flow/routing/risk.ts";
import { runConfidence } from "../examples/control-flow/routing/confidence.ts";
import { type Runner, type RecordValue } from "../examples/control-flow/routing/shared.ts";
const { values } = parseArgs({ options: {
provider: { type: "string" }, python: { type: "string" },
report: { type: "string", default: ".examples-routing-tasks-live-results.json" },
"output-dir": { type: "string", default: ".examples-routing-tasks" },
} });
const config = loadRuntimeConfigFile("ditto.yaml", process.env);
const provider = values.provider ?? config.model?.provider;
assert.ok(provider && config.providers[provider], "Select a configured HTTP provider");
const modelName = config.providers[provider].model ?? (config.model?.provider === provider ? config.model.model : undefined);
assert.ok(modelName);
const model = { provider, model: modelName };
const python = resolve(values.python ?? process.env.DITTO_EXAMPLE_TOOLS_PYTHON ?? "examples/_shared/tools/.venv/bin/python");
await mkdir(resolve(values["output-dir"]!), { recursive: true });
const directory = await mkdtemp(join(resolve(values["output-dir"]!), "run-"));
const inputDirectory = join(directory, "input");
await promisify(execFile)(python, [fileURLToPath(new URL("./fixtures/routing-files.py", import.meta.url)), inputDirectory], { timeout: 120_000 });
console.log(`Task artifacts: ${directory}`);
const store = new PickupTaskStore(directory);
const tools = [...createFileTools({ root: inputDirectory, python,
pdftotext: process.env.DITTO_EXAMPLE_PDFTOTEXT ?? "pdftotext", tesseract: process.env.DITTO_EXAMPLE_TESSERACT ?? "tesseract",
asrModel: process.env.DITTO_EXAMPLE_ASR_MODEL ?? "tiny.en",
}), store.tool];
const toolCalls: { name: string; status: string }[] = [];
const runtime = createDitto({ config, sandbox: { ...config.sandbox, tools: tools.map(tool => tool.name) }, workers: [
createInferWorker(), createInteractionWorker({ output: store.output, tools: tools.map(tool => ({ ...tool,
async execute(args, context) {
try { const result = await tool.execute(args, context); toolCalls.push({ name: tool.name, status: result.status }); return result; }
catch (error) { toolCalls.push({ name: tool.name, status: "threw" }); throw error; }
},
})) }),
] });
const graphs: string[] = [];
const samples: NodeResult<SampleOutput>[] = [];
const observed: Runner = { async run(plan, input, options) {
graphs.push(plan.id);
const result = await runtime.run(plan, input, options);
for (const output of Object.values(result)) if (output && typeof output === "object" && "node" in output && output.node === "INFER.REASONING.SAMPLE") samples.push(output as NodeResult<SampleOutput>);
return result;
} };
const results: Record<string, unknown>[] = [];
const expectedTasks = new Map<string, { status: string; content: unknown; pickup?: RecordValue }>();
const startedAt = new Date().toISOString();
async function artifact(id: string, expected: unknown, status = "completed") {
const file = JSON.parse(await readFile(join(directory, "deliveries", `${id}.json`), "utf8"));
assert.deepEqual(file, { taskId: id, status, content: expected });
assert.equal(store.task(id).status, status);
assert.deepEqual(JSON.parse(store.task(id).outcome!), expected);
expectedTasks.set(id, { status, content: expected });
}
async function check(name: string, run: () => Promise<Record<string, unknown>>) {
graphs.length = 0; samples.length = 0; toolCalls.length = 0;
const start = Date.now();
const report: Record<string, unknown> = { name, status: "failed" };
console.log(JSON.stringify({ name, event: "started" }));
try { Object.assign(report, await run()); report.status = "passed"; }
catch (error) { report.error = error instanceof Error ? error.message : "Unknown failure"; }
finally {
Object.assign(report, { durationMs: Date.now() - start, graphs: [...graphs], toolCalls: [...toolCalls],
modelCalls: samples.length, executions: samples.map(sample => ({ executionId: sample.executionId, status: sample.status, usage: sample.output?.usage })) });
results.push(report);
await writeFile(values.report!, JSON.stringify({ startedAt, checkedAt: new Date().toISOString(), provider, model: modelName,
directory, database: join(directory, "tasks.sqlite"), reviewActor: "automated acceptance reviewer (not a real human)",
passed: results.filter(row => row.status === "passed").length, failed: results.filter(row => row.status !== "passed").length, results }, null, 2) + "\n");
console.log(JSON.stringify(report));
}
}
async function request(id: string) {
const value = { code: `PICKUP-${randomBytes(4).toString("hex")}`, quantity: randomInt(2, 30) };
const path = join(inputDirectory, `${id}.json`);
await writeFile(path, JSON.stringify({ id, text: `Pickup code ${value.code}; quantity ${value.quantity}.` }));
const persisted = JSON.parse(await readFile(path, "utf8")) as { id: string; text: string };
return { input: { ...persisted, model }, value, path };
}
try {
for (const task of ["extract", "count"] as const) await check(`task-condition-${task}`, async () => {
const id = `condition-${task}`;
const { input, value, path } = await request(id);
store.create(id, "condition");
await runStateRouting(observed, { ...input, task, state: "ready" });
const expected = task === "extract" ? { route: task, ...value } : { route: task, quantity: value.quantity };
await artifact(id, expected);
assert.deepEqual(graphs, [`condition-${task}`, "routing-delivery"]);
return { taskId: id, inputFile: path, expected, taskStatus: store.task(id).status };
});
await check("task-blocked-resume", async () => {
const id = "blocked-resume";
const { input, value } = await request(id);
store.create(id, "condition");
await runStateRouting(observed, { ...input, task: "extract", state: "blocked" });
await artifact(id, { route: "defer", status: "blocked" }, "blocked");
assert.equal(samples.length, 0);
await runStateRouting(observed, { ...input, task: "extract", state: "ready" });
await artifact(id, { route: "extract", ...value });
return { taskId: id, transitions: ["queued", "blocked", "completed"] };
});
for (const [level, route, queue] of [["urgent", "priority", "expedited"], ["routine", "standard", "normal"]] as const) await check(`task-branch-${route}`, async () => {
const id = `branch-${route}`;
const { input, value } = await request(id);
store.create(id, "dispatch");
await runConditional(observed, { ...input, text: `Service level: ${level}. ${input.text}` });
await artifact(id, { route, queue, ...value });
assert.deepEqual(graphs, ["branch-decision", `branch-${route}`, "routing-delivery"]);
return { taskId: id, dispatchQueue: queue, taskStatus: "completed" };
});
await check("task-branch-merge", async () => {
const id = "handoff";
const owner = `TEAM-${randomBytes(4).toString("hex")}`;
const deadline = "2027-04-17";
const path = join(inputDirectory, "handoff.txt");
await writeFile(path, `Owner: ${owner}. Deadline: ${deadline}.`);
store.create(id, "handoff");
await runBranchMerge(observed, { id, model, text: await readFile(path, "utf8") });
await artifact(id, { route: "merged", owner, deadline });
assert.equal(samples.length, 2);
return { taskId: id, artifact: join(directory, "deliveries/handoff.json"), expected: { owner, deadline } };
});
const files = JSON.parse(await readFile(join(inputDirectory, "manifest.json"), "utf8")) as { name: string; mediaType: string; code: string; quantity: number }[];
for (const file of files) await check(`task-file-${file.name}`, async () => {
const id = `file-${file.name.replaceAll(".", "-")}`;
const path = join(inputDirectory, file.name);
store.create(id, "file-extraction");
// Expected answers remain in this runner; only path/MIME/name reach the actual parser.
const result = await runFileTask(observed, { id, model, file: { path, name: file.name, mediaType: file.mediaType } });
const digest = createHash("sha256").update(await readFile(path)).digest("hex");
assert.equal((result.parsed as Record<string, unknown>).sourceSha256, digest);
const route = file.mediaType === "application/pdf" ? "pdf" : file.mediaType.startsWith("image/") ? "image" : file.mediaType.startsWith("audio/") ? "audio" : "spreadsheet";
await artifact(id, { route, file: file.name, code: file.code, quantity: file.quantity });
assert.deepEqual(graphs, ["file-ingestion", `file-${route}`, "routing-delivery"]);
assert.equal(toolCalls.length, 1);
return { taskId: id, inputFile: path, sourceSha256: digest, parsed: result.parsed, taskStatus: "completed" };
});
for (const [id, risk, verified, decision] of [
["risk-direct", 0, true, null], ["risk-verify", 30, true, null], ["risk-verify-review", 30, false, "approve"],
["risk-confirm", 60, true, "approve"], ["risk-rejected", 60, true, "reject"], ["risk-human", 90, true, "approve"],
] as const) await check(id, async () => {
const { input, value } = await request(id);
store.create(id, "risk", risk);
if (verified) store.addReference(value);
const result = await runRisk(observed, { ...input, risk, verify: actual => store.verify(actual) });
if (decision) {
const pendingStatus = risk >= 50 && risk < 75 ? "pending_confirmation" : "pending_human";
await artifact(id, { route: result.content.route, status: pendingStatus, ...value }, pendingStatus);
assert.equal(store.pickup(id), undefined);
assert.equal(toolCalls.length, 0);
await assert.rejects(runtime.invoke("INTERACTION.ACT.TOOL", { call: { id, name: "record_pickup", arguments: { id, ...value } } }), /not authorized/);
await store.review(id, decision, "acceptance-reviewer");
assert.equal(store.authorized(id, { ...value, quantity: value.quantity + 1 }), false, "Approval must bind the exact payload");
const resume = () => resumeRisk(observed, { id, value, route: result.content.route as "confirm" | "human" | "verify", authorize: (id, value) => store.authorized(id, value) });
if (decision === "reject") {
await assert.rejects(resume(), /authorization/);
assert.equal(store.pickup(id), undefined);
assert.equal(store.task(id).status, "rejected");
await artifact(id, { route: result.content.route, status: "rejected", ...value }, "rejected");
return { taskId: id, review: decision, finalStatus: "rejected", businessWrites: 0 };
}
await resume();
// A retry after successful completion must not insert another business record.
await resume();
}
const expected = { route: result.content.route, status: "executed", ...value };
await artifact(id, expected);
assert.deepEqual({ ...store.pickup(id) }, value);
expectedTasks.set(id, { status: "completed", content: expected, pickup: value });
return { taskId: id, review: decision, finalStatus: "completed", persistedPickup: value };
});
for (const [id, initialSources, resolvedSources, route] of [
["confidence-return", 5, 5, "return"], ["confidence-analyze", 4, 5, "analyze"],
["confidence-retry", 2, 5, "retry"], ["confidence-escalate", 1, 1, "escalate"],
["confidence-analysis-unresolved", 4, 4, "analyze"], ["confidence-retry-unresolved", 2, 2, "retry"],
] as const) await check(id, async () => {
const { input, value } = await request(id);
store.create(id, "reconciliation");
const references: string[] = [];
for (let index = 0; index < 5; index++) {
const path = join(inputDirectory, `${id}-reference-${index}.json`);
await writeFile(path, JSON.stringify(index < resolvedSources ? value : { ...value, quantity: value.quantity + 1 }));
references.push(path);
}
const assessments: unknown[] = [];
await runConfidence(observed, { ...input, evidence: await readFile(references[0]!, "utf8"), assess: async (candidate, attempt) => {
const paths = attempt === 0 ? references.slice(0, initialSources) : references;
const checked = await Promise.all(paths.map(async path => JSON.parse(await readFile(path, "utf8")) as RecordValue));
const matches = checked.filter(reference => reference.code === candidate.code && reference.quantity === candidate.quantity).length;
const score = matches / references.length;
assessments.push({ attempt, checkedFiles: paths, matches, requiredSources: references.length, score });
return score;
} });
const score = (route === "return" || route === "escalate" ? initialSources : resolvedSources) / 5;
const status = score >= 0.9 ? "returned" : "pending_human";
await artifact(id, { route, status, score, ...value }, status === "returned" ? "completed" : "pending_human");
assert.equal(samples.length, route === "analyze" || route === "retry" ? 2 : 1);
return { taskId: id, assessments, finalStatus: store.task(id).status };
});
await check("task-corrupt-file", async () => {
const id = "corrupt-file";
const path = join(inputDirectory, "corrupt.pdf");
await writeFile(path, "This is not a PDF file");
store.create(id, "file-extraction");
await assert.rejects(runFileTask(observed, { id, model, file: { path, name: "corrupt.pdf", mediaType: "application/pdf" } }));
assert.equal(samples.length, 0);
await store.fail(id, "FILE_PARSE_FAILED");
await artifact(id, { route: "error", status: "failed", error: "FILE_PARSE_FAILED" }, "failed");
return { taskId: id, parserRejected: true, modelCalls: 0, finalStatus: "failed" };
});
await runtime.close();
store.close();
await check("task-reopen-persistence", async () => {
const reopened = new PickupTaskStore(directory);
try {
for (const [id, expected] of expectedTasks) {
const task = reopened.task(id);
assert.equal(task.status, expected.status);
assert.deepEqual(JSON.parse(task.outcome!), expected.content);
if (expected.pickup) assert.deepEqual({ ...reopened.pickup(id) }, expected.pickup);
}
const count = reopened.db.prepare("SELECT COUNT(*) AS count FROM pickups").get()?.count;
assert.equal(count, 5, "Only authorized tasks produce durable business rows, including retries");
return { reopenedTasks: expectedTasks.size, durablePickups: count, database: join(directory, "tasks.sqlite") };
} finally { reopened.close(); }
});
} finally {
await runtime.close();
// The store may already have been closed for the restart check.
if (store.db.isOpen) store.close();
}
assert.equal(results.length, 27);
if (results.some(row => row.status !== "passed")) process.exitCode = 1;