polyc-query 2026.10.1

The Query plane's read model: a DataFusion engine over signed projection artifacts, behind a verified credential.
//! Proves the fixed SQL statements `ListRoutineFires` and `GetRoutineFire`
//! send to the Fleet Query plane against a `DataFusion` `MemTable` built
//! from the `routine-lifecycle/v1` `fires` schema.
//!
//! `$1` and `$2` are typed parameters in production. This proof substitutes
//! literals so `DataFusion` can plan the rest of the statement against the
//! real schema. The production files still contain the placeholders.

#![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");
}