#![allow(clippy::unwrap_used)]
use std::sync::Arc;
use arrow::array::{BooleanArray, RecordBatch, StringArray, UInt64Array};
use datafusion::datasource::MemTable;
use datafusion::prelude::SessionContext;
use polyc_projection::family::{OBSERVED_ROUTINES_ENTRY, OBSERVED_ROUTINES_TABLE, ROUTINE_FIRES};
const ROUTINE_OVERVIEW_SQL: &str = polyc_query_model::statements::ROUTINE_OVERVIEW_SQL;
fn observed_routines_schema() -> arrow::datatypes::SchemaRef {
use polyc_projection::family::LogicalType;
let declared = OBSERVED_ROUTINES_ENTRY
.table(OBSERVED_ROUTINES_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 fires_schema() -> arrow::datatypes::SchemaRef {
use polyc_projection::family::LogicalType;
let declared = polyc_projection::family::ROUTINE_LIFECYCLE_ENTRY
.table(ROUTINE_FIRES)
.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 fixed_incarnation(len: usize) -> arrow::array::FixedSizeBinaryArray {
let mut builder = arrow::array::FixedSizeBinaryBuilder::with_capacity(len, 32);
for _ in 0..len {
builder.append_value([0_u8; 32]).unwrap();
}
builder.finish()
}
#[allow(
clippy::too_many_lines,
reason = "one fixture batch per table; splitting it hides which columns each table carries"
)]
fn context() -> SessionContext {
let ctx = SessionContext::new();
let routines_schema = observed_routines_schema();
let n = 3;
let routines = RecordBatch::try_new(
Arc::clone(&routines_schema),
vec![
Arc::new(StringArray::from(vec!["conv-a", "conv-b", "conv-a"])), Arc::new(fixed_incarnation(n)), Arc::new(UInt64Array::from(vec![0_u64, 0, 1])), Arc::new(StringArray::from(vec!["default", "default", "default"])), Arc::new(UInt64Array::from(vec![1_u64, 1, 2])), Arc::new(StringArray::from(vec!["1", "1", "2"])), Arc::new(UInt64Array::from(vec![0_u64, 0, 0])), Arc::new(UInt64Array::from(vec![0_u64, 0, 0])), Arc::new(StringArray::from(vec!["obs", "obs", "obs"])), Arc::new(StringArray::from(vec!["r-1", "r-2", "r-1"])), Arc::new(StringArray::from(vec!["u1", "u2", "u1"])), Arc::new(StringArray::from(vec!["1", "1", "2"])), Arc::new(UInt64Array::from(vec![1_u64, 1, 1])), Arc::new(StringArray::from(vec!["conv-a", "conv-b", "conv-a"])), Arc::new(StringArray::from(vec![
"persona-1",
"persona-2",
"persona-1",
])), Arc::new(StringArray::from(vec!["conv-a", "conv-b", "conv-a"])), Arc::new(StringArray::from(vec!["private", "public", "private"])), Arc::new(StringArray::from(vec!["R1", "R2", "R1-newer"])), Arc::new(StringArray::from(vec!["", "", ""])), Arc::new(StringArray::from(vec!["sched-1", "sched-2", "sched-1"])), Arc::new(StringArray::from(vec!["UTC", "UTC", "UTC"])), Arc::new(StringArray::from(vec!["[]", "[]", "[]"])), Arc::new(StringArray::from(vec!["do it", "do it", "do it"])), Arc::new(BooleanArray::from(vec![false, false, false])), Arc::new(BooleanArray::from(vec![true, true, true])), Arc::new(BooleanArray::from(vec![true, true, true])), Arc::new(StringArray::from(vec!["Ready", "Ready", "Ready"])), Arc::new(BooleanArray::from(vec![false, false, false])), Arc::new(StringArray::from(vec!["", "", ""])), Arc::new(BooleanArray::from(vec![false, false, false])), Arc::new(UInt64Array::from(vec![0_u64, 0, 0])), Arc::new(BooleanArray::from(vec![false, false, false])), Arc::new(UInt64Array::from(vec![0_u64, 0, 0])), Arc::new(StringArray::from(vec!["[]", "[]", "[]"])), Arc::new(BooleanArray::from(vec![false, false, false])), Arc::new(BooleanArray::from(vec![false, false, false])), Arc::new(StringArray::from(vec!["", "", ""])), Arc::new(BooleanArray::from(vec![false, false, false])), Arc::new(UInt64Array::from(vec![0_u64, 0, 0])), Arc::new(BooleanArray::from(vec![false, false, false])), Arc::new(StringArray::from(vec!["", "", ""])), Arc::new(BooleanArray::from(vec![false, false, false])), Arc::new(BooleanArray::from(vec![true, true, true])), ],
)
.unwrap();
ctx.register_table(
"observed_routines",
Arc::new(MemTable::try_new(routines_schema, vec![vec![routines]]).unwrap()),
)
.unwrap();
let fires_schema = fires_schema();
let fires = RecordBatch::try_new(
Arc::clone(&fires_schema),
vec![
Arc::new(StringArray::from(vec![
"routine-scheduler",
"routine-scheduler",
"routine-scheduler",
])),
Arc::new(fixed_incarnation(3)),
Arc::new(UInt64Array::from(vec![0_u64, 1, 2])),
Arc::new(StringArray::from(vec!["r-1", "r-1", "r-1"])),
Arc::new(StringArray::from(vec!["u1", "u1", "u1"])),
Arc::new(StringArray::from(vec![
"occurrence-1",
"occurrence-2",
"occurrence-3",
])),
Arc::new(UInt64Array::from(vec![0_u64, 0, 0])),
Arc::new(UInt64Array::from(vec![1_000_u64, 2_000, 3_000])),
Arc::new(BooleanArray::from(vec![true, true, true])),
Arc::new(StringArray::from(vec!["ok", "stopped_ungranted", "ok"])),
Arc::new(BooleanArray::from(vec![false, false, false])),
Arc::new(StringArray::from(vec!["", "", ""])),
],
)
.unwrap();
ctx.register_table(
"fires",
Arc::new(MemTable::try_new(fires_schema, vec![vec![fires]]).unwrap()),
)
.unwrap();
ctx
}
#[tokio::test]
#[allow(
clippy::too_many_lines,
reason = "asserts every column plus the duplicate-uid dedup in one proof; splitting it hides which assertion covers which fixture row"
)]
async fn routine_overview_sql_never_fans_out_and_derives_the_right_columns() {
let ctx = context();
let df = ctx.sql(ROUTINE_OVERVIEW_SQL).await.unwrap();
let batches = df.collect().await.unwrap();
let rows: usize = batches.iter().map(RecordBatch::num_rows).sum();
assert_eq!(
rows, 2,
"one row per DISTINCT routine — r-1's three fires AND its duplicate \
observed_routines row must never fan it out into more than one row"
);
assert_eq!(
batches[0]
.schema()
.fields()
.iter()
.map(|field| field.name().as_str())
.collect::<Vec<_>>(),
vec![
"uid",
"name",
"creator_persona",
"scope",
"prompt",
"schedule_json",
"fire_conversation_id",
"suspended",
"paused_by",
"paused_at_ms",
"pause_reason",
"display_name",
"description",
"schedule_timezone",
"orphaned",
"setup_completed",
"fire_count",
"last_fire_at_ms",
"last_fire_outcome",
"denial_count",
],
"must match fleet_routine_views's own column-name lookups, in order"
);
let mut by_uid: std::collections::HashMap<String, (u64, u64, String, u64, String)> =
std::collections::HashMap::new();
for batch in &batches {
let uid = batch
.column_by_name("uid")
.unwrap()
.as_any()
.downcast_ref::<StringArray>()
.unwrap();
let display_name = batch
.column_by_name("display_name")
.unwrap()
.as_any()
.downcast_ref::<StringArray>()
.unwrap();
let fire_count = batch
.column_by_name("fire_count")
.unwrap()
.as_any()
.downcast_ref::<UInt64Array>()
.unwrap();
let denial_count = batch
.column_by_name("denial_count")
.unwrap()
.as_any()
.downcast_ref::<UInt64Array>()
.unwrap();
let last_fire_outcome = batch
.column_by_name("last_fire_outcome")
.unwrap()
.as_any()
.downcast_ref::<StringArray>()
.unwrap();
let last_fire_at_ms = batch
.column_by_name("last_fire_at_ms")
.unwrap()
.as_any()
.downcast_ref::<UInt64Array>()
.unwrap();
for row in 0..batch.num_rows() {
by_uid.insert(
uid.value(row).to_owned(),
(
fire_count.value(row),
denial_count.value(row),
last_fire_outcome.value(row).to_owned(),
last_fire_at_ms.value(row),
display_name.value(row).to_owned(),
),
);
}
}
let (fire_count, denial_count, last_fire_outcome, last_fire_at_ms, display_name) =
&by_uid["u1"];
assert_eq!(*fire_count, 3, "r-1's fire count sums all three fires");
assert_eq!(
*denial_count, 1,
"exactly one of r-1's three fires was stopped_ungranted"
);
assert_eq!(
last_fire_outcome, "ok",
"the LATEST fire by fired_at_ms (3000) is 'ok', not the stopped_ungranted one"
);
assert_eq!(*last_fire_at_ms, 3_000);
assert_eq!(
display_name, "R1-newer",
"the higher-observation_ordinal duplicate observed_routines row wins, \
proving latest_routine's own dedup — never a fan-out from the duplicate"
);
let (fire_count, denial_count, last_fire_outcome, last_fire_at_ms, _) = &by_uid["u2"];
assert_eq!(*fire_count, 0, "r-2 has never fired");
assert_eq!(*denial_count, 0);
assert_eq!(last_fire_outcome, "");
assert_eq!(*last_fire_at_ms, 0);
}