#![allow(clippy::unwrap_used)]
use std::sync::Arc;
use arrow::array::{BooleanArray, FixedSizeBinaryBuilder, RecordBatch, StringArray, UInt64Array};
use datafusion::datasource::MemTable;
use datafusion::prelude::SessionContext;
use polyc_projection::family::{ROUTINE_FIRES, ROUTINE_LIFECYCLE_ENTRY};
const FIRES_SQL: &str = polyc_query_model::statements::ROUTINE_FIRES_SQL;
const FIRE_SQL: &str = polyc_query_model::statements::ROUTINE_FIRE_SQL;
fn schema() -> arrow::datatypes::SchemaRef {
use polyc_projection::family::LogicalType;
let declared = ROUTINE_LIFECYCLE_ENTRY
.table(ROUTINE_FIRES)
.expect("fires table declared");
let fields: Vec<arrow::datatypes::Field> = declared
.fields()
.iter()
.map(|field| {
let data_type = match field.logical_type() {
LogicalType::Utf8 => arrow::datatypes::DataType::Utf8,
LogicalType::FixedBytes { len } => {
arrow::datatypes::DataType::FixedSizeBinary(i32::try_from(len).unwrap())
}
LogicalType::UInt64 => arrow::datatypes::DataType::UInt64,
LogicalType::Boolean => arrow::datatypes::DataType::Boolean,
};
arrow::datatypes::Field::new(field.name(), data_type, field.nullable())
})
.collect();
Arc::new(arrow::datatypes::Schema::new(fields))
}
fn context() -> SessionContext {
let ctx = SessionContext::new();
let fires_schema = schema();
let incarnation = {
let mut builder = FixedSizeBinaryBuilder::with_capacity(3, 32);
for _ in 0..3 {
builder.append_value([3_u8; 32]).unwrap();
}
builder.finish()
};
let fires = RecordBatch::try_new(
Arc::clone(&fires_schema),
vec![
Arc::new(StringArray::from(vec!["routine-scheduler"; 3])),
Arc::new(incarnation),
Arc::new(UInt64Array::from(vec![0_u64, 1, 2])),
Arc::new(StringArray::from(vec!["r-1", "r-1", "r-2"])),
Arc::new(StringArray::from(vec!["uid-1", "uid-1", "uid-2"])),
Arc::new(StringArray::from(vec!["r-1-a", "r-1-b", "r-2-a"])),
Arc::new(UInt64Array::from(vec![100_u64, 200, 150])),
Arc::new(UInt64Array::from(vec![110_u64, 210, 160])),
Arc::new(BooleanArray::from(vec![true, false, true])),
Arc::new(StringArray::from(vec!["ok", "", "paused"])),
Arc::new(BooleanArray::from(vec![true, false, true])),
Arc::new(StringArray::from(vec!["[]", "", r#"["fs_write"]"#])),
],
)
.unwrap();
ctx.register_table(
"fires",
Arc::new(MemTable::try_new(fires_schema, vec![vec![fires]]).unwrap()),
)
.unwrap();
ctx
}
#[test]
fn the_list_statement_binds_uid_not_name() {
assert!(FIRES_SQL.contains("WHERE routine_uid = $1 AND routine_uid <> ''"));
assert!(!FIRES_SQL.contains("WHERE routine = $1"));
assert!(FIRES_SQL.contains("FROM fires"));
}
#[test]
fn the_get_statement_binds_uid_and_occurrence() {
assert!(FIRE_SQL.contains("WHERE routine_uid = $1 AND routine_uid <> '' AND occurrence = $2"));
assert!(!FIRE_SQL.contains("WHERE routine = $1"));
}
#[tokio::test]
async fn list_sql_restores_nulls_filters_and_orders_newest_first() {
let ctx = context();
let sql = FIRES_SQL.replace("$1", "'uid-1'");
let df = ctx.sql(&sql).await.unwrap();
let batches = df.collect().await.unwrap();
assert_eq!(batches.len(), 1);
let batch = &batches[0];
assert_eq!(batch.num_rows(), 2);
let occurrences = batch
.column_by_name("occurrence")
.unwrap()
.as_any()
.downcast_ref::<StringArray>()
.unwrap();
let outcomes = batch
.column_by_name("outcome")
.unwrap()
.as_any()
.downcast_ref::<StringArray>()
.unwrap();
let ordered: Vec<Option<&str>> = occurrences.iter().collect();
assert_eq!(ordered, vec![Some("r-1-b"), Some("r-1-a")]);
let ordered_outcomes: Vec<Option<&str>> = outcomes.iter().collect();
assert_eq!(ordered_outcomes, vec![None, Some("ok")]);
}
#[tokio::test]
async fn get_sql_selects_one_occurrence() {
let ctx = context();
let sql = FIRE_SQL.replace("$1", "'uid-1'").replace("$2", "'r-1-a'");
let df = ctx.sql(&sql).await.unwrap();
let batches = df.collect().await.unwrap();
assert_eq!(batches.len(), 1);
let batch = &batches[0];
assert_eq!(batch.num_rows(), 1);
let occurrences = batch
.column_by_name("occurrence")
.unwrap()
.as_any()
.downcast_ref::<StringArray>()
.unwrap();
assert_eq!(occurrences.value(0), "r-1-a");
}