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 `crate::control_plane::dashboard`
//! (`crates/control-plane/src/forensics.rs`'s `DASHBOARD_ATTRIBUTION_SQL`)
//! sends to the Fleet Query plane for `/api/dashboard`'s caller-attribution
//! stats (POLY-350) — against real `DataFusion` `MemTable`s built from the
//! exact `conversation-security/v1` schema the registry declares, not a live
//! projected stack.
//!
//! Mirrors `financial_dashboard_sql.rs`'s own shape and disclaimer: this is a
//! syntax-and-semantics proof for the SQL text itself, not the Query plane's
//! authority, resolution, or transport.

#![allow(clippy::unwrap_used)]

use std::sync::Arc;

use arrow::array::{RecordBatch, StringArray, UInt64Array};
use datafusion::datasource::MemTable;
use datafusion::prelude::SessionContext;
use polyc_projection::family::{
    CONVERSATION_SECURITY_ENTRY, SECURITY_ATTRIBUTION, SECURITY_ATTRIBUTION_PROVENANCE,
};

/// The statement under test — the exact constant `crate::forensics`
/// (`crates/control-plane/src/forensics.rs`) runs in production.
const ATTRIBUTION_SQL: &str = polyc_query_model::statements::DASHBOARD_ATTRIBUTION_SQL;

/// The same logical-to-Arrow mapping `polyc_projector::artifact::arrow_schema`
/// applies, restated here rather than imported — see
/// `financial_dashboard_sql.rs`'s own `schema` for why a Component cannot
/// depend on `polyc-projector`, a Container.
fn schema(table: polyc_projection::family::TableId) -> arrow::datatypes::SchemaRef {
    use polyc_projection::family::LogicalType;

    let declared = CONVERSATION_SECURITY_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 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()
}

/// `conv-a` carries THREE attribution records — two participants sharing one
/// edge and an initiator asserted by a second edge — proving the statement
/// still returns exactly one row for it (proof 1's fixture). `conv-b` carries
/// a single initiator record with a display name. `conv-c` has no
/// attribution at all and must be absent from the result, the same way
/// `dashboard_conversations.sql`'s `conv-c` (no committed turn) is absent
/// from that statement.
fn context() -> SessionContext {
    let ctx = SessionContext::new();

    let attribution_schema = schema(SECURITY_ATTRIBUTION);
    let attribution = RecordBatch::try_new(
        Arc::clone(&attribution_schema),
        vec![
            Arc::new(StringArray::from(vec![
                "conv-a", "conv-a", "conv-a", "conv-b",
            ])),
            Arc::new(fixed_incarnation(4)),
            Arc::new(UInt64Array::from(vec![0_u64, 1, 2, 0])),
            Arc::new(StringArray::from(vec![
                "turn-1", "turn-1", "turn-1", "turn-1",
            ])),
            Arc::new(StringArray::from(vec![
                "persona-a",
                "persona-b",
                "persona-c",
                "persona-d",
            ])),
            Arc::new(StringArray::from(vec![
                "participant",
                "participant",
                "initiator",
                "initiator",
            ])),
        ],
    )
    .unwrap();
    ctx.register_table(
        "attribution",
        Arc::new(MemTable::try_new(attribution_schema, vec![vec![attribution]]).unwrap()),
    )
    .unwrap();

    let provenance_schema = schema(SECURITY_ATTRIBUTION_PROVENANCE);
    let provenance = RecordBatch::try_new(
        Arc::clone(&provenance_schema),
        vec![
            Arc::new(StringArray::from(vec![
                "conv-a", "conv-a", "conv-a", "conv-b",
            ])),
            Arc::new(fixed_incarnation(4)),
            Arc::new(UInt64Array::from(vec![0_u64, 1, 2, 0])),
            Arc::new(StringArray::from(vec![
                "slack", "slack", "trigger", "slack",
            ])),
            Arc::new(StringArray::from(vec!["team", "team", "svc", "team"])),
            Arc::new(StringArray::from(vec!["u1", "u2", "svc-1", "u3"])),
            Arc::new(StringArray::from(vec!["Alice", "Bob", "Carol", "Dana"])),
            // conv-a's two participants share `edge-1`; its initiator posted
            // through a SECOND, distinct edge (`edge-2`) — two distinct
            // edges for conv-a. conv-b's one record posted through `edge-3`.
            Arc::new(StringArray::from(vec![
                "edge-1", "edge-1", "edge-2", "edge-3",
            ])),
            Arc::new(StringArray::from(vec!["", "", "", ""])),
            Arc::new(StringArray::from(vec!["", "", "", ""])),
        ],
    )
    .unwrap();
    ctx.register_table(
        "attribution_provenance",
        Arc::new(MemTable::try_new(provenance_schema, vec![vec![provenance]]).unwrap()),
    )
    .unwrap();

    ctx
}

#[tokio::test]
async fn one_row_per_conversation_with_correct_edge_count_and_initiator() {
    let ctx = context();
    let df = ctx.sql(ATTRIBUTION_SQL).await.unwrap();
    let batches = df.collect().await.unwrap();
    let rows: usize = batches.iter().map(RecordBatch::num_rows).sum();
    // conv-a (three attributions) and conv-b (one) each collapse to exactly
    // one row; conv-c never appears at all — a fan-out bug would report 3
    // rows for conv-a instead of 1.
    assert_eq!(
        rows, 2,
        "one row per conversation with any attribution, whatever its row count"
    );

    let batch = &batches[0];
    let ids = batch
        .column_by_name("id")
        .unwrap()
        .as_any()
        .downcast_ref::<StringArray>()
        .unwrap();
    let edge_counts = batch
        .column_by_name("edge_count")
        .unwrap()
        .as_any()
        .downcast_ref::<UInt64Array>()
        .unwrap();
    let initiators = batch
        .column_by_name("initiator_display_name")
        .unwrap()
        .as_any()
        .downcast_ref::<StringArray>()
        .unwrap();
    let fleet_edge_counts = batch
        .column_by_name("fleet_edge_count")
        .unwrap()
        .as_any()
        .downcast_ref::<UInt64Array>()
        .unwrap();

    let a = ids
        .iter()
        .position(|id| id == Some("conv-a"))
        .expect("conv-a present");
    assert_eq!(
        edge_counts.value(a),
        2,
        "conv-a's three attribution rows post through two DISTINCT edges — a fanout bug \
         would report 3, one per attribution row"
    );
    assert_eq!(
        initiators.value(a),
        "Carol",
        "conv-a's initiator is persona-c"
    );

    let b = ids
        .iter()
        .position(|id| id == Some("conv-b"))
        .expect("conv-b present");
    assert_eq!(edge_counts.value(b), 1);
    assert_eq!(initiators.value(b), "Dana");

    // Every row carries the SAME fleet-wide distinct-edge total: edge-1,
    // edge-2, edge-3 — three, never conv-a's own count leaking into it.
    assert_eq!(fleet_edge_counts.value(a), 3);
    assert_eq!(fleet_edge_counts.value(b), 3);
}