Dataflow-rs
A high-performance rules engine for IFTTT-style automation in Rust with zero-overhead JSONLogic evaluation.
Dataflow-rs is a lightweight, embeddable rules engine that lets you define IF → THEN → THAT automation in JSON. Rules are evaluated using pre-compiled JSONLogic for zero runtime overhead, and actions execute asynchronously for high throughput. Whether you're routing events, validating data, or building complex automation pipelines, Dataflow-rs gives you enterprise-grade performance with minimal complexity.
⚡ Blazing Fast Performance
Dataflow-rs is built for high-throughput hot paths. By compiling all JSONLogic expressions once at engine startup, runtime evaluation runs with zero allocations, zero parsing overhead, and predictable latency.
A multi-threaded benchmark (1,000,000 concurrent events) on a 10-core machine yields:
- Throughput: ~640,000 messages/sec
- Median (P50) Latency: 6 μs
- Tail (P99) Latency: 51 μs
- Tail (P99.9) Latency: 93 μs
🧩 Full-Stack Ecosystem
Go beyond backend microservices. Use the same rule definitions across your entire stack:
- Rust Backend: Run natively with maximum speed and concurrency using
dataflow-rs. - Browser & Edge: Run client-side validations or edge routing using WebAssembly bindings via @goplasmatic/dataflow-wasm.
- React UI Admin Portal: Let users and developers visualize, edit, and step-by-step debug rules using @goplasmatic/dataflow-ui.
How It Works: IF → THEN → THAT
┌─────────────────────────────────────────────────────────────────┐
│ Rule (Workflow) │
│ │
│ IF condition matches → JSONLogic against any field │
│ THEN execute actions (tasks) → map, validate, custom logic │
│ THAT chain more rules → priority-ordered execution │
└─────────────────────────────────────────────────────────────────┘
Example: IF order.total > 1000 THEN apply_discount AND notify_manager
Core Concepts
| Rules Engine | Workflow Engine | Description |
|---|---|---|
| Rule | Workflow | A condition + actions bundle — IF condition THEN execute actions |
| Action | Task | An individual processing step (map, validate, or custom function) |
| RulesEngine | Engine | Evaluates rules against messages and executes matching actions |
Both naming conventions are fully supported — use whichever fits your mental model.
Why dataflow-rs?
If you need dynamic business rules or user-customizable workflows, writing manual if/else checks makes your code rigid, while running full orchestrators (like Temporal or Zeebe) adds heavy infrastructure overhead and milliseconds of network latency. Dataflow-rs gives you the best of both worlds:
| Capability | Hardcoded Rust | dataflow-rs | Heavy Orchestrators (Temporal/Zeebe) |
|---|---|---|---|
| Hot Reload Rules | ❌ Recompile & redeploy | Instant JSON update | ❌ Deploy new worker code |
| Execution Overhead | None | Zero (pre-compiled JSONLogic) | ❌ DB reads/writes (tens of ms) |
| Browser execution | ❌ Compile full app to WASM | Run same rules in JS via WASM | ❌ Network round-trip required |
| Visual Debugger | ❌ Build your own UI | Included React UI components | Included Dashboard |
| Infrastructure | None | None (embeddable library) | ❌ Requires server clusters & DBs |
Getting Started
1. Add to Cargo.toml
[]
= "3.0"
= { = "1", = ["rt-multi-thread", "macros"] }
= "1.0"
2. Define Rules in JSON
A message arrives with its body in payload. Conditions and mappings are
evaluated against data, so the first rule loads the payload into data, and
the second rule acts on it. This is the chaining in IF → THEN → THAT: rules
run in order, and each one sees what the previous rules wrote.
A rule's condition is evaluated before any of its own tasks run. A condition can only read what earlier rules produced — never what its own tasks are about to write. That is why the parse lives in its own rule here rather than as a first task on
premium_order.
3. Run the Engine
use ;
use Message;
use json;
async
Handling Errors — Two Channels
process_message reports errors through two complementary channels:
Result::Errsignals that the engine stopped early (a task failed withoutcontinue_on_error, or an engine-level error occurred).message.errors()always contains every error encountered, including errors from tasks that ran withcontinue_on_error = trueand so didn't short-circuit the workflow.
A short-circuit ? will surface only the first kind. For full coverage:
use ;
use Message;
use json;
async
Using Rules Engine Aliases
use ;
// These are type aliases — same types, rules-engine terminology
let rule = from_json?;
let engine = builder.with_workflow.build?;
Key Features
- IF → THEN → THAT Model: Define rules with JSONLogic conditions, execute actions, chain with priority ordering.
- Zero Runtime Compilation: All JSONLogic expressions pre-compiled at startup for optimal performance.
- Full Context Access: Conditions can access any field —
data,metadata,temp_data. - Async-First Architecture: Native async/await support with Tokio for high-throughput processing.
- Execution Tracing: Step-by-step debugging with message snapshots after each action.
- Built-in Functions: Parse, Map, Validate, Filter, Log, and Publish for complete data pipelines.
- Pipeline Control Flow: Filter/gate function to halt workflows or skip tasks based on conditions.
- Channel Routing: Route messages to specific workflow channels with O(1) lookup.
- Workflow Lifecycle: Manage workflow status (active/paused/archived), versioning, and tagging.
- Hot Reload: Swap workflows at runtime without re-registering custom functions.
- Extensible: Add custom async actions by implementing the
AsyncFunctionHandlertrait. - Typed Integration Configs: Pre-validated configs for HTTP, Enrich, and Kafka integrations.
- WebAssembly Support: Run rules in the browser with
@goplasmatic/dataflow-wasm. - React UI Components: Visualize and debug rules with
@goplasmatic/dataflow-ui. - Auditing: Full audit trail of all changes as data flows through the pipeline.
Architecture
Compilation Phase (Startup)
- All JSONLogic expressions compiled once when the Engine is created
- Compiled logic cached with Arc for zero-copy sharing
- Validates all expressions early, failing fast on errors
Execution Phase (Runtime)
- Engine evaluates each rule's condition against the message context
- Matching rules execute their actions with pre-compiled logic (zero compilation overhead)
process_message()for normal execution,process_message_with_trace()for debugging- Each action can be async, enabling I/O operations without blocking
Performance
On a 10-core machine processing 1,000,000 messages concurrently (Tokio multi-threaded runtime, --release; per message: 1 parse + 6 mappings + 3 validations):
| Metric | Value |
|---|---|
| Throughput | ~640,000 msg/sec |
| Avg Latency | 10 μs |
| P50 Latency | 6 μs |
| P99 Latency | 51 μs |
| P99.9 Latency | 93 μs |
Why it's fast:
- Pre-Compilation: All JSONLogic compiled at startup, zero runtime parsing
- Arc-Wrapped Logic: Zero-copy sharing of compiled expressions across threads
- Arena Evaluation: Consecutive sync tasks evaluate against one bump-arena view of the context; map writes are spliced into it in place instead of re-cloning the written subtree
- Precomputed Paths: Mapping, parse, and publish target paths are split and interned at compile time — the hot path never re-parses a path string
- Async I/O: Non-blocking operations for external services via Tokio
Tuning tip: if you never read audit trails, build messages with
Message::builder().capture_changes(false) — skipping the per-mapping
old/new value snapshots is the largest single lever in mapping-heavy
workloads. See the performance guide for more.
Run the benchmarks and examples yourself:
Targeted microbenchmarks for profiling a specific hot path. The micro_* ones
run a tight current_thread loop so the signal isn't buried under Tokio
scheduling; the last two measure throughput on a multi-threaded runtime:
Custom Functions
Extend the engine with your own async actions. Each handler declares a typed
Input (deserialized once at engine init), receives a TaskContext that
records audit-trail changes automatically, and returns a TaskOutcome:
use async_trait;
use ;
use OwnedDataValue;
use Deserialize;
use json;
/// Typed config for the handler — fails at `Engine::new()` if malformed,
/// not on first message.
;
// Register handlers via the builder. `.register("name", h)` accepts any
// `AsyncFunctionHandler` and boxes it internally.
Built-in Functions
| Function | Purpose | Modifies Data |
|---|---|---|
parse_json |
Parse JSON from payload into data context | Yes |
parse_xml |
Parse XML string into JSON data structure | Yes |
map |
Data transformation using JSONLogic | Yes |
validation |
Rule-based data validation | No (read-only) |
filter |
Pipeline control flow — halt workflow or skip task | No |
log |
Structured logging with JSONLogic expressions | No |
publish_json |
Serialize data to JSON string | Yes |
publish_xml |
Serialize data to XML string | Yes |
Filter (Pipeline Control Flow)
The filter function evaluates a JSONLogic condition and controls pipeline execution:
on_reject: "halt"— stops the entire workflow when the condition is falseon_reject: "skip"— skips just the current task and continues
Log (Structured Logging)
The log function outputs structured log messages using the log crate:
Log levels: trace, debug, info, warn, error. Messages and fields support JSONLogic expressions.
Channel Routing
Route messages to specific workflow channels for efficient O(1) dispatch:
// Workflows define their channel
// { "id": "order_rule", "channel": "orders", "status": "active", ... }
// Process only workflows on a specific channel
engine.process_message_for_channel.await?;
Only active workflows are included in channel routing. Workflows default to the "default" channel.
Workflow Lifecycle
Workflows support lifecycle management fields:
| Field | Type | Default | Description |
|---|---|---|---|
channel |
string | "default" |
Channel for message routing |
version |
number | 1 |
Workflow version |
status |
string | "active" |
active, paused, or archived |
tags |
array | [] |
Arbitrary tags for organization |
created_at |
datetime | null |
Creation timestamp (ISO 8601) |
updated_at |
datetime | null |
Last update timestamp (ISO 8601) |
All fields are optional and backward-compatible with existing configurations.
Engine Hot Reload
Swap workflows at runtime without losing custom function registrations:
let new_workflows = vec!;
let new_engine = engine.with_new_workflows;
// Old engine remains valid for in-flight messages
Visualize & Debug Rules
Because every rule is plain JSON, the React UI can render it: JSONLogic expressions become readable flow diagrams, and the debugger steps through execution with a message diff after every task.
Ecosystem
| Package | Description | Install |
|---|---|---|
| dataflow-rs | Async rules engine in Rust (this crate) | cargo add dataflow-rs |
| @goplasmatic/dataflow-wasm | WebAssembly bindings — run rules in browser or Node.js | npm i @goplasmatic/dataflow-wasm |
| @goplasmatic/dataflow-ui | React components for rule visualization, editing, and step-by-step debugging | npm i @goplasmatic/dataflow-ui |
| datalogic-rs | JSONLogic compiler/evaluator used internally | cargo add datalogic-rs |
📖 Documentation: User Guide & API Reference · Interactive Playground · Visual Debugger
Contributing
We welcome contributions! Here's how to get started:
- Fork the repository and clone your fork
- Run tests:
cargo testto ensure everything passes - Make changes and add tests for any new features
- Run the benchmark before and after:
cargo run --example benchmark --release - Submit a pull request with a clear description of your changes
See the CHANGELOG for recent changes and release history.
About Plasmatic
Dataflow-rs is developed by the team at Plasmatic. We're passionate about building open-source tools for data processing and automation.
License
This project is licensed under the Apache License, Version 2.0. See the LICENSE file for more details.