Skip to content
gitmono-devPublic

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Repository files navigation

Dagrs

Dagrs is an asynchronous Rust runtime for dependency graphs and an optional, durable Harness step loop.

Scope

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.

Graph quick start

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.

Harness quick start

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.

Harness execution contract

  1. A StateSchema declares typed keys, versions, and reducers. A participant declares its read and write sets.
  2. All participants selected for a super-step receive the same immutable StateSnapshot.
  3. The barrier validates access declarations and reduces deltas in stable participant-ID order.
  4. The runtime commits the new state revision, work receipts, checkpoint, and durable event drafts through one RunStore transaction.
  5. 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.

Durability and external effects

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.

Feature map

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

Default Harness limits

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.

Packages and features

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.

Build and test

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_store

The 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-features

For 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.

Repository map for maintainers and coding agents

Read in this order when changing behavior:

  1. Cargo.toml and src/lib.rs for package features and public exports.
  2. The implementation module under src/ or adapters/.
  3. Focused tests under tests/ or the package's tests/ directory.
  4. The matching current guide under docs/api/ and the migration guide.
  5. docs/CHANGELOG.md for 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.

Contributing

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.

License

Dual licensed under MIT and Apache-2.0; see LICENSE-MIT and LICENSE-APACHE.

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages