Runtime, Graph and Loop API
简体中文 · Worker API · Runnable examples
Import Runtime APIs from @codesoul-co/ditto or @codesoul-co/ditto/runtime. Define the graph's nodes and dependencies, define loop state transitions, register concrete Workers, then execute. Model, database and tool plugins stay inside their Workers.
Placement and communication
| Worker placement | Execution | Registration |
|---|---|---|
| Same process | Direct call; no serialization | register(definition) |
| Same machine, separate process | Native Node IPC; no TCP socket | createIpcTransport + registerRemote; receiver uses serveWorkerIpc |
| Separate machines | HTTP(S) request/response | createHttpTransport + registerRemote; receiver uses createWorkerHttpHandler |
Routing filters public capabilities, availability and capacity, then prefers direct execution, same host, and other hosts, with round robin among peers. Each external registration names an installed transport. Runtime does not launch processes or discover endpoints. Use the same hostId for one machine and distinct processId values for its runtimes; Worker IDs must be unique within a caller's directory.
Local IPC / network HTTP are deployment modes. Separately, invoke means request/response and emit means asynchronous event acceptance. Events stay within the injected EventFabric; invoking a Worker through IPC/HTTP does not broadcast events. See communication for credentials and Artifact storage.
Creation, registration and cleanup
import { createDitto, createContextWorker, loadRuntimeConfigFile } from "@codesoul-co/ditto";
const config = loadRuntimeConfigFile("ditto.yaml", process.env);
const runtime = createDitto({ config, hostId: "machine-a", processId: "agent" });
const context = runtime.register(createContextWorker({ policy: config.context.policy ?? {} }), "context-a");
try {
console.log(runtime.workers());
console.log(await runtime.invoke("CONTEXT.LOAD", {
sources: [{ role: "user", content: "hello" }],
}, { workerId: context.address.workerId }));
} finally { await runtime.close(); }| API | Arguments and behavior |
|---|---|
createDitto(options?) / new DittoRuntime(options?) | config/providers/sandbox/sandboxExecutor build default services; workers registers definitions; hostId/processId identify the runtime; transports/events/artifacts/inlineLimitBytes inject communication resources |
register(worker, idOrOptions?) | Returns a WorkerHandle; second argument is an ID or { id?, services? }; resources are instantiated per registration |
registerRemote({ address, capabilities, transportId, concurrency? }) | Registers an external Worker; transport must be installed; caller capacity defaults to unlimited, receiver limits still apply |
workers() | Snapshots address, capabilities, available, active and concurrency without exposing SDKs or credentials |
handle.setAvailable(false/true) | Pause/resume routing; accepted invocations continue |
handle.unregister() | Removes directory entry, returning whether it existed; cleanup remains owned by the handle or Runtime |
handle.close() | Unregister, drain accepted calls, dispose once; repeated calls share the promise |
runtime.close() | Reject new top-level work, remove this Runtime's subscriptions, drain accepted graphs/loops/calls/event handlers, then dispose Workers; repeated calls share the promise |
context.setAvailable(false);
context.setAvailable(true);
context.unregister();
await context.close();Applications own shared EventFabric instances, transports, child processes, HTTP servers and injected database connections. Stop ingress, close Runtime, then release external resources. Worker-owned resources can use defineWorker({ resources, dispose, ... }). Do not await Runtime.close from its own handler: that would wait for itself.
Graph construction and execution
import { graph } from "@codesoul-co/ditto";
const prepare = graph<string>("prepare")
.node("load", "CONTEXT.LOAD", [], content => ({ sources: [{ role: "user", content }] }))
.node("select", "CONTEXT.SELECT", ["load"], (_input, { load }) => ({ context: load, purpose: "infer" }));
const result = await runtime.run(prepare, "Find relevant information", {
concurrency: 2,
workers: { load: "context-a", select: "context-a" },
signal: AbortSignal.timeout(30_000),
});
console.log(result.select.context);graph<I>(id?), runtime.graph<I>(id?) and ExecutionGraph.create<I>(id) create immutable graphs. .node(id, nodeType, dependencies, bind) returns a new graph; retain the return value. IDs must be unique and dependencies must name earlier nodes. Empty dependencies allow parallel execution. bind(input, outputs) sees only declared dependencies and returns the selected Node's typed input. Worker IDs, clients, credentials and endpoints belong to execution options rather than graph definitions.
runtime.run(plan, input, options?) returns results keyed by task ID. workers maps task IDs to Worker IDs, leaving unbound tasks on automatic routing. Invalid bindings and capability mismatches reject before execution; a pinned unavailable Worker does not silently fall back. Concurrency precedence is call option → YAML runtime.graphConcurrency → unlimited.
The scheduler builds an O(V+E) dependency index and ready queue. On an exception it stops starting tasks, drains started branches, and rejects with the first observed error. A business result with status: "failed" remains ordinary data: check it in bind and throw or forward deliberately. Output freezing is shallow; large payloads are not copied.
Worker concurrency limits each replica; graph concurrency limits each run. Full replicas cause NoWorkerAvailableError; no additional capacity queue is created. For parallel tasks pinned to a single-capacity replica, set graph concurrency to 1, add replicas, or change dependencies. Applications also control admission across concurrent graph runs.
Step dependencies and input binding
The four .node(id, nodeType, dependencies, bind) arguments are a unique step ID, public Node type, direct dependency IDs and a synchronous input mapper. The first bind argument is always the original runtime.run input; the second contains only the declared dependency outputs. The return value must match the selected Node's input contract. Results retain each Node's output type rather than being uniformly wrapped in NodeResult.
In loaded → updated → selected → delivered, a delivery binding that also reads the original loaded result must declare ["loaded", "selected"]. Declaring only ["selected"] does not expose transitive dependencies. This restriction applies both to TypeScript types and the object passed at runtime.
Use explicit dependency chains for sequential work; no concurrency override or separate sequential executor is required. Bindings map data; Nodes perform asynchronous work. Running a Graph once does not require a Loop.
See sequencing and step execution for complete calls. The examples export Graphs, input types, default inputs and run functions. They use public package entrypoints without TypeScript path aliases; real-model entrypoints explicitly load application YAML and environment variables:
npm run example:sequence:pipeline
npm run example:sequence:dependencies
npm run example:sequence:stages
npm run example:sequence:batchRun npm run check:examples:sequence:live for a real-model path. It covers all four execution types by extending the fixed-step and dependency Graphs with inference and delivery, and directly invoking the staged and batch functions to verify upstream data reaches the models. Check INFER NodeResult status and finishReason explicitly before delivery; an unsuccessful result is not automatically a Graph exception. See end-to-end validation for configuration, commands and reports.
Loop: repetition and graph selection
import { loop } from "@codesoul-co/ditto";
const alternate = graph<string>("alternate")
.node("load", "CONTEXT.LOAD", [], content => ({ sources: [{ role: "user", content }] }))
.node("select", "CONTEXT.SELECT", ["load"], (_input, { load }) => ({ context: load, purpose: "infer" }));
const workflow = loop({
graph: (state: { round: number; text: string }) => state.round % 2 ? alternate : prepare,
maxIterations: 4,
bind: state => state.text,
update: (state, output) => ({ round: state.round + 1, text: String(output.select.context.items[0]?.content ?? "") }),
done: state => state.round >= 4,
});
const state = await runtime.loop(workflow, { round: 0, text: "hello" }, {
concurrency: 1,
workers: { prepare: { load: "context-a" }, alternate: { load: "context-a" } },
});
console.log(state.round); // 4loop(definition) freezes a definition; runtime.loop also accepts a plain definition. graph is a fixed graph or synchronous selector based on state. Each iteration selects → binds → runs → updates → checks done. It runs at least once; done receives updated state and current output. Candidate graphs share an input/output contract; tagged unions can represent different business states.
The positive integer iteration limit comes from the definition, YAML runtime.loopMaxIterations, or 32. Reaching the limit without done throws. Execution worker bindings are grouped by graph ID; concurrency and signal cover the whole loop. State stays on the caller; only Node requests cross transports. There is no implicit persistence, retry or checkpointing.
Composing stages and batches
For complete Agent workflows, define a Graph per stage and compose them through a Loop plan: const a = yield* graphStep(stageA, input); const b = yield* graphStep(stageB, a);. Each stage retains its own input/output types. Loop owns execution, branching, repetition and the total Graph budget. See Graph composition through Loop for a runnable example and cancellation/recovery contracts.
For batches, use a fixed Graph with Loop: bind selects the current object, update validates its result and appends it to the collection, and done checks whether all objects are processed. Set maxIterations to the batch size and run a separate summary Graph afterwards. Loop executes at least once, so an empty batch must bypass it rather than setting maxIterations: 0. A Graph failure or an exception from update stops later objects and prevents summary execution. See batch.ts.
Worker-specific services and Sandbox
import { createRuntimeServices, createInteractionWorker } from "@codesoul-co/ditto";
const services = createRuntimeServices({
config, sandbox: { tools: ["inspect_text"], read: true },
// sandboxExecutor: applicationContainerExecutor,
});
const reader = runtime.register(createInteractionWorker({ tools: [inspectTextTool] }), {
id: "file-reader", services,
});inspectTextTool is an application RegisteredTool; see the complete example. Omit services to share Runtime defaults. Independent services can supply a different workspace (config.workspace), permission policy, ProviderRegistry and SandboxExecutor. HTTP providers are constructed against their corresponding Sandbox; reusing a broader provider while replacing only sandbox would not isolate its network access. Custom injected providers/SDKs must honor application policy themselves.
ctx.invoke uses the destination Worker's services. Remote hosts supply their own configuration; envelopes cannot override deployment settings. Sandbox is a cooperative capability guard. For arbitrary JS or SDK isolation, deploy the Worker in a separate process/container and inject the appropriate executor; Runtime does not create containers.
Worker context and cancellation
import { defineWorker } from "@codesoul-co/ditto";
const loader = defineWorker({ type: "CONTEXT", nodes: {
"CONTEXT.LOAD": async (input, ctx) => {
ctx.signal?.throwIfAborted();
console.log(ctx.worker.workerId, ctx.execution?.graphId, ctx.execution?.runId, ctx.execution?.nodeId);
const result = await ctx.invoke("CONTEXT.LOAD", input, { workerId: "context-a" });
await ctx.emit({ type: "context.loaded", payload: { count: result.items.length } });
return result;
},
} });ctx.resources/config hold Worker resources/config; ctx.services holds its execution environment and ctx.artifacts the optional store. ctx.invoke(node, input, { workerId?, signal? }) invokes public capabilities with inherited execution scope and cancellation. ctx.run(internalGraph, input) runs on the current replica, including unexposed internal nodes, without consuming another top-level capacity slot. Await internal calls before the handler returns.
Cancellation is cooperative. Local handlers receive ctx.signal and can forward it to SDKs or ctx.services.sandbox.run(command, ctx.signal). Graph/Loop stop subsequent tasks and drain started local work; arbitrary JS is not forcibly interrupted. IPC/HTTP cancellation and timeout stop the caller's wait, not remote execution. Use business idempotency keys for retried side effects.
Events
const unsubscribe = runtime.subscribe("context.loaded", event => { console.log(event.payload); });
await runtime.emit({ type: "context.loaded", payload: { count: 2 } });
const failures = await runtime.drainEvents();
for (const failure of failures) console.error(failure.event.type, failure.error);
unsubscribe();Emit acknowledges acceptance rather than consumer completion. Subscribe returns an independently owned removal function. Drain waits for accepted Fabric handlers and consumes failure records; a shared Fabric includes shared consumers. Default LocalEventFabric dispatches asynchronously within the process. Middleware adapters implement EventFabric.emit/subscribe/drain and are explicitly injected.
IPC and HTTP adapters
import { createIpcTransport, serveWorkerIpc } from "@codesoul-co/ditto";
// Parent: child is a node:child_process.fork result; exchange address at startup.
const ipc = createIpcTransport({ id: "local-ipc", channel: child, timeoutMs: config.timeoutMs });
const app = createDitto({ hostId: "machine-a", transports: [ipc] });
app.registerRemote({ address: childAddress, capabilities: ["CONTEXT.LOAD"], transportId: ipc.id, concurrency: 4 });
await app.invoke("CONTEXT.LOAD", { sources: [] });
await app.close(); ipc.close(); // Application separately shuts down the child.
// Child: register local Workers before mounting the receiver.
const receiver = serveWorkerIpc(workerRuntime, process);
await receiver.close(); // During shutdown, detach and drain accepted replies.
await workerRuntime.close();IpcChannel accepts Node ChildProcess or an IPC-enabled process. It requires an existing parent/child channel. Timeout defaults to 30000 ms, valid range 1–2147483647. Invocation IDs correlate responses; timeout, abort, disconnect and close clean up pending callers. Failures return a generic message. IPC is for trusted application channels, without a network token.
import { createHttpTransport, createWorkerHttpHandler } from "@codesoul-co/ditto";
const http = createHttpTransport({ id: "remote", url: endpoint, token, timeoutMs: config.timeoutMs });
const app = createDitto({ hostId: "machine-a", transports: [http] });
app.registerRemote({ address: remoteAddress, capabilities: ["CONTEXT.LOAD"], transportId: http.id });
// Server: createServer(createWorkerHttpHandler(workerRuntime, { token })).listen(...)HTTP transport requires a nonempty id; timeoutMs defaults to 30000 and accepts 1–2147483647, validated at construction. The HTTP endpoint is /ditto/invoke; default request/response limit is 1 MiB, configurable with maxBodyBytes. Both ends must use the same token; use HTTPS between production hosts. placement.ts includes startup, address exchange, graph execution and cleanup. receive(envelope) is the adapter-facing receiver with target/capability validation; business code should use typed invoke/run.
For InMemoryArtifactStore, PayloadCodec and external stores, see communication. Existing runRagFlow/runSkillFlow/runToolCallFlow/runMcpFlow/runReactFlow remain available; see Interaction API.
Cancellation through built-in Workers
For local execution, the Runtime signal reaches INFER model providers, MEMORY database adapters, CONTEXT services, RETRIEVAL pipelines and INTERACTION tools/MCP clients. SDK adapters must forward it to an underlying implementation that supports cancellation. Cancellation does not undo completed side effects. Shutdown drains already accepted Graph/Loop work, including nested ctx.invoke delegation. HTTP and IPC currently cancel the caller's wait without a remote task cancellation protocol.
Sandbox API and local execution
Import from @codesoul-co/ditto/runtime/sandbox or the root package. Implementations and replaceable ports require no additional process management dependency.
| API | Arguments and behavior |
|---|---|
| new Sandbox(workspace, policy?, executor?) | Resolves an absolute workspace; snapshots and freezes permissions; deny by default. read/write/execute must be booleans; tools/mcp/skills/network are arrays of nonempty strings |
| allows(kind, name) | Synchronous boolean; named permissions support exact matches or * |
| assert(kind, name) | Throws PermissionDeniedError if denied |
| readText(path, signal?) | UTF-8 read; requires read permission; rejects paths/symlinks outside the workspace |
| writeText(path, content, signal?) | UTF-8 overwrite/create; requires write permission and an existing parent directory; rejects outside/dangling symlinks; does not create directories |
| run({ command, args }, signal?) | Requires execute=true and an injected executor; snapshots arguments, resolves the workspace, checks cancellation before/after work and returns stdout/stderr/exitCode |
| createLocalSandboxExecutor(options) | Explicit trusted host executor for Sandbox or Runtime services |
LocalSandboxExecutorOptions requires commands: exact executable names or absolute paths. An empty array disables commands; relative executable paths are rejected. timeoutMs defaults to 5000 and accepts 1–2147483647 ms. maxOutputBytes defaults to 65536 and accepts 1–16777216, counting combined stdout/stderr UTF-8 bytes. env supplies explicit string values, defaulting to PATH=/usr/bin:/bin without inherited parent credentials or implicit Node coverage settings; the OS may add environment variables. Commands, environment and limits are snapshotted at construction.
Execution uses spawn with shell=false, separate literal arguments and closed stdin. Nonzero exits return their exitCode. Startup errors reject; timeout, cancellation or excess output kills the directly spawned process, closes output streams and rejects after process close. Oversized output never becomes a truncated success. Cancellation cannot roll back completed writes or external side effects.
import { Sandbox, createLocalSandboxExecutor } from "@codesoul-co/ditto/runtime/sandbox";
import { loadRuntimeConfigFile } from "@codesoul-co/ditto";
const config = loadRuntimeConfigFile("ditto.yaml", process.env);
const executor = createLocalSandboxExecutor({
commands: ["uname", "printf"], ...config.sandboxExecution,
});
const sandbox = new Sandbox(config.workspace, {
read: true, write: true, execute: true, tools: ["inspect"],
}, executor);
console.log(sandbox.allows("tools", "inspect")); // true
sandbox.assert("tools", "inspect");
const signal = AbortSignal.timeout(3000);
await sandbox.writeText("example.txt", "hello", signal);
console.log(await sandbox.readText("example.txt", signal));
const result = await sandbox.run({ command: "printf", args: ["%s", "$(uname) stays literal"] }, signal);
console.log(result); // { stdout: "$(uname) stays literal", stderr: "", exitCode: 0 }Inject with createDitto({ config, sandbox: { execute: true }, sandboxExecutor: executor }), or provide independent executors per Worker using createRuntimeServices. Container/remote adapters implement SandboxExecutor.run(command, { workspace, signal? }) without changing Graphs/Tools.
This is a cooperative capability boundary. Local commands and their descendants have no OS-level filesystem/network isolation; executable allowlists do not constrain files accessible through arguments. Terminating the direct child does not guarantee termination of its entire process tree. Untrusted code needs an external OS/container isolation executor. See the runnable Linux/macOS tool composition.
See composition for custom nodes/Workers, resources, events and Artifacts; see flows for predefined compositions.
Conditions and routing with public APIs: application-selected Graphs, dependency joins, risk/confidence policies and tool injection.
Parallel execution with public APIs: per-Graph concurrency, model-planned Graphs, strict joins and partial failures across independent Graphs.
Loops with public APIs: synchronous state transitions, asynchronous tool checks, model-planned Graph selection, budget/round termination and durable artifacts.
Recovery with public APIs: bounded Loop retries, Graph cancellation boundaries, application checkpoints, external effect queries and compensation.
Human intervention APIs and examples: approval before execution, intermediate confirmation, edited continuations, reviewed publication and human handoff.
Task lifecycle APIs and examples: state tracking, guarded execution, safe stopping, scheduled triggers and event triggers.
Public API boundaries and unified package acceptance: capability mapping for 38 examples, strict external consumer types and real task verification.
Request understanding and interaction through public APIs: goals, constraints, clarification, conversation, choices, and intent handling.
Planning and task management APIs: decomposition, dependency graphs, budget admission, and tool selection with Redis Context and database Memory.
Graph / Loop composition: Loop owns stage Graphs, branches, repetition, budgets and recovery.
