Expand description
A scheduled, dependency-ordered orchestrator of asset graphs.
swale runs on the taquba durable execution crates. Every task instance runs as a single workflow run, and the record of its outcome and the event for the scheduler commit in one transaction.
§Definition file
One TOML file per graph. An asset node declares the asset it produces and
the assets it consumes, and a task node declares its upstream nodes with
after.
[graph]
name = "orders_daily"
schedule = "0 2 * * *"
partition = "daily"
[[node]]
name = "extract"
produces = "orders_raw"
operator = "shell"
[node.params]
command = "dbt run --select orders_raw >&2 && printf '{\"table\": \"orders_raw\"}'"
[[node]]
name = "load"
produces = "orders_warehouse"
consumes = ["orders_raw"]
operator = "http"
[node.params]
method = "POST"
url = "https://warehouse.example/load"
body = '{"table": "{{ upstream.extract.table }}"}'The definition module documents every field, and the loader reports
every fault of a definition at once. From Rust, load_path returns the
checked graph:
use std::path::Path;
use swale::OperatorSet;
let graph = swale::load_path(
Path::new("examples/orders_daily.toml"),
&OperatorSet::builtin(),
)?;
println!("{}: {} nodes", graph.name(), graph.nodes().len());§Operators
An operator runs a node, and its parameters are the params table of the
node. A parameter string refers to {{ partition }}, {{ run.<field> }}
(id or summary), {{ upstream.<node>.<path> }} and {{ env.<NAME> }}.
The upstream reference is a path into the JSON output of an upstream node,
and the env reference is an environment variable of the daemon.
subprocessrunsargvwith a JSON document on stdin (the identity, the parameters and the upstream outputs) and reads the output from stdout.shellrunscommandwithsh -c, with the identity inSWALE_*variables, the upstream outputs inSWALE_INPUTSand theenvtable in the environment, and reads the output from stdout.httpsendsmethod,url,headersandbodywithintimeoutand outputs{"status": <code>, "body": <value>}, with a JSON body parsed.object_existswaits until an object exists aturl, a store URL, with a poll everyintervaland a failure aftertimeout. No worker is held between polls, and the output is the URL, the size and the last-modified time of the object.
Exit code 0 and a 2xx status are success. Exit code 75, a 5xx status, a 429
status, a connection failure and a timeout are transient errors, retried up
to retries times. Any other exit code or status is a permanent error,
which dead-letters the task instance. A program or a request must be
idempotent per attempt.
§Schedule
The daemon fires the schedule of a graph, a cron expression with five
fields in UTC. The partition of a firing is the start of the schedule
interval that ends at the firing time: a daily schedule at 02:00 that
fires on 2026-09-16 runs the partition 20260915. A graph with a schedule
declares partition as daily or hourly.
The first adoption of a graph runs the partitions of the firings within
its catchup window. After downtime the daemon replays the missed firings
within the same window.
§Command
The swale binary is enabled by the default cli feature. A program that
embeds the library can set default-features = false.
swale validate loads a definition and prints each fault, one per line.
swale run runs a graph for one partition on a store, prints each node’s
record as it is written and waits for the run to settle. A second run
for the same partition resumes the existing graph run, and a process
interrupted mid-run resumes without repeating a completed node.
swale publish writes a definition to the store as the current definition
of its graph. swale daemon runs the published graphs until it is
interrupted. It adopts each published definition within one sync interval,
fires each schedule and runs the task instances. A graph run keeps the
definition it started from, so an edit applies from the next graph run.
The default pool always exists, and --pool name=steps adds a pool.
The daemon keeps the records of a settled graph run and the request
records for --retention, ninety days by default. A start of a partition
whose records were removed runs the graph again.
Every command that opens a store takes --store: a directory or an object
store URL (s3://bucket/prefix, gs://bucket/prefix,
az://container/prefix). A cloud scheme needs the matching cargo feature
(aws, gcp or azure) and reads the provider’s environment variables
for its credentials. The default store is ~/.swale/store.
swale run and swale daemon open the store as its only writer. A
swale run on the store of a running daemon opens a second writer, and
the store then refuses the writes of the daemon.
swale status lists the graphs of a store with the count of their graph
runs in each state. swale status <graph> lists the latest graph runs of
the graph, and swale status <graph> <partition> lists the nodes of one
graph run. A node without a record is ready, waiting or blocked. A
blocked node does not run until a rerun changes the record of an upstream.
swale queues lists the job counts of every queue, and swale queues <queue> lists the dead jobs of one queue. Both commands only read from the
store, so they can run alongside a daemon.
swale start <graph> <partition>... starts the graph run of each
partition, swale rerun <graph> <partition> <node> runs a node again, and
swale cancel <graph> <partition> cancels an active graph run. Each
command writes a request to the store and does not open the queue, and
the daemon applies the request at its next sync pass. With --wait the
command waits for the outcome and prints it.
A rerun of a failed or cancelled node runs the node again, and its downstream nodes follow as their trigger rules allow. A rerun of a succeeded node runs the node again and, after it, every node downstream of it through an asset edge or an all-succeeded edge, each with the new outputs. A task node with the all-done or one-failed rule keeps its record. Until such a node runs again, the status view shows it as ready, waiting or blocked with its last record.
$ swale validate examples/orders_daily.toml
orders_daily: 5 nodes, 3 edges
$ swale run examples/local.toml --store ./swale-store
local/none: started, 1 root node(s) submitted
first: succeeded (local-none-first-r0)
second: succeeded (local-none-second-r0)
local/none: complete
$ swale status local none --store ./swale-store
local/none: complete, requested 2026-09-17T13:48:19Z, definition 9450f82715de
NODE POOL STATE RUN TERMINATED
first default succeeded local-none-first-r0 2026-09-17T13:48:19Z
second default succeeded local-none-second-r0 2026-09-17T13:48:20Z
$ swale publish examples/orders_daily.toml --store s3://bucket/swale
orders_daily: published 36e831ff15e026ad45115374cb98edec6b617dffb0d4ef40f9fd351ea36bca9a
$ swale daemon --store s3://bucket/swale --pool warehouse=2
$ swale rerun orders_daily 20260915 transform --store s3://bucket/swale --wait
request 01K5ARRE8ZW9K8XTJTQ6PVYJK7: submitted orders_daily-20260915-transform-r1A command exits with status 0, 1 or 2:
- 0. The command succeeded. An interrupt ends
swale daemonwith this status. - 1. A definition has a fault, a graph run failed, a request was refused or another error occurred.
- 2. The arguments are not valid.
Re-exports§
pub use daemon::Daemon;pub use daemon::DaemonOptions;pub use daemon::RequestReport;pub use definition::load_path;pub use definition::load_str;pub use definition_store::DefinitionError;pub use definition_store::DefinitionStore;pub use definition_store::Published;pub use graph::Graph;pub use graph::GraphSpec;pub use graph::Node;pub use graph::NodeKind;pub use graph::NodeSpec;pub use graph::Partitioning;pub use graph::Problem;pub use graph::TriggerRule;pub use hook::EVENTS_QUEUE;pub use hook::Event;pub use hook::RecordHook;pub use operator::Operator;pub use operator::OperatorSet;pub use operator::Outcome;pub use operator::Task;pub use partition::Partition;pub use pools::Pools;pub use readiness::NodeState;pub use records::GraphRecord;pub use records::GraphRunRecord;pub use records::GraphRunState;pub use records::JsonBytes;pub use records::NodeRecord;pub use records::RecordStatus;pub use records::RequestOutcome;pub use records::RequestRecord;pub use request::Request;pub use request::RequestId;pub use request::RequestStore;pub use scheduler::RerunOutcome;pub use scheduler::Scheduler;pub use scheduler::SchedulerOptions;pub use scheduler::StartOutcome;pub use scheduler::TRIGGERS_QUEUE;pub use status::GraphRunStatus;pub use status::GraphStatus;pub use status::NodeStatus;pub use status::RunCounts;pub use status::RunSummary;pub use status::StatusReader;pub use task::TaskIdentity;pub use template::Template;
Modules§
- daemon
- The daemon: the long-running process of a deployment. It runs the pools
and the
Scheduler, adopts the published definitions and schedules the cron entry of every adopted graph with a schedule. - definition
- The definition file: one TOML document per graph.
- definition_
store - The definitions of a deployment, stored in the object store by content hash.
- dispatch
- The step runner of every pool: it reads the task input from the payload
and the identity from the headers, renders the parameters and runs the
operator. An
{{ env.<NAME> }}reference renders from the environment of the process, read once when the runner is built. AnOutcome::Continueenqueues the next step of the run after its delay, with the input and the state as the payload. - duration
- Durations in the definition file: an integer and a unit, such as
30s,5m,6hor7d. - graph
- The internal graph model: nodes, their edges and the checks a graph passes
before it runs. Every front-end builds a
Graphfrom aGraphSpecthroughGraph::build. - hook
- The terminal hook of every pool: it writes the node’s record and enqueues the event for the scheduler, and both commit with the notification’s acknowledgement.
- input
- The input of a task instance: the operator, the node’s parameters before rendering and the records of its upstreams. The input is the run’s payload, so a definition edit that changes a node’s parameters while its run is active fails the resubmission with an input mismatch. The payload of a later step of the run is the input with the state of the previous step.
- operator
- The operators a definition can name: the check of their parameters at load time and their execution at run time.
- partition
- The partition key of an asset, in the charset of a run id.
- pools
- The pools: one
WorkflowRuntimeper pool, all over one queue and one store. - readiness
- The readiness rule: when a node runs, and when a graph run is finished.
- records
- The KV records of the orchestrator and their keys. Every key has the
swale/prefix, and every value is JSON written as an absolute value. - request
- The request: an operation on a running daemon, written to the object store by a command and applied by the daemon.
- retention
- The retention of the records: the pass that removes the records of a settled graph run and the request records past the retention window.
- scheduler
- The scheduler: it starts graph runs, submits a node when its upstreams satisfy its trigger rule, and settles the state of a graph run.
- status
- The view of a deployment from another process: the graphs, the graph runs, the state of every node of a graph run and the queues.
- store
- The layout of the object store: every path the process writes is within
the store prefix of the
--storeURL.store_pathjoins the prefix the one way, and anObjectPrefixis the objects within one such path. - task
- The identity of a task instance: the graph, the partition, the node and
the rerun count, which together form the run id and the
swale.headers of the run. - template
- The substitution language of a parameter string. A reference is written
between
{{and}}, and the language is substitution only.
Enums§
- Error
- The failure of a definition load.