polyc-query 2026.10.2

The Query plane's read model: a DataFusion engine over signed projection artifacts, behind a verified credential.
//! Proves the fixed SQL statement `GET /api/routines/lifecycle`
//! (`crates/control-plane/src/forensics.rs`) sends to the Fleet Query plane
//! — against a real `DataFusion` `MemTable` built from the exact schema
//! `routine-lifecycle/v1` declares, not a live projected stack.
//!
//! This is a syntax-and-semantics proof for the SQL text itself: given rows
//! shaped exactly like the family publishes, the statement restores the
//! nullable `reason`/`scope` shape the route's readers carry and orders by
//! recorded time with journal position as the tiebreak. It does not
//! exercise the Query plane's authority, resolution, or transport —
//! `crates/query-service`'s own harness covers that.
//!
//! See `financial_dashboard_sql.rs` for the same proof's shape over the
//! dashboard statements.

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

/// The statement under test — the exact constant `crate::forensics` runs
/// in production, read from the shared statement registry rather than
/// duplicated as a second literal, so a change to it is caught without two
/// copies to keep in sync.
const LIFECYCLE_SQL: &str = polyc_query_model::statements::ROUTINE_LIFECYCLE_SQL;

/// The same logical-to-Arrow mapping `polyc_projector::artifact::arrow_schema`
/// applies, restated here rather than imported: this crate is a Component and
/// `polyc-projector` is a Container, so a test in this crate must not depend
/// on it (the layer rule points inward, not sideways into a container).
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))
}

/// Four lifecycle rows, deliberately out of `at_ms` order and mixing every
/// `has_` combination the fold emits:
///
/// - `routine_created` at position 0 — `has_reason`/`has_scope`/`has_channel`
///   all false;
/// - `routine_paused` at position 1 — `has_reason` true, the reason text a
///   reader must see;
/// - `routine_resumed` at position 2 — same `at_ms` as the pause, so journal
///   position is what orders the pair;
/// - `routine_scope_changed` at position 3 — `has_scope` true;
/// - a second `routine_paused` at position 4 — `has_reason` false, the
///   pauser gave none.
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![
            // `partition` — the fixed scheduler partition name.
            Arc::new(StringArray::from(vec!["routine-scheduler"; 5])),
            // `source_incarnation`.
            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",
            ])),
            // `at_ms` — positions 1 and 2 deliberately share one clock, and
            // position 4's clock precedes them all so the statement's ORDER
            // BY, not the row order, sets the timeline.
            Arc::new(UInt64Array::from(vec![100_u64, 300, 300, 400, 200])),
            // `has_reason` / `reason`.
            Arc::new(BooleanArray::from(vec![false, true, false, false, false])),
            Arc::new(StringArray::from(vec!["", "maintenance", "", "", ""])),
            // `has_scope` / `scope`.
            Arc::new(BooleanArray::from(vec![false, false, false, true, false])),
            Arc::new(StringArray::from(vec!["", "", "", "public", ""])),
            // `has_channel` / `channel`.
            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();

    // `ORDER BY at_ms, position`: the reason-less pause (position 4, at_ms
    // 200) precedes the shared-clock pair (positions 1 then 2, at_ms 300).
    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)]
    );

    // `reason` reads as a value only where `has_reason` stood — every other
    // row, including the recorded-but-empty-string rows, reads NULL.
    let ordered_reasons: Vec<Option<&str>> = reasons.iter().collect();
    assert_eq!(
        ordered_reasons,
        vec![None, None, Some("maintenance"), None, None]
    );

    // `scope` likewise — only the scope-changed row carries one.
    let ordered_scopes: Vec<Option<&str>> = scopes.iter().collect();
    assert_eq!(ordered_scopes, vec![None, None, None, None, Some("public")]);
}