operon 0.7.0

A workflow engine for parallel, incremental scheduling of DAG-defined multiplex tasks.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
# Operon

[![arXiv](https://img.shields.io/badge/arXiv-2511.16080-b31b1b.svg)](https://arxiv.org/abs/2511.16080)
[![crates.io](https://img.shields.io/badge/crates.io-v0.7.0-blue.svg)](https://crates.io/crates/operon)
[![License](https://img.shields.io/badge/license-MIT%2FApache--2.0-yellow.svg)](LICENSE-MIT)
[![MSRV](https://img.shields.io/badge/MSRV-1.91+-lightgray.svg)](https://blog.rust-lang.org/2025/10/30/Rust-1.91.0/)

A Rust-native workflow engine designed for parallel, incremental scheduling of [DAG-defined](#running-dag-defined-tasks) [multiplex](#multiplexing) tasks.
Powered by a PostgreSQL-based transactional backend, Operon specializes in orchestrating complex and long-running workflows with minimal downtime, flexible recovery, and high parallelism.

## Table of Contents

- [Operon]#operon
  - [Table of Contents]#table-of-contents
  - [Examples \& Demo](#examples--demo)
  - [Prerequisites]#prerequisites
  - [Quick Start]#quick-start
  - [Key Features]#key-features
    - [Running DAG-Defined Tasks]#running-dag-defined-tasks
    - [Multiplexing]#multiplexing
    - [Incremental Scheduling]#incremental-scheduling
    - [Transactional Backend]#transactional-backend
    - [Interactive UI \& Workflow Control](#interactive-ui--workflow-control)
    - [Per-Task Parallelism]#per-task-parallelism
  - [Usage]#usage
    - [Installation]#installation
    - [Defining Entities]#defining-entities
    - [Defining the Pipeline]#defining-the-pipeline
    - [Implementing the Service]#implementing-the-service
    - [Implementing the Storage (Optional)]#implementing-the-storage-optional
    - [Running Operon]#running-operon
      - [Operon TUI]#operon-tui
  - [Roadmap]#roadmap
  - [License]#license
    - [Contribution]#contribution

## Examples & Demo

![Demo 1](docs/assets/demo1.gif)

▲ Animation of running [ex2](operon/examples/ex2.rs) with Operon. _(log level `Info`)_

![Demo 2](docs/assets/demo2.gif)

▲ Animation of recovering from a poisoned run of ex2.

You can find more examples in the [examples](operon/examples/) directory of this repository.

## Prerequisites

You will need the following to run Operon:

- [Rust]https://www.rust-lang.org/tools/install (tested with Rust 1.91+)
  - An async runtime configured via [`tokio`]https://crates.io/crates/tokio
- A working [PostgreSQL]https://www.postgresql.org/download/ database (version 14 or later), to use the PostgreSQL backends
  - You will need a full [connection URI]https://www.postgresql.org/docs/current/libpq-connect.html#LIBPQ-CONNSTRING-URIS that can access the database.
  - A pipeline can also run without a database, keeping everything in process. See [Running Operon]#running-operon.

## Quick Start

If you want to try out Operon, you can clone the repository and run the provided examples:

```bash
git clone https://github.com/Asteromorph-Corp/operon
cd operon
# Make sure the URI points to a running PostgreSQL database.
export POSTGRES_URI=<your_postgres_uri>
cargo run --release --example ex1
```

We recommend reading the source code of [ex1](operon/examples/ex1.rs) to get a hang of how everything works.

If you do not have a database at hand, [ex5](operon/examples/ex5.rs) runs a whole pipeline in memory and needs nothing but `cargo run --release --example ex5`.

## Key Features

### Running DAG-Defined Tasks

Operon's primary use case is best described as _a known pipeline of an unknown number of jobs_.
It executes these jobs in parallel until all possible jobs have completed.

Here, a _task_ is one kind of work the pipeline does, and a _job_ is one execution of a task.
A task's outputs ([_entities_](#defining-entities)) can serve as inputs for other tasks.
Tasks and their dependencies must be predefined, forming a directed acyclic graph (DAG).
This DAG's validity is checked at macro-expansion time.
How many jobs each task runs, on the other hand, is discovered as the run proceeds.

### Multiplexing

Tasks in Operon are _multiplex_, meaning that one job may produce multiple entities of the same type (as a Rust `Vec`).
From another perspective, allowing multiplexing means that a single task may run many jobs, each using different input entities.
In this sense, Operon's DAG could also be viewed as a dynamic graph of _jobs_ that evolves as the run progresses.

The number of jobs a task runs cannot be known until upstream tasks produce the necessary entities.
Due to this, the number of jobs is quantified using an abstraction called _named dimensions_ instead of a simple count.

### Incremental Scheduling

Operon utilizes incremental scheduling, which means the scheduler never needs to know the entire task graph up front, saving memory and startup time.
As an event-driven system, each individual task runner is only aware of the jobs it can execute immediately, enabling efficient resource usage and pooling.

### Transactional Backend

Operon keeps two kinds of state: the _metadata_ that drives scheduling and recovery, and the _entity data_ that tasks produce and consume.
Both are backed by a PostgreSQL implementation that ships with the engine.
Its transactional design allows for atomic updates to job states, and by extension, reliable recovery from failures.

As a tradeoff, the PostgreSQL backend often requires heavy database access, which may become a bottleneck for systems with high-throughput workloads.
Both halves are swappable: an in-memory metadata backend ships alongside the PostgreSQL one, and the entity storage is an interface you may implement over any store you like.
A run kept entirely in memory gives up recovery in exchange for dropping the database.

### Interactive UI & Workflow Control

Operon provides a terminal-based TUI for real-time monitoring and control.
Users can track task progress, browse past logs, and interact with the workflow through shell-like commands — including pausing, resuming, or gracefully shutting down the engine.

### Per-Task Parallelism

Operon supports per-task parallelism, meaning that each task maintains its own thread pool for the jobs it runs.
This is particularly useful for tasks that benefit from internal parallel execution or must adhere to external concurrency limits (e.g., database connections or API rate limits).
The pipeline definition sizes each pool, and may also fix the order in which a task's jobs are picked up.
See [the `define_operon!` documentation](docs/define_operon_dsl.md).

## Usage

### Installation

Add Operon to your project's dependencies by adding `operon` in `Cargo.toml`:

```bash
cargo add operon
```

Alternatively, clone this repository:

```bash
git clone https://github.com/Asteromorph-Corp/operon
```

Once that's done, add the following to your project's `Cargo.toml`:

```toml
[dependencies]
# Assuming you cloned the repository to your home directory:
operon = { path = "~/operon/operon" }
```

### Defining Entities

_Entities_ are typed values that are produced and consumed by tasks in Operon.
Any valid Rust type with a `PascalCase` name can be used as an entity, given that it implements `Debug + Clone + Serialize + DeserializeOwned + Send + Sync + 'static`.
The serde bounds hold even for a pipeline that never touches a database, since the PostgreSQL storage is generated for every pipeline.
An example of entity declarations is as follows:

```rust
// In operon/examples/ex1.rs:

use serde::{Serialize, Deserialize};

// Strings already implement all the necessary traits,
// so using a type alias of `String` is sufficient for our `Input` type.
type Input = String;

// For composite types, we need to implement or derive the necessary traits.
#[derive(Debug, Clone, Serialize, Deserialize)]
struct Intermediate(String);

#[derive(Debug, Clone, Serialize, Deserialize)]
struct Output(char);
```

_Note 1_. These entity types must be directly accessible (without module scoping) in the scope where the [pipeline definition](#defining-the-pipeline) `define_operon!` macro is used.

_Note 2_. The engine can only recognize type names that are in `PascalCase` (a.k.a. `UpperCamelCase`) as defined in the [`heck` crate](https://docs.rs/heck/latest/heck).
Here are some examples of valid and invalid entity names:

| Invalid Name          | Valid Name                           |
| --------------------- | ------------------------------------ |
| `a`                   | `A`                                  |
| `lowerCamel`          | `UpperCamel`                         |
| `APIResponse`         | `ApiResponse`                        |
| `XYCoordinates`       | `XyCoordinates` / `XAndYCoordinates` |
| `_String` / `String_` | `StringEntity`                       |

### Defining the Pipeline

The _pipeline_ is the skeleton of the Operon workflow, defining _how_ the entities will be produced and consumed.
More specifically, the pipeline consists of the following components:

- **Name**: A unique identifier for the pipeline.
- **Tasks**: A listing of tasks that will be executed in the pipeline.
  Each task introduces a new type of entity to the pipeline, which can be used as input for subsequent tasks.

Additionally, entities in Operon are paired with _named dimensions_ that represent the way you can iterate over the entities.
Simply put, these dimensions can be understood as _directions_ the entities repeat in.
For example, if you have an `Intermediate` entity that has two dimensions, `input_no` and `word_no`, you can think of it as a 2D grid where each cell is an `Intermediate` entity.

![Figure 1](docs/assets/figure1.svg)

▲ An example of a 2D grid of `Intermediate` entities.

The following is an example of a pipeline definition using the `define_operon!` macro:

```rust
// In operon/examples/ex1.rs:

operon::define_operon! {
    splitter = {
        Input<input_no> = get_inputs();
        Intermediate<word_no> = get_words(Input) for input_no;
        Output<char_no> = get_chars(Intermediate) for input_no, word_no;
    }
}
```

![Figure 2](docs/assets/figure2.svg)

▲ A visual representation of the "splitter" pipeline.

Invoking the `define_operon!` macro brings several utilities into scope:

- a `schema` module that contains the metadata of the pipeline;
- a `{PipelineName}Service` trait that provides the parsed tasks [you would need to implement]#implementing-the-service;
- a `{PipelineName}Storage` trait that exposes [the storage interface]#implementing-the-storage-optional for the entities;
- a `Psql{PipelineName}Storage` struct that serves as a default implementation of the storage interface using PostgreSQL.

The pipeline must follow a few rules that are enforced at macro-expansion time:

- Each task takes a list of "arguments" or "inputs" that must be entities that were defined earlier in the pipeline.
  Each task input must be either a single entity (`EntityType`) or a slice across dimensions (`EntityType<dim1, dim2, ...>`).
- Each task must return one of the following two options:
  - A single entity, denoted `SpawnedEntityType`.
  - A 1D vector of entities, denoted `SpawnedEntityType<spawned_dimension_name>`.
    In this case, this task spawns a dimension that can be iterated over in subsequent tasks.
- The dimension specifications must be "well-formed," as thoroughly described in [our technical report]https://arxiv.org/abs/2511.16080.
  - For illustration, take the list of `Intermediate`s as shown in [the "splitter" pipeline example]docs/assets/figure1.svg: `[["Good", "morning"], ["Bonjour"], ["Buenos", "días"]]`.
  - Writing `Intermediate<word_no>` represents a vector/slice of `Intermediate` entities indexed by `word_no`, which we will have for each `input_no` "coordinate."
  `["Good", "morning"]` or `["Bonjour"]` would be a valid example of such a vector.
  - However, writing `Intermediate<input_no>` would not be feasible.
  If we apply the same logic with above, we need a vector of `Intermediate` entities indexed by `input_no` "for each `word_no` coordinate."
  When `word_no` is `0`, we would have `["Good", "Bonjour", "Buenos"]`, but when `word_no` is `1`, what would we have — `["morning", ???, "días"]`?
  The range of `word_no` is unknown until the coordinate of `input_no` is fixed, so we cannot implicitly iterate over `word_no` while collapsing `input_no`.

We provide brief diagnostics for violations of these rules.

A task may additionally carry an `#[operon(...)]` attribute that sizes its worker pool or fixes the order its jobs run in:

```rust
operon::define_operon! {
    splitter = {
        Input<input_no> = get_inputs();
        #[operon(concurrency = 8, ord = (-input_no))]
        Intermediate<word_no> = get_words(Input) for input_no;
        Output<char_no> = get_chars(Intermediate) for input_no, word_no;
    }
}
```

Here `get_words` gets a pool of 8 workers and takes its jobs in descending `input_no`, so the last input is split first.

If you need further information, refer to the [`define_operon!` documentation](docs/define_operon_dsl.md) and the [technical report](https://arxiv.org/abs/2511.16080) for more details on the system.

### Implementing the Service

The pipeline definition serves as a blueprint for the tasks that will be executed — now you would need to implement the actual logic of these tasks.
This is done by deriving `OperonService` on a type of your own and providing an `impl` for the `{PipelineName}Service` trait that was generated by the `define_operon!` macro.
Continuing with the previous example, you would implement the `splitter` pipeline as follows:

```rust
// In operon/examples/ex1.rs (slightly modified):

use async_trait::async_trait;
use operon::OperonService;

#[derive(OperonService)]
#[operon(error = std::convert::Infallible)]
struct MySplitterService;

#[async_trait]
impl SplitterService for MySplitterService {
    async fn get_inputs(&self) -> Result<Vec<Input>, Self::Error> {
        Ok(vec![
            Input::from("Good morning"),
            Input::from("Bonjour"),
            Input::from("Buenos días"),
        ])
    }
    async fn get_words(&self, input: Input) -> Result<Vec<Intermediate>, Self::Error> {
        Ok(input
            .split_whitespace()
            .map(|s| Intermediate(s.to_string()))
            .collect())
    }
    async fn get_chars(&self, intermediate: Intermediate) -> Result<Vec<Output>, Self::Error> {
        Ok(intermediate.0.chars().map(Output).collect())
    }
}
```

The exact signature of each task function is parsed from the pipeline definition, and will be provided in a docstring of the generated `{PipelineName}Service` trait.

A task's methods report failure as `Self::Error`, which defaults to a boxed `std::error::Error`.
You can name a concrete error type instead with `#[operon(error = MyError)]` on the derive.
In the above example, we wrote `#[operon(error = std::convert::Infallible)]` because the service is infallible.
A job that returns an error puts its own task into an error state and reports it to the UI.
The run keeps going elsewhere and ends as aborted, leaving what did complete available to a later recovery.

### Implementing the Storage (Optional)

The Operon engine assumes all entities are accessible through a storage interface — we call this interface the `{PipelineName}Storage` trait.
We provide a struct `Psql{PipelineName}Storage` that already implements this trait using PostgreSQL, which you build from `PsqlStorageOptions`.
The `build` method needs to be told which storage it is building, either by a turbofish or by annotating the binding.

```rust
// In operon/examples/ex1.rs (slightly modified):

use operon::options::PsqlStorageOptions;

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
    // ...
    let storage = PsqlStorageOptions::new("postgres://username:password@hostname:port/dbname")
        .with_schema("data")
        .build::<PsqlSplitterStorage>()?;
    // ...
}
```

You may also choose to implement your own storage by providing an `impl` for two traits `OperonStorage` and `{PipelineName}Storage`.

The `OperonStorage` half covers the pipeline-independent interface: namely, the error type, the backend's lifecycle, and footprint operations.
The generated half asks for a `get`/`put` pair per entity type and provides default implementations for batch operations.
The batch operations will, by default, iterate over the `get`/`put` methods you provide, but you may override them if your backend supports more efficient bulk operations.
Analogous to the service implementation, signatures of the functions you need to implement are parsed from the pipeline definition.
The signatures will be provided in a generated docstring on the `{PipelineName}Storage` trait.

Having an alternative storage backend may be useful if you want to use a different database or have a quick in-memory storage for testing purposes.
For a concrete example, [ex5](operon/examples/ex5.rs) implements one over `DashMap`.
However, note that the engine will not provide recoverability if the storage is volatile or you leave the footprint methods of `OperonStorage` at their defaults.

### Running Operon

Once you have all the pieces in place, you build the metadata backend, construct an `Operon` instance from your service, storage, and that backend, then call the `.run()` method.
Surface-level run settings (UI mode, logging) live in a separate `OperonOptions` you attach with `.with_options(...)`.

```rust
// In operon/examples/ex1.rs (slightly modified):

use operon::options::{PsqlMetaStorageOptions, PsqlStorageOptions};
use operon::Operon;

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
    //# ——————————————————— Initializing Settings ————————————————————— #//
    let database_uri = "postgres://username:password@hostname:port/dbname";

    let service = MySplitterService;
    let storage = PsqlStorageOptions::new(&database_uri)
        .with_schema("ex1_data")
        .build::<PsqlSplitterStorage>()?;
    let meta = PsqlMetaStorageOptions::new(&database_uri)
        .with_schema("ex1_meta")
        .build()?;

    //# ———————————————————————— Running Operon ——————————————————————— #//
    Operon::new(service, storage, meta)
        .run()
        .await?;

    Ok(())
}
```

The `.run()` method will start the Operon engine that will execute the defined pipeline using the provided service and storage implementations.
If the execution is successful, the results will be stored in the storage, and you can retrieve them using the storage interface after the `.run().await?` call.

The metadata backend is a separate choice from where the entities live.
Swapping `PsqlMetaStorageOptions` for `MemMetaStorageOptions` keeps the metadata in process, so a run needs no database at all once the entity storage is also database-free:

```rust
use operon::options::MemMetaStorageOptions;

let meta = MemMetaStorageOptions::new().build();
```

If you use an in-memory metadata backend, the metadata will be dropped at the end of the run, which amounts to opting out of recoverability.
To defer the choice to runtime, build the backend through `MetaBackendOptions` instead, which resolves either backend into a single `AnyBackend` type:

```rust
use operon::AnyBackend;
use operon::options::MetaBackendOptions;

let meta: AnyBackend = match std::env::var("POSTGRES_URI") {
    Ok(uri) => MetaBackendOptions::psql(uri),
    Err(_) => MetaBackendOptions::mem(),
}
.build()?;
```

#### Operon TUI

When run in Interactive mode (`.with_ui_mode(UiMode::Interactive)` in `OperonOptions`, which is the default behavior), the Operon engine takes over the terminal and launches a text user interface (TUI).
The UI allows you to interact with the engine and control the workflow using shell-like commands.

```text
Navigation keys:
    Ctrl+C              Clear input.
    Ctrl+D              Exit.
    Ctrl+L              Clear logs.
    Left, Right         Scroll progress bars.
    Alt+Up, Alt+Down    Scroll logs 1 line.
    Up, Down            Scroll logs 5 lines.
    PgUp, PgDn          Scroll logs 20 lines.
    Esc                 Show most recent logs.

Commands:
    run [OPTIONS]       Start a new run using the best available restoration
                        (unless specified by options).
                        --fresh, --rebuild, and --redo are mutually exclusive.
        -f, --fresh         Start a fresh run, ignoring any existing data.
        -r, --rebuild       Rebuild the run from trusted data before starting.
        -s, --skip <TASK>[ ...]
                            With --rebuild, do not rebuild the given 1 or more task(s).
        -R, --redo <TASK>[ ...]
                            Shorthand for --rebuild --skip <...>.
        -i, --redo-inconsistent-tasks
                            Rebuild the run even on a failed check,
                            ignoring tasks with corrupt data and their downstream tasks.
                            Cannot be used with --fresh.
                            Note that --redo <INCONSISTENT_TASKS> will NOT allow a rebuild
                            on a failed check without this flag.
    check [OPTIONS]     Check the consistency of the data from the last run.
        -m, --mode [MODE]   Mode of the consistency check. Defaults to "quick". Options:
            trust-all           Assume all data is trustworthy, skipping checks.
            metadata-only       Check only metadata consistency.
            quick               Perform a metadata check plus data validation only at boundaries.
            exhaustive          Perform a full consistency check of all data. (Can be very slow.)
    exit                Exit the UI.
    clear               Clear the log buffer.
    quit [OPTIONS]      Stop all jobs and exit the UI. Defaults to graceful shutdown.
        -f, --force         Force quit.
        -n, --no-exit       Don't exit the UI.
    pause [OPTIONS] [<TASK>[ ...]]
                        Pause executing new jobs.
        -c, --cascade       Cascade the pause command to dependent tasks.
    resume [<TASK>[ ...]]
                        Resume paused tasks.
    help                Print this help message.
```

You may disable the UI by setting `.with_ui_mode(UiMode::Headless)` in `OperonOptions`.
Certain features, such as real-time monitoring, pause/resume functionality, and recovery options, will not be available in Headless mode.

## Roadmap

Operon is under active development.
Please check the [issues](https://github.com/Asteromorph-Corp/operon/issues) page for the full list of planned features and known problems.
You can also reach out via opening an issue if you [experienced bugs](https://github.com/Asteromorph-Corp/operon/issues/new?template=bug-report.yml) or [have any suggestions](https://github.com/Asteromorph-Corp/operon/issues/new?template=feature-request.yml).

## License

This project is licensed under either the [MIT License](LICENSE-MIT) or the [Apache License (Version 2.0)](LICENSE-APACHE), at your option.

### Contribution

Unless you explicitly state otherwise, any contribution intentionally submitted for inclusion in the work by you shall be dual licensed as above, without any additional terms or conditions.
Please read our [CONTRIBUTING.md](CONTRIBUTING.md) file for further information.