#![cfg(feature = "store-postgres")]
use std::collections::HashMap;
use std::env::var;
use chrono::{TimeDelta, Utc};
use ironflow_store::postgres::PostgresStore;
use ironflow_store::prelude::*;
use ironflow_store::store::RunStore;
use serde_json::json;
use tokio::sync::Mutex;
use uuid::Uuid;
static SERIAL: Mutex<()> = Mutex::const_new(());
async fn get_store() -> PostgresStore {
let url = var("DATABASE_URL").expect("DATABASE_URL must be set");
PostgresStore::new(&url)
.await
.expect("failed to connect to PostgreSQL")
}
fn new_run(priority: i16) -> NewRun {
NewRun {
workflow_name: "priority".to_string(),
trigger: TriggerKind::Manual,
payload: json!({}),
max_retries: 0,
handler_version: None,
labels: HashMap::new(),
scheduled_at: None,
created_by: None,
idempotency_key: None,
concurrency_key: None,
priority,
concurrency_limits: Vec::new(),
max_cost_usd: None,
worker_tags: Vec::new(),
}
}
async fn create(store: &PostgresStore, req: NewRun) -> Run {
store.create_run(req).await.unwrap().into_run()
}
async fn pick_id(store: &PostgresStore) -> Option<Uuid> {
store.pick_next_pending(None).await.unwrap().map(|r| r.id)
}
async fn drain_pending(store: &PostgresStore) {
while store.pick_next_pending(None).await.unwrap().is_some() {}
}
async fn finish(store: &PostgresStore, ids: &[Uuid]) {
for id in ids {
let run = store.get_run(*id).await.unwrap().unwrap();
let target = match run.status.state {
RunStatus::Pending | RunStatus::Sleeping => RunStatus::Cancelled,
RunStatus::Running => RunStatus::Completed,
_ => continue,
};
store.update_run_status(*id, target).await.unwrap();
}
}
#[tokio::test]
#[ignore = "requires DATABASE_URL"]
async fn priority_orders_the_queue_then_fifo() {
let _serial = SERIAL.lock().await;
let store = get_store().await;
drain_pending(&store).await;
let default_old = create(&store, new_run(0)).await;
let background = create(&store, new_run(-10)).await;
let urgent = create(&store, new_run(50)).await;
let default_young = create(&store, new_run(0)).await;
let mut order = Vec::new();
for _ in 0..4 {
order.push(pick_id(&store).await.expect("a run is due"));
}
assert_eq!(
order,
vec![urgent.id, default_old.id, default_young.id, background.id]
);
assert_eq!(pick_id(&store).await, None);
finish(&store, &order).await;
}
#[tokio::test]
#[ignore = "requires DATABASE_URL"]
async fn priority_does_not_bypass_scheduled_at() {
let _serial = SERIAL.lock().await;
let store = get_store().await;
drain_pending(&store).await;
let deferred = create(
&store,
NewRun {
scheduled_at: Some(Utc::now() + TimeDelta::seconds(3600)),
..new_run(100)
},
)
.await;
let due = create(&store, new_run(-100)).await;
assert_eq!(pick_id(&store).await, Some(due.id));
assert_eq!(pick_id(&store).await, None);
finish(&store, &[deferred.id, due.id]).await;
}
#[tokio::test]
#[ignore = "requires DATABASE_URL"]
async fn priority_round_trips_and_defaults_to_zero() {
let _serial = SERIAL.lock().await;
let store = get_store().await;
let default = create(&store, new_run(0)).await;
let urgent = create(&store, new_run(MAX_PRIORITY)).await;
let background = create(&store, new_run(MIN_PRIORITY)).await;
assert_eq!(
store.get_run(default.id).await.unwrap().unwrap().priority,
0
);
assert_eq!(
store.get_run(urgent.id).await.unwrap().unwrap().priority,
MAX_PRIORITY
);
assert_eq!(
store
.get_run(background.id)
.await
.unwrap()
.unwrap()
.priority,
MIN_PRIORITY
);
finish(&store, &[default.id, urgent.id, background.id]).await;
}
#[tokio::test]
#[ignore = "requires DATABASE_URL"]
async fn priority_check_constraint_rejects_out_of_range() {
let _serial = SERIAL.lock().await;
let store = get_store().await;
for priority in [MAX_PRIORITY + 1, MIN_PRIORITY - 1] {
let err = store.create_run(new_run(priority)).await.unwrap_err();
assert!(matches!(err, StoreError::Database(_)), "{err:?}");
}
}
#[tokio::test]
#[ignore = "requires DATABASE_URL"]
async fn priority_filter_lists_exact_matches() {
let _serial = SERIAL.lock().await;
let store = get_store().await;
let marked = create(&store, new_run(-37)).await;
let other = create(&store, new_run(37)).await;
let page = store
.list_runs(
RunFilter {
priority: Some(-37),
..RunFilter::default()
},
1,
100,
)
.await
.unwrap();
assert!(page.items.iter().any(|r| r.id == marked.id));
assert!(page.items.iter().all(|r| r.priority == -37));
finish(&store, &[marked.id, other.id]).await;
}
#[tokio::test]
#[ignore = "requires DATABASE_URL"]
async fn priority_schedule_firing_carries_the_schedule_priority() {
let _serial = SERIAL.lock().await;
let store = get_store().await;
let schedule = store
.create_schedule(NewSchedule {
workflow_name: format!("priority-schedule-{}", Uuid::now_v7()),
cron_expression: "0 0 * * * *".to_string(),
inputs: json!({}),
source: ScheduleSource::Api,
priority: 25,
created_by_user_id: None,
next_trigger_at: Some(Utc::now() - TimeDelta::seconds(60)),
})
.await
.unwrap();
assert_eq!(schedule.priority, 25);
let occurrence = schedule.next_trigger_at.expect("due");
let firing = store
.fire_due_schedule(
schedule.id,
occurrence,
ScheduleNext::At(Utc::now() + TimeDelta::seconds(3600)),
)
.await
.unwrap()
.expect("due occurrence fires");
let run_id = firing.run.run().id;
assert_eq!(firing.run.run().priority, 25);
let updated = store
.update_schedule(
schedule.id,
ScheduleUpdate {
priority: Some(-5),
..ScheduleUpdate::default()
},
)
.await
.unwrap();
assert_eq!(updated.priority, -5);
finish(&store, &[run_id]).await;
store.delete_schedule(schedule.id).await.unwrap();
}