use crate::meta_storage::mem::MemMetaStorage;
use crate::meta_storage::tests::utils::{
TicketView, dim_i, psql_backend, task_alpha, task_beta, task_delta, task_gamma, task_over_j,
views,
};
use crate::meta_storage::{
MetaBackend, MetaClientApi, MetaConnApi, MetaResolutionApi, MetaTicketApi,
};
use crate::schema::{Job, Resolution, Ticket, TicketStatus};
const SCHEMA: &str = "operon_differential";
#[derive(Debug, PartialEq, Eq)]
enum Observed {
Tickets(Vec<TicketView>),
Status((i64, i64, i64)),
Resolution(Option<(Vec<usize>, usize)>),
}
type Step = (&'static str, Observed);
async fn exercise<MSto: MetaBackend>(backend: &MSto) -> Vec<Step> {
let conn = backend.scheduler_conn().await.expect("scheduler conn");
let client = conn.as_client();
let mut log = Vec::new();
client.init_schema().await.expect("init_schema");
client
.init_dimension_hash()
.await
.expect("init_dimension_hash");
client.init_ticket_hash().await.expect("init_ticket_hash");
let _ = client
.resolution(dim_i())
.init()
.await
.expect("resolution init");
client
.init_ticket_summary()
.await
.expect("init_ticket_summary");
client
.init_ticket_status_type()
.await
.expect("init_ticket_status_type");
let _ = client
.ticket(task_alpha())
.init()
.await
.expect("alpha init");
let _ = client.ticket(task_beta()).init().await.expect("beta init");
let _ = client
.ticket(task_gamma())
.init()
.await
.expect("gamma init");
let _ = client
.ticket(task_delta())
.init()
.await
.expect("delta init");
let _ = client.init_footprint().await.expect("init_footprint");
client
.resolution(dim_i())
.clear()
.await
.expect("resolution clear");
client
.ticket(task_alpha())
.clear()
.await
.expect("alpha clear");
client
.ticket(task_beta())
.clear()
.await
.expect("beta clear");
client
.ticket(task_gamma())
.clear()
.await
.expect("gamma clear");
client
.ticket(task_delta())
.clear()
.await
.expect("delta clear");
client.clear_footprint().await.expect("clear_footprint");
client
.ticket(task_alpha())
.put(Ticket::new(0))
.await
.expect("alpha put");
client
.ticket(task_beta())
.put(Ticket::new(1))
.await
.expect("beta put");
log.push((
"alpha status after put",
Observed::Status(
client
.ticket(task_alpha())
.get_status()
.await
.expect("status"),
),
));
log.push((
"beta status after put",
Observed::Status(
client
.ticket(task_beta())
.get_status()
.await
.expect("status"),
),
));
client
.ticket(task_beta())
.put(Ticket::new(99))
.await
.expect("beta duplicate put");
log.push((
"beta waiting after duplicate put",
Observed::Tickets(views(
&client
.ticket(task_beta())
.get_all(TicketStatus::Waiting)
.await
.expect("get_all"),
)),
));
client
.resolution(dim_i())
.put(Resolution {
coordinate: [],
ub: 3,
})
.await
.expect("resolution put");
let resolution = client
.resolution(dim_i())
.get([])
.await
.expect("resolution get");
log.push((
"resolution i",
Observed::Resolution(resolution.map(|res| (res.coordinate.to_vec(), res.ub))),
));
let popped = client
.ticket(task_beta())
.explode::<0, 0>(
dim_i(),
Resolution {
coordinate: [],
ub: 3,
},
)
.await
.expect("beta explode");
log.push(("beta explode returned", Observed::Tickets(views(&popped))));
log.push((
"beta status after explode",
Observed::Status(
client
.ticket(task_beta())
.get_status()
.await
.expect("status"),
),
));
log.push((
"beta waiting after explode",
Observed::Tickets(views(
&client
.ticket(task_beta())
.get_all(TicketStatus::Waiting)
.await
.expect("get_all"),
)),
));
let raised = client
.ticket(task_beta())
.raise_deps_quota::<0>(task_alpha(), Ticket::new(0), &[], 2)
.await
.expect("beta raise_deps_quota");
log.push((
"beta raise_deps_quota returned",
Observed::Tickets(views(&raised)),
));
log.push((
"beta waiting after raise_deps_quota",
Observed::Tickets(views(
&client
.ticket(task_beta())
.get_all(TicketStatus::Waiting)
.await
.expect("get_all"),
)),
));
let first = client
.ticket(task_beta())
.raise_deps_done::<0>(task_alpha(), Job { coordinate: [] }, &[])
.await
.expect("beta raise_deps_done");
log.push((
"beta raise_deps_done returned (short of quota)",
Observed::Tickets(views(&first)),
));
let second = client
.ticket(task_beta())
.raise_deps_done::<0>(task_alpha(), Job { coordinate: [] }, &[])
.await
.expect("beta raise_deps_done");
log.push((
"beta raise_deps_done returned (quota met)",
Observed::Tickets(views(&second)),
));
log.push((
"beta status after raise_deps_done",
Observed::Status(
client
.ticket(task_beta())
.get_status()
.await
.expect("status"),
),
));
client
.ticket(task_beta())
.mark_done(Job { coordinate: [0] })
.await
.expect("beta mark_done");
log.push((
"beta status after mark_done",
Observed::Status(
client
.ticket(task_beta())
.get_status()
.await
.expect("status"),
),
));
log.push((
"beta done after mark_done",
Observed::Tickets(views(
&client
.ticket(task_beta())
.get_all(TicketStatus::Done)
.await
.expect("get_all"),
)),
));
for coord in 0..3 {
client
.ticket(task_gamma())
.put(Ticket::new(1).with_coordinate::<0>(coord))
.await
.expect("gamma put");
}
let pinned = client
.ticket(task_gamma())
.raise_deps_done::<1>(task_beta(), Job { coordinate: [0] }, &[])
.await
.expect("gamma raise_deps_done");
log.push((
"gamma raise_deps_done returned (pinned to i = 0)",
Observed::Tickets(views(&pinned)),
));
log.push((
"gamma status after pinned raise_deps_done",
Observed::Status(
client
.ticket(task_gamma())
.get_status()
.await
.expect("status"),
),
));
log.push((
"gamma waiting after pinned raise_deps_done",
Observed::Tickets(views(
&client
.ticket(task_gamma())
.get_all(TicketStatus::Waiting)
.await
.expect("get_all"),
)),
));
for i in 0..2 {
for j in 0..2 {
client
.ticket(task_delta())
.put(
Ticket::new(2)
.with_coordinate::<0>(i)
.with_coordinate::<1>(j),
)
.await
.expect("delta put");
}
}
let over_i = client
.ticket(task_delta())
.raise_deps_done::<1>(task_gamma(), Job { coordinate: [0] }, &[])
.await
.expect("delta raise_deps_done over i");
log.push((
"delta raise_deps_done returned (pinned to i = 0, j free)",
Observed::Tickets(views(&over_i)),
));
log.push((
"delta waiting after raise_deps_done over i",
Observed::Tickets(views(
&client
.ticket(task_delta())
.get_all(TicketStatus::Waiting)
.await
.expect("get_all"),
)),
));
let over_j = client
.ticket(task_delta())
.raise_deps_done::<1>(task_over_j(), Job { coordinate: [0] }, &[])
.await
.expect("delta raise_deps_done over j");
log.push((
"delta raise_deps_done returned (pinned to j = 0, i free)",
Observed::Tickets(views(&over_j)),
));
log.push((
"delta status after both pinned raises",
Observed::Status(
client
.ticket(task_delta())
.get_status()
.await
.expect("status"),
),
));
log.push((
"delta queued after both pinned raises",
Observed::Tickets(views(
&client
.ticket(task_delta())
.get_all(TicketStatus::Queued)
.await
.expect("get_all"),
)),
));
client
.ticket(task_alpha())
.mark_done(Job { coordinate: [] })
.await
.expect("alpha mark_done");
log.push((
"alpha status after mark_done",
Observed::Status(
client
.ticket(task_alpha())
.get_status()
.await
.expect("status"),
),
));
log
}
#[tokio::test]
async fn mem_matches_psql() {
let Some(psql) = psql_backend(SCHEMA) else {
eprintln!("skipping: POSTGRES_URI is not set, so there is no Postgres to compare against");
return;
};
let from_psql = exercise(&psql).await;
let from_mem = exercise(&MemMetaStorage::default()).await;
for (psql_step, mem_step) in from_psql.iter().zip(&from_mem) {
pretty_assertions::assert_eq!(psql_step, mem_step, "backends drifted at: {}", psql_step.0);
}
pretty_assertions::assert_eq!(from_psql.len(), from_mem.len());
}