swale/lib.rs
1//! A scheduled, dependency-ordered orchestrator of asset graphs.
2//!
3//! swale runs on the taquba durable execution crates. Every task instance
4//! runs as a single workflow run, and the record of its outcome and the event
5//! for the scheduler commit in one transaction.
6//!
7//! # Definition file
8//!
9//! One TOML file per graph. An asset node declares the asset it produces and
10//! the assets it consumes, and a task node declares its upstream nodes with
11//! `after`.
12//!
13//! ```toml
14//! [graph]
15//! name = "orders_daily"
16//! schedule = "0 2 * * *"
17//! partition = "daily"
18//!
19//! [[node]]
20//! name = "extract"
21//! produces = "orders_raw"
22//! operator = "shell"
23//! [node.params]
24//! command = "dbt run --select orders_raw >&2 && printf '{\"table\": \"orders_raw\"}'"
25//!
26//! [[node]]
27//! name = "load"
28//! produces = "orders_warehouse"
29//! consumes = ["orders_raw"]
30//! operator = "http"
31//! [node.params]
32//! method = "POST"
33//! url = "https://warehouse.example/load"
34//! body = '{"table": "{{ upstream.extract.table }}"}'
35//! ```
36//!
37//! The [`definition`] module documents every field, and the loader reports
38//! every fault of a definition at once. From Rust, [`load_path`] returns the
39//! checked graph:
40//!
41//! ```
42//! use std::path::Path;
43//! use swale::OperatorSet;
44//!
45//! # fn main() -> Result<(), Box<dyn std::error::Error>> {
46//! let graph = swale::load_path(
47//! Path::new("examples/orders_daily.toml"),
48//! &OperatorSet::builtin(),
49//! )?;
50//! println!("{}: {} nodes", graph.name(), graph.nodes().len());
51//! # Ok(()) }
52//! ```
53//!
54//! # Operators
55//!
56//! An operator runs a node, and its parameters are the `params` table of the
57//! node. A parameter string refers to `{{ partition }}`, `{{ run.<field> }}`
58//! (`id` or `summary`), `{{ upstream.<node>.<path> }}` and `{{ env.<NAME> }}`.
59//! The upstream reference is a path into the JSON output of an upstream node,
60//! and the env reference is an environment variable of the daemon.
61//!
62//! - **`subprocess`** runs `argv` with a JSON document on stdin (the identity,
63//! the parameters and the upstream outputs) and reads the output from stdout.
64//! - **`shell`** runs `command` with `sh -c`, with the identity in `SWALE_*`
65//! variables, the upstream outputs in `SWALE_INPUTS` and the `env` table in
66//! the environment, and reads the output from stdout.
67//! - **`http`** sends `method`, `url`, `headers` and `body` within `timeout`
68//! and outputs `{"status": <code>, "body": <value>}`, with a JSON body
69//! parsed.
70//! - **`object_exists`** waits until an object exists at `url`, a store URL,
71//! with a poll every `interval` and a failure after `timeout`. No worker
72//! is held between polls, and the output is the URL, the size and the
73//! last-modified time of the object.
74//!
75//! Exit code 0 and a 2xx status are success. Exit code 75, a 5xx status, a 429
76//! status, a connection failure and a timeout are transient errors, retried up
77//! to `retries` times. Any other exit code or status is a permanent error,
78//! which dead-letters the task instance. A program or a request must be
79//! idempotent per attempt.
80//!
81//! # Schedule
82//!
83//! The daemon fires the `schedule` of a graph, a cron expression with five
84//! fields in UTC. The partition of a firing is the start of the schedule
85//! interval that ends at the firing time: a daily schedule at 02:00 that
86//! fires on 2026-09-16 runs the partition `20260915`. A graph with a schedule
87//! declares `partition` as `daily` or `hourly`.
88//!
89//! The first adoption of a graph runs the partitions of the firings within
90//! its `catchup` window. After downtime the daemon replays the missed firings
91//! within the same window.
92//!
93//! # Command
94//!
95//! The `swale` binary is enabled by the default `cli` feature. A program that
96//! embeds the library can set `default-features = false`.
97//!
98//! `swale validate` loads a definition and prints each fault, one per line.
99//!
100//! `swale run` runs a graph for one partition on a store, prints each node's
101//! record as it is written and waits for the run to settle. A second `run`
102//! for the same partition resumes the existing graph run, and a process
103//! interrupted mid-run resumes without repeating a completed node.
104//!
105//! `swale publish` writes a definition to the store as the current definition
106//! of its graph. `swale daemon` runs the published graphs until it is
107//! interrupted. It adopts each published definition within one sync interval,
108//! fires each schedule and runs the task instances. A graph run keeps the
109//! definition it started from, so an edit applies from the next graph run.
110//! The `default` pool always exists, and `--pool name=steps` adds a pool.
111//!
112//! The daemon keeps the records of a settled graph run and the request
113//! records for `--retention`, ninety days by default. A start of a partition
114//! whose records were removed runs the graph again.
115//!
116//! Every command that opens a store takes `--store`: a directory or an object
117//! store URL (`s3://bucket/prefix`, `gs://bucket/prefix`,
118//! `az://container/prefix`). A cloud scheme needs the matching cargo feature
119//! (`aws`, `gcp` or `azure`) and reads the provider's environment variables
120//! for its credentials. The default store is `~/.swale/store`.
121//!
122//! `swale run` and `swale daemon` open the store as its only writer. A
123//! `swale run` on the store of a running daemon opens a second writer, and
124//! the store then refuses the writes of the daemon.
125//!
126//! `swale status` lists the graphs of a store with the count of their graph
127//! runs in each state. `swale status <graph>` lists the latest graph runs of
128//! the graph, and `swale status <graph> <partition>` lists the nodes of one
129//! graph run. A node without a record is `ready`, `waiting` or `blocked`. A
130//! blocked node does not run until a rerun changes the record of an upstream.
131//!
132//! `swale queues` lists the job counts of every queue, and `swale queues
133//! <queue>` lists the dead jobs of one queue. Both commands only read from the
134//! store, so they can run alongside a daemon.
135//!
136//! `swale start <graph> <partition>...` starts the graph run of each
137//! partition, `swale rerun <graph> <partition> <node>` runs a node again, and
138//! `swale cancel <graph> <partition>` cancels an active graph run. Each
139//! command writes a request to the store and does not open the queue, and
140//! the daemon applies the request at its next sync pass. With `--wait` the
141//! command waits for the outcome and prints it.
142//!
143//! A rerun of a failed or cancelled node runs the node again, and its
144//! downstream nodes follow as their trigger rules allow. A rerun of a
145//! succeeded node runs the node again and, after it, every node downstream
146//! of it through an asset edge or an all-succeeded edge, each with the new
147//! outputs. A task node with the all-done or one-failed rule keeps its
148//! record. Until such a node runs again, the status view shows it as ready,
149//! waiting or blocked with its last record.
150//!
151//! ```console
152//! $ swale validate examples/orders_daily.toml
153//! orders_daily: 5 nodes, 3 edges
154//! $ swale run examples/local.toml --store ./swale-store
155//! local/none: started, 1 root node(s) submitted
156//! first: succeeded (local-none-first-r0)
157//! second: succeeded (local-none-second-r0)
158//! local/none: complete
159//! $ swale status local none --store ./swale-store
160//! local/none: complete, requested 2026-09-17T13:48:19Z, definition 9450f82715de
161//! NODE POOL STATE RUN TERMINATED
162//! first default succeeded local-none-first-r0 2026-09-17T13:48:19Z
163//! second default succeeded local-none-second-r0 2026-09-17T13:48:20Z
164//! $ swale publish examples/orders_daily.toml --store s3://bucket/swale
165//! orders_daily: published 36e831ff15e026ad45115374cb98edec6b617dffb0d4ef40f9fd351ea36bca9a
166//! $ swale daemon --store s3://bucket/swale --pool warehouse=2
167//! $ swale rerun orders_daily 20260915 transform --store s3://bucket/swale --wait
168//! request 01K5ARRE8ZW9K8XTJTQ6PVYJK7: submitted orders_daily-20260915-transform-r1
169//! ```
170//!
171//! A command exits with status 0, 1 or 2:
172//!
173//! - **0.** The command succeeded. An interrupt ends `swale daemon` with this
174//! status.
175//! - **1.** A definition has a fault, a graph run failed, a request was
176//! refused or another error occurred.
177//! - **2.** The arguments are not valid.
178
179pub mod daemon;
180pub mod definition;
181pub mod definition_store;
182pub mod dispatch;
183pub mod duration;
184mod error;
185pub mod graph;
186pub mod hook;
187pub mod input;
188pub mod operator;
189pub mod partition;
190pub mod pools;
191pub mod readiness;
192pub mod records;
193pub mod request;
194pub mod retention;
195pub mod scheduler;
196pub mod status;
197pub mod store;
198pub mod task;
199pub mod template;
200
201pub use daemon::{Daemon, DaemonOptions, RequestReport};
202pub use definition::{load_path, load_str};
203pub use definition_store::{DefinitionError, DefinitionStore, Published};
204pub use error::Error;
205pub use graph::{Graph, GraphSpec, Node, NodeKind, NodeSpec, Partitioning, Problem, TriggerRule};
206pub use hook::{EVENTS_QUEUE, Event, RecordHook};
207pub use operator::{Operator, OperatorSet, Outcome, Task};
208pub use partition::Partition;
209pub use pools::Pools;
210pub use readiness::NodeState;
211pub use records::{
212 GraphRecord, GraphRunRecord, GraphRunState, JsonBytes, NodeRecord, RecordStatus,
213 RequestOutcome, RequestRecord,
214};
215pub use request::{Request, RequestId, RequestStore};
216pub use scheduler::{RerunOutcome, Scheduler, SchedulerOptions, StartOutcome, TRIGGERS_QUEUE};
217pub use status::{GraphRunStatus, GraphStatus, NodeStatus, RunCounts, RunSummary, StatusReader};
218pub use task::TaskIdentity;
219pub use template::Template;