polyc-query 2026.10.2

The Query plane's read model: a DataFusion engine over signed projection artifacts, behind a verified credential.
//! Proves `GetMyPayments`'s two fixed statements
//! (`crates/control-plane/src/payments_read.rs`, 8C-2's Q10/Q11 move)
//! against real `DataFusion` `MemTable`s built from the exact schemas the
//! registry declares.
//!
//! `FixtureProjectedQuery` (the test double `crates/control-plane/src/
//! forensics.rs` uses for its own route tests) re-implements the
//! `subject = $1` filter in Rust rather than executing this SQL text, so a
//! mutation that dropped `WHERE subject = $1` from either statement would
//! pass every control-plane test unchanged. This file is the missing proof:
//! the actual statement text, against a real planner, with a bound
//! parameter, proving a co-participant's row is excluded — not a
//! Rust-level reimplementation of the same filter.
//!
//! Same posture as `financial_dashboard_sql.rs`'s own `PAYMENTS_UNION_SQL`:
//! a hand-copied literal, kept in sync by hand (this crate is a Component
//! and cannot depend on the Container crate that owns the real constant), so
//! a change to either copy without the other is this file's own risk, not
//! one it can catch. What it DOES catch is a regression in the filter logic
//! itself, which no other test in the workspace exercises against a planner.

#![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_FINANCIAL_ENTRY, FINANCIAL_OUTBOUND_PAYMENTS, FINANCIAL_PAYMENTS,
    FINANCIAL_REFUSALS,
};

/// The exact statements `crate::control_plane::payments_read`
/// (`crates/control-plane/src/payments_read.rs`) sends to the Fleet Query
/// plane for `GetMyPayments`. See this file's own header for why these are
/// hand-copied rather than `include_str!`-shared.
const PAYMENTS_SQL: &str = "SELECT partition, position, turn_id, direction, reference, \
amount_base_units, asset, recipient, method, tool_call_id, payer_kind, timestamp_unix \
FROM payments WHERE subject = $1 \
UNION ALL \
SELECT partition, position, turn_id, direction, reference, amount_base_units, asset, \
recipient, method, tool_call_id, payer_kind, timestamp_unix FROM outbound_payments \
WHERE subject = $1 \
ORDER BY timestamp_unix DESC, partition DESC, position DESC LIMIT 500";

const RECENT_REFUSALS_SQL: &str = "SELECT partition, position, turn_id, reason, reason_detail, \
merchant_host, requested_base_units, permitted_base_units, tool_call_id, timestamp_unix \
FROM refusals WHERE subject = $1 \
ORDER BY timestamp_unix DESC, partition DESC, position DESC LIMIT 500";

fn schema(table: polyc_projection::family::TableId) -> arrow::datatypes::SchemaRef {
    use polyc_projection::family::LogicalType;

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

/// One shared conversation, TWO subjects (`persona-a` and `persona-b`, a
/// co-participant): one inbound and one outbound payment each, and one
/// refusal each. Every row in this fixture belongs to a subject other than
/// the one under test in at least one assertion, so a filter that leaked
/// would be caught.
#[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 payments_schema = schema(FINANCIAL_PAYMENTS);
    let payments = RecordBatch::try_new(
        Arc::clone(&payments_schema),
        vec![
            Arc::new(StringArray::from(vec!["conv-shared", "conv-shared"])),
            Arc::new(fixed_incarnation(2)),
            Arc::new(UInt64Array::from(vec![0_u64, 1])),
            Arc::new(StringArray::from(vec!["persona-a", "persona-b"])),
            Arc::new(UInt64Array::from(vec![1_000_u64, 2_000])),
            Arc::new(StringArray::from(vec!["USDC", "USDC"])),
            Arc::new(StringArray::from(vec!["ref-a-in", "ref-b-in"])),
            Arc::new(StringArray::from(vec!["0xRa", "0xRb"])),
            Arc::new(StringArray::from(vec!["tempo", "tempo"])),
            Arc::new(StringArray::from(vec!["call-a-in", "call-b-in"])),
            Arc::new(StringArray::from(vec!["turn-a1", "turn-b1"])),
            Arc::new(StringArray::from(vec!["inbound", "inbound"])),
            Arc::new(StringArray::from(vec!["linked_wallet", "linked_wallet"])),
            Arc::new(UInt64Array::from(vec![1_000_u64, 2_000])),
        ],
    )
    .unwrap();
    ctx.register_table(
        "payments",
        Arc::new(MemTable::try_new(payments_schema, vec![vec![payments]]).unwrap()),
    )
    .unwrap();

    let outbound_schema = schema(FINANCIAL_OUTBOUND_PAYMENTS);
    let outbound = RecordBatch::try_new(
        Arc::clone(&outbound_schema),
        vec![
            Arc::new(StringArray::from(vec!["conv-shared", "conv-shared"])),
            Arc::new(fixed_incarnation(2)),
            Arc::new(UInt64Array::from(vec![2_u64, 3])),
            Arc::new(StringArray::from(vec!["persona-a", "persona-b"])),
            Arc::new(UInt64Array::from(vec![500_u64, 700])),
            Arc::new(StringArray::from(vec!["0xToken", "0xToken"])),
            Arc::new(StringArray::from(vec!["ref-a-out", "ref-b-out"])),
            Arc::new(StringArray::from(vec!["0xRc", "0xRd"])),
            Arc::new(StringArray::from(vec!["tempo", "tempo"])),
            Arc::new(StringArray::from(vec!["call-a-out", "call-b-out"])),
            Arc::new(StringArray::from(vec!["turn-a2", "turn-b2"])),
            Arc::new(StringArray::from(vec!["outbound", "outbound"])),
            Arc::new(StringArray::from(vec!["deployment", "deployment"])),
            Arc::new(UInt64Array::from(vec![3_000_u64, 4_000])),
        ],
    )
    .unwrap();
    ctx.register_table(
        "outbound_payments",
        Arc::new(MemTable::try_new(outbound_schema, vec![vec![outbound]]).unwrap()),
    )
    .unwrap();

    let refusals_schema = schema(FINANCIAL_REFUSALS);
    let refusals = RecordBatch::try_new(
        Arc::clone(&refusals_schema),
        vec![
            Arc::new(StringArray::from(vec!["conv-shared", "conv-shared"])),
            Arc::new(fixed_incarnation(2)),
            Arc::new(UInt64Array::from(vec![0_u64, 1])),
            Arc::new(StringArray::from(vec!["persona-a", "persona-b"])),
            Arc::new(StringArray::from(vec!["over_spend_cap", "over_spend_cap"])),
            Arc::new(StringArray::from(vec!["m-a.example", "m-b.example"])),
            Arc::new(UInt64Array::from(vec![500_u64, 600])),
            Arc::new(UInt64Array::from(vec![100_u64, 200])),
            Arc::new(StringArray::from(vec!["call-a-blocked", "call-b-blocked"])),
            Arc::new(UInt64Array::from(vec![1_500_u64, 2_500])),
            Arc::new(StringArray::from(vec!["turn-a3", "turn-b3"])),
            Arc::new(StringArray::from(vec!["detail-a", "detail-b"])),
        ],
    )
    .unwrap();
    ctx.register_table(
        "refusals",
        Arc::new(MemTable::try_new(refusals_schema, vec![vec![refusals]]).unwrap()),
    )
    .unwrap();

    ctx
}

#[tokio::test]
async fn payments_sql_subject_filter_excludes_the_co_participant_and_binds_both_union_arms() {
    let ctx = context();
    // Textual substitution, the same pattern `routine_fires_sql.rs` already
    // uses for its own `$1`/`$2` — proves the SAME placeholder bound TWICE
    // across a `UNION ALL` (once per arm) filters both arms, not just the
    // first one a naive single-bind implementation might reach.
    let sql = PAYMENTS_SQL.replace("$1", "'persona-a'");
    let df = ctx.sql(&sql).await.unwrap();
    let batches = df.collect().await.unwrap();
    let rows: usize = batches.iter().map(RecordBatch::num_rows).sum();
    assert_eq!(
        rows, 2,
        "persona-a's own inbound AND outbound row — never persona-b's, \
         proving `$1` filters BOTH arms of the UNION ALL"
    );

    let mut references = Vec::new();
    for batch in &batches {
        let reference = batch
            .column_by_name("reference")
            .unwrap()
            .as_any()
            .downcast_ref::<StringArray>()
            .unwrap();
        for row in 0..batch.num_rows() {
            references.push(reference.value(row).to_owned());
        }
    }
    assert!(
        references.contains(&"ref-a-in".to_owned()) && references.contains(&"ref-a-out".to_owned()),
        "persona-a's own rows must both be present: {references:?}"
    );
    assert!(
        !references.contains(&"ref-b-in".to_owned())
            && !references.contains(&"ref-b-out".to_owned()),
        "persona-b's rows must never appear on persona-a's own read: {references:?}"
    );
    // `ORDER BY timestamp_unix DESC`: the outbound row (3_000) sorts before
    // the inbound one (1_000).
    assert_eq!(
        references,
        vec!["ref-a-out".to_owned(), "ref-a-in".to_owned()]
    );
}

#[tokio::test]
async fn recent_refusals_sql_subject_filter_excludes_the_co_participant() {
    let ctx = context();
    let sql = RECENT_REFUSALS_SQL.replace("$1", "'persona-a'");
    let df = ctx.sql(&sql).await.unwrap();
    let batches = df.collect().await.unwrap();
    let rows: usize = batches.iter().map(RecordBatch::num_rows).sum();
    assert_eq!(rows, 1, "only persona-a's own blocked attempt");

    let batch = &batches[0];
    let tool_call_id = batch
        .column_by_name("tool_call_id")
        .unwrap()
        .as_any()
        .downcast_ref::<StringArray>()
        .unwrap();
    assert_eq!(
        tool_call_id.value(0),
        "call-a-blocked",
        "never persona-b's blocked attempt in the same conversation"
    );
}