Dagrs is an asynchronous Rust runtime for dependency graphs and an optional, durable Harness step loop.
Dagrs schedules work. It does not provide a model SDK, prompt format, tokenizer, MCP client, approval UI, policy service, or operating-system sandbox. Applications inject those integrations through Dagrs traits. Operation policy metadata is not a security boundary.
Choose the execution API by workload:
| API | Use it for | State and recovery |
|---|---|---|
Graph |
A static flow graph of Nodes and bounded channels |
EnvVar, graph events, and the legacy CheckpointStore contract |
HarnessRun |
A sequence of participant super-steps with durable work, budgets, interrupts, or child runs | Typed State, versioned checkpoints, RunStore, and durable RunEvents |
Graph supports loop/conditional/router nodes, bounded async channels, graph events, and the separate legacy checkpoint/hook APIs. The Harness API is additive; existing Graph::async_start() applications do not need to migrate.
use std::{error::Error, sync::Arc};
use dagrs::{Action, DefaultNode, EnvVar, Graph, InChannels, NodeTable, OutChannels, Output};
struct NoOp;
#[dagrs::async_trait::async_trait]
impl Action for NoOp {
async fn run(
&self,
_input: &mut InChannels,
_output: &mut OutChannels,
_env: Arc<EnvVar>,
) -> Output {
Output::empty()
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> {
let mut nodes = NodeTable::new();
let worker = DefaultNode::with_action("worker".to_owned(), NoOp, &mut nodes);
let mut graph = Graph::new();
graph.add_node(worker)?;
let report = graph.async_start().await?;
println!("{:?}", report.status);
Ok(())
}Graph construction and execution are asynchronous. The application owns the Tokio runtime. Graph::add_node, Graph::add_edge, and loop graph construction return DagrsResult; handle those errors rather than relying on panics.
A participant reads one immutable state snapshot and returns a delta plus an explicit transition. Finish is the only successful terminal transition.
use std::{collections::BTreeSet, sync::Arc};
use dagrs::{
DagrsResult, HarnessRun, MemoryRunStore, ParticipantId, ParticipantRegistry, RunConfig,
RunDependencies, RunStatus, State, StateDelta, StateSchema, StateSnapshot, StepOutcome,
StepParticipant, StepTransition, WorkResults,
};
struct Finish {
id: ParticipantId,
}
#[dagrs::async_trait::async_trait]
impl StepParticipant for Finish {
fn id(&self) -> ParticipantId { self.id.clone() }
fn read_keys(&self) -> BTreeSet<String> { BTreeSet::new() }
fn write_keys(&self) -> BTreeSet<String> { BTreeSet::new() }
async fn run(
&self,
_snapshot: StateSnapshot,
_work_results: WorkResults,
) -> DagrsResult<StepOutcome> {
Ok(StepOutcome::new(StateDelta::new(), Vec::new(), StepTransition::finish()))
}
}
#[tokio::main]
async fn main() -> DagrsResult<()> {
let schema = StateSchema::new(1)?;
let initial_state: State = schema.empty_state();
let registry = ParticipantRegistry::new();
registry.register(Finish { id: ParticipantId::new("finish")? })?;
let store = Arc::new(MemoryRunStore::new());
let dependencies = RunDependencies::new(registry, store, schema, initial_state)?;
let report = HarnessRun::start(
RunConfig::new(),
[ParticipantId::new("finish")?],
dependencies,
).await?;
assert_eq!(report.status, RunStatus::Succeeded);
Ok(())
}This minimal example uses in-memory storage. Use FileRunStore or the optional SeaORM adapter when a run must survive process restart.
- A
StateSchemadeclares typed keys, versions, and reducers. A participant declares its read and write sets. - All participants selected for a super-step receive the same immutable
StateSnapshot. - The barrier validates access declarations and reduces deltas in stable participant-ID order.
- The runtime commits the new state revision, work receipts, checkpoint, and durable event drafts through one
RunStoretransaction. - Only after the commit does the next participant wave run.
Transitions are explicit: Continue, Dispatch, Interrupt, or Finish. Work is identified by stable IDs and versioned opaque payloads. Dispatch and interrupt handlers are injected by the application.
A store transaction can make Dagrs state durable; it cannot make an external API call part of that database transaction. Before invoking an operation, Dagrs records it as started. If recovery cannot determine whether the effect happened, it records outcome_unknown and will not retry it automatically. The embedding application must reconcile the external result before continuing.
RunEvent is the replayable journal. RunNotice is an in-process live update and can be lost on disconnect. Event drafts must already be redacted by the caller; the store does not discover secrets in arbitrary payloads.
| Capability | Main API | Guide |
|---|---|---|
| Typed state, reducers, participant barriers | StateSchema, StepParticipant, StateStepExecutor |
Harness run and state |
| Run limits, cancellation, checkpoint/resume | RunConfig, RunDependencies, HarnessRun |
Harness run and state |
| Durable events and cursor replay | RunEvent, RunEventSource, RunNotice |
Harness run and state |
| Interrupt approval and validated edits | InterruptRequest, InterruptDecision |
Harness run and state |
| Dynamic operations, retry, and policy | OperationSpec, OperationRuntime, OperationPolicy |
Harness run and state |
| Child runs, joins, and persisted scheduling | ChildRunSpec, HarnessRun child APIs |
Harness run and state |
| Large results and step middleware | ArtifactStore, ContextMiddleware |
Artifacts and middleware |
| SQL RunStore | dagrs-store-sea-orm |
SeaORM adapter |
#[auto_node] and dependencies! |
dagrs-derive |
Derive crate |
| Provider-neutral integration example | fake model/tool adapters | Agent Harness example |
RunConfig::new() applies finite defaults. The runtime validates required limits. Explicit numeric ceilings are 256 operations per step, 256 buffered join results, and 16 MiB per inline result; other fields have no additional public cap beyond their serialized integer/duration representation.
| Limit | Default |
|---|---|
| Committed super-steps | 100 |
| Active execution time | 30 minutes |
| Interrupt TTL | 30 minutes |
| Concurrent operations | 16 |
| Operation requests per step | 32 |
| Buffered join results | 32 |
| Inline work result | 256 KiB |
| Concurrent children per root | 8 |
| Children per parent / total descendants per root | 16 / 64 |
| Child depth | 4 |
| Cancellation drain / terminal finalization | 5 seconds / 5 seconds |
RunStoreConfig::default() retains up to 10,000 checkpoints per thread and 100,000 events per run. The hard cap is 1,000,000 for each. Page reads accept at most 100 records.
| Package | Version | Notes |
|---|---|---|
dagrs |
0.9.0 |
Core graph and Harness runtime. Default feature derive re-exports dagrs-derive. |
dagrs-derive |
0.6.0 |
Procedural macros; also usable directly. |
dagrs-store-sea-orm |
0.1.0 |
Separate workspace package; driver features sqlite, postgres, mysql, none enabled by default. MySQL driver is also used for MariaDB. |
The filesystem artifact backend requires Unix inter-process file locking. On non-Unix platforms, FileArtifactStore::new returns a structured error. Memory and RunStore APIs remain available.
Run from the repository root:
cargo fmt --all -- --check
cargo test --workspace --all-targets --all-features
cargo test -p dagrs-store-sea-orm --features sqlite --test sea_orm_storeThe examples are separate Cargo workspaces:
cargo test --manifest-path examples/agent-harness/Cargo.toml --all-targets --all-features
cargo test --manifest-path examples/dagrs-sklearn/Cargo.toml --all-targets --all-featuresFor PostgreSQL, MySQL, or MariaDB conformance, enable the matching adapter feature and set DAGRS_TEST_DATABASE_URL to a disposable database. Do not run destructive conformance tests against production data.
Read in this order when changing behavior:
Cargo.tomlandsrc/lib.rsfor package features and public exports.- The implementation module under
src/oradapters/. - Focused tests under
tests/or the package'stests/directory. - The matching current guide under
docs/api/and the migration guide. docs/CHANGELOG.mdfor user-visible changes.
The completed plan and files under docs/plan/evidence/ are dated design/verification records. They are not a substitute for checking the current source and tests.
Clone the canonical repository with git clone https://github.com/gitmono-dev/dagrs.git. Run formatting and focused tests for changed behavior, then the workspace checks above. Public API, persisted schema, error-code, or recovery changes must update the matching guide and docs/CHANGELOG.md.
Dual licensed under MIT and Apache-2.0; see LICENSE-MIT and LICENSE-APACHE.