Overview
High-throughput streaming data processor. A long-running Rust runner spawns tools (any language) as child processes, routes NDJSON envelopes between them over OS pipes, and coordinates the whole pipeline from a declarative YAML/JSON graph.
Five-minute tour
Section titled “Five-minute tour”A pipeline is a folder with one or more variants:
my-pipeline/├── tools/ # pipeline-local tools (optional; override global tools)├── configs/ # prompts, dictionaries, anything tools read├── storage/ # caches that survive across runs├── temp/ # intermediaries (gates, checkpoints, spools)├── sessions/<id>/ # per-run artefacts (auto-created)└── variants/ └── main.yamlA variant declares the DAG:
pipeline: my-pipelinevariant: mainstages: scan: tool: scan-fs settings: { include: "*.txt", hash: blake2b } input: $input read: tool: read-file-stream settings: { format: lines } input: scan sink: tool: write-file-stream settings: { default_file: "$output/_index.ndjson", format: ndjson } input: readThen run it:
dpe run my-pipeline:main -i input/ -o output/That’s it. Each stage is a separate OS process; the runner wires stdout→stdin between them; backpressure comes free from the kernel.
Start here
Section titled “Start here”| Doc | Read if you want to… |
|---|---|
| Installing | get dpe + standard tools on your machine (npm / script / Docker) |
| Configuration | config.toml schema + resolution order |
| Concepts | envelopes, stages, DAG topology, the runner’s job |
| Writing pipelines | author YAML variants: linear, fan-in, route, filter, replicas, dedup |
| Expressions | route / filter / condition expressions in the DSL |
| Path prefixes | $input / $output / $configs / $storage / $temp / $session |
| CLI reference | dpe init / run / check / test / coverage / tools / install / config / log / logs / monitor — full flag reference + run –json –stats wire format |
| Testing | dpe test — per-stage snapshot tests with expected/data.ndjson diff; dpe coverage — per-variant matrix + skip-list |
| Session artefacts | what lands in sessions/<id>/: trace, errors, journal, stages.json |
| Editor / programmatic integration | dpe run --json --stats event stream, dpe log --stage, dpe check --plan — what consumers can rely on |
| Caching | ctx.cached(...) — skip expensive work (LLM calls, big parses) when the same inputs have been seen before |
| Tools overview | standard-tool catalogue + tool contract |
| Frameworks | write your own tool in Rust / Python / TypeScript |
| Authoring a tool | one-command scaffold + headless build from a spec.yaml |
| Docker | dpe-base + dpe-dev images + client pipeline patterns |
| Examples | worked variants walked through end-to-end |
What’s in this monorepo
Section titled “What’s in this monorepo”| Component | Path | Purpose |
|---|---|---|
| Runner + CLI | runner/ |
dpe binary — spawns tools, wires pipes, traces, controls sessions |
| Dev CLI | dpe-dev/ |
Tool authoring: scaffold / build / test / check / setup |
| Rust framework | frameworks/rust/ |
SDK for writing tools in Rust (combycode-dpe crate) |
| TypeScript framework | frameworks/ts/ |
SDK for Bun/TS (@combycode/dpe-framework-ts) |
| Python framework | frameworks/python/ |
SDK for Python (combycode-dpe package) |
| scan-fs | tools/scan-fs/ |
filesystem scanner (files / dirs / both) |
| read-file-stream | tools/read-file-stream/ |
stream text-file rows (NDJSON / lines / CSV) |
| write-file-stream | tools/write-file-stream/ |
write envelopes to files, LRU-bounded |
| write-file-stream-hashed | tools/write-file-stream-hashed/ |
same with per-file content dedup |
| normalize | tools/normalize/ |
row-level normaliser (dict / parse / rename / compute / template / require) |
| gate | tools/gate/ |
stateful pass-through that publishes progress |
| checkpoint | tools/checkpoint/ |
spool-then-release on gate(s) met |
| Built-ins | in-runner | route / filter / dedup / group-by |
| dev-workspace-template | dev-workspace-template/ |
Claude skill pack + fixtures — embedded into dpe-dev |
| Docker | docker/ |
Multi-stage Dockerfile: dpe-base + dpe-dev images |
| Test pipeline | test-pipeline/ |
Regression suite — runs every standard tool against synthesised inputs |
| Catalogue | catalog.json |
Manifest of standard tools |
Design rules (read these before arguing with the code)
Section titled “Design rules (read these before arguing with the code)”- Each tool is an independent program. Any language. Receives settings as one JSON argument on argv[1]. Reads NDJSON envelopes line-by-line from stdin. Writes NDJSON to stdout. Writes typed events to stderr. Tools never self-exit — the runner owns their lifecycle.
- Envelopes are the contract.
{"t":"d","id":"...","src":"...","v":{...}}for data;{"t":"m","v":{...}}for meta. Onlyvis business-meaningful;id/srcchain the provenance. - Filesystem-first coordination. Gates, checkpoints, dedup indices — all live on disk. The runner observes files; it doesn’t own state across tools.
- Framework does the merged trace event. On every
ctx.output(), the framework emits one{type:"trace",id,src,labels}to stderr. Runner just appends to$session/trace/trace.N.ndjson. No sniffing of stdout. - IPC is cross-platform local sockets, never TCP. Named pipe on Windows, UDS on Unix, via
interprocess. Runner writes$session/control.addr; CLI reads it to connect. - DAG is acyclic. If you need iteration, loop inside a tool. Cycles in the DAG are explicitly rejected.
License
Section titled “License”DPE is licensed under AGPL-3.0-or-later. See LICENSE for the full text. Commercial licensing is available — contact CombyCode for details.