#![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_LIFECYCLE_ENTRY, ROUTINE_LIFECYCLE_EVENTS};
const LIFECYCLE_SQL: &str = polyc_query_model::statements::ROUTINE_LIFECYCLE_SQL;
fn schema(table: polyc_projection::family::TableId) -> arrow::datatypes::SchemaRef {
use polyc_projection::family::LogicalType;
let declared = ROUTINE_LIFECYCLE_ENTRY
.table(table)
.expect("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 lifecycle_schema = schema(ROUTINE_LIFECYCLE_EVENTS);
let signer = {
let mut builder = FixedSizeBinaryBuilder::with_capacity(5, 32);
for _ in 0..5 {
builder.append_value([7_u8; 32]).unwrap();
}
builder.finish()
};
let lifecycle = RecordBatch::try_new(
Arc::clone(&lifecycle_schema),
vec![
Arc::new(StringArray::from(vec!["routine-scheduler"; 5])),
Arc::new(signer.clone()),
Arc::new(UInt64Array::from(vec![0_u64, 1, 2, 3, 4])),
Arc::new(StringArray::from(vec![
"routine_created",
"routine_paused",
"routine_resumed",
"routine_scope_changed",
"routine_paused",
])),
Arc::new(StringArray::from(vec!["r-1", "r-1", "r-1", "r-1", "r-1"])),
Arc::new(StringArray::from(vec![
"persona-a",
"persona-a",
"persona-a",
"persona-a",
"persona-b",
])),
Arc::new(StringArray::from(vec![
"conv-1", "conv-2", "conv-3", "conv-4", "conv-5",
])),
Arc::new(UInt64Array::from(vec![100_u64, 300, 300, 400, 200])),
Arc::new(BooleanArray::from(vec![false, true, false, false, false])),
Arc::new(StringArray::from(vec!["", "maintenance", "", "", ""])),
Arc::new(BooleanArray::from(vec![false, false, false, true, false])),
Arc::new(StringArray::from(vec!["", "", "", "public", ""])),
Arc::new(BooleanArray::from(vec![false, true, true, true, true])),
Arc::new(StringArray::from(vec!["", "chat", "rpc", "chat", "rpc"])),
Arc::new(signer),
],
)
.unwrap();
ctx.register_table(
"routine_lifecycle",
Arc::new(MemTable::try_new(lifecycle_schema, vec![vec![lifecycle]]).unwrap()),
)
.unwrap();
ctx
}
#[tokio::test]
async fn lifecycle_sql_restores_nulls_and_orders_by_clock_then_position() {
let ctx = context();
let df = ctx.sql(LIFECYCLE_SQL).await.unwrap();
let batches = df.collect().await.unwrap();
assert_eq!(batches.len(), 1);
let batch = &batches[0];
assert_eq!(batch.num_rows(), 5);
let phases = batch
.column_by_name("phase")
.unwrap()
.as_any()
.downcast_ref::<StringArray>()
.unwrap();
let at_ms = batch
.column_by_name("at_ms")
.unwrap()
.as_any()
.downcast_ref::<UInt64Array>()
.unwrap();
let reasons = batch
.column_by_name("reason")
.unwrap()
.as_any()
.downcast_ref::<StringArray>()
.unwrap();
let scopes = batch
.column_by_name("scope")
.unwrap()
.as_any()
.downcast_ref::<StringArray>()
.unwrap();
let ordered_phases: Vec<Option<&str>> = phases.iter().collect();
assert_eq!(
ordered_phases,
vec![
Some("routine_created"),
Some("routine_paused"),
Some("routine_paused"),
Some("routine_resumed"),
Some("routine_scope_changed"),
]
);
let ordered_clocks: Vec<Option<u64>> = at_ms.iter().collect();
assert_eq!(
ordered_clocks,
vec![Some(100), Some(200), Some(300), Some(300), Some(400)]
);
let ordered_reasons: Vec<Option<&str>> = reasons.iter().collect();
assert_eq!(
ordered_reasons,
vec![None, None, Some("maintenance"), None, None]
);
let ordered_scopes: Vec<Option<&str>> = scopes.iter().collect();
assert_eq!(ordered_scopes, vec![None, None, None, None, Some("public")]);
}