Skip to content

order-files.ts ​

Source: examples/_shared/tools/order-files.ts · Download source · Example guide

ts
import { createHash, randomUUID } from "node:crypto";
import { link, mkdir, readFile, realpath, rm, stat, writeFile } from "node:fs/promises";
import { isAbsolute, join, relative, resolve } from "node:path";
import { pathToFileURL } from "node:url";
import type { JsonObject, JsonValue } from "@codesoul-co/ditto/contracts";
import type { OutputSink, RegisteredTool } from "@codesoul-co/ditto/worker/interaction";

export interface OrderRecord { code: string; quantity: number; unitPriceCents: number }
export interface SavedOrder extends OrderRecord { sourceId: string; sourceSha256: string; totalCents: number }
export interface OrderFailure { sourceId: string; code: string; message: string }
export interface OrderReport {
  batchId: string; status: "completed" | "partial" | "failed";
  successes: SavedOrder[]; failures: OrderFailure[];
  totals: { orders: number; quantity: number; totalCents: number };
}
export function object(value: unknown): Record<string, unknown> {
  if (!value || typeof value !== "object" || Array.isArray(value)) throw new Error("Expected an object");
  return value as Record<string, unknown>;
}
export function identifier(value: unknown): string {
  if (typeof value !== "string" || !/^[A-Za-z0-9_-]{1,64}$/.test(value)) throw new Error("Invalid identifier");
  return value;
}
export function orderRecord(value: unknown): OrderRecord {
  const data = object(value);
  if (typeof data.code !== "string" || !data.code.trim() || typeof data.quantity !== "number"
    || !Number.isSafeInteger(data.quantity) || data.quantity < 1 || typeof data.unitPriceCents !== "number"
    || !Number.isSafeInteger(data.unitPriceCents) || data.unitPriceCents < 0
    || !Number.isSafeInteger(data.quantity * data.unitPriceCents)) throw new Error("Invalid order record");
  return { code: data.code, quantity: data.quantity, unitPriceCents: data.unitPriceCents };
}
export async function readJson(path: string): Promise<unknown> { return JSON.parse(await readFile(path, "utf8")); }

class ArtifactConflict extends Error {}

/** Atomically publish an immutable artifact. Identical retries succeed; conflicting writes fail. */
async function persist(path: string, value: unknown): Promise<void> {
  const temporary = `${path}.${randomUUID()}.tmp`;
  const text = JSON.stringify(value, null, 2) + "\n";
  await writeFile(temporary, text, { flag: "wx" });
  try {
    try { await link(temporary, path); }
    catch (error) {
      if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error;
      if (await readFile(path, "utf8") !== text) throw new ArtifactConflict("Artifact write conflict");
    }
  } finally { await rm(temporary, { force: true }); }
}
function jsonValue(value: unknown): JsonValue { return JSON.parse(JSON.stringify(value)) as JsonValue; }

/** Application-owned filesystem integration; no vendor or business configuration enters Core. */
export function createOrderFiles(options: { inputDirectory: string; outputDirectory: string }) {
  const outputDirectory = resolve(options.outputDirectory);
  const batchDirectory = (id: string) => join(outputDirectory, identifier(id));
  const pathFor = (id: string, sourceId: string) => join(batchDirectory(id), `${identifier(sourceId)}.order.json`);
  const tools: RegisteredTool[] = [
    {
      name: "read_order_source", effects: ["read"], inputSchema: { type: "object", required: ["sourceId", "path"], properties: { sourceId: { type: "string" }, path: { type: "string" } } },
      validate(args) { identifier(args.sourceId); if (typeof args.path !== "string" || !args.path.trim()) throw new Error("Source path required"); },
      async execute(args) {
        const root = await realpath(options.inputDirectory);
        let path: string;
        try { path = await realpath(String(args.path)); }
        catch (error) {
          if ((error as NodeJS.ErrnoException).code === "ENOENT") return { status: "failed", error: { code: "SOURCE_NOT_FOUND", message: "Order source was not found" } };
          throw error;
        }
        const rel = relative(root, path);
        if (rel === ".." || rel.startsWith("../") || rel.startsWith("..\\") || isAbsolute(rel)) throw new Error("Source outside input directory");
        const info = await stat(path);
        if (!info.isFile() || info.size > 1024 * 1024) return { status: "failed", error: { code: "INVALID_SOURCE", message: "Order source must be a regular file of at most one MiB" } };
        const bytes = await readFile(path);
        const text = new TextDecoder("utf-8", { fatal: true }).decode(bytes);
        if (!text.trim()) return { status: "failed", error: { code: "EMPTY_SOURCE", message: "Order source is empty" } };
        return { status: "success", structuredContent: { sourceId: args.sourceId!, text, sourceSha256: createHash("sha256").update(bytes).digest("hex") } };
      },
    },
    {
      name: "save_order", effects: ["write"], inputSchema: { type: "object", required: ["batchId", "sourceId", "sourceSha256", "record"] },
      validate(args) {
        identifier(args.batchId); identifier(args.sourceId); orderRecord(args.record);
        if (typeof args.sourceSha256 !== "string" || !/^[a-f0-9]{64}$/.test(args.sourceSha256)) throw new Error("Invalid source digest");
      },
      async execute(args) {
        const record = orderRecord(args.record);
        const saved: SavedOrder = { sourceId: String(args.sourceId), sourceSha256: String(args.sourceSha256), ...record, totalCents: record.quantity * record.unitPriceCents };
        await mkdir(batchDirectory(String(args.batchId)), { recursive: true });
        try { await persist(pathFor(String(args.batchId), saved.sourceId), saved); }
        catch (error) {
          if (error instanceof ArtifactConflict) return { status: "failed", error: { code: "ARTIFACT_CONFLICT", message: "An existing order artifact differs from this result" } };
          throw error;
        }
        return { status: "success", structuredContent: jsonValue(saved) };
      },
    },
    {
      name: "save_parallel_plan", effects: ["write"], inputSchema: { type: "object", required: ["batchId", "plan"] },
      validate(args) { identifier(args.batchId); object(args.plan); },
      async execute(args) {
        const directory = batchDirectory(String(args.batchId));
        await mkdir(directory, { recursive: true });
        await persist(join(directory, "plan.json"), args.plan);
        return { status: "success", structuredContent: args.plan! };
      },
    },
    {
      name: "build_order_report", effects: ["read", "write"], inputSchema: { type: "object", required: ["batchId", "successIds", "failures"] },
      validate(args) {
        identifier(args.batchId);
        if (!Array.isArray(args.successIds) || !Array.isArray(args.failures)) throw new Error("Expected result arrays");
        const ids = args.successIds.map(identifier);
        for (const failure of args.failures) {
          const value = object(failure); ids.push(identifier(value.sourceId));
          if (typeof value.code !== "string" || !/^[A-Z_]{1,64}$/.test(value.code) || typeof value.message !== "string" || !value.message.trim()) throw new Error("Invalid failure details");
        }
        if (new Set(ids).size !== ids.length) throw new Error("Duplicate or overlapping result IDs");
      },
      async execute(args) {
        const batchId = String(args.batchId);
        const successes: SavedOrder[] = [];
        for (const id of args.successIds as readonly string[]) {
          const stored = object(await readJson(pathFor(batchId, id)));
          const record = orderRecord(stored);
          if (stored.sourceId !== id || typeof stored.sourceSha256 !== "string" || !/^[a-f0-9]{64}$/.test(stored.sourceSha256)
            || stored.totalCents !== record.quantity * record.unitPriceCents) throw new Error("Invalid persisted order");
          successes.push({ sourceId: id, sourceSha256: stored.sourceSha256, ...record, totalCents: stored.totalCents });
        }
        const failures = args.failures as unknown as OrderFailure[];
        const quantity = successes.reduce((sum, item) => sum + item.quantity, 0);
        const totalCents = successes.reduce((sum, item) => sum + item.totalCents, 0);
        if (!Number.isSafeInteger(quantity) || !Number.isSafeInteger(totalCents)) throw new Error("Report totals exceed safe integers");
        const report: OrderReport = { batchId, status: failures.length ? successes.length ? "partial" : "failed" : "completed",
          successes, failures, totals: { orders: successes.length, quantity, totalCents } };
        const directory = batchDirectory(batchId);
        await mkdir(directory, { recursive: true });
        await persist(join(directory, "report.json"), report);
        return { status: "success", structuredContent: jsonValue(report) };
      },
    },
  ];
  const output: OutputSink = { async deliver(input) {
    const id = identifier(input.deliveryId);
    const report = await readJson(join(batchDirectory(id), "report.json"));
    if (JSON.stringify(report) !== JSON.stringify(input.message.content)) throw new Error("Delivery does not match the persisted report");
    const path = join(batchDirectory(id), "delivery.json");
    await persist(path, { deliveryId: id, report });
    return { deliveryId: id, status: "accepted", artifacts: [{ name: "order-report.json", reference: { uri: pathToFileURL(path).href, mediaType: "application/json" } }] };
  } };
  return { tools, output, batchDirectory, pathFor };
}
export function toJsonObject(value: unknown): JsonObject { return object(jsonValue(value)) as JsonObject; }

Ditto · @codesoul-co/ditto · Node.js 24+