use std::sync::Arc;
use std::time::Duration;
use swale::JsonBytes;
use swale::records::{
EXPIRY_PREFIX, Expiring, GraphRunRecord, GraphRunState, NodeRecord, RecordStatus,
RequestOutcome, graph_key, graph_run_key, parse_expiry_key, request_key,
};
use swale::retention::ExpireReport;
use swale::scheduler::{Error, firing_headers};
use swale::{
Daemon, DaemonOptions, DefinitionStore, GraphRecord, OperatorSet, Partition, Pools, RecordHook,
Request, RequestId, RequestStore, Scheduler, SchedulerOptions, StatusReader, TRIGGERS_QUEUE,
};
use taquba::object_store::path::Path as ObjectPath;
use taquba::object_store::{ObjectStore, ObjectStoreExt};
use taquba::{MockClock, Queue};
use taquba_cron::{Backfill, BackfillStart, Schedule};
use tokio_util::sync::CancellationToken;
mod common;
fn definition_text(graph: &str, schedule: &str, asset: &str, pool: &str) -> String {
definition_with_catchup(graph, schedule, asset, pool, "3d")
}
fn definition_with_catchup(
graph: &str,
schedule: &str,
asset: &str,
pool: &str,
catchup: &str,
) -> String {
format!(
r#"
[graph]
name = "{graph}"
schedule = "{schedule}"
catchup = "{catchup}"
partition = "daily"
[[node]]
name = "extract"
produces = "{asset}"
operator = "subprocess"
pool = "{pool}"
[node.params]
argv = ["sh", "-c", "cat >/dev/null; printf '{{}}'"]
"#
)
}
fn ms(text: &str) -> u64 {
let time: chrono::DateTime<chrono::Utc> = text.parse().unwrap();
time.timestamp_millis() as u64
}
struct Harness {
objects: Arc<dyn ObjectStore>,
queue: Arc<Queue>,
clock: MockClock,
definitions: Arc<DefinitionStore>,
operators: Arc<OperatorSet>,
}
impl Harness {
async fn start(now: &str) -> Harness {
let clock = MockClock::new(ms(now));
let (objects, queue) = common::open_queue(clock.clone()).await;
let operators = Arc::new(OperatorSet::builtin());
let definitions = Arc::new(DefinitionStore::new(objects.clone(), "", operators.clone()));
Harness {
objects,
queue,
clock,
definitions,
operators,
}
}
fn requests(&self) -> RequestStore {
RequestStore::new(self.objects.clone(), "")
}
fn daemon(&self) -> (Arc<Daemon>, Arc<Scheduler>, Arc<Pools>) {
let hook = RecordHook::new(self.queue.clock());
let pools = Arc::new(
Pools::builder(
self.queue.clone(),
self.objects.clone(),
self.operators.clone(),
hook,
)
.poll_interval(Duration::from_millis(10))
.pool("default", 4)
.build(),
);
let scheduler = Arc::new(Scheduler::new(
self.queue.clone(),
self.definitions.clone(),
pools.clone(),
));
let daemon = Arc::new(Daemon::new(
self.queue.clone(),
scheduler.clone(),
pools.clone(),
self.requests(),
));
(daemon, scheduler, pools)
}
fn spawn(
&self,
) -> (
CancellationToken,
tokio::task::JoinHandle<Result<(), Error>>,
) {
self.spawn_with_retention(Duration::from_secs(90 * 86_400))
}
fn spawn_with_retention(
&self,
retention: Duration,
) -> (
CancellationToken,
tokio::task::JoinHandle<Result<(), Error>>,
) {
let (daemon, _, _) = self.daemon();
let stop = CancellationToken::new();
let shutdown = stop.clone().cancelled_owned();
let options = DaemonOptions {
scheduler: SchedulerOptions {
concurrency: 2,
poll_interval: Duration::from_millis(10),
reconcile_interval: Duration::from_secs(3600),
},
sync_interval: Duration::from_millis(20),
retention,
};
let handle = tokio::spawn(async move { daemon.run(options, shutdown).await });
(stop, handle)
}
async fn graph_run(&self, partition: &str) -> Option<GraphRunRecord> {
let key = graph_run_key("orders", &Partition::new(partition).unwrap());
self.queue
.view()
.kv_get(&key)
.await
.unwrap()
.map(|b| GraphRunRecord::from_bytes(&b).unwrap())
}
async fn wait_for_complete(&self, partition: &str) {
common::wait_until(
&format!("the graph run of {partition} never completed"),
async || {
self.graph_run(partition)
.await
.filter(|run| run.state == GraphRunState::Complete)
},
)
.await;
}
async fn graph_record(&self, graph: &str) -> Option<GraphRecord> {
self.queue
.view()
.kv_get(&graph_key(graph))
.await
.unwrap()
.map(|b| GraphRecord::from_bytes(&b).unwrap())
}
async fn wait_for_extract(&self, partition: &str, status: RecordStatus) -> NodeRecord {
let key = format!("swale/assets/orders_raw/{partition}");
common::wait_until(
&format!("extract of {partition} never reached {status}"),
async || {
let bytes = self.queue.view().kv_get(key.as_bytes()).await.unwrap()?;
Some(NodeRecord::from_bytes(&bytes).unwrap()).filter(|r| r.status == status)
},
)
.await
}
}
#[tokio::test(flavor = "multi_thread")]
async fn a_new_graph_catches_up_firings_run_their_interval_start_and_downtime_is_backfilled() {
let h = Harness::start("2026-09-16T01:59:59.950Z").await;
h.definitions
.publish(&definition_text(
"orders",
"0 2 * * *",
"orders_raw",
"default",
))
.await
.unwrap();
h.definitions
.publish(&definition_with_catchup(
"far",
"0 2 * * *",
"far_raw",
"default",
"99999999999d",
))
.await
.unwrap();
let (stop, handle) = h.spawn();
for partition in ["20260912", "20260913", "20260914"] {
h.wait_for_complete(partition).await;
}
assert!(h.graph_run("20260911").await.is_none());
assert!(h.graph_run("20260915").await.is_none());
h.clock.advance(Duration::from_secs(1));
h.wait_for_complete("20260915").await;
let edit = h
.definitions
.publish(&definition_text(
"orders",
"0 */12 * * *",
"orders_raw",
"default",
))
.await
.unwrap();
h.clock.advance(Duration::from_secs(10 * 3600));
h.wait_for_complete("20260916").await;
assert_eq!(h.graph_run("20260916").await.unwrap().definition, edit.hash);
stop.cancel();
handle.await.unwrap().unwrap();
h.clock.advance(Duration::from_secs(2 * 86_400));
assert!(h.graph_run("20260917").await.is_none());
let (stop, handle) = h.spawn();
for partition in ["20260917", "20260918"] {
h.wait_for_complete(partition).await;
}
stop.cancel();
handle.await.unwrap().unwrap();
}
#[tokio::test]
async fn a_sync_pass_adopts_a_changed_pointer_and_refuses_a_definition_that_cannot_run() {
let h = Harness::start("2026-09-16T12:00:00Z").await;
let (daemon, scheduler, _) = h.daemon();
let schedule = |expression: &str| {
Schedule::new(
"orders",
expression.parse().unwrap(),
TRIGGERS_QUEUE,
Vec::new(),
)
.headers(firing_headers("orders"))
.backfill(Some(Backfill {
lookback: Duration::from_secs(3 * 86_400),
start: BackfillStart::Lookback,
}))
};
assert!(matches!(
scheduler.handle_trigger("orders", Some(0)).await,
Err(Error::UnknownGraph(graph)) if graph == "orders"
));
let v1 = h
.definitions
.publish(&definition_text(
"orders",
"0 2 * * *",
"orders_raw",
"default",
))
.await
.unwrap();
let report = daemon.sync().await.unwrap();
assert_eq!(report.adopted, [("orders".to_string(), v1.hash.clone())]);
assert_eq!(report.schedules, [schedule("0 2 * * *")]);
let record = h.graph_record("orders").await.unwrap();
assert_eq!(record.definition, v1.hash);
assert_eq!(record.adopted_at_ms, ms("2026-09-16T12:00:00Z"));
let report = daemon.sync().await.unwrap();
assert!(report.adopted.is_empty());
assert_eq!(report.schedules, [schedule("0 2 * * *")]);
assert!(matches!(
scheduler.handle_trigger("orders", None).await,
Err(Error::NoPartition(graph)) if graph == "orders"
));
h.definitions
.publish(&definition_text("orders", "0 4 * * *", "orders_raw", "gpu"))
.await
.unwrap();
let report = daemon.sync().await.unwrap();
assert!(report.adopted.is_empty());
assert_eq!(report.refused.len(), 1);
assert!(report.refused[0].1.contains("pool `gpu`"), "{report:?}");
assert_eq!(report.schedules, [schedule("0 2 * * *")]);
assert_eq!(h.graph_record("orders").await.unwrap().definition, v1.hash);
h.clock.advance(Duration::from_secs(3600));
let v3 = h
.definitions
.publish(&definition_text(
"orders",
"0 3 * * *",
"orders_raw",
"default",
))
.await
.unwrap();
let report = daemon.sync().await.unwrap();
assert_eq!(report.adopted, [("orders".to_string(), v3.hash.clone())]);
assert_eq!(report.schedules, [schedule("0 3 * * *")]);
let record = h.graph_record("orders").await.unwrap();
assert_eq!(record.adopted_at_ms, ms("2026-09-16T13:00:00Z"));
let partitions = [Partition::new("20260901").unwrap()];
let started = scheduler.start_runs("orders", &partitions).await.unwrap();
assert_eq!(started, partitions);
let started = scheduler.start_runs("orders", &partitions).await.unwrap();
assert!(started.is_empty());
assert_eq!(h.graph_run("20260901").await.unwrap().definition, v3.hash);
let (hash, _) = h
.definitions
.put(&definition_text(
"other",
"0 5 * * *",
"orders_raw",
"default",
))
.await
.unwrap();
h.objects
.put(
&ObjectPath::from("definitions/current/other"),
hash.into_bytes().into(),
)
.await
.unwrap();
let report = daemon.sync().await.unwrap();
assert!(
report.refused[0].1.contains("produced by graph `orders`"),
"{report:?}"
);
assert!(h.graph_record("other").await.is_none());
}
#[tokio::test(flavor = "multi_thread")]
async fn a_request_in_the_store_is_applied_once_and_its_record_is_readable() {
let h = Harness::start("2026-09-16T12:00:00Z").await;
let dir = std::path::Path::new(env!("CARGO_TARGET_TMPDIR")).join("requests");
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
let marker = dir.join("marker");
let definition = format!(
r#"
[graph]
name = "orders"
partition = "daily"
[[node]]
name = "extract"
produces = "orders_raw"
operator = "subprocess"
[node.params]
argv = ["sh", "-c", "cat >/dev/null; test -f {} || exit 3; printf '{{}}'"]
"#,
marker.display()
);
h.definitions.publish(&definition).await.unwrap();
let (daemon, _, pools) = h.daemon();
let stop = CancellationToken::new();
let pool_handles = pools.spawn(&stop);
daemon.sync().await.unwrap();
let requests = h.requests();
let partition = |key: &str| Partition::new(key).unwrap();
let id = |n: u64| RequestId::new(format!("req-{n}")).unwrap();
let start = Request::Start {
graph: "orders".into(),
partitions: vec![partition("20260915"), partition("20260916")],
};
requests.submit(&id(1), &start).await.unwrap();
let report = daemon.apply_requests().await.unwrap();
assert_eq!(
report.applied,
[(
id(1),
RequestOutcome::Started {
partitions: vec![partition("20260915"), partition("20260916")]
}
)]
);
assert!(requests.list().await.unwrap().is_empty());
for key in ["20260915", "20260916"] {
assert_eq!(h.graph_run(key).await.unwrap().state, GraphRunState::Active);
h.wait_for_extract(key, RecordStatus::Failed).await;
}
requests
.submit(
&id(2),
&Request::Start {
graph: "orders".into(),
partitions: vec![partition("20260915"), partition("20260917")],
},
)
.await
.unwrap();
let report = daemon.apply_requests().await.unwrap();
assert_eq!(
report.applied[0].1,
RequestOutcome::Started {
partitions: vec![partition("20260917")]
}
);
std::fs::write(&marker, b"").unwrap();
h.clock.advance(Duration::from_secs(60));
let rerun = Request::Rerun {
graph: "orders".into(),
partition: partition("20260915"),
node: "extract".into(),
};
requests.submit(&id(3), &rerun).await.unwrap();
let report = daemon.apply_requests().await.unwrap();
assert_eq!(
report.applied,
[(
id(3),
RequestOutcome::Rerun {
run_id: "orders-20260915-extract-r1".into()
}
)]
);
let record = h
.wait_for_extract("20260915", RecordStatus::Succeeded)
.await;
assert_eq!(record.rerun, 1);
let bytes = h
.queue
.view()
.kv_get(&request_key(&id(3)))
.await
.unwrap()
.unwrap();
let recorded = swale::RequestRecord::from_bytes(&bytes).unwrap();
assert_eq!(recorded.request, rerun);
assert_eq!(recorded.handled_at_ms, ms("2026-09-16T12:01:00Z"));
h.clock.advance(Duration::from_secs(60));
requests.submit(&id(3), &rerun).await.unwrap();
let report = daemon.apply_requests().await.unwrap();
assert!(report.applied.is_empty());
assert!(requests.list().await.unwrap().is_empty());
let bytes = h
.queue
.view()
.kv_get(&request_key(&id(3)))
.await
.unwrap()
.unwrap();
assert_eq!(swale::RequestRecord::from_bytes(&bytes).unwrap(), recorded);
tokio::time::sleep(Duration::from_millis(200)).await;
assert_eq!(
h.wait_for_extract("20260915", RecordStatus::Succeeded)
.await
.rerun,
1
);
requests.submit(&id(4), &rerun).await.unwrap();
let report = daemon.apply_requests().await.unwrap();
assert_eq!(
report.applied,
[(
id(4),
RequestOutcome::Rerun {
run_id: "orders-20260915-extract-r2".into()
}
)]
);
let refused = [
(
id(5),
Request::Rerun {
graph: "orders".into(),
partition: partition("20260915"),
node: "nope".into(),
},
),
(
id(6),
Request::Start {
graph: "nope".into(),
partitions: vec![partition("20260915")],
},
),
(
id(7),
Request::Cancel {
graph: "orders".into(),
partition: partition("20260901"),
},
),
];
for (id, request) in &refused {
requests.submit(id, request).await.unwrap();
}
h.objects
.put(
&ObjectPath::from("requests/not-a-request"),
b"nope".to_vec().into(),
)
.await
.unwrap();
let report = daemon.apply_requests().await.unwrap();
let reasons: Vec<String> = report
.applied
.iter()
.map(|(_, outcome)| match outcome {
RequestOutcome::Refused { reason } => reason.clone(),
other => panic!("{other:?}"),
})
.collect();
assert_eq!(
reasons,
[
"graph `orders` does not have a node `nope`",
"graph `nope` does not have an adopted definition",
"graph `orders` does not have an active run for partition `20260901`",
]
);
assert!(requests.list().await.unwrap().is_empty());
assert!(
h.queue
.view()
.kv_get(b"swale/requests/not-a-request")
.await
.unwrap()
.is_none()
);
requests
.submit(
&id(8),
&Request::Start {
graph: "orders".into(),
partitions: vec![partition("20260918")],
},
)
.await
.unwrap();
requests
.submit(
&id(9),
&Request::Cancel {
graph: "orders".into(),
partition: partition("20260918"),
},
)
.await
.unwrap();
let report = daemon.apply_requests().await.unwrap();
assert_eq!(report.applied[1], (id(9), RequestOutcome::Cancelled));
assert_eq!(
h.graph_run("20260918").await.unwrap().state,
GraphRunState::Cancelled
);
let read = common::wait_until("no reader saw the request record", async || {
let reader =
StatusReader::open(h.objects.clone(), common::QUEUE_PATH, h.definitions.clone())
.await
.unwrap();
let record = reader.request(&id(3)).await.unwrap();
reader.close().await.unwrap();
record
})
.await;
assert_eq!(read, recorded);
stop.cancel();
for handle in pool_handles {
let _ = handle.wait().await;
}
}
#[tokio::test(flavor = "multi_thread")]
async fn a_retention_pass_removes_the_records_of_settled_runs_and_of_requests_past_the_window() {
let h = Harness::start("2026-09-16T12:00:00Z").await;
h.definitions
.publish(&definition_text(
"orders",
"0 2 * * *",
"orders_raw",
"default",
))
.await
.unwrap();
let (daemon, scheduler, pools) = h.daemon();
let stop = CancellationToken::new();
let pool_handles = pools.spawn(&stop);
let scheduler_handle = scheduler.clone().spawn(
SchedulerOptions {
concurrency: 2,
poll_interval: Duration::from_millis(10),
reconcile_interval: Duration::from_secs(3600),
},
stop.clone().cancelled_owned(),
);
daemon.sync().await.unwrap();
let requests = h.requests();
let partition = |key: &str| Partition::new(key).unwrap();
let id = |n: u64| RequestId::new(format!("req-{n}")).unwrap();
let start = |keys: &[&str]| Request::Start {
graph: "orders".into(),
partitions: keys.iter().map(|key| partition(key)).collect(),
};
let retention = Duration::from_secs(90 * 86_400);
let record_key = |key: &str| format!("swale/assets/orders_raw/{key}");
let has_key = async |key: &[u8]| h.queue.view().kv_get(key).await.unwrap().is_some();
let entries = async || {
let page = h
.queue
.view()
.kv_scan(EXPIRY_PREFIX.as_bytes(), .., 100)
.await
.unwrap();
page.entries
.into_iter()
.map(|(key, _)| parse_expiry_key(&key).unwrap())
.collect::<Vec<_>>()
};
let t0 = ms("2026-09-16T12:00:00Z");
let day = 86_400_000;
let run_entry = |time: u64, key: &str| {
(
time,
Expiring::Run {
graph: "orders".into(),
partition: partition(key),
},
)
};
let request_entry = |time: u64, n: u64| (time, Expiring::Request(id(n)));
requests
.submit(&id(1), &start(&["20260910", "20260911"]))
.await
.unwrap();
daemon.apply_requests().await.unwrap();
for key in ["20260910", "20260911"] {
h.wait_for_complete(key).await;
assert_eq!(h.graph_run(key).await.unwrap().settled_at_ms, Some(t0));
}
assert_eq!(
entries().await,
[
request_entry(t0, 1),
run_entry(t0, "20260910"),
run_entry(t0, "20260911"),
]
);
h.clock.advance(Duration::from_secs(30 * 86_400));
requests
.submit(&id(2), &start(&["20260912"]))
.await
.unwrap();
daemon.apply_requests().await.unwrap();
h.wait_for_complete("20260912").await;
h.clock.advance(Duration::from_secs(30 * 86_400));
requests
.submit(
&id(3),
&Request::Rerun {
graph: "orders".into(),
partition: partition("20260912"),
node: "extract".into(),
},
)
.await
.unwrap();
daemon.apply_requests().await.unwrap();
let run = h.graph_run("20260912").await.unwrap();
assert_eq!(run.state, GraphRunState::Active);
assert_eq!(run.settled_at_ms, None);
h.wait_for_complete("20260912").await;
assert_eq!(
h.graph_run("20260912").await.unwrap().settled_at_ms,
Some(t0 + 60 * day)
);
assert_eq!(
entries().await[3..],
[
request_entry(t0 + 30 * day, 2),
run_entry(t0 + 30 * day, "20260912"),
request_entry(t0 + 60 * day, 3),
run_entry(t0 + 60 * day, "20260912"),
]
);
h.clock.advance(Duration::from_secs(30 * 86_400));
let report = scheduler.expire(retention).await.unwrap();
assert_eq!(
report,
ExpireReport {
runs: 2,
requests: 1
}
);
for key in ["20260910", "20260911"] {
assert!(h.graph_run(key).await.is_none(), "{key}");
assert!(!has_key(record_key(key).as_bytes()).await, "{key}");
}
assert!(h.graph_run("20260912").await.is_some());
assert!(has_key(record_key("20260912").as_bytes()).await);
assert!(!has_key(&request_key(&id(1))).await);
assert!(has_key(&request_key(&id(2))).await);
assert_eq!(entries().await.len(), 4);
assert_eq!(
scheduler.expire(retention).await.unwrap(),
ExpireReport::default()
);
h.clock.advance(Duration::from_secs(30 * 86_400));
let report = scheduler.expire(retention).await.unwrap();
assert_eq!(
report,
ExpireReport {
runs: 0,
requests: 1
}
);
assert!(h.graph_run("20260912").await.is_some());
assert!(!has_key(&request_key(&id(2))).await);
assert_eq!(
entries().await,
[
request_entry(t0 + 60 * day, 3),
run_entry(t0 + 60 * day, "20260912"),
]
);
requests
.submit(&id(4), &start(&["20260910"]))
.await
.unwrap();
let report = daemon.apply_requests().await.unwrap();
assert_eq!(
report.applied,
[(
id(4),
RequestOutcome::Started {
partitions: vec![partition("20260910")]
}
)]
);
h.wait_for_complete("20260910").await;
assert_eq!(
h.wait_for_extract("20260910", RecordStatus::Succeeded)
.await
.rerun,
0
);
stop.cancel();
for handle in pool_handles {
let _ = handle.wait().await;
}
scheduler_handle.wait().await.unwrap();
h.clock.advance(Duration::from_secs(3600));
let (stop, handle) = h.spawn_with_retention(Duration::from_secs(3600));
common::wait_until("the daemon never expired the run", async || {
h.graph_run("20260912").await.is_none().then_some(())
})
.await;
stop.cancel();
handle.await.unwrap().unwrap();
}