use arrow::array::{Array as _, Int32Array};
use arrow::datatypes::{DataType, Field, Schema};
use buffa::Message as _;
use commonware_runtime::{Runner, deterministic};
use datafusion::execution::memory_pool::FairSpillPool;
use datafusion::execution::runtime_env::RuntimeEnvBuilder;
use polyc_crypto::approval::{
ApprovalSigner, ReceiptPayload, RefusalPayload, WalletLinkLifecyclePayload, receipt_payload,
refusal_payload, routine_created_payload, routine_deleted_payload, routine_paused_payload,
routine_resumed_payload, routine_scope_changed_payload, wallet_link_lifecycle_payload,
};
use polyc_eventlog::{EventLog, EventLogConfig};
use polyc_proto::kinds;
use polyc_proto::proto::polychrome::agent::v1::{
Content, FunctionCallContent, FunctionResultContent, Message as WireMessage, TextContent,
ToolCallContent, ToolResultContent, content, function_result_content, tool_call_content,
tool_result_content,
};
use polyc_proto::proto::polychrome::events::v1::{
AttributionEvent, ModelCallEvent, RoutineFireOutcome, RoutineFireOutcomeEvent,
RoutineFiredEvent, SummaryEvent, TurnDispatchedEvent, TurnFailedEvent, UsageEvent,
};
use polyc_proto::proto::polychrome::harness::v1::TurnFailureKind;
use polyc_proto::proto::polychrome::persona::v1::{
ExternalIdentity, PersonaCredential, SpendPolicy, UsageRollup, WalletLink,
};
use uuid::Uuid;
use super::*;
impl QueryEngine {
pub(crate) async fn execute(
&self,
sql: &str,
allow_explain: bool,
) -> Result<QueryOutput, QueryEngineError> {
self.execute_with_params(sql, &[], allow_explain).await
}
}
fn test_base_state() -> SessionState {
let runtime = RuntimeEnvBuilder::new()
.with_memory_pool(Arc::new(FairSpillPool::new(64 * 1024 * 1024)))
.build_arc()
.expect("build test runtime");
SessionStateBuilder::new()
.with_runtime_env(runtime)
.with_default_features()
.build()
}
fn fixture_events(committed_turn: Uuid, orphaned_turn: Uuid) -> Vec<(u64, Event)> {
let usage = UsageEvent {
input_tokens: 11,
output_tokens: 22,
..Default::default()
};
vec![
(
0,
Event::new(
kinds::tagged(kinds::TURN_START, &committed_turn),
Vec::new(),
),
),
(
1,
Event::trusted(
kinds::tagged(kinds::USER_MSG, &committed_turn),
b"hello".to_vec(),
),
),
(
2,
Event::new(
kinds::tagged(kinds::USAGE, &committed_turn),
usage.encode_to_vec(),
),
),
(
3,
Event::new(
kinds::tagged(kinds::TURN_COMPLETE, &committed_turn),
Vec::new(),
),
),
(
4,
Event::new(kinds::tagged(kinds::TURN_START, &orphaned_turn), Vec::new()),
),
(
5,
Event::trusted(
kinds::tagged(kinds::USER_MSG, &orphaned_turn),
b"orphaned".to_vec(),
),
),
(
6,
Event::new(
kinds::tagged(kinds::OUTPUT_MSG, &orphaned_turn),
b"partial reply".to_vec(),
),
),
]
}
fn one_partition(name: &str, committed_turn: Uuid, orphaned_turn: Uuid) -> Vec<PartitionEvents> {
vec![PartitionEvents {
partition: name.to_string(),
events: fixture_events(committed_turn, orphaned_turn),
}]
}
fn single_count(output: &QueryOutput) -> i64 {
assert_eq!(output.batches.len(), 1, "expected exactly one batch");
let batch = &output.batches[0];
assert_eq!(batch.num_rows(), 1, "expected exactly one row");
batch
.column(0)
.as_any()
.downcast_ref::<arrow::array::Int64Array>()
.expect("COUNT(*) is Int64")
.value(0)
}
#[test]
fn real_journal_round_trip_reproduces_the_spikes_3_vs_0() {
let committed_turn = Uuid::from_u128(0x0195_abcd_ef01_2345_6789_abcd_ef01_2345);
let orphaned_turn = Uuid::from_u128(0x0195_abcd_ef01_2345_6789_abcd_ef01_9999);
let executor = deterministic::Runner::default();
let replayed: Vec<(u64, Event)> = executor.start(|context| async move {
let log = EventLog::open(context, EventLogConfig::for_partition("conv-p3-real"))
.await
.expect("open log");
for (_position, event) in fixture_events(committed_turn, orphaned_turn) {
log.append(&event).await.expect("append");
}
log.commit().await.expect("commit");
log.replay_with_positions().await.expect("replay")
});
assert_eq!(replayed.len(), 7, "the fixture's own seven events");
let tokio_rt = tokio::runtime::Runtime::new().expect("build a tokio runtime");
tokio_rt.block_on(async move {
let partitions = vec![PartitionEvents {
partition: "conv-p3-real".to_string(),
events: replayed,
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let raw = engine
.execute(
&format!("SELECT COUNT(*) AS c FROM events_raw WHERE turn_id = '{orphaned_turn}'"),
false,
)
.await
.expect("raw count query");
assert_eq!(
single_count(&raw),
3,
"raw replay includes the orphaned turn's rows"
);
let committed = engine
.execute(
&format!("SELECT COUNT(*) AS c FROM events WHERE turn_id = '{orphaned_turn}'"),
false,
)
.await
.expect("committed count query");
assert_eq!(
single_count(&committed),
0,
"the committed-turns view excludes the orphaned turn entirely"
);
});
}
#[tokio::test]
async fn fleet_scope_events_view_excludes_orphaned_turn_raw_includes_it() {
let committed_turn = Uuid::from_u128(1);
let orphaned_turn = Uuid::from_u128(2);
let partitions = one_partition("conv-a", committed_turn, orphaned_turn);
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let raw = engine
.execute(
&format!("SELECT COUNT(*) AS c FROM events_raw WHERE turn_id = '{orphaned_turn}'"),
false,
)
.await
.expect("raw query");
assert_eq!(single_count(&raw), 3);
let committed = engine
.execute(
&format!("SELECT COUNT(*) AS c FROM events WHERE turn_id = '{orphaned_turn}'"),
false,
)
.await
.expect("committed query");
assert_eq!(single_count(&committed), 0);
}
#[tokio::test]
async fn conversations_scope_hides_raw_table_by_every_name() {
let committed_turn = Uuid::from_u128(3);
let orphaned_turn = Uuid::from_u128(4);
let partitions = one_partition("conv-b", committed_turn, orphaned_turn);
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-b".to_string()]),
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let unqualified = engine.execute("SELECT * FROM events_raw", false).await;
assert!(
matches!(unqualified, Err(QueryEngineError::DataFusion(_))),
"events_raw must not resolve for a non-Fleet scope: {unqualified:?}"
);
let qualified = engine
.execute("SELECT * FROM datafusion.public.events_raw", false)
.await;
assert!(
matches!(qualified, Err(QueryEngineError::DataFusion(_))),
"a fully-qualified name must not reach the raw rows either: {qualified:?}"
);
let attribution_raw_unqualified = engine.execute("SELECT * FROM attribution_raw", false).await;
assert!(
matches!(
attribution_raw_unqualified,
Err(QueryEngineError::DataFusion(_))
),
"attribution_raw must not resolve for a non-Fleet scope: {attribution_raw_unqualified:?}"
);
let payments_raw_unqualified = engine.execute("SELECT * FROM payments_raw", false).await;
assert!(
matches!(
payments_raw_unqualified,
Err(QueryEngineError::DataFusion(_))
),
"payments_raw must not resolve for a non-Fleet scope: {payments_raw_unqualified:?}"
);
}
#[tokio::test]
async fn fleet_scope_reference_tables_return_rows_and_participation_joins_to_its_conversation() {
let committed_turn = Uuid::from_u128(20);
let orphaned_turn = Uuid::from_u128(21);
let partitions = one_partition("conv-ref", committed_turn, orphaned_turn);
let profile = PersonaProfile {
persona_id: "persona-ref-1".to_string(),
display_name: "Ada".to_string(),
created_at_ms: 5_000,
status: "linked".to_string(),
..Default::default()
};
let participation = Participation {
conversation_id: "ref".to_string(),
role: "initiator".to_string(),
first_at_ms: 6_000,
..Default::default()
};
let reference = ReferenceData {
personas: vec![profile],
participations: vec![("persona-ref-1".to_string(), participation)],
..ReferenceData::empty()
};
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
reference,
&[],
QueryLimits::default(),
)
.await
.expect("build");
let personas = engine
.execute("SELECT COUNT(*) AS c FROM personas", false)
.await
.expect("personas query");
assert_eq!(single_count(&personas), 1);
let participations = engine
.execute("SELECT COUNT(*) AS c FROM participations", false)
.await
.expect("participations query");
assert_eq!(single_count(&participations), 1);
let joined = engine
.execute(
"SELECT DISTINCT e.partition FROM participations p \
JOIN events_raw e ON e.partition = 'conv-' || p.conversation_id \
WHERE p.persona_id = 'persona-ref-1'",
false,
)
.await
.expect("join query");
assert_eq!(
joined.batches[0].num_rows(),
1,
"the join matches exactly one distinct conversation partition"
);
let partition_col = joined.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("partition is Utf8");
assert_eq!(partition_col.value(0), "conv-ref");
}
#[tokio::test]
#[allow(clippy::too_many_lines)] async fn qry8_bulk_export_query_returns_parity_columns_and_hides_key_material() {
let committed_turn = Uuid::from_u128(28);
let orphaned_turn = Uuid::from_u128(29);
let partitions = one_partition("conv-parity", committed_turn, orphaned_turn);
let profile = PersonaProfile {
persona_id: "persona-parity-1".to_string(),
display_name: "Ada".to_string(),
created_at_ms: 1_000,
status: "linked".to_string(),
..Default::default()
};
let wallet = WalletLink {
persona_id: "persona-parity-1".to_string(),
wallet_address: "0xParity".to_string(),
currency: "0xUSDC".to_string(),
expiry_unix: 123,
key_ref: "keychain://persona-parity-1-secret".to_string(),
created_at_ms: 1_000,
revoked: false,
revoked_at_ms: 0,
..Default::default()
};
let spend_policy = SpendPolicy {
persona_id: "persona-parity-1".to_string(),
limit: "5".to_string(),
period_secs: 86_400,
max_lifetime_secs: 604_800,
allowed_hosts: vec!["api.example.com".to_string(), "pay.example.com".to_string()],
updated_at_ms: 2_000,
set_by: "admin-1".to_string(),
..Default::default()
};
let credential = PersonaCredential {
persona_id: "persona-parity-1".to_string(),
credential_id: vec![9, 9, 9],
p256_public_key_sec1: vec![8, 8, 8],
signing_public_key: vec![7, 7, 7],
rp_id: "polychrome.example".to_string(),
origin: "https://polychrome.example".to_string(),
created_at_ms: 3_000,
revoked: true,
revoked_at_ms: 3_500,
..Default::default()
};
let usage_rollup = UsageRollup {
persona_id: "persona-parity-1".to_string(),
committed_turns: 5,
input_tokens: 10,
output_tokens: 20,
last_active_ms: 4_000,
conversation_ids: vec!["conv-parity".to_string()],
updated_at_ms: 4_500,
..Default::default()
};
let reference = ReferenceData {
personas: vec![profile],
participations: Vec::new(),
wallets: vec![("persona-parity-1".to_string(), wallet)],
spend_policies: vec![("persona-parity-1".to_string(), spend_policy)],
credentials: vec![("persona-parity-1".to_string(), credential)],
usage_rollups: vec![("persona-parity-1".to_string(), usage_rollup)],
routines: Vec::new(),
dashboard: Vec::new(),
};
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
reference,
&[],
QueryLimits::default(),
)
.await
.expect("build");
let joined = engine
.execute(
"SELECT w.wallet_address, sp.\"limit\", sp.allowed_hosts, c.rp_id, c.origin, \
c.revoked, u.committed_turns \
FROM personas p \
JOIN persona_wallets w ON w.persona_id = p.persona_id \
JOIN persona_spend_policies sp ON sp.persona_id = p.persona_id \
JOIN persona_credentials c ON c.persona_id = p.persona_id \
JOIN persona_usage u ON u.persona_id = p.persona_id \
WHERE p.persona_id = 'persona-parity-1'",
false,
)
.await
.expect("bulk-export join query");
assert_eq!(joined.batches[0].num_rows(), 1);
let wallet_address = joined.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("wallet_address is Utf8")
.value(0);
assert_eq!(wallet_address, "0xParity", "wallet_address is unredacted");
let limit = joined.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("limit is Utf8")
.value(0);
assert_eq!(limit, "5", "spend limit is unredacted");
let allowed_hosts = joined.batches[0]
.column(2)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("allowed_hosts is Utf8")
.value(0);
assert_eq!(
allowed_hosts, r#"["api.example.com","pay.example.com"]"#,
"the full allowed_hosts list is unredacted"
);
let rp_id = joined.batches[0]
.column(3)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("rp_id is Utf8")
.value(0);
assert_eq!(rp_id, "polychrome.example");
let origin = joined.batches[0]
.column(4)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("origin is Utf8")
.value(0);
assert_eq!(origin, "https://polychrome.example");
let revoked = joined.batches[0]
.column(5)
.as_any()
.downcast_ref::<arrow::array::BooleanArray>()
.expect("revoked is Boolean")
.value(0);
assert!(revoked, "credential revocation state is unredacted");
let committed_turns = joined.batches[0]
.column(6)
.as_any()
.downcast_ref::<arrow::array::UInt64Array>()
.expect("committed_turns is UInt64")
.value(0);
assert_eq!(committed_turns, 5);
for (table, column) in [
("persona_wallets", "key_ref"),
("persona_credentials", "credential_id"),
("persona_credentials", "p256_public_key_sec1"),
("persona_credentials", "signing_public_key"),
] {
let result = engine
.execute(&format!("SELECT {column} FROM {table}"), false)
.await;
assert!(
matches!(result, Err(QueryEngineError::DataFusion(_))),
"{table}.{column} must not resolve: {result:?}"
);
}
let columns = engine
.execute(
"SELECT column_name FROM information_schema.columns WHERE table_name IN \
('persona_wallets', 'persona_spend_policies', 'persona_credentials', \
'persona_usage')",
false,
)
.await
.expect("information_schema.columns query");
let banned = [
"key_ref",
"credential_id",
"p256_public_key_sec1",
"signing_public_key",
];
for batch in &columns.batches {
let column_names = batch
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("column_name is Utf8");
for i in 0..column_names.len() {
let name = column_names.value(i);
assert!(
!banned.contains(&name),
"key-material column {name} must not appear in any new reference table's schema"
);
}
}
}
#[tokio::test]
async fn conversations_scope_hides_reference_tables() {
let committed_turn = Uuid::from_u128(22);
let orphaned_turn = Uuid::from_u128(23);
let partitions = one_partition("conv-ref-hidden", committed_turn, orphaned_turn);
let profile = PersonaProfile {
persona_id: "persona-hidden-1".to_string(),
display_name: "Hidden".to_string(),
created_at_ms: 1_000,
status: "linked".to_string(),
identities: vec![sample_identity("slack", "team-1", "U-HIDDEN", "Hidden")],
..Default::default()
};
let participation = Participation {
conversation_id: "ref-hidden".to_string(),
role: "initiator".to_string(),
first_at_ms: 2_000,
..Default::default()
};
let wallet = WalletLink {
persona_id: "persona-hidden-1".to_string(),
wallet_address: "0xhidden".to_string(),
currency: "0xusdc".to_string(),
created_at_ms: 1_500,
..Default::default()
};
let spend_policy = SpendPolicy {
persona_id: "persona-hidden-1".to_string(),
limit: "5".to_string(),
allowed_hosts: vec!["pay.example.com".to_string()],
..Default::default()
};
let credential = PersonaCredential {
persona_id: "persona-hidden-1".to_string(),
rp_id: "polychrome.example".to_string(),
origin: "https://polychrome.example".to_string(),
created_at_ms: 1_600,
..Default::default()
};
let usage_rollup = UsageRollup {
persona_id: "persona-hidden-1".to_string(),
committed_turns: 3,
..Default::default()
};
let reference = ReferenceData {
personas: vec![profile],
participations: vec![("persona-hidden-1".to_string(), participation)],
wallets: vec![("persona-hidden-1".to_string(), wallet)],
spend_policies: vec![("persona-hidden-1".to_string(), spend_policy)],
credentials: vec![("persona-hidden-1".to_string(), credential)],
usage_rollups: vec![("persona-hidden-1".to_string(), usage_rollup)],
routines: Vec::new(),
dashboard: Vec::new(),
};
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-ref-hidden".to_string()]),
partitions,
reference,
&[],
QueryLimits::default(),
)
.await
.expect("build");
for table in [
"personas",
"participations",
"persona_identities",
"persona_wallets",
"persona_spend_policies",
"persona_credentials",
"persona_usage",
] {
let bare = engine
.execute(&format!("SELECT * FROM {table}"), false)
.await;
assert!(
matches!(bare, Err(QueryEngineError::DataFusion(_))),
"{table} must not resolve for a non-Fleet scope: {bare:?}"
);
let qualified = engine
.execute(&format!("SELECT * FROM datafusion.public.{table}"), false)
.await;
assert!(
matches!(qualified, Err(QueryEngineError::DataFusion(_))),
"a fully-qualified name must not reach {table} either: {qualified:?}"
);
}
}
#[tokio::test]
async fn fleet_scope_dashboard_table_returns_rows_with_summary_text_and_json_columns() {
let committed_turn = Uuid::from_u128(30);
let orphaned_turn = Uuid::from_u128(31);
let partitions = one_partition("conv-dash-ref", committed_turn, orphaned_turn);
let row = crate::dashboard::DashboardRow {
conversation_id: "dash-ref".to_string(),
total_events: 9,
committed_turns: 1,
input_tokens: 10,
output_tokens: 5,
summary_text: Some("condensed so far".to_string()),
last_turn_id: Some(committed_turn.as_simple().to_string()),
edges: vec!["web".to_string()],
settlements: vec![crate::dashboard::DashboardSettlement {
persona: "persona-a".to_string(),
spend_base_units: 1_500,
charged_base_units: 250,
}],
created_at_ms: Some(1_700_000_000_000),
last_activity_ms: Some(1_700_000_000_000),
persona_id: Some("persona-a".to_string()),
first_message_preview: Some("hello there".to_string()),
};
let reference = ReferenceData {
dashboard: vec![row],
..ReferenceData::empty()
};
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
reference,
&[],
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute(
"SELECT conversation_id, summary_active, summary_text, edges_json, settlements_json \
FROM dashboard",
false,
)
.await
.expect("dashboard query");
assert_eq!(output.batches[0].num_rows(), 1);
let conversation_id = output.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("conversation_id is Utf8");
assert_eq!(conversation_id.value(0), "dash-ref");
let summary_active = output.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::BooleanArray>()
.expect("summary_active is Boolean");
assert!(summary_active.value(0));
let summary_text = output.batches[0]
.column(2)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("summary_text is Utf8");
assert_eq!(
summary_text.value(0),
"condensed so far",
"a Fleet session sees summary_text — the same posture the `summary` table's own \
QRY-2-D Fleet-only gate applies"
);
let edges_json = output.batches[0]
.column(3)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("edges_json is Utf8");
assert_eq!(edges_json.value(0), r#"["web"]"#);
let settlements_json = output.batches[0]
.column(4)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("settlements_json is Utf8");
let decoded: serde_json::Value = serde_json::from_str(settlements_json.value(0)).unwrap();
assert_eq!(decoded[0]["persona"], "persona-a");
assert_eq!(decoded[0]["spend_base_units"], "1500");
assert_eq!(decoded[0]["charged_base_units"], "250");
}
#[tokio::test]
async fn conversations_scope_hides_dashboard_table_and_its_summary_text() {
let committed_turn = Uuid::from_u128(32);
let orphaned_turn = Uuid::from_u128(33);
let partitions = one_partition("conv-dash-hidden", committed_turn, orphaned_turn);
let row = crate::dashboard::DashboardRow {
conversation_id: "dash-hidden".to_string(),
summary_text: Some("must never leak to a non-admin scope".to_string()),
first_message_preview: Some("also must never leak".to_string()),
..Default::default()
};
let reference = ReferenceData {
dashboard: vec![row],
..ReferenceData::empty()
};
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-dash-hidden".to_string()]),
partitions,
reference,
&[],
QueryLimits::default(),
)
.await
.expect("build");
let dashboard = engine.execute("SELECT * FROM dashboard", false).await;
assert!(
matches!(dashboard, Err(QueryEngineError::DataFusion(_))),
"dashboard must not resolve for a non-Fleet scope: {dashboard:?}"
);
let qualified = engine
.execute("SELECT * FROM datafusion.public.dashboard", false)
.await;
assert!(
matches!(qualified, Err(QueryEngineError::DataFusion(_))),
"a fully-qualified name must not reach dashboard either: {qualified:?}"
);
}
#[tokio::test]
async fn fleet_scope_personas_merged_into_column_resolves_the_survivor() {
let committed_turn = Uuid::from_u128(24);
let orphaned_turn = Uuid::from_u128(25);
let partitions = one_partition("conv-merge", committed_turn, orphaned_turn);
let survivor = PersonaProfile {
persona_id: "persona-survivor".to_string(),
display_name: "Ada".to_string(),
created_at_ms: 1_000,
status: "linked".to_string(),
..Default::default()
};
let absorbed = PersonaProfile {
persona_id: "persona-absorbed".to_string(),
display_name: "Ada (old)".to_string(),
created_at_ms: 500,
status: "merged".to_string(),
merged_into: "persona-survivor".to_string(),
..Default::default()
};
let reference = ReferenceData {
personas: vec![survivor, absorbed],
participations: Vec::new(),
..ReferenceData::empty()
};
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
reference,
&[],
QueryLimits::default(),
)
.await
.expect("build");
let rows = engine
.execute(
"SELECT persona_id, merged_into FROM personas ORDER BY persona_id",
false,
)
.await
.expect("personas query");
assert_eq!(rows.batches[0].num_rows(), 2);
let persona_id_col = rows.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("persona_id is Utf8");
let merged_into_col = rows.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("merged_into is Utf8");
assert_eq!(persona_id_col.value(0), "persona-absorbed");
assert_eq!(merged_into_col.value(0), "persona-survivor");
assert_eq!(persona_id_col.value(1), "persona-survivor");
assert_eq!(
merged_into_col.value(1),
"",
"an unmerged survivor's merged_into is empty, not null"
);
}
#[tokio::test]
async fn fleet_scope_persona_identities_are_queryable_and_resolve_by_identity() {
let committed_turn = Uuid::from_u128(26);
let orphaned_turn = Uuid::from_u128(27);
let partitions = one_partition("conv-idents", committed_turn, orphaned_turn);
let profile = PersonaProfile {
persona_id: "persona-idents-1".to_string(),
display_name: "Ada".to_string(),
created_at_ms: 1_000,
status: "linked".to_string(),
identities: vec![
sample_identity("slack", "team-1", "U-1", "Ada"),
sample_identity("telegram", "", "T-1", "Ada"),
],
..Default::default()
};
let reference = ReferenceData {
personas: vec![profile],
participations: Vec::new(),
..ReferenceData::empty()
};
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
reference,
&[],
QueryLimits::default(),
)
.await
.expect("build");
let count = engine
.execute("SELECT COUNT(*) AS c FROM persona_identities", false)
.await
.expect("persona_identities query");
assert_eq!(single_count(&count), 2);
let joined = engine
.execute(
"SELECT DISTINCT p.display_name FROM persona_identities i \
JOIN personas p ON p.persona_id = i.persona_id \
WHERE i.persona_id = 'persona-idents-1'",
false,
)
.await
.expect("join query");
assert_eq!(
joined.batches[0].num_rows(),
1,
"the join matches the one persona exactly once"
);
let display_name_col = joined.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("display_name is Utf8");
assert_eq!(display_name_col.value(0), "Ada");
let resolved = engine
.execute(
"SELECT persona_id FROM persona_identities \
WHERE provider = 'slack' AND scope = 'team-1' AND external_id = 'U-1'",
false,
)
.await
.expect("resolve query");
assert_eq!(
resolved.batches[0].num_rows(),
1,
"the resolve-shape query matches exactly one identity"
);
let resolved_persona_id = resolved.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("persona_id is Utf8");
assert_eq!(resolved_persona_id.value(0), "persona-idents-1");
}
#[tokio::test]
async fn conversations_scope_payload_column_is_unreachable() {
let committed_turn = Uuid::from_u128(17);
let orphaned_turn = Uuid::from_u128(18);
let partitions = one_partition("conv-redact", committed_turn, orphaned_turn);
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-redact".to_string()]),
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let events_payload = engine.execute("SELECT payload FROM events", false).await;
assert!(
matches!(events_payload, Err(QueryEngineError::DataFusion(_))),
"`payload` must not resolve as a column of `events`: {events_payload:?}"
);
let usage_payload = engine.execute("SELECT payload FROM usage", false).await;
assert!(
matches!(usage_payload, Err(QueryEngineError::DataFusion(_))),
"`payload` must not resolve as a column of `usage`: {usage_payload:?}"
);
let events_payload_json = engine
.execute("SELECT payload_json FROM events", false)
.await;
assert!(
matches!(events_payload_json, Err(QueryEngineError::DataFusion(_))),
"`payload_json` must not resolve as a column of `events`: {events_payload_json:?}"
);
let usage_payload_json = engine
.execute("SELECT payload_json FROM usage", false)
.await;
assert!(
matches!(usage_payload_json, Err(QueryEngineError::DataFusion(_))),
"`payload_json` must not resolve as a column of `usage`: {usage_payload_json:?}"
);
let events_output = engine
.execute("SELECT * FROM events", false)
.await
.expect("events view");
let events_schema = events_output.batches[0].schema();
let events_columns: Vec<&str> = events_schema
.fields()
.iter()
.map(|f| f.name().as_str())
.collect();
assert_eq!(
events_columns,
vec!["partition", "position", "kind", "turn_id"],
"events' schema for a non-Fleet scope must be exactly the committed-turns \
projection, with no payload or payload_json column"
);
let usage_output = engine
.execute("SELECT * FROM usage", false)
.await
.expect("usage table");
let usage_schema = usage_output.batches[0].schema();
let usage_columns: Vec<&str> = usage_schema
.fields()
.iter()
.map(|f| f.name().as_str())
.collect();
assert_eq!(
usage_columns,
vec![
"partition",
"position",
"turn_id",
"input_tokens",
"output_tokens"
],
"usage's schema for a non-Fleet scope carries its uniform keys but no \
payload or payload_json column"
);
}
#[tokio::test]
async fn conversations_scope_model_call_payload_column_is_unreachable() {
let committed_turn = Uuid::from_u128(19);
let orphaned_turn = Uuid::from_u128(20);
let mut partitions = one_partition("conv-redact-mc", committed_turn, orphaned_turn);
partitions[0].events.push((
7,
Event::new(
kinds::tagged(kinds::MODEL_CALL, &committed_turn),
ModelCallEvent {
provider: "stub-provider".to_string(),
model: "stub-model".to_string(),
captured_clock_unix_ms: 1_700_000_000_000,
..Default::default()
}
.encode_to_vec(),
),
));
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-redact-mc".to_string()]),
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let payload = engine
.execute("SELECT payload FROM model_call", false)
.await;
assert!(
matches!(payload, Err(QueryEngineError::DataFusion(_))),
"`payload` must not resolve as a column of `model_call`: {payload:?}"
);
let payload_json = engine
.execute("SELECT payload_json FROM model_call", false)
.await;
assert!(
matches!(payload_json, Err(QueryEngineError::DataFusion(_))),
"`payload_json` must not resolve as a column of `model_call`: {payload_json:?}"
);
let output = engine
.execute("SELECT * FROM model_call", false)
.await
.expect("model_call table");
let schema = output.batches[0].schema();
let columns: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
assert_eq!(
columns,
vec![
"partition",
"position",
"turn_id",
"provider",
"model",
"captured_clock_unix_ms",
],
"model_call's schema for a non-Fleet scope carries its uniform keys but no \
payload or payload_json column"
);
}
#[tokio::test]
async fn conversations_scope_attribution_payload_column_is_unreachable() {
let committed_turn = Uuid::from_u128(21);
let orphaned_turn = Uuid::from_u128(22);
let mut partitions = one_partition("conv-redact-attr", committed_turn, orphaned_turn);
let caller = sample_attribution(
"persona-payload-check",
"initiator",
sample_identity("slack", "team-payload", "U-PAYLOAD", "Payload Checker"),
);
partitions[0].events.push((
7,
Event::new(
kinds::tagged(kinds::CALLER, &committed_turn),
caller.encode_to_vec(),
),
));
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-redact-attr".to_string()]),
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let payload = engine
.execute("SELECT payload FROM attribution", false)
.await;
assert!(
matches!(payload, Err(QueryEngineError::DataFusion(_))),
"`payload` must not resolve as a column of `attribution`: {payload:?}"
);
let payload_json = engine
.execute("SELECT payload_json FROM attribution", false)
.await;
assert!(
matches!(payload_json, Err(QueryEngineError::DataFusion(_))),
"`payload_json` must not resolve as a column of `attribution`: {payload_json:?}"
);
let output = engine
.execute("SELECT * FROM attribution", false)
.await
.expect("attribution table");
let schema = output.batches[0].schema();
let columns: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
assert_eq!(
columns,
vec!["partition", "position", "turn_id", "persona_id", "role"],
"attribution's schema for a non-Fleet scope carries its uniform keys and the \
resolved persona_id/role, but no payload, payload_json, or identity_* column \
(see conversations_scope_attribution_identity_columns_are_unreachable for the \
identity_* redaction proof itself)"
);
}
#[tokio::test]
async fn conversations_scope_attribution_identity_columns_are_unreachable() {
let turn = Uuid::from_u128(2_100);
let identity = sample_identity("slack", "team-other", "U-OTHER", "Other Participant");
let participant = sample_attribution("persona-other", "participant", identity);
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
),
(
1,
Event::new(
kinds::tagged(kinds::PARTICIPANT, &turn),
participant.encode_to_vec(),
),
),
(
2,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-identity-redact".to_string(),
events,
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-identity-redact".to_string()]),
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
for column in [
"identity_provider",
"identity_scope",
"identity_external_id",
"identity_display_name",
] {
let result = engine
.execute(&format!("SELECT {column} FROM attribution"), false)
.await;
assert!(
matches!(result, Err(QueryEngineError::DataFusion(_))),
"`{column}` must not resolve as a column of `attribution` for a non-Fleet \
scope — a conversation-scoped session must never read another \
participant's raw external identity: {result:?}"
);
}
let output = engine
.execute("SELECT * FROM attribution", false)
.await
.expect("attribution table");
let schema = output.batches[0].schema();
let columns: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
assert_eq!(
columns,
vec!["partition", "position", "turn_id", "persona_id", "role"],
"a non-Fleet attribution schema must carry exactly the redacted column set, \
never any identity_* column"
);
let resolved = engine
.execute("SELECT persona_id, role FROM attribution", false)
.await
.expect("persona_id/role must still resolve for a non-Fleet scope");
assert_eq!(resolved.batches[0].num_rows(), 1);
let persona_id_col = resolved.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("persona_id is Utf8");
let role_col = resolved.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("role is Utf8");
assert_eq!(persona_id_col.value(0), "persona-other");
assert_eq!(role_col.value(0), "participant");
}
#[tokio::test]
async fn conversations_scope_events_usage_model_call_attribution_and_turn_failed_still_work() {
let committed_turn = Uuid::from_u128(5);
let orphaned_turn = Uuid::from_u128(6);
let failed_turn = Uuid::from_u128(7);
let mut events = fixture_events(committed_turn, orphaned_turn);
events.push((
7,
Event::new(kinds::tagged(kinds::TURN_START, &failed_turn), Vec::new()),
));
events.push((
8,
Event::new(
kinds::tagged(kinds::TURN_FAILED, &failed_turn),
TurnFailedEvent {
kind: TurnFailureKind::RateLimit.into(),
message: "provider throttled the request".to_string(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
.encode_to_vec(),
),
));
events.push((
9,
Event::new(
kinds::tagged(kinds::TURN_COMPLETE, &failed_turn),
Vec::new(),
),
));
let partitions = vec![PartitionEvents {
partition: "conv-c".to_string(),
events,
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-c".to_string()]),
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let events = engine
.execute("SELECT COUNT(*) AS c FROM events", false)
.await
.expect("events view");
assert_eq!(
single_count(&events),
7,
"the committed turn's four rows plus the now-committed failed turn's three \
(turn_start, turn_failed, turn_complete)"
);
let usage = engine
.execute("SELECT COUNT(*) AS c FROM usage", false)
.await
.expect("usage table");
assert_eq!(single_count(&usage), 1);
let model_call = engine
.execute("SELECT COUNT(*) AS c FROM model_call", false)
.await
.expect("model_call table");
assert_eq!(
single_count(&model_call),
0,
"the fixture carries no model_call events"
);
let attribution = engine
.execute("SELECT COUNT(*) AS c FROM attribution", false)
.await
.expect("attribution table");
assert_eq!(
single_count(&attribution),
0,
"the fixture carries no caller/participant events"
);
let turn_failed = engine
.execute("SELECT failure_kind, message FROM turn_failed", false)
.await
.expect("turn_failed table");
assert_eq!(turn_failed.batches[0].num_rows(), 1);
let failure_kind_col = turn_failed.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("failure_kind is Utf8");
let message_col = turn_failed.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("message is Utf8");
assert_eq!(failure_kind_col.value(0), "rate_limit");
assert_eq!(message_col.value(0), "provider throttled the request");
}
#[allow(clippy::too_many_lines)]
#[tokio::test]
async fn conversations_scope_hides_uncommitted_turn_usage_payments_tool_calls_approvals_and_turn_failed()
{
let committed_turn = Uuid::from_u128(9_001);
let orphaned_turn = Uuid::from_u128(9_002);
let signer = ApprovalSigner::from_seed(9_004);
let committed_receipt =
signed_payment_receipt(&signer, kinds::PAYMENT_RECEIPT, "call-committed");
let orphaned_receipt = signed_payment_receipt(&signer, kinds::PAYMENT_RECEIPT, "call-orphaned");
let committed_approval_request = polyc_crypto::approval::request_payload(
"call-committed",
"rm",
r#"{"path":"/tmp/committed"}"#,
"default",
"",
&[],
"",
"",
"",
&[],
false,
);
let committed_approval_response = signed_approval_response(
&signer,
"call-committed",
"rm",
r#"{"path":"/tmp/committed"}"#,
"caller-1",
"ok",
"conv-uncommitted-hidden",
"nonce-committed",
);
let orphaned_approval_request = polyc_crypto::approval::request_payload(
"call-orphaned",
"rm",
r#"{"path":"/tmp/orphaned"}"#,
"default",
"",
&[],
"",
"",
"",
&[],
false,
);
let orphaned_approval_response = signed_approval_response(
&signer,
"call-orphaned",
"rm",
r#"{"path":"/tmp/orphaned"}"#,
"caller-1",
"ok",
"conv-uncommitted-hidden",
"nonce-orphaned",
);
let committed_call = wire_tool_call_message("call-committed", "search", r#"{"q":"a"}"#);
let orphaned_call = wire_tool_call_message("call-orphaned", "search", r#"{"q":"b"}"#);
let events = vec![
(
0,
Event::new(
kinds::tagged(kinds::TURN_START, &committed_turn),
Vec::new(),
),
),
(
1,
Event::new(
kinds::tagged(kinds::USAGE, &committed_turn),
UsageEvent {
input_tokens: 5,
output_tokens: 6,
..Default::default()
}
.encode_to_vec(),
),
),
(
2,
Event::new(
kinds::tagged(kinds::PAYMENT_RECEIPT, &committed_turn),
committed_receipt,
),
),
(
3,
Event::new(
kinds::tagged(kinds::APPROVAL_REQUEST, &committed_turn),
committed_approval_request,
),
),
(
4,
Event::new(
kinds::tagged(kinds::APPROVAL_RESPONSE, &committed_turn),
committed_approval_response,
),
),
(
5,
Event::new(
kinds::tagged(kinds::OUTPUT_MSG, &committed_turn),
committed_call.encode_to_vec(),
),
),
(
6,
Event::new(
kinds::tagged(kinds::TURN_FAILED, &committed_turn),
TurnFailedEvent {
kind: TurnFailureKind::RateLimit.into(),
message: "provider throttled the request".to_string(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
.encode_to_vec(),
),
),
(
7,
Event::new(
kinds::tagged(kinds::TURN_COMPLETE, &committed_turn),
Vec::new(),
),
),
(
8,
Event::new(kinds::tagged(kinds::TURN_START, &orphaned_turn), Vec::new()),
),
(
9,
Event::new(
kinds::tagged(kinds::USAGE, &orphaned_turn),
UsageEvent {
input_tokens: 100,
output_tokens: 200,
..Default::default()
}
.encode_to_vec(),
),
),
(
10,
Event::new(
kinds::tagged(kinds::PAYMENT_RECEIPT, &orphaned_turn),
orphaned_receipt,
),
),
(
11,
Event::new(
kinds::tagged(kinds::APPROVAL_REQUEST, &orphaned_turn),
orphaned_approval_request,
),
),
(
12,
Event::new(
kinds::tagged(kinds::APPROVAL_RESPONSE, &orphaned_turn),
orphaned_approval_response,
),
),
(
13,
Event::new(
kinds::tagged(kinds::OUTPUT_MSG, &orphaned_turn),
orphaned_call.encode_to_vec(),
),
),
(
14,
Event::new(
kinds::tagged(kinds::TURN_FAILED, &orphaned_turn),
TurnFailedEvent {
kind: TurnFailureKind::Timeout.into(),
message: "deadline exceeded before the turn could complete".to_string(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
.encode_to_vec(),
),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-uncommitted-hidden".to_string(),
events,
}];
let trusted_signers = vec![signer.public_key_bytes()];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-uncommitted-hidden".to_string()]),
partitions,
ReferenceData::empty(),
&trusted_signers,
QueryLimits::default(),
)
.await
.expect("build");
let usage = engine
.execute(
&format!("SELECT COUNT(*) AS c FROM usage WHERE turn_id = '{committed_turn}'"),
false,
)
.await
.expect("usage committed query");
assert_eq!(
single_count(&usage),
1,
"the committed turn's usage row is visible"
);
let payments = engine
.execute(
"SELECT COUNT(*) AS c FROM payments WHERE tool_call_id = 'call-committed'",
false,
)
.await
.expect("payments committed query");
assert_eq!(
single_count(&payments),
1,
"the committed turn's payment row is visible"
);
let approvals = engine
.execute(
"SELECT COUNT(*) AS c FROM approvals WHERE request_id = 'call-committed'",
false,
)
.await
.expect("approvals committed query");
assert_eq!(
single_count(&approvals),
2,
"the committed turn's request/response rows are both visible"
);
let tool_calls = engine
.execute(
"SELECT COUNT(*) AS c FROM tool_calls WHERE tool_call_id = 'call-committed'",
false,
)
.await
.expect("tool_calls committed query");
assert_eq!(
single_count(&tool_calls),
1,
"the committed turn's tool call row is visible"
);
let turn_failed_committed = engine
.execute(
&format!("SELECT COUNT(*) AS c FROM turn_failed WHERE turn_id = '{committed_turn}'"),
false,
)
.await
.expect("turn_failed committed query");
assert_eq!(
single_count(&turn_failed_committed),
1,
"the committed turn's turn_failed row is visible — turn_complete lands in the same \
atomic batch as turn_failed, so this turn is committed like any other"
);
let orphaned_usage = engine
.execute(
&format!("SELECT COUNT(*) AS c FROM usage WHERE turn_id = '{orphaned_turn}'"),
false,
)
.await
.expect("usage orphaned query");
assert_eq!(
single_count(&orphaned_usage),
0,
"the orphaned turn's usage row must be hidden"
);
let orphaned_payments = engine
.execute(
"SELECT COUNT(*) AS c FROM payments WHERE tool_call_id = 'call-orphaned'",
false,
)
.await
.expect("payments orphaned query");
assert_eq!(
single_count(&orphaned_payments),
0,
"the orphaned turn's payment row must be hidden"
);
let orphaned_approvals = engine
.execute(
"SELECT COUNT(*) AS c FROM approvals WHERE request_id = 'call-orphaned'",
false,
)
.await
.expect("approvals orphaned query");
assert_eq!(
single_count(&orphaned_approvals),
0,
"the orphaned turn's request/response rows must both be hidden"
);
let orphaned_tool_calls = engine
.execute(
"SELECT COUNT(*) AS c FROM tool_calls WHERE tool_call_id = 'call-orphaned'",
false,
)
.await
.expect("tool_calls orphaned query");
assert_eq!(
single_count(&orphaned_tool_calls),
0,
"the orphaned turn's tool call row must be hidden"
);
let orphaned_turn_failed = engine
.execute(
&format!("SELECT COUNT(*) AS c FROM turn_failed WHERE turn_id = '{orphaned_turn}'"),
false,
)
.await
.expect("turn_failed orphaned query");
assert_eq!(
single_count(&orphaned_turn_failed),
0,
"the orphaned turn's turn_failed row must be hidden too — no turn_complete ever \
landed for it, so it is exactly as uncommitted as its usage/payment/approval/ \
tool_call siblings"
);
}
#[tokio::test]
async fn usage_sum_matches_encoded_usage_events() {
let committed_turn = Uuid::from_u128(150);
let partitions = vec![PartitionEvents {
partition: "conv-usage".to_string(),
events: vec![
(
0,
Event::new(
kinds::tagged(kinds::TURN_START, &committed_turn),
Vec::new(),
),
),
(
1,
Event::new(
kinds::tagged(kinds::USAGE, &committed_turn),
UsageEvent {
input_tokens: 10,
output_tokens: 3,
..Default::default()
}
.encode_to_vec(),
),
),
(
2,
Event::new(
kinds::tagged(kinds::USAGE, &committed_turn),
UsageEvent {
input_tokens: 40,
output_tokens: 7,
..Default::default()
}
.encode_to_vec(),
),
),
(
3,
Event::new(
kinds::tagged(kinds::TURN_COMPLETE, &committed_turn),
Vec::new(),
),
),
],
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute(
"SELECT SUM(input_tokens) AS i, SUM(output_tokens) AS o FROM usage",
false,
)
.await
.expect("sum query");
let batch = &output.batches[0];
let input_sum = batch
.column(0)
.as_any()
.downcast_ref::<arrow::array::UInt64Array>()
.expect("UInt64 sum");
let output_sum = batch
.column(1)
.as_any()
.downcast_ref::<arrow::array::UInt64Array>()
.expect("UInt64 sum");
assert_eq!(input_sum.value(0), 50);
assert_eq!(output_sum.value(0), 10);
}
#[tokio::test]
async fn model_call_rows_are_queryable_and_join_to_their_conversation() {
let committed_turn = Uuid::from_u128(200);
let events = vec![
(
0,
Event::new(
kinds::tagged(kinds::TURN_START, &committed_turn),
Vec::new(),
),
),
(
1,
Event::new(
kinds::tagged(kinds::MODEL_CALL, &committed_turn),
ModelCallEvent {
provider: "stub-provider".to_string(),
model: "stub-model".to_string(),
decode_params: buffa::MessageField::default(),
system_config_ref: String::new(),
captured_clock_unix_ms: 1_700_000_000_000,
clear_trigger_bytes: 0,
clear_keep_recent: 0,
clear_marker: String::new(),
persona_block_placement: 0,
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
.encode_to_vec(),
),
),
(
2,
Event::new(
kinds::tagged(kinds::TURN_COMPLETE, &committed_turn),
Vec::new(),
),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-model-call".to_string(),
events,
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let count = engine
.execute("SELECT COUNT(*) AS c FROM model_call", false)
.await
.expect("model_call count query");
assert_eq!(single_count(&count), 1);
let joined = engine
.execute(
"SELECT DISTINCT mc.provider FROM model_call mc \
JOIN events e ON e.partition = mc.partition \
WHERE mc.partition = 'conv-model-call'",
false,
)
.await
.expect("join query");
assert_eq!(
joined.batches[0].num_rows(),
1,
"the join matches the one model_call row's provider exactly once"
);
let provider_col = joined.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("provider is Utf8");
assert_eq!(provider_col.value(0), "stub-provider");
}
#[tokio::test]
async fn turn_failed_rows_are_queryable_and_join_to_their_conversation() {
let failed_turn = Uuid::from_u128(300);
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &failed_turn), Vec::new()),
),
(
1,
Event::new(
kinds::tagged(kinds::TURN_FAILED, &failed_turn),
TurnFailedEvent {
kind: TurnFailureKind::RateLimit.into(),
message: "provider throttled the request".to_string(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
.encode_to_vec(),
),
),
(
2,
Event::new(
kinds::tagged(kinds::TURN_COMPLETE, &failed_turn),
Vec::new(),
),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-turn-failed".to_string(),
events,
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let count = engine
.execute("SELECT COUNT(*) AS c FROM turn_failed", false)
.await
.expect("turn_failed count query");
assert_eq!(single_count(&count), 1);
let joined = engine
.execute(
"SELECT DISTINCT tf.failure_kind, tf.message FROM turn_failed tf \
JOIN events_raw e ON e.partition = tf.partition \
WHERE tf.partition = 'conv-turn-failed'",
false,
)
.await
.expect("join query");
assert_eq!(
joined.batches[0].num_rows(),
1,
"the join matches the one turn_failed row's failure_kind/message exactly once"
);
let failure_kind_col = joined.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("failure_kind is Utf8");
let message_col = joined.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("message is Utf8");
assert_eq!(failure_kind_col.value(0), "rate_limit");
assert_eq!(message_col.value(0), "provider throttled the request");
}
#[tokio::test]
async fn summary_rows_are_queryable_and_join_to_their_conversation() {
let committed_turn = Uuid::from_u128(400);
let summary_id = Uuid::from_u128(401);
let events = vec![
(
0,
Event::new(
kinds::tagged(kinds::TURN_START, &committed_turn),
Vec::new(),
),
),
(
1,
Event::new(
kinds::tagged(kinds::TURN_COMPLETE, &committed_turn),
Vec::new(),
),
),
(
2,
Event::new(
kinds::tagged(kinds::SUMMARY, &summary_id),
SummaryEvent {
text: "the older prefix, condensed".to_string(),
covers_through_position: 2,
..Default::default()
}
.encode_to_vec(),
),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-summary".to_string(),
events,
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let count = engine
.execute("SELECT COUNT(*) AS c FROM summary", false)
.await
.expect("summary count query");
assert_eq!(single_count(&count), 1);
let joined = engine
.execute(
"SELECT DISTINCT s.text, s.covers_through_position, s.turn_id \
FROM summary s \
JOIN events_raw e ON e.partition = s.partition \
WHERE s.partition = 'conv-summary'",
false,
)
.await
.expect("join query");
assert_eq!(
joined.batches[0].num_rows(),
1,
"the join matches the one summary row's text/covers_through_position exactly once"
);
let text_col = joined.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("text is Utf8");
let covers_col = joined.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::UInt64Array>()
.expect("covers_through_position is UInt64");
let turn_id_col = joined.batches[0]
.column(2)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("turn_id is Utf8");
assert_eq!(text_col.value(0), "the older prefix, condensed");
assert_eq!(covers_col.value(0), 2);
assert_eq!(
turn_id_col.value(0),
summary_id.to_string(),
"summary.turn_id is the summary's OWN synthetic tag, not the committed turn's id"
);
}
#[tokio::test]
async fn conversations_scope_hides_summary_rows() {
let summary_id = Uuid::from_u128(410);
let events = vec![(
0,
Event::new(
kinds::tagged(kinds::SUMMARY, &summary_id),
SummaryEvent {
text: "fleet-only content".to_string(),
covers_through_position: 9,
..Default::default()
}
.encode_to_vec(),
),
)];
let conversations_engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-summary-hidden".to_string()]),
vec![PartitionEvents {
partition: "conv-summary-hidden".to_string(),
events,
}],
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build conversations engine");
let unqualified = conversations_engine
.execute("SELECT * FROM summary", false)
.await;
assert!(
matches!(unqualified, Err(QueryEngineError::DataFusion(_))),
"summary must not resolve for a non-Fleet scope: {unqualified:?}"
);
let qualified = conversations_engine
.execute("SELECT * FROM datafusion.public.summary", false)
.await;
assert!(
matches!(qualified, Err(QueryEngineError::DataFusion(_))),
"a fully-qualified name must not reach summary either: {qualified:?}"
);
let raw_unqualified = conversations_engine
.execute("SELECT * FROM summary_raw", false)
.await;
assert!(
matches!(raw_unqualified, Err(QueryEngineError::DataFusion(_))),
"summary_raw must not resolve for a non-Fleet scope either: {raw_unqualified:?}"
);
}
#[allow(clippy::too_many_lines)]
#[tokio::test]
async fn turn_failed_filter_by_failure_kind() {
let turn_a = Uuid::from_u128(301);
let turn_b = Uuid::from_u128(302);
let turn_c = Uuid::from_u128(303);
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn_a), Vec::new()),
),
(
1,
Event::new(
kinds::tagged(kinds::TURN_FAILED, &turn_a),
TurnFailedEvent {
kind: TurnFailureKind::Timeout.into(),
message: "deadline exceeded".to_string(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
.encode_to_vec(),
),
),
(
2,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn_a), Vec::new()),
),
(
3,
Event::new(kinds::tagged(kinds::TURN_START, &turn_b), Vec::new()),
),
(
4,
Event::new(
kinds::tagged(kinds::TURN_FAILED, &turn_b),
TurnFailedEvent {
kind: TurnFailureKind::Timeout.into(),
message: "second deadline exceeded".to_string(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
.encode_to_vec(),
),
),
(
5,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn_b), Vec::new()),
),
(
6,
Event::new(kinds::tagged(kinds::TURN_START, &turn_c), Vec::new()),
),
(
7,
Event::new(
kinds::tagged(kinds::TURN_FAILED, &turn_c),
TurnFailedEvent {
kind: TurnFailureKind::Auth.into(),
message: "credential rejected".to_string(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
.encode_to_vec(),
),
),
(
8,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn_c), Vec::new()),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-turn-failed-filter".to_string(),
events,
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let timeouts = engine
.execute(
"SELECT turn_id, message FROM turn_failed \
WHERE failure_kind = 'timeout' ORDER BY turn_id",
false,
)
.await
.expect("failure_kind filter query");
assert_eq!(
timeouts.batches[0].num_rows(),
2,
"exactly the two timeout rows, not the auth row"
);
let message_col = timeouts.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("message is Utf8");
let messages: Vec<&str> = (0..2).map(|i| message_col.value(i)).collect();
assert!(messages.contains(&"deadline exceeded"));
assert!(messages.contains(&"second deadline exceeded"));
let auths = engine
.execute(
"SELECT COUNT(*) AS c FROM turn_failed WHERE failure_kind = 'auth'",
false,
)
.await
.expect("auth filter query");
assert_eq!(single_count(&auths), 1);
}
fn sample_identity(
provider: &str,
scope: &str,
external_id: &str,
display_name: &str,
) -> ExternalIdentity {
ExternalIdentity {
provider: provider.to_string(),
scope: scope.to_string(),
external_id: external_id.to_string(),
display_name: display_name.to_string(),
time_zone: String::new(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn sample_attribution(
persona_id: &str,
role: &str,
identity: ExternalIdentity,
) -> AttributionEvent {
AttributionEvent {
persona_id: persona_id.to_string(),
identity: buffa::MessageField::some(identity),
role: role.to_string(),
asserting_edge_id: String::new(),
signer_pk_hex: String::new(),
signature_hex: String::new(),
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
#[tokio::test]
async fn attribution_rows_are_queryable_and_join_to_their_conversation() {
let committed_turn = Uuid::from_u128(210);
let caller = sample_attribution(
"persona-ada",
"initiator",
sample_identity("slack", "team-caller", "U-CALLER", "Ada"),
);
let participant = sample_attribution(
"persona-bob",
"participant",
sample_identity("telegram", "", "T-PARTICIPANT", "Bob"),
);
let events = vec![
(
0,
Event::new(
kinds::tagged(kinds::TURN_START, &committed_turn),
Vec::new(),
),
),
(
1,
Event::new(
kinds::tagged(kinds::CALLER, &committed_turn),
caller.encode_to_vec(),
),
),
(
2,
Event::new(
kinds::tagged(kinds::PARTICIPANT, &committed_turn),
participant.encode_to_vec(),
),
),
(
3,
Event::new(
kinds::tagged(kinds::TURN_COMPLETE, &committed_turn),
Vec::new(),
),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-attribution".to_string(),
events,
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let count = engine
.execute("SELECT COUNT(*) AS c FROM attribution", false)
.await
.expect("attribution count query");
assert_eq!(single_count(&count), 2, "both caller and participant rows");
let joined = engine
.execute(
"SELECT a.role, a.identity_provider, a.identity_external_id \
FROM attribution a \
JOIN events e ON e.partition = a.partition \
WHERE a.partition = 'conv-attribution' \
GROUP BY a.role, a.identity_provider, a.identity_external_id \
ORDER BY a.role",
false,
)
.await
.expect("join query");
assert_eq!(
joined.batches[0].num_rows(),
2,
"the join matches both attribution rows exactly once each"
);
let role_col = joined.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("role is Utf8");
let provider_col = joined.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("identity_provider is Utf8");
let external_id_col = joined.batches[0]
.column(2)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("identity_external_id is Utf8");
assert_eq!(role_col.value(0), "initiator");
assert_eq!(provider_col.value(0), "slack");
assert_eq!(external_id_col.value(0), "U-CALLER");
assert_eq!(role_col.value(1), "participant");
assert_eq!(provider_col.value(1), "telegram");
assert_eq!(external_id_col.value(1), "T-PARTICIPANT");
}
#[tokio::test]
async fn fleet_scope_attribution_returns_full_identity_columns() {
let committed_turn = Uuid::from_u128(212);
let caller = sample_attribution(
"persona-ada",
"initiator",
sample_identity("slack", "team-caller", "U-CALLER", "Ada"),
);
let events = vec![
(
0,
Event::new(
kinds::tagged(kinds::TURN_START, &committed_turn),
Vec::new(),
),
),
(
1,
Event::new(
kinds::tagged(kinds::CALLER, &committed_turn),
caller.encode_to_vec(),
),
),
(
2,
Event::new(
kinds::tagged(kinds::TURN_COMPLETE, &committed_turn),
Vec::new(),
),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-attribution-fleet".to_string(),
events,
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let full = engine
.execute(
"SELECT identity_scope, identity_display_name FROM attribution \
WHERE persona_id = 'persona-ada'",
false,
)
.await
.expect("identity_scope/identity_display_name must resolve for Fleet");
assert_eq!(full.batches[0].num_rows(), 1);
let scope_col = full.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("identity_scope is Utf8");
let display_name_col = full.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("identity_display_name is Utf8");
assert_eq!(scope_col.value(0), "team-caller");
assert_eq!(display_name_col.value(0), "Ada");
let star_output = engine
.execute("SELECT * FROM attribution", false)
.await
.expect("attribution table");
let star_schema = star_output.batches[0].schema();
let star_columns: Vec<&str> = star_schema
.fields()
.iter()
.map(|f| f.name().as_str())
.collect();
assert_eq!(
star_columns,
vec![
"partition",
"position",
"turn_id",
"persona_id",
"role",
"identity_provider",
"identity_scope",
"identity_external_id",
"identity_display_name",
],
"a Fleet session's attribution schema must carry every identity_* column"
);
}
#[tokio::test]
async fn attribution_persona_id_joins_to_personas_reference_table() {
let committed_turn = Uuid::from_u128(211);
let caller = sample_attribution(
"persona-join-1",
"initiator",
sample_identity("slack", "team-1", "U-JOIN", "Ada"),
);
let events = vec![
(
0,
Event::new(
kinds::tagged(kinds::TURN_START, &committed_turn),
Vec::new(),
),
),
(
1,
Event::new(
kinds::tagged(kinds::CALLER, &committed_turn),
caller.encode_to_vec(),
),
),
(
2,
Event::new(
kinds::tagged(kinds::TURN_COMPLETE, &committed_turn),
Vec::new(),
),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-attr-join".to_string(),
events,
}];
let profile = PersonaProfile {
persona_id: "persona-join-1".to_string(),
display_name: "Ada".to_string(),
created_at_ms: 3_000,
status: "linked".to_string(),
..Default::default()
};
let reference = ReferenceData {
personas: vec![profile],
participations: Vec::new(),
..ReferenceData::empty()
};
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
reference,
&[],
QueryLimits::default(),
)
.await
.expect("build");
let joined = engine
.execute(
"SELECT p.display_name FROM attribution a \
JOIN personas p ON p.persona_id = a.persona_id \
WHERE a.partition = 'conv-attr-join'",
false,
)
.await
.expect("persona join query");
assert_eq!(
joined.batches[0].num_rows(),
1,
"the join matches the one attribution row's persona exactly once"
);
let display_name_col = joined.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("display_name is Utf8");
assert_eq!(display_name_col.value(0), "Ada");
}
fn signed_payment_receipt(
signer: &ApprovalSigner,
kind: &'static str,
tool_call_id: &str,
) -> Vec<u8> {
let (payload, sig, pk) = receipt_payload(
&ReceiptPayload {
kind,
reference: "tx-engine-1",
amount: "2.75",
currency: "USDC",
recipient: "0xengine-recipient",
method: "tempo",
timestamp: "2026-07-21T00:00:00Z",
tool_call_id,
approval_pos: "1",
approved_args_hash: "hash-engine",
subject: "persona-payer",
payer_kind: "linked_wallet",
paying_account: "0xengine-payer",
},
signer,
);
let _ = (sig, pk);
payload
}
fn signed_refusal(signer: &ApprovalSigner, tool_call_id: &str) -> Vec<u8> {
let (payload, sig, pk) = refusal_payload(
&RefusalPayload {
kind: kinds::PAYMENT_REFUSAL,
reason: "over_spend_cap",
reason_detail: "payment of 500 base units exceeds the cap of 100",
merchant_host: "merchant.example",
requested_base_units: "500",
permitted_base_units: "100",
tool_call_id,
subject: "persona-payer",
timestamp: "1780000000",
},
signer,
);
let _ = (sig, pk);
payload
}
#[tokio::test]
async fn conversations_scope_refusals_signer_public_key_is_unreachable() {
let turn = Uuid::from_u128(602);
let signer = ApprovalSigner::from_seed(204);
let bytes = signed_refusal(&signer, "call-redact");
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
),
(
1,
Event::new(kinds::tagged(kinds::PAYMENT_REFUSAL, &turn), bytes),
),
(
2,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-refusals-redact".to_string(),
events,
}];
let trusted_signers = vec![signer.public_key_bytes()];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-refusals-redact".to_string()]),
partitions,
ReferenceData::empty(),
&trusted_signers,
QueryLimits::default(),
)
.await
.expect("build");
let result = engine
.execute("SELECT signer_public_key FROM refusals", false)
.await;
assert!(
matches!(result, Err(QueryEngineError::DataFusion(_))),
"`signer_public_key` must not resolve as a column of `refusals` for a \
non-Fleet scope: {result:?}"
);
let raw_result = engine.execute("SELECT * FROM refusals_raw", false).await;
assert!(
matches!(raw_result, Err(QueryEngineError::DataFusion(_))),
"`refusals_raw` must not resolve at all for a non-Fleet scope: {raw_result:?}"
);
let output = engine
.execute("SELECT * FROM refusals", false)
.await
.expect("refusals table");
let schema = output.batches[0].schema();
let columns: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
assert_eq!(
columns,
vec![
"partition",
"position",
"turn_id",
"reason",
"reason_detail",
"merchant_host",
"requested_base_units",
"permitted_base_units",
"tool_call_id",
"subject",
"timestamp",
],
"a non-Fleet refusals schema must carry every column except signer_public_key"
);
let resolved = engine
.execute(
"SELECT reason, requested_base_units, subject FROM refusals WHERE tool_call_id = 'call-redact'",
false,
)
.await
.expect("reason/requested_base_units/subject must still resolve for a non-Fleet scope");
assert_eq!(resolved.batches[0].num_rows(), 1);
let reason_col = resolved.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("reason is Utf8");
let requested_col = resolved.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("requested_base_units is Utf8");
let subject_col = resolved.batches[0]
.column(2)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("subject is Utf8");
assert_eq!(reason_col.value(0), "over_spend_cap");
assert_eq!(requested_col.value(0), "500");
assert_eq!(subject_col.value(0), "persona-payer");
}
#[tokio::test]
async fn fleet_scope_refusals_returns_signer_public_key() {
let turn = Uuid::from_u128(603);
let signer = ApprovalSigner::from_seed(205);
let bytes = signed_refusal(&signer, "call-fleet");
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
),
(
1,
Event::new(kinds::tagged(kinds::PAYMENT_REFUSAL, &turn), bytes),
),
(
2,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-refusals-fleet".to_string(),
events,
}];
let trusted_signers = vec![signer.public_key_bytes()];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&trusted_signers,
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute(
"SELECT signer_public_key FROM refusals WHERE tool_call_id = 'call-fleet'",
false,
)
.await
.expect("signer_public_key must resolve for Fleet");
assert_eq!(output.batches[0].num_rows(), 1);
let signer_col = output.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::BinaryArray>()
.expect("signer_public_key is Binary");
assert_eq!(signer_col.value(0), signer.public_key_bytes());
}
fn signed_wallet_link_lifecycle_event(signer: &ApprovalSigner, conversation_id: &str) -> Vec<u8> {
let (payload, sig, pk) = wallet_link_lifecycle_payload(
&WalletLinkLifecyclePayload {
kind: kinds::WALLET_LINK_LIFECYCLE,
transition: "linked",
subject: "persona-payer",
wallet_address: "0x1111111111111111111111111111111111111111",
currency: "0x2222222222222222222222222222222222222222",
chain_id: "42431",
limit_base_units: "5000000",
limit_human: "5",
period_secs: "86400",
expiry_unix: "1780086400",
recipients: "",
conversation_id,
timestamp: "1780000000",
},
signer,
);
let _ = (sig, pk);
payload
}
#[tokio::test]
async fn conversations_scope_wallet_link_lifecycle_signer_public_key_is_unreachable() {
let signer = ApprovalSigner::from_seed(304);
let bytes = signed_wallet_link_lifecycle_event(&signer, "conv-wll-redact");
let events = vec![(0, Event::new(kinds::WALLET_LINK_LIFECYCLE, bytes))];
let partitions = vec![PartitionEvents {
partition: "conv-wll-redact".to_string(),
events,
}];
let trusted_signers = vec![signer.public_key_bytes()];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-wll-redact".to_string()]),
partitions,
ReferenceData::empty(),
&trusted_signers,
QueryLimits::default(),
)
.await
.expect("build");
let result = engine
.execute("SELECT signer_public_key FROM wallet_link_lifecycle", false)
.await;
assert!(
matches!(result, Err(QueryEngineError::DataFusion(_))),
"`signer_public_key` must not resolve as a column of `wallet_link_lifecycle` for a \
non-Fleet scope: {result:?}"
);
let raw_result = engine
.execute("SELECT * FROM wallet_link_lifecycle_raw", false)
.await;
assert!(
matches!(raw_result, Err(QueryEngineError::DataFusion(_))),
"`wallet_link_lifecycle_raw` must not resolve at all for a non-Fleet scope: {raw_result:?}"
);
let output = engine
.execute("SELECT * FROM wallet_link_lifecycle", false)
.await
.expect("wallet_link_lifecycle table");
let schema = output.batches[0].schema();
let columns: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
assert_eq!(
columns,
vec![
"partition",
"position",
"turn_id",
"transition",
"subject",
"wallet_address",
"currency",
"chain_id",
"limit_base_units",
"limit_human",
"period_secs",
"expiry_unix",
"recipients",
"conversation_id",
"timestamp",
],
"a non-Fleet wallet_link_lifecycle schema must carry every column except signer_public_key"
);
let resolved = engine
.execute(
"SELECT transition, subject, wallet_address FROM wallet_link_lifecycle WHERE conversation_id = 'conv-wll-redact'",
false,
)
.await
.expect("transition/subject/wallet_address must still resolve for a non-Fleet scope");
assert_eq!(resolved.batches[0].num_rows(), 1);
let transition_col = resolved.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("transition is Utf8");
assert_eq!(transition_col.value(0), "linked");
}
#[tokio::test]
async fn fleet_scope_wallet_link_lifecycle_returns_signer_public_key() {
let signer = ApprovalSigner::from_seed(305);
let bytes = signed_wallet_link_lifecycle_event(&signer, "conv-wll-fleet");
let events = vec![(0, Event::new(kinds::WALLET_LINK_LIFECYCLE, bytes))];
let partitions = vec![PartitionEvents {
partition: "conv-wll-fleet".to_string(),
events,
}];
let trusted_signers = vec![signer.public_key_bytes()];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&trusted_signers,
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute(
"SELECT signer_public_key FROM wallet_link_lifecycle WHERE conversation_id = 'conv-wll-fleet'",
false,
)
.await
.expect("signer_public_key must resolve for Fleet");
assert_eq!(output.batches[0].num_rows(), 1);
let signer_col = output.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::BinaryArray>()
.expect("signer_public_key is Binary");
assert_eq!(signer_col.value(0), signer.public_key_bytes());
}
#[tokio::test]
async fn payments_rows_are_queryable_and_join_to_their_conversation() {
let committed_turn = Uuid::from_u128(400);
let signer = ApprovalSigner::from_seed(101);
let inbound = signed_payment_receipt(&signer, kinds::PAYMENT_RECEIPT, "call-in");
let outbound = signed_payment_receipt(&signer, kinds::OUTBOUND_PAYMENT_RECEIPT, "call-out");
let events = vec![
(
0,
Event::new(
kinds::tagged(kinds::TURN_START, &committed_turn),
Vec::new(),
),
),
(
1,
Event::new(
kinds::tagged(kinds::PAYMENT_RECEIPT, &committed_turn),
inbound,
),
),
(
2,
Event::new(
kinds::tagged(kinds::OUTBOUND_PAYMENT_RECEIPT, &committed_turn),
outbound,
),
),
(
3,
Event::new(
kinds::tagged(kinds::TURN_COMPLETE, &committed_turn),
Vec::new(),
),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-payments".to_string(),
events,
}];
let trusted_signers = vec![signer.public_key_bytes()];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&trusted_signers,
QueryLimits::default(),
)
.await
.expect("build");
let count = engine
.execute("SELECT COUNT(*) AS c FROM payments", false)
.await
.expect("payments count query");
assert_eq!(single_count(&count), 2, "both inbound and outbound rows");
let joined = engine
.execute(
"SELECT p.direction, p.tool_call_id FROM payments p \
JOIN events e ON e.partition = p.partition \
WHERE p.partition = 'conv-payments' \
GROUP BY p.direction, p.tool_call_id \
ORDER BY p.direction",
false,
)
.await
.expect("join query");
assert_eq!(
joined.batches[0].num_rows(),
2,
"the join matches both payments rows exactly once each"
);
let direction_col = joined.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("direction is Utf8");
let tool_call_id_col = joined.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("tool_call_id is Utf8");
assert_eq!(direction_col.value(0), "inbound");
assert_eq!(tool_call_id_col.value(0), "call-in");
assert_eq!(direction_col.value(1), "outbound");
assert_eq!(tool_call_id_col.value(1), "call-out");
}
#[tokio::test]
async fn payments_excludes_a_forged_or_untrusted_signer_receipt() {
let turn = Uuid::from_u128(401);
let trusted = ApprovalSigner::from_seed(102);
let untrusted = ApprovalSigner::from_seed(103);
let legit = signed_payment_receipt(&trusted, kinds::PAYMENT_RECEIPT, "call-legit");
let forged = signed_payment_receipt(&untrusted, kinds::PAYMENT_RECEIPT, "call-forged");
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
),
(
1,
Event::new(kinds::tagged(kinds::PAYMENT_RECEIPT, &turn), legit),
),
(
2,
Event::new(kinds::tagged(kinds::PAYMENT_RECEIPT, &turn), forged),
),
(
3,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-payments-forged".to_string(),
events,
}];
let trusted_signers = vec![trusted.public_key_bytes()];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&trusted_signers,
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute("SELECT tool_call_id FROM payments", false)
.await
.expect("payments query");
assert_eq!(
output.batches[0].num_rows(),
1,
"only the legitimately-signed receipt appears — the forged one is dropped by \
the shared fold, not by this table's own logic"
);
let tool_call_id = output.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("tool_call_id is Utf8");
assert_eq!(tool_call_id.value(0), "call-legit");
}
#[tokio::test]
#[allow(clippy::too_many_lines)] async fn conversations_scope_payments_signer_public_key_is_unreachable() {
let turn = Uuid::from_u128(402);
let signer = ApprovalSigner::from_seed(104);
let bytes = signed_payment_receipt(&signer, kinds::PAYMENT_RECEIPT, "call-redact");
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
),
(
1,
Event::new(kinds::tagged(kinds::PAYMENT_RECEIPT, &turn), bytes),
),
(
2,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-payments-redact".to_string(),
events,
}];
let trusted_signers = vec![signer.public_key_bytes()];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-payments-redact".to_string()]),
partitions,
ReferenceData::empty(),
&trusted_signers,
QueryLimits::default(),
)
.await
.expect("build");
let result = engine
.execute("SELECT signer_public_key FROM payments", false)
.await;
assert!(
matches!(result, Err(QueryEngineError::DataFusion(_))),
"`signer_public_key` must not resolve as a column of `payments` for a \
non-Fleet scope: {result:?}"
);
let output = engine
.execute("SELECT * FROM payments", false)
.await
.expect("payments table");
let schema = output.batches[0].schema();
let columns: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
assert_eq!(
columns,
vec![
"partition",
"position",
"turn_id",
"direction",
"reference",
"amount",
"currency",
"recipient",
"method",
"timestamp",
"version",
"tool_call_id",
"approval_pos",
"approved_args_hash",
"subject",
"payer_kind",
"paying_account",
],
"a non-Fleet payments schema must carry every column except signer_public_key"
);
let resolved = engine
.execute(
"SELECT direction, amount, subject, payer_kind, paying_account FROM payments WHERE tool_call_id = 'call-redact'",
false,
)
.await
.expect("direction/amount/subject/payer_kind/paying_account must still resolve for a non-Fleet scope");
assert_eq!(resolved.batches[0].num_rows(), 1);
let direction_col = resolved.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("direction is Utf8");
let amount_col = resolved.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("amount is Utf8");
let subject_col = resolved.batches[0]
.column(2)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("subject is Utf8");
let payer_kind_col = resolved.batches[0]
.column(3)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("payer_kind is Utf8");
let paying_account_col = resolved.batches[0]
.column(4)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("paying_account is Utf8");
assert_eq!(direction_col.value(0), "inbound");
assert_eq!(amount_col.value(0), "2.75");
assert_eq!(subject_col.value(0), "persona-payer");
assert_eq!(
payer_kind_col.value(0),
"linked_wallet",
"payer_kind must resolve with its real value at a non-Fleet scope — same non-redacted \
class as subject/recipient, not the raw signer key"
);
assert_eq!(paying_account_col.value(0), "0xengine-payer");
}
#[tokio::test]
async fn fleet_scope_payments_returns_signer_public_key() {
let turn = Uuid::from_u128(403);
let signer = ApprovalSigner::from_seed(105);
let bytes = signed_payment_receipt(&signer, kinds::PAYMENT_RECEIPT, "call-fleet");
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
),
(
1,
Event::new(kinds::tagged(kinds::PAYMENT_RECEIPT, &turn), bytes),
),
(
2,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-payments-fleet".to_string(),
events,
}];
let trusted_signers = vec![signer.public_key_bytes()];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&trusted_signers,
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute(
"SELECT signer_public_key, payer_kind, paying_account FROM payments WHERE tool_call_id = 'call-fleet'",
false,
)
.await
.expect("signer_public_key/payer_kind/paying_account must resolve for Fleet");
assert_eq!(output.batches[0].num_rows(), 1);
let signer_col = output.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::BinaryArray>()
.expect("signer_public_key is Binary");
assert_eq!(signer_col.value(0), signer.public_key_bytes());
let payer_kind_col = output.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("payer_kind is Utf8");
let paying_account_col = output.batches[0]
.column(2)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("paying_account is Utf8");
assert_eq!(payer_kind_col.value(0), "linked_wallet");
assert_eq!(paying_account_col.value(0), "0xengine-payer");
}
#[allow(clippy::too_many_arguments)] fn signed_approval_response(
signer: &ApprovalSigner,
request_id: &str,
tool_name: &str,
args_json: &str,
caller: &str,
reason: &str,
conversation_id: &str,
nonce: &str,
) -> Vec<u8> {
polyc_crypto::approval::response_payload(
request_id,
tool_name,
args_json,
"",
true,
false,
&[],
caller,
"",
"default",
reason,
"",
conversation_id,
nonce,
"",
signer,
)
.0
}
#[tokio::test]
async fn approvals_rows_are_queryable_and_join_to_their_conversation() {
let committed_turn = Uuid::from_u128(500);
let signer = ApprovalSigner::from_seed(106);
let request_bytes = polyc_crypto::approval::request_payload(
"call-approve",
"rm",
r#"{"path":"/tmp"}"#,
"default",
"",
&[],
"",
"",
"",
&[],
false,
);
let response_bytes = signed_approval_response(
&signer,
"call-approve",
"rm",
r#"{"path":"/tmp"}"#,
"caller-1",
"ok",
"conv-approvals",
"nonce-1",
);
let events = vec![
(
0,
Event::new(
kinds::tagged(kinds::TURN_START, &committed_turn),
Vec::new(),
),
),
(
1,
Event::new(
kinds::tagged(kinds::APPROVAL_REQUEST, &committed_turn),
request_bytes,
),
),
(
2,
Event::new(
kinds::tagged(kinds::APPROVAL_RESPONSE, &committed_turn),
response_bytes,
),
),
(
3,
Event::new(
kinds::tagged(kinds::TURN_COMPLETE, &committed_turn),
Vec::new(),
),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-approvals".to_string(),
events,
}];
let trusted_signers = vec![signer.public_key_bytes()];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&trusted_signers,
QueryLimits::default(),
)
.await
.expect("build");
let count = engine
.execute("SELECT COUNT(*) AS c FROM approvals", false)
.await
.expect("approvals count query");
assert_eq!(single_count(&count), 2, "both request and response rows");
let joined = engine
.execute(
"SELECT a.phase, a.request_id FROM approvals a \
JOIN events e ON e.partition = a.partition \
WHERE a.partition = 'conv-approvals' \
GROUP BY a.phase, a.request_id \
ORDER BY a.phase",
false,
)
.await
.expect("join query");
assert_eq!(
joined.batches[0].num_rows(),
2,
"the join matches both approvals rows exactly once each"
);
let phase_col = joined.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("phase is Utf8");
assert_eq!(phase_col.value(0), "request");
assert_eq!(phase_col.value(1), "response");
}
#[allow(clippy::too_many_lines)]
#[tokio::test]
async fn conversations_scope_approvals_signer_public_key_is_unreachable() {
let turn = Uuid::from_u128(501);
let signer = ApprovalSigner::from_seed(107);
let response_bytes = signed_approval_response(
&signer,
"call-redact",
"rm",
"{}",
"caller-redact",
"ok",
"conv-approvals-redact",
"nonce-1",
);
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
),
(
1,
Event::new(
kinds::tagged(kinds::APPROVAL_RESPONSE, &turn),
response_bytes,
),
),
(
2,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-approvals-redact".to_string(),
events,
}];
let trusted_signers = vec![signer.public_key_bytes()];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-approvals-redact".to_string()]),
partitions,
ReferenceData::empty(),
&trusted_signers,
QueryLimits::default(),
)
.await
.expect("build");
let result = engine
.execute("SELECT signer_public_key FROM approvals", false)
.await;
assert!(
matches!(result, Err(QueryEngineError::DataFusion(_))),
"`signer_public_key` must not resolve as a column of `approvals` for a \
non-Fleet scope: {result:?}"
);
let approvals_raw_unqualified = engine.execute("SELECT * FROM approvals_raw", false).await;
assert!(
matches!(
approvals_raw_unqualified,
Err(QueryEngineError::DataFusion(_))
),
"approvals_raw must not resolve for a non-Fleet scope: {approvals_raw_unqualified:?}"
);
let output = engine
.execute("SELECT * FROM approvals", false)
.await
.expect("approvals table");
let schema = output.batches[0].schema();
let columns: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
assert_eq!(
columns,
vec![
"partition",
"position",
"turn_id",
"phase",
"request_id",
"tool_name",
"args_json",
"approved",
"response_reason",
"signature_status",
"routine_grant",
"tool_descriptor_hash",
"grant_scope",
],
"a non-Fleet approvals schema must carry exactly the participant-visible column set"
);
let resolved = engine
.execute(
"SELECT approved, response_reason, signature_status FROM approvals \
WHERE request_id = 'call-redact'",
false,
)
.await
.expect(
"approved/response_reason/signature_status must still resolve for a non-Fleet scope",
);
assert_eq!(resolved.batches[0].num_rows(), 1);
let approved_col = resolved.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::BooleanArray>()
.expect("approved is Boolean");
let reason_col = resolved.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("response_reason is Utf8");
let signature_col = resolved.batches[0]
.column(2)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("signature_status is Utf8");
assert!(approved_col.value(0));
assert_eq!(reason_col.value(0), "ok");
assert_eq!(signature_col.value(0), "verified");
}
#[tokio::test]
async fn fleet_scope_approvals_returns_signer_public_key() {
let turn = Uuid::from_u128(502);
let signer = ApprovalSigner::from_seed(108);
let response_bytes = signed_approval_response(
&signer,
"call-fleet-approval",
"rm",
"{}",
"caller-fleet",
"ok",
"conv-approvals-fleet",
"nonce-1",
);
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
),
(
1,
Event::new(
kinds::tagged(kinds::APPROVAL_RESPONSE, &turn),
response_bytes,
),
),
(
2,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-approvals-fleet".to_string(),
events,
}];
let trusted_signers = vec![signer.public_key_bytes()];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&trusted_signers,
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute(
"SELECT signer_public_key FROM approvals WHERE request_id = 'call-fleet-approval'",
false,
)
.await
.expect("signer_public_key must resolve for Fleet");
assert_eq!(output.batches[0].num_rows(), 1);
let signer_col = output.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::BinaryArray>()
.expect("signer_public_key is Binary");
assert_eq!(signer_col.value(0), signer.public_key_bytes());
}
fn signed_handoff_bytes(
signer: &polyc_crypto::signing_role::HandoffSigner,
child_conversation_id: &str,
child_agent_id: &str,
carried_count: u32,
reason: &str,
) -> Vec<u8> {
let mut h = polyc_proto::proto::polychrome::handoff::v1::Handoff {
child_conversation_id: child_conversation_id.to_owned(),
child_agent_id: child_agent_id.to_owned(),
carried_count,
reason: reason.to_owned(),
..Default::default()
};
polyc_crypto::handoff::sign_handoff_into(signer, &mut h);
h.encode_to_vec()
}
#[allow(clippy::too_many_arguments)] fn signed_handoff_denied_bytes(
signer: &polyc_crypto::signing_role::HandoffSigner,
parent_conversation_id: &str,
parent_agent_id: &str,
child_agent_id: &str,
reason: &str,
denial_reason: &str,
allowed: &[&str],
) -> Vec<u8> {
let mut d = polyc_proto::proto::polychrome::handoff::v1::HandoffDenied {
parent_conversation_id: parent_conversation_id.to_owned(),
parent_agent_id: parent_agent_id.to_owned(),
child_agent_id: child_agent_id.to_owned(),
reason: reason.to_owned(),
denial_reason: denial_reason.to_owned(),
allowed: allowed.iter().map(|s| (*s).to_owned()).collect(),
..Default::default()
};
polyc_crypto::handoff::sign_handoff_denied_into(signer, &mut d);
d.encode_to_vec()
}
#[tokio::test]
async fn handoffs_rows_are_queryable_and_join_to_their_conversation() {
let committed_turn = Uuid::from_u128(600);
let signer = super::fixture_handoff_signer();
let handoff_bytes = signed_handoff_bytes(&signer, "child-x", "researcher", 2, "delegate");
let denied_bytes = signed_handoff_denied_bytes(
&signer,
"conv-handoffs",
"assistant",
"banned-agent",
"delegate weird task",
"this agent can't hand off to that agent",
&["coding", "research"],
);
let events = vec![
(
0,
Event::new(
kinds::tagged(kinds::TURN_START, &committed_turn),
Vec::new(),
),
),
(
1,
Event::new(
kinds::tagged(kinds::HANDOFF, &committed_turn),
handoff_bytes,
),
),
(
2,
Event::new(
kinds::tagged(kinds::HANDOFF_DENIED, &committed_turn),
denied_bytes,
),
),
(
3,
Event::new(
kinds::tagged(kinds::TURN_COMPLETE, &committed_turn),
Vec::new(),
),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-handoffs".to_string(),
events,
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let count = engine
.execute("SELECT COUNT(*) AS c FROM handoffs", false)
.await
.expect("handoffs count query");
assert_eq!(single_count(&count), 2, "handoff and handoff_denied rows");
let joined = engine
.execute(
"SELECT h.phase, h.child_conversation_id FROM handoffs h \
JOIN events e ON e.partition = h.partition \
WHERE h.partition = 'conv-handoffs' \
GROUP BY h.phase, h.child_conversation_id \
ORDER BY h.phase",
false,
)
.await
.expect("join query");
assert_eq!(
joined.batches[0].num_rows(),
2,
"the join matches every handoffs row exactly once each"
);
let phase_col = joined.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("phase is Utf8");
assert_eq!(phase_col.value(0), "handoff");
assert_eq!(phase_col.value(1), "handoff_denied");
}
#[tokio::test]
async fn independently_committed_handoffs_are_queryable() {
let handoff_id = Uuid::from_u128(603);
let denied_id = Uuid::from_u128(604);
let signer = super::fixture_handoff_signer();
let events = vec![
(
0,
Event::new(
kinds::tagged(kinds::HANDOFF, &handoff_id),
signed_handoff_bytes(&signer, "child-live", "researcher", 1, "delegate"),
),
),
(
1,
Event::new(
kinds::tagged(kinds::HANDOFF_DENIED, &denied_id),
signed_handoff_denied_bytes(
&signer,
"conv-handoffs-live",
"assistant",
"blocked",
"delegate",
"not allowed",
&["researcher"],
),
),
),
];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
vec![PartitionEvents {
partition: "conv-handoffs-live".to_owned(),
events,
}],
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let count = engine
.execute("SELECT COUNT(*) AS c FROM handoffs", false)
.await
.expect("handoffs count query");
assert_eq!(
single_count(&count),
2,
"independently committed handoff records must remain visible"
);
}
#[allow(clippy::too_many_lines)]
#[tokio::test]
async fn conversations_scope_handoffs_signed_by_is_unreachable() {
let handoff_id = Uuid::from_u128(601);
let denied_id = Uuid::from_u128(605);
let signer = super::fixture_handoff_signer();
let handoff_bytes =
signed_handoff_bytes(&signer, "child-redact", "researcher", 1, "delegate-redact");
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::HANDOFF, &handoff_id), handoff_bytes),
),
(
1,
Event::new(
kinds::tagged(kinds::HANDOFF_DENIED, &denied_id),
signed_handoff_denied_bytes(
&signer,
"conv-handoffs-redact",
"assistant",
"blocked",
"delegate-redact",
"not allowed",
&["researcher"],
),
),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-handoffs-redact".to_string(),
events,
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-handoffs-redact".to_string()]),
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let result = engine
.execute("SELECT signed_by FROM handoffs", false)
.await;
assert!(
matches!(result, Err(QueryEngineError::DataFusion(_))),
"`signed_by` must not resolve as a column of `handoffs` for a non-Fleet scope: {result:?}"
);
let handoffs_raw_unqualified = engine.execute("SELECT * FROM handoffs_raw", false).await;
assert!(
matches!(
handoffs_raw_unqualified,
Err(QueryEngineError::DataFusion(_))
),
"handoffs_raw must not resolve for a non-Fleet scope: {handoffs_raw_unqualified:?}"
);
let output = engine
.execute("SELECT * FROM handoffs", false)
.await
.expect("handoffs table");
assert_eq!(
output.batches[0].num_rows(),
2,
"the redacted view keeps independent handoff commands"
);
let schema = output.batches[0].schema();
let columns: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
assert_eq!(
columns,
vec![
"partition",
"position",
"turn_id",
"phase",
"child_conversation_id",
"child_agent_id",
"carried_count",
"reason",
"parent_agent_id",
"denial_reason",
"allowed",
"signature_status",
],
"a non-Fleet handoffs schema must carry every column except signed_by"
);
let resolved = engine
.execute(
"SELECT child_agent_id, carried_count, reason FROM handoffs \
WHERE child_conversation_id = 'child-redact'",
false,
)
.await
.expect("child_agent_id/carried_count/reason must still resolve for a non-Fleet scope");
assert_eq!(resolved.batches[0].num_rows(), 1);
let child_agent_id_col = resolved.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("child_agent_id is Utf8");
let carried_count_col = resolved.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::UInt32Array>()
.expect("carried_count is UInt32");
let reason_col = resolved.batches[0]
.column(2)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("reason is Utf8");
assert_eq!(child_agent_id_col.value(0), "researcher");
assert_eq!(carried_count_col.value(0), 1);
assert_eq!(reason_col.value(0), "delegate-redact");
}
#[tokio::test]
async fn fleet_scope_handoffs_returns_signed_by() {
let turn = Uuid::from_u128(602);
let signer = super::fixture_handoff_signer();
let handoff_bytes =
signed_handoff_bytes(&signer, "child-fleet", "researcher", 3, "delegate-fleet");
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
),
(
1,
Event::new(kinds::tagged(kinds::HANDOFF, &turn), handoff_bytes),
),
(
2,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-handoffs-fleet".to_string(),
events,
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute(
"SELECT signed_by FROM handoffs WHERE child_conversation_id = 'child-fleet'",
false,
)
.await
.expect("signed_by must resolve for Fleet");
assert_eq!(output.batches[0].num_rows(), 1);
let signed_by_col = output.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::BinaryArray>()
.expect("signed_by is Binary");
assert_eq!(signed_by_col.value(0), signer.public_key_bytes());
}
fn signed_grant_replay(
signer: &ApprovalSigner,
conversation_id: &str,
turn_id: &str,
tool: &str,
grant_ref: &str,
covered_capabilities: &[String],
coverage_hash: &str,
) -> Vec<u8> {
polyc_crypto::approval::test_util::grant_replay_payload(
conversation_id,
turn_id,
tool,
grant_ref,
covered_capabilities,
coverage_hash,
signer,
)
.0
}
#[tokio::test]
async fn grant_replays_rows_are_queryable_and_join_to_their_conversation() {
let committed_turn = Uuid::from_u128(700);
let signer = ApprovalSigner::from_seed(120);
let payload_bytes = signed_grant_replay(
&signer,
"conv-grant-replays",
&committed_turn.to_string(),
"read_file",
"grant-ref-engine",
&["fs.read".to_owned()],
"sha256:coverage-engine",
);
let events = vec![
(
0,
Event::new(
kinds::tagged(kinds::TURN_START, &committed_turn),
Vec::new(),
),
),
(
1,
Event::new(
kinds::tagged(kinds::GRANT_REPLAY, &committed_turn),
payload_bytes,
),
),
(
2,
Event::new(
kinds::tagged(kinds::TURN_COMPLETE, &committed_turn),
Vec::new(),
),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-grant-replays".to_string(),
events,
}];
let trusted_signers = vec![signer.public_key_bytes()];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&trusted_signers,
QueryLimits::default(),
)
.await
.expect("build");
let count = engine
.execute("SELECT COUNT(*) AS c FROM grant_replays", false)
.await
.expect("grant_replays count query");
assert_eq!(single_count(&count), 1);
let joined = engine
.execute(
"SELECT g.tool, g.grant_ref FROM grant_replays g \
JOIN events e ON e.partition = g.partition \
WHERE g.partition = 'conv-grant-replays' \
GROUP BY g.tool, g.grant_ref",
false,
)
.await
.expect("join query");
assert_eq!(
joined.batches[0].num_rows(),
1,
"the join matches the grant_replays row exactly once"
);
let tool_col = joined.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("tool is Utf8");
assert_eq!(tool_col.value(0), "read_file");
}
#[allow(clippy::too_many_lines)]
#[tokio::test]
async fn conversations_scope_grant_replays_signer_public_key_is_unreachable() {
let turn = Uuid::from_u128(701);
let signer = ApprovalSigner::from_seed(121);
let payload_bytes = signed_grant_replay(
&signer,
"conv-grant-replays-redact",
&turn.to_string(),
"read_file",
"grant-ref-redact",
&["fs.read".to_owned()],
"sha256:coverage-redact",
);
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
),
(
1,
Event::new(kinds::tagged(kinds::GRANT_REPLAY, &turn), payload_bytes),
),
(
2,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-grant-replays-redact".to_string(),
events,
}];
let trusted_signers = vec![signer.public_key_bytes()];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-grant-replays-redact".to_string()]),
partitions,
ReferenceData::empty(),
&trusted_signers,
QueryLimits::default(),
)
.await
.expect("build");
let result = engine
.execute("SELECT signer_public_key FROM grant_replays", false)
.await;
assert!(
matches!(result, Err(QueryEngineError::DataFusion(_))),
"`signer_public_key` must not resolve as a column of `grant_replays` for a \
non-Fleet scope: {result:?}"
);
let grant_replays_raw_unqualified = engine
.execute("SELECT * FROM grant_replays_raw", false)
.await;
assert!(
matches!(
grant_replays_raw_unqualified,
Err(QueryEngineError::DataFusion(_))
),
"grant_replays_raw must not resolve for a non-Fleet scope: {grant_replays_raw_unqualified:?}"
);
let output = engine
.execute("SELECT * FROM grant_replays", false)
.await
.expect("grant_replays table");
let schema = output.batches[0].schema();
let columns: Vec<&str> = schema.fields().iter().map(|f| f.name().as_str()).collect();
assert_eq!(
columns,
vec![
"partition",
"position",
"turn_id",
"tool",
"grant_ref",
"covered_capabilities",
"coverage_hash",
"signature_status",
],
"a non-Fleet grant_replays schema must carry exactly the participant-visible column set"
);
let resolved = engine
.execute(
"SELECT tool, grant_ref, covered_capabilities, coverage_hash, signature_status \
FROM grant_replays WHERE grant_ref = 'grant-ref-redact'",
false,
)
.await
.expect(
"tool/grant_ref/covered_capabilities/coverage_hash/signature_status must still \
resolve for a non-Fleet scope",
);
assert_eq!(resolved.batches[0].num_rows(), 1);
let tool_col = resolved.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("tool is Utf8");
let signature_col = resolved.batches[0]
.column(4)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("signature_status is Utf8");
assert_eq!(tool_col.value(0), "read_file");
assert_eq!(signature_col.value(0), "verified");
}
#[tokio::test]
async fn fleet_scope_grant_replays_returns_signer_public_key() {
let turn = Uuid::from_u128(702);
let signer = ApprovalSigner::from_seed(122);
let payload_bytes = signed_grant_replay(
&signer,
"conv-grant-replays-fleet",
&turn.to_string(),
"read_file",
"grant-ref-fleet",
&["fs.read".to_owned()],
"sha256:coverage-fleet",
);
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
),
(
1,
Event::new(kinds::tagged(kinds::GRANT_REPLAY, &turn), payload_bytes),
),
(
2,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-grant-replays-fleet".to_string(),
events,
}];
let trusted_signers = vec![signer.public_key_bytes()];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&trusted_signers,
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute(
"SELECT signer_public_key FROM grant_replays WHERE grant_ref = 'grant-ref-fleet'",
false,
)
.await
.expect("signer_public_key must resolve for Fleet");
assert_eq!(output.batches[0].num_rows(), 1);
let signer_col = output.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::BinaryArray>()
.expect("signer_public_key is Binary");
assert_eq!(signer_col.value(0), signer.public_key_bytes());
}
fn wire_text_message(role: &str, text: &str, internal_only: bool) -> WireMessage {
WireMessage {
role: role.to_string(),
content: buffa::MessageField::some(Content {
r#type: Some(content::Type::Text(Box::new(TextContent {
text: text.to_string(),
..Default::default()
}))),
..Default::default()
}),
internal_only,
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn wire_tool_call_message(id: &str, name: &str, args_json: &str) -> WireMessage {
let arguments = serde_json::from_str::<buffa_types::google::protobuf::Struct>(args_json)
.map(buffa::MessageField::some)
.unwrap_or_default();
WireMessage {
role: "model".to_string(),
content: buffa::MessageField::some(Content {
r#type: Some(content::Type::ToolCall(Box::new(ToolCallContent {
id: id.to_string(),
r#type: Some(tool_call_content::Type::FunctionCall(Box::new(
FunctionCallContent {
name: name.to_string(),
arguments,
..Default::default()
},
))),
..Default::default()
}))),
..Default::default()
}),
internal_only: false,
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
fn wire_tool_result_message(
id: &str,
name: &str,
result_json: &str,
first_party: bool,
) -> WireMessage {
let response = serde_json::from_str::<buffa_types::google::protobuf::Struct>(result_json)
.ok()
.map(|s| function_result_content::Result::Response(Box::new(s)));
WireMessage {
role: "tool".to_string(),
content: buffa::MessageField::some(Content {
r#type: Some(content::Type::ToolResult(Box::new(ToolResultContent {
call_id: id.to_string(),
first_party,
r#type: Some(tool_result_content::Type::FunctionResult(Box::new(
FunctionResultContent {
name: name.to_string(),
result: response,
..Default::default()
},
))),
..Default::default()
}))),
..Default::default()
}),
internal_only: false,
__buffa_unknown_fields: buffa::UnknownFields::default(),
}
}
#[tokio::test]
async fn message_content_rows_are_queryable_and_join_to_their_conversation() {
let committed_turn = Uuid::from_u128(310);
let text = wire_text_message("user", "please search", false);
let call = wire_tool_call_message("call-1", "search", r#"{"q":"rust"}"#);
let result = wire_tool_result_message("call-1", "search", r#"{"hits":3}"#, true);
let events = vec![
(
0,
Event::new(
kinds::tagged(kinds::TURN_START, &committed_turn),
Vec::new(),
),
),
(
1,
Event::trusted(
kinds::tagged(kinds::USER_MSG, &committed_turn),
text.encode_to_vec(),
),
),
(
2,
Event::new(
kinds::tagged(kinds::OUTPUT_MSG, &committed_turn),
call.encode_to_vec(),
),
),
(
3,
Event::new(
kinds::tagged(kinds::USER_MSG, &committed_turn),
result.encode_to_vec(),
),
),
(
4,
Event::new(
kinds::tagged(kinds::TURN_COMPLETE, &committed_turn),
Vec::new(),
),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-message-content".to_string(),
events,
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let messages_count = engine
.execute("SELECT COUNT(*) AS c FROM messages", false)
.await
.expect("messages count query");
assert_eq!(single_count(&messages_count), 1, "one text block");
let tool_calls_count = engine
.execute("SELECT COUNT(*) AS c FROM tool_calls", false)
.await
.expect("tool_calls count query");
assert_eq!(
single_count(&tool_calls_count),
2,
"one call row and one result row"
);
let joined = engine
.execute(
"SELECT m.role, m.text FROM messages m \
JOIN events e ON e.partition = m.partition \
WHERE m.partition = 'conv-message-content' \
GROUP BY m.role, m.text",
false,
)
.await
.expect("messages join query");
assert_eq!(joined.batches[0].num_rows(), 1);
let role_col = joined.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("role is Utf8");
let text_col = joined.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("text is Utf8");
assert_eq!(role_col.value(0), "user");
assert_eq!(text_col.value(0), "please search");
}
#[tokio::test]
async fn tool_call_and_its_result_share_a_tool_call_id() {
let turn = Uuid::from_u128(311);
let call = wire_tool_call_message("call-1", "search", r#"{"q":"rust"}"#);
let result = wire_tool_result_message("call-1", "search", r#"{"hits":3}"#, true);
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
),
(
1,
Event::new(
kinds::tagged(kinds::OUTPUT_MSG, &turn),
call.encode_to_vec(),
),
),
(
2,
Event::new(
kinds::tagged(kinds::USER_MSG, &turn),
result.encode_to_vec(),
),
),
(
3,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
),
];
let partitions = vec![PartitionEvents {
partition: "conv-pairing".to_string(),
events,
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let paired = engine
.execute(
"SELECT c.arguments, r.result, r.first_party \
FROM tool_calls c \
JOIN tool_calls r ON r.tool_call_id = c.tool_call_id AND r.block_type = 'result' \
WHERE c.block_type = 'call' AND c.tool_call_id = 'call-1'",
false,
)
.await
.expect("call/result self-join");
assert_eq!(paired.batches[0].num_rows(), 1);
let arguments_col = paired.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("arguments is Utf8");
let result_col = paired.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("result is Utf8");
let first_party_col = paired.batches[0]
.column(2)
.as_any()
.downcast_ref::<arrow::array::BooleanArray>()
.expect("first_party is Boolean");
assert_eq!(arguments_col.value(0), r#"{"q":"rust"}"#);
assert_eq!(result_col.value(0), r#"{"hits":3.0}"#);
assert!(first_party_col.value(0));
let call_row = engine
.execute(
"SELECT result, first_party FROM tool_calls \
WHERE block_type = 'call' AND tool_call_id = 'call-1'",
false,
)
.await
.expect("call row query");
assert!(
call_row.batches[0].column(0).is_null(0),
"call row has no result"
);
assert!(
call_row.batches[0].column(1).is_null(0),
"call row has no first_party, not false"
);
}
#[tokio::test]
async fn conversations_scope_hides_a_paused_turns_narration() {
let turn = Uuid::from_u128(743);
let narration = wire_text_message("model", "Removing it now (pending your approval).", false);
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
),
(
1,
Event::new(
kinds::tagged(kinds::USER_MSG, &turn),
wire_text_message("user", "remove the export", false).encode_to_vec(),
),
),
(
2,
Event::new(
kinds::tagged(kinds::OUTPUT_MSG, &turn),
narration.encode_to_vec(),
),
),
(
3,
Event::new(kinds::tagged(kinds::TURN_TEXT_WITHHELD, &turn), Vec::new()),
),
(
4,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
),
];
let partition_name = "conv-paused-withheld".to_owned();
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec![partition_name.clone()]),
vec![PartitionEvents {
partition: partition_name,
events,
}],
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("conversations build");
let messages = engine
.execute("SELECT text FROM messages", false)
.await
.expect("conversations messages query");
let texts: Vec<String> = messages
.batches
.iter()
.flat_map(|batch| {
let col = batch
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("text is Utf8");
(0..batch.num_rows())
.map(|i| col.value(i).to_owned())
.collect::<Vec<_>>()
})
.collect();
assert!(
texts.iter().any(|t| t == "remove the export"),
"the person's own message stays visible: {texts:?}"
);
assert!(
!texts
.iter()
.any(|t| t == "Removing it now (pending your approval)."),
"a paused turn's narration must not reach a conversation-scoped \
reader: {texts:?}"
);
}
fn internal_only_redaction_fixture() -> (String, Vec<(u64, Event)>) {
let turn = Uuid::from_u128(312);
let visible = wire_text_message("user", "visible to the participant", false);
let hidden = wire_text_message("model", "internal ground-truth note", true);
let hidden_call = {
let mut call = wire_tool_call_message("call-hidden", "search", r#"{"q":"x"}"#);
call.internal_only = true;
call
};
let events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
),
(
1,
Event::new(
kinds::tagged(kinds::USER_MSG, &turn),
visible.encode_to_vec(),
),
),
(
2,
Event::new(
kinds::tagged(kinds::OUTPUT_MSG, &turn),
hidden.encode_to_vec(),
),
),
(
3,
Event::new(
kinds::tagged(kinds::OUTPUT_MSG, &turn),
hidden_call.encode_to_vec(),
),
),
(
4,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
),
];
("conv-internal-redact".to_string(), events)
}
#[tokio::test]
async fn fleet_scope_message_content_includes_internal_only_rows() {
let (partition_name, events) = internal_only_redaction_fixture();
let partitions = vec![PartitionEvents {
partition: partition_name,
events,
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("fleet build");
let messages = engine
.execute("SELECT COUNT(*) AS c FROM messages", false)
.await
.expect("fleet messages count");
assert_eq!(
single_count(&messages),
2,
"Fleet sees both the visible and the internal_only message"
);
let tool_calls = engine
.execute("SELECT COUNT(*) AS c FROM tool_calls", false)
.await
.expect("fleet tool_calls count");
assert_eq!(
single_count(&tool_calls),
1,
"Fleet sees the internal_only tool call too"
);
}
#[tokio::test]
async fn conversations_scope_hides_internal_only_messages_and_tool_calls() {
let (partition_name, events) = internal_only_redaction_fixture();
let partitions = vec![PartitionEvents {
partition: partition_name.clone(),
events,
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec![partition_name]),
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("conversations build");
let messages = engine
.execute("SELECT text FROM messages", false)
.await
.expect("conversations messages query");
assert_eq!(
messages.batches[0].num_rows(),
1,
"the internal_only row must not reach a non-Fleet scope"
);
let text_col = messages.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("text is Utf8");
assert_eq!(text_col.value(0), "visible to the participant");
let tool_calls = engine
.execute("SELECT COUNT(*) AS c FROM tool_calls", false)
.await
.expect("conversations tool_calls count");
assert_eq!(
single_count(&tool_calls),
0,
"the internal_only tool call must not reach a non-Fleet scope either"
);
for table in ["messages_raw", "tool_calls_raw"] {
let result = engine
.execute(&format!("SELECT * FROM {table}"), false)
.await;
assert!(
matches!(result, Err(QueryEngineError::DataFusion(_))),
"`{table}` must not resolve for a non-Fleet scope: {result:?}"
);
}
}
#[tokio::test]
async fn row_cap_truncates_to_exactly_the_cap() {
let events: Vec<(u64, Event)> = (0..25)
.map(|i| (i, Event::new(format!("k{i}"), Vec::new())))
.collect();
let partitions = vec![PartitionEvents {
partition: "conv-cap".to_string(),
events,
}];
let limits = QueryLimits {
row_cap: 10,
..QueryLimits::default()
};
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
limits,
)
.await
.expect("build");
let output = engine
.execute("SELECT * FROM events_raw ORDER BY position", false)
.await
.expect("query");
let total_rows: usize = output.batches.iter().map(RecordBatch::num_rows).sum();
assert_eq!(total_rows, 10);
assert!(output.truncated);
}
#[tokio::test]
async fn row_cap_not_truncated_when_rows_are_at_or_below_cap() {
let events: Vec<(u64, Event)> = (0..5)
.map(|i| (i, Event::new(format!("k{i}"), Vec::new())))
.collect();
let partitions = vec![PartitionEvents {
partition: "conv-under-cap".to_string(),
events,
}];
let limits = QueryLimits {
row_cap: 5,
..QueryLimits::default()
};
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
limits,
)
.await
.expect("build");
let output = engine
.execute("SELECT * FROM events_raw", false)
.await
.expect("query");
let total_rows: usize = output.batches.iter().map(RecordBatch::num_rows).sum();
assert_eq!(total_rows, 5);
assert!(!output.truncated);
}
#[tokio::test]
async fn row_cap_usize_max_does_not_panic_building_the_limit() {
let events: Vec<(u64, Event)> = (0..3)
.map(|i| (i, Event::new(format!("k{i}"), Vec::new())))
.collect();
let partitions = vec![PartitionEvents {
partition: "conv-max-cap".to_string(),
events,
}];
let limits = QueryLimits {
row_cap: usize::MAX,
..QueryLimits::default()
};
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
limits,
)
.await
.expect("build");
let result = engine.execute("SELECT * FROM events_raw", false).await;
assert!(
matches!(result, Err(QueryEngineError::DataFusion(_))),
"row_cap == usize::MAX must fail gracefully through DataFusion's own limit-conversion \
error, never panic: {result:?}"
);
}
fn int_batch(start: i32, rows: i32) -> RecordBatch {
let schema = Arc::new(Schema::new(vec![Field::new("n", DataType::Int32, false)]));
let values: Vec<i32> = (start..start + rows).collect();
RecordBatch::try_new(schema, vec![Arc::new(Int32Array::from(values))]).expect("batch build")
}
fn collect_n_column(batches: &[RecordBatch]) -> Vec<i32> {
batches
.iter()
.flat_map(|batch| {
batch
.column(0)
.as_any()
.downcast_ref::<Int32Array>()
.expect("Int32 column")
.values()
.iter()
.copied()
})
.collect()
}
#[test]
fn cap_rows_truncates_mid_batch_across_multiple_physical_batches() {
let batches = vec![int_batch(0, 3), int_batch(3, 4), int_batch(7, 5)];
let (kept, truncated) = cap_rows(batches, 6);
assert!(truncated, "12 total rows exceeds the cap of 6");
let total_rows: usize = kept.iter().map(RecordBatch::num_rows).sum();
assert_eq!(total_rows, 6, "exactly cap rows are returned");
assert_eq!(
kept.len(),
2,
"the third batch is dropped once the cap is exhausted"
);
assert_eq!(kept[0].num_rows(), 3, "first batch kept whole");
assert_eq!(kept[1].num_rows(), 3, "second batch sliced from 4 to 3");
assert_eq!(
collect_n_column(&kept),
vec![0, 1, 2, 3, 4, 5],
"returned rows are the first `cap` in original order"
);
}
#[test]
fn cap_rows_exact_total_is_not_truncated() {
let batches = vec![int_batch(0, 3), int_batch(3, 3)];
let (kept, truncated) = cap_rows(batches, 6);
assert!(!truncated, "total rows exactly equals the cap");
let total_rows: usize = kept.iter().map(RecordBatch::num_rows).sum();
assert_eq!(total_rows, 6, "all rows are returned");
assert_eq!(kept.len(), 2, "both batches pass through unmodified");
assert_eq!(collect_n_column(&kept), vec![0, 1, 2, 3, 4, 5]);
}
#[tokio::test]
async fn information_schema_resolves_under_fleet() {
let partitions = one_partition("conv-info-fleet", Uuid::from_u128(7), Uuid::from_u128(8));
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute(
"SELECT table_name FROM information_schema.tables WHERE table_name = 'usage'",
false,
)
.await
.expect("information_schema must resolve under Fleet");
let total_rows: usize = output.batches.iter().map(RecordBatch::num_rows).sum();
assert_eq!(total_rows, 1);
}
#[tokio::test]
async fn information_schema_fails_under_conversations() {
let partitions = one_partition("conv-info-conv", Uuid::from_u128(9), Uuid::from_u128(10));
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-info-conv".to_string()]),
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let result = engine
.execute("SELECT * FROM information_schema.tables", false)
.await;
assert!(
matches!(result, Err(QueryEngineError::DataFusion(_))),
"information_schema must not resolve outside Fleet: {result:?}"
);
}
#[tokio::test]
async fn registry_typed_kinds_match_the_tables_query_engine_build_actually_registers() {
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
Vec::new(),
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute("SELECT table_name FROM information_schema.tables", false)
.await
.expect("information_schema query");
let table_name_col = output.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("table_name is Utf8");
let mut registered: std::collections::BTreeSet<String> = (0..output.batches[0].num_rows())
.map(|i| table_name_col.value(i).to_string())
.collect();
registered.retain(|name| {
name != "events"
&& !name.ends_with("_raw")
&& name != "personas"
&& name != "participations"
&& name != "persona_identities"
&& name != "persona_wallets"
&& name != "persona_spend_policies"
&& name != "persona_credentials"
&& name != "persona_usage"
&& name != "dashboard"
&& name != "routine_grants"
&& name != "routine_active_grants"
&& name != "routine_overview"
&& name != "routine_refusals"
&& !matches!(
name.as_str(),
"tables"
| "views"
| "columns"
| "df_settings"
| "schemata"
| "routines"
| "parameters"
)
});
let mut expected: std::collections::BTreeSet<String> = std::collections::BTreeSet::new();
for (_kind_base, decision) in crate::decode::REGISTRY {
if let crate::decode::Decode::Typed(label) = decision {
if *label == "message_content" {
expected.insert("messages".to_string());
expected.insert("tool_calls".to_string());
} else {
expected.insert((*label).to_string());
}
}
}
assert_eq!(
registered, expected,
"REGISTRY's typed kinds and QueryEngine::build's actual table registrations have \
diverged — every REGISTRY Typed kind must have a matching registered table, and \
every registered typed table must have a matching REGISTRY entry"
);
}
#[tokio::test]
async fn url_table_reference_fails_under_every_scope() {
for scope in [
QueryScope::Fleet,
QueryScope::Conversations(vec!["conv-url".to_string()]),
] {
let partitions = one_partition("conv-url", Uuid::from_u128(11), Uuid::from_u128(12));
let engine = QueryEngine::build(
&test_base_state(),
&scope,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let result = engine.execute("SELECT * FROM 'x.csv'", false).await;
assert!(
matches!(result, Err(QueryEngineError::DataFusion(_))),
"ad hoc external tables must never resolve ({scope:?}): {result:?}"
);
}
}
#[tokio::test]
async fn explain_honors_allow_explain_both_ways() {
let partitions = one_partition("conv-explain", Uuid::from_u128(13), Uuid::from_u128(14));
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let allowed = engine.execute("EXPLAIN SELECT * FROM events", true).await;
assert!(
allowed.is_ok(),
"opted-in EXPLAIN must succeed: {allowed:?}"
);
assert!(!allowed.expect("checked above").batches.is_empty());
let rejected = engine.execute("EXPLAIN SELECT * FROM events", false).await;
assert!(
matches!(rejected, Err(QueryEngineError::Rejected(_))),
"EXPLAIN without opt-in must be rejected before planning: {rejected:?}"
);
}
#[tokio::test]
async fn explain_analyze_is_always_rejected() {
let partitions = one_partition("conv-analyze", Uuid::from_u128(15), Uuid::from_u128(16));
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
for allow_explain in [true, false] {
let result = engine
.execute("EXPLAIN ANALYZE SELECT * FROM events", allow_explain)
.await;
assert!(
matches!(result, Err(QueryEngineError::Rejected(_))),
"EXPLAIN ANALYZE executes the query and stays banned everywhere \
(allow_explain={allow_explain}): {result:?}"
);
}
}
#[tokio::test]
async fn multi_partition_tables_union_every_supplied_partition() {
let partitions = vec![
PartitionEvents {
partition: "conv-x".to_string(),
events: vec![(0, Event::new(kinds::TURN_START, Vec::new()))],
},
PartitionEvents {
partition: "conv-y".to_string(),
events: vec![
(0, Event::new(kinds::TURN_START, Vec::new())),
(1, Event::new(kinds::TURN_START, Vec::new())),
],
},
];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute(
"SELECT partition, COUNT(*) AS c FROM events_raw GROUP BY partition ORDER BY partition",
false,
)
.await
.expect("group-by query");
let batch = &output.batches[0];
assert_eq!(batch.num_rows(), 2, "one row per unioned partition");
let partition_col = batch
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("partition is Utf8");
assert_eq!(partition_col.value(0), "conv-x");
assert_eq!(partition_col.value(1), "conv-y");
}
#[tokio::test]
async fn output_to_json_round_trips_the_raw_payload_column() {
let payload = vec![0xDE, 0xAD, 0xBE, 0xEF, 0x00, 0x01];
let partitions = vec![PartitionEvents {
partition: "conv-payload".to_string(),
events: vec![(0, Event::trusted(kinds::USER_MSG, payload.clone()))],
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute(
"SELECT position, payload FROM events_raw ORDER BY position",
false,
)
.await
.expect("payload query");
let json = crate::output::output_to_json(&output, 0).expect("convert to wire JSON");
assert_eq!(json.columns, vec!["position", "payload"]);
assert_eq!(json.rows.len(), 1, "one event, one row");
assert_eq!(
json.rows[0][1],
serde_json::json!("deadbeef0001"),
"Binary encodes as lowercase hex with no 0x prefix"
);
}
#[tokio::test]
async fn events_raw_payload_json_is_queryable_for_a_schemaless_kind() {
let partitions = vec![PartitionEvents {
partition: "conv-json".to_string(),
events: vec![(
0,
Event::new(
kinds::APPROVAL_DEFERRED,
br#"{"reason":"awaiting reviewer","attempt":2}"#.to_vec(),
),
)],
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute(
"SELECT payload_json FROM events_raw \
WHERE kind_base = 'approval_deferred' AND payload_json LIKE '%\"attempt\":2%'",
false,
)
.await
.expect("payload_json query");
assert_eq!(
output.batches[0].num_rows(),
1,
"the LIKE predicate matches the one row"
);
let payload_json = output.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("payload_json is Utf8")
.value(0);
let parsed: serde_json::Value =
serde_json::from_str(payload_json).expect("payload_json is valid JSON text");
assert_eq!(
parsed["reason"], "awaiting reviewer",
"a field is extractable from the returned JSON text"
);
}
#[tokio::test]
async fn events_raw_payload_json_is_null_for_a_non_json_payload() {
let partitions = vec![PartitionEvents {
partition: "conv-nonjson".to_string(),
events: vec![(
0,
Event::trusted(kinds::USER_MSG, vec![0xFF, 0xFE, 0x00, 0x01]),
)],
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute("SELECT payload_json FROM events_raw", false)
.await
.expect("payload_json query");
let array = output.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("payload_json is Utf8");
assert_eq!(array.len(), 1);
assert!(
array.is_null(0),
"a non-JSON/binary payload must decode to NULL, not a query error"
);
}
#[tokio::test]
async fn events_raw_payload_json_is_queryable_via_json_get_for_a_schemaless_kind() {
let partitions = vec![PartitionEvents {
partition: "conv-json-get".to_string(),
events: vec![(
0,
Event::new(
kinds::APPROVAL_DEFERRED,
br#"{"reason":"awaiting reviewer","attempt":2}"#.to_vec(),
),
)],
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute(
"SELECT payload_json ->> 'reason' AS reason, \
json_get_int(payload_json, 'attempt') AS attempt \
FROM events_raw WHERE kind_base = 'approval_deferred'",
false,
)
.await
.expect("json_get query over payload_json");
assert_eq!(output.batches[0].num_rows(), 1);
let reason = output.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("reason (->>) is Utf8")
.value(0);
assert_eq!(reason, "awaiting reviewer");
let attempt = output.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::Int64Array>()
.expect("attempt (json_get_int) is Int64")
.value(0);
assert_eq!(attempt, 2);
}
#[tokio::test]
async fn tool_calls_arguments_is_queryable_via_json_get() {
let turn = Uuid::from_u128(320);
let call = wire_tool_call_message("call-1", "search", r#"{"q":"rust"}"#);
let partitions = vec![PartitionEvents {
partition: "conv-json-tool-call".to_string(),
events: vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
),
(
1,
Event::new(
kinds::tagged(kinds::OUTPUT_MSG, &turn),
call.encode_to_vec(),
),
),
(
2,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
),
],
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute(
"SELECT arguments ->> 'q' AS q FROM tool_calls WHERE block_type = 'call'",
false,
)
.await
.expect("json_get query over tool_calls.arguments");
assert_eq!(output.batches[0].num_rows(), 1);
let q = output.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("q (->>) is Utf8")
.value(0);
assert_eq!(q, "rust");
}
#[tokio::test]
async fn events_raw_payload_json_extraction_is_null_for_a_typed_payload() {
let usage = UsageEvent {
input_tokens: 1,
output_tokens: 2,
..Default::default()
};
let partitions = vec![PartitionEvents {
partition: "conv-json-null".to_string(),
events: vec![(0, Event::new(kinds::USAGE, usage.encode_to_vec()))],
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute(
"SELECT payload_json ->> 'anything' AS extracted, \
json_get_str(payload_json, 'anything') AS typed_get \
FROM events_raw",
false,
)
.await
.expect("json functions over a NULL payload_json must not error");
assert_eq!(output.batches[0].num_rows(), 1);
assert!(
output.batches[0].column(0).is_null(0),
"->> over NULL payload_json is NULL, not an error"
);
assert!(
output.batches[0].column(1).is_null(0),
"json_get_str over NULL payload_json is NULL, not an error"
);
}
#[tokio::test]
#[allow(clippy::similar_names)] async fn multi_partition_committed_view_includes_each_partitions_committed_rows_only() {
let committed_x = Uuid::from_u128(101);
let orphaned_x = Uuid::from_u128(102);
let committed_y = Uuid::from_u128(103);
let orphaned_y = Uuid::from_u128(104);
let partitions = vec![
PartitionEvents {
partition: "conv-mp-x".to_string(),
events: fixture_events(committed_x, orphaned_x),
},
PartitionEvents {
partition: "conv-mp-y".to_string(),
events: fixture_events(committed_y, orphaned_y),
},
];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let output = engine
.execute(
"SELECT partition, turn_id FROM events ORDER BY partition, position",
false,
)
.await
.expect("committed view query");
let batch = &output.batches[0];
assert_eq!(
batch.num_rows(),
8,
"each partition's four committed rows (turn_start, user_msg, usage, \
turn_complete), no orphaned rows from either partition"
);
let partition_col = batch
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("partition is Utf8");
let turn_id_col = batch
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("turn_id is Utf8");
let orphaned_x_str = orphaned_x.to_string();
let orphaned_y_str = orphaned_y.to_string();
let mut x_rows = 0;
let mut y_rows = 0;
for i in 0..batch.num_rows() {
let turn_id = turn_id_col.value(i);
assert_ne!(
turn_id, orphaned_x_str,
"conv-mp-x's orphaned turn must not appear"
);
assert_ne!(
turn_id, orphaned_y_str,
"conv-mp-y's orphaned turn must not appear"
);
match partition_col.value(i) {
"conv-mp-x" => x_rows += 1,
"conv-mp-y" => y_rows += 1,
other => panic!("unexpected partition {other}"),
}
}
assert_eq!(x_rows, 4, "conv-mp-x's committed turn's four rows");
assert_eq!(y_rows, 4, "conv-mp-y's committed turn's four rows");
}
#[allow(clippy::too_many_lines)]
#[tokio::test]
async fn qry5_shared_turn_id_across_partitions_does_not_cross_contaminate_commit_status() {
let shared_turn = Uuid::from_u128(0x0195_abcd_ef01_2345_6789_abcd_ef01_5005);
let committed_events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &shared_turn), Vec::new()),
),
(
1,
Event::new(
kinds::tagged(kinds::USAGE, &shared_turn),
UsageEvent {
input_tokens: 7,
output_tokens: 8,
..Default::default()
}
.encode_to_vec(),
),
),
(
2,
Event::new(
kinds::tagged(kinds::TURN_COMPLETE, &shared_turn),
Vec::new(),
),
),
];
let orphaned_events = vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &shared_turn), Vec::new()),
),
(
1,
Event::new(
kinds::tagged(kinds::USAGE, &shared_turn),
UsageEvent {
input_tokens: 999,
output_tokens: 999,
..Default::default()
}
.encode_to_vec(),
),
),
];
let partitions = vec![
PartitionEvents {
partition: "conv-qry5-committed".to_string(),
events: committed_events,
},
PartitionEvents {
partition: "conv-qry5-orphaned".to_string(),
events: orphaned_events,
},
];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let committed_count = engine
.execute(
&format!(
"SELECT COUNT(*) AS c FROM events \
WHERE partition = 'conv-qry5-committed' AND turn_id = '{shared_turn}'"
),
false,
)
.await
.expect("committed partition events query");
assert_eq!(
single_count(&committed_count),
3,
"the committed partition's own three rows under the shared turn_id are visible"
);
let orphaned_count = engine
.execute(
&format!(
"SELECT COUNT(*) AS c FROM events \
WHERE partition = 'conv-qry5-orphaned' AND turn_id = '{shared_turn}'"
),
false,
)
.await
.expect("orphaned partition events query");
assert_eq!(
orphaned_count.batches[0].num_rows(),
1,
"COUNT(*) still returns exactly one row"
);
assert_eq!(
single_count(&orphaned_count),
0,
"the orphaned partition's rows under the SAME turn_id must stay hidden — the \
committed partition's turn_complete marker must never vouch for a different \
conversation's uncommitted turn under the identical turn_id"
);
let usage_committed = engine
.execute(
&format!(
"SELECT COUNT(*) AS c FROM usage \
WHERE partition = 'conv-qry5-committed' AND turn_id = '{shared_turn}'"
),
false,
)
.await
.expect("committed partition usage query");
assert_eq!(single_count(&usage_committed), 1);
let usage_orphaned = engine
.execute(
&format!(
"SELECT COUNT(*) AS c FROM usage \
WHERE partition = 'conv-qry5-orphaned' AND turn_id = '{shared_turn}'"
),
false,
)
.await
.expect("orphaned partition usage query");
assert_eq!(
single_count(&usage_orphaned),
0,
"the orphaned partition's usage row must stay hidden even though another \
conversation committed the identical turn_id"
);
}
#[tokio::test]
async fn null_turn_id_typed_row_is_dropped_by_the_committed_turn_filter() {
let bare_usage = UsageEvent {
input_tokens: 3,
output_tokens: 4,
..Default::default()
};
let partitions = vec![PartitionEvents {
partition: "conv-null-turn-id".to_string(),
events: vec![(
0,
Event::new(kinds::USAGE.to_owned(), bare_usage.encode_to_vec()),
)],
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let raw = engine
.execute("SELECT COUNT(*) AS c FROM usage_raw", false)
.await
.expect("usage_raw query");
assert_eq!(
single_count(&raw),
1,
"the bare-kind row is visible through the Fleet-only raw table, turn_id NULL and all"
);
let raw_turn_id = engine
.execute("SELECT turn_id FROM usage_raw", false)
.await
.expect("usage_raw turn_id query");
assert!(
raw_turn_id.batches[0].column(0).is_null(0),
"the row's own turn_id is genuinely NULL, not an empty string"
);
let filtered = engine
.execute("SELECT COUNT(*) AS c FROM usage", false)
.await
.expect("usage query");
assert_eq!(
single_count(&filtered),
0,
"a NULL turn_id row must be dropped by the committed-turn filter, exactly like an \
orphaned/uncommitted turn's row"
);
}
#[tokio::test]
async fn ddl_is_rejected_through_the_engine() {
let partitions = one_partition("conv-ddl", Uuid::from_u128(17), Uuid::from_u128(18));
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let result = engine.execute("CREATE TABLE t (a INT)", false).await;
assert!(
matches!(result, Err(QueryEngineError::Rejected(_))),
"DDL must be rejected before it ever reaches planning: {result:?}"
);
}
#[tokio::test]
async fn sql_options_reject_ddl_independent_of_the_statement_gate() {
let ctx = SessionContext::new();
let sql_options = SQLOptions::new()
.with_allow_ddl(false)
.with_allow_dml(false)
.with_allow_statements(false);
let result = ctx
.sql_with_options("CREATE TABLE t (a INT)", sql_options)
.await;
assert!(
result.is_err(),
"DataFusion's own SQLOptions guard must reject DDL even without the statement \
gate's involvement: {result:?}"
);
}
#[tokio::test]
async fn fires_and_routines_rows_are_queryable_and_join_by_routine_name() {
let fired = RoutineFiredEvent {
routine: "daily-standup".to_string(),
occurrence: "daily-standup-28461600".to_string(),
scheduled_at_ms: 1_753_300_800_000,
fired_at_ms: 1_753_300_805_000,
..Default::default()
};
let partitions = vec![PartitionEvents {
partition: "routine-scheduler".to_string(),
events: vec![(
0,
Event::trusted(kinds::ROUTINE_FIRED, fired.encode_to_vec()),
)],
}];
let reference = ReferenceData {
routines: vec![crate::routine_catalog::RoutineStatusRecord {
name: "daily-standup".to_string(),
uid: "uid-daily-standup".to_string(),
fire_conversation_id: "fire-conv-fixture".to_string(),
ready: true,
phase: Some("Ready".to_string()),
message: None,
last_fire_time_ms: Some(1_784_797_200_000),
next_fire_time_ms: Some(1_784_883_600_000),
conditions_json: "[]".to_string(),
creator_persona: "persona-1".to_string(),
provenance_conversation_id: "conv-1".to_string(),
schedule_json: r#"{"kind":"cron","expression":"0 9 * * *","timezone":null}"#
.to_string(),
next_fires_json: r"[1784883600000]".to_string(),
suspended: false,
paused_by: None,
paused_at_ms: None,
pause_reason: None,
prompt: "post the morning standup".to_string(),
scope: "private".to_string(),
orphaned: false,
display_name: String::new(),
description: String::new(),
schedule_timezone: "UTC".to_owned(),
}],
..ReferenceData::empty()
};
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
reference,
&[],
QueryLimits::default(),
)
.await
.expect("build");
let fires_count = engine
.execute("SELECT COUNT(*) AS c FROM fires", false)
.await
.expect("fires count query");
assert_eq!(single_count(&fires_count), 1);
let joined = engine
.execute(
"SELECT f.occurrence, f.scheduled_at_ms, f.fired_at_ms, r.ready, r.next_fire_time_ms \
FROM fires f JOIN routines r ON r.name = f.routine \
WHERE f.routine = 'daily-standup'",
false,
)
.await
.expect("join query");
assert_eq!(joined.batches[0].num_rows(), 1);
let occurrence_col = joined.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("occurrence is Utf8");
assert_eq!(occurrence_col.value(0), "daily-standup-28461600");
let scheduled_col = joined.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::UInt64Array>()
.expect("scheduled_at_ms is UInt64");
assert_eq!(scheduled_col.value(0), 1_753_300_800_000);
let ready_col = joined.batches[0]
.column(3)
.as_any()
.downcast_ref::<arrow::array::BooleanArray>()
.expect("ready is Boolean");
assert!(ready_col.value(0));
let next_fire_col = joined.batches[0]
.column(4)
.as_any()
.downcast_ref::<arrow::array::Int64Array>()
.expect("next_fire_time_ms is Int64");
assert_eq!(next_fire_col.value(0), 1_784_883_600_000);
}
#[tokio::test]
async fn fires_outcome_column_is_selectable_through_the_view() {
let fired = RoutineFiredEvent {
routine: "daily-standup".to_string(),
occurrence: "daily-standup-28461600".to_string(),
scheduled_at_ms: 1_753_300_800_000,
fired_at_ms: 1_753_300_805_000,
..Default::default()
};
let outcome = RoutineFireOutcomeEvent {
routine: "daily-standup".to_string(),
occurrence: "daily-standup-28461600".to_string(),
outcome: RoutineFireOutcome::Ok.into(),
fired_at_ms: 1_753_300_805_000,
..Default::default()
};
let partitions = vec![PartitionEvents {
partition: "routine-scheduler".to_string(),
events: vec![
(
0,
Event::trusted(kinds::ROUTINE_FIRED, fired.encode_to_vec()),
),
(
1,
Event::trusted(kinds::ROUTINE_FIRE_OUTCOME, outcome.encode_to_vec()),
),
],
}];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let result = engine
.execute(
"SELECT outcome FROM fires WHERE routine = 'daily-standup'",
false,
)
.await
.expect("outcome query");
assert_eq!(result.batches[0].num_rows(), 1);
let outcome_col = result.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("outcome is Utf8");
assert_eq!(outcome_col.value(0), "ok");
}
#[tokio::test]
async fn fires_joins_turn_dispatch_by_occurrence_where_a_fired_turn_exists() {
let turn = Uuid::from_u128(0x0195_abcd_ef01_2345_6789_abcd_ef01_7777);
let fired = RoutineFiredEvent {
routine: "daily-standup".to_string(),
occurrence: "daily-standup-28461600".to_string(),
scheduled_at_ms: 1_753_300_800_000,
fired_at_ms: 1_753_300_805_000,
..Default::default()
};
let dispatched = TurnDispatchedEvent {
occurrence: "daily-standup-28461600".to_string(),
..Default::default()
};
let orphan_fired = RoutineFiredEvent {
routine: "daily-standup".to_string(),
occurrence: "daily-standup-28465200".to_string(),
scheduled_at_ms: 1_753_300_800_000 + 3_600_000,
fired_at_ms: 1_753_300_800_000 + 3_600_005,
..Default::default()
};
let partitions = vec![
PartitionEvents {
partition: "routine-scheduler".to_string(),
events: vec![
(
0,
Event::trusted(kinds::ROUTINE_FIRED, fired.encode_to_vec()),
),
(
1,
Event::trusted(kinds::ROUTINE_FIRED, orphan_fired.encode_to_vec()),
),
],
},
PartitionEvents {
partition: "conv-standup-1".to_string(),
events: vec![
(
0,
Event::new(kinds::tagged(kinds::TURN_START, &turn), Vec::new()),
),
(
1,
Event::new(
kinds::tagged(kinds::TURN_DISPATCHED, &turn),
dispatched.encode_to_vec(),
),
),
(
2,
Event::new(kinds::tagged(kinds::TURN_COMPLETE, &turn), Vec::new()),
),
],
},
];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let joined = engine
.execute(
"SELECT f.occurrence, t.partition, t.turn_id \
FROM fires f JOIN turn_dispatch t ON t.occurrence = f.occurrence",
false,
)
.await
.expect("fires-joins-turn_dispatch query");
assert_eq!(
joined.batches[0].num_rows(),
1,
"only the fire with a matching dispatched turn must survive the join"
);
let occurrence_col = joined.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("occurrence is Utf8");
assert_eq!(occurrence_col.value(0), "daily-standup-28461600");
let partition_col = joined.batches[0]
.column(1)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("partition is Utf8");
assert_eq!(partition_col.value(0), "conv-standup-1");
}
#[tokio::test]
#[allow(clippy::too_many_lines)] async fn routine_lifecycle_decodes_the_five_signed_kinds_end_to_end() {
let signer = ApprovalSigner::from_seed(42);
let (created, _, _) = routine_created_payload(
"daily-standup",
"persona-1",
"conv-1",
"call-1",
"hash-1",
1_000,
&signer,
);
let (paused, _, _) = routine_paused_payload(
"daily-standup",
"persona-1",
"conv-2",
"call-2",
"hash-2",
"chat",
2_000,
Some("rotating out old announcements"),
&signer,
);
let (resumed, _, _) = routine_resumed_payload(
"daily-standup",
"persona-1",
"conv-3",
"call-3",
"hash-3",
"chat",
3_000,
&signer,
);
let (deleted, _, _) = routine_deleted_payload(
"daily-standup",
"persona-1",
"conv-4",
"call-4",
"hash-4",
"chat",
4_000,
&signer,
);
let (scope_changed, _, _) = routine_scope_changed_payload(
"daily-standup",
"persona-1",
"conv-5",
"call-5",
"hash-5",
"public",
5_000,
&signer,
);
let partitions = vec![PartitionEvents {
partition: "routine-scheduler".to_string(),
events: vec![
(0, Event::new(kinds::ROUTINE_CREATED, created)),
(1, Event::new(kinds::ROUTINE_PAUSED, paused)),
(2, Event::new(kinds::ROUTINE_RESUMED, resumed)),
(3, Event::new(kinds::ROUTINE_DELETED, deleted)),
(4, Event::new(kinds::ROUTINE_SCOPE_CHANGED, scope_changed)),
],
}];
let trusted_signers = vec![signer.public_key_bytes()];
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&trusted_signers,
QueryLimits::default(),
)
.await
.expect("build");
let result = engine
.execute(
"SELECT phase, routine, actor_persona, at_ms, reason, scope FROM routine_lifecycle \
ORDER BY at_ms",
false,
)
.await
.expect("routine_lifecycle query");
assert_eq!(result.batches[0].num_rows(), 5);
let phase = result.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("phase is Utf8");
assert_eq!(phase.value(0), "routine_created");
assert_eq!(phase.value(1), "routine_paused");
assert_eq!(phase.value(2), "routine_resumed");
assert_eq!(phase.value(3), "routine_deleted");
assert_eq!(phase.value(4), "routine_scope_changed");
let reason = result.batches[0]
.column(4)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("reason is Utf8");
assert!(reason.is_null(0), "routine_created carries no reason");
assert_eq!(reason.value(1), "rotating out old announcements");
assert!(reason.is_null(2), "routine_resumed carries no reason");
assert!(reason.is_null(3), "routine_deleted carries no reason");
assert!(reason.is_null(4), "routine_scope_changed carries no reason");
let scope = result.batches[0]
.column(5)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("scope is Utf8");
assert!(scope.is_null(0), "routine_created carries no scope");
assert!(scope.is_null(1), "routine_paused carries no scope");
assert!(scope.is_null(2), "routine_resumed carries no scope");
assert!(scope.is_null(3), "routine_deleted carries no scope");
assert_eq!(scope.value(4), "public");
}
#[allow(clippy::too_many_lines)]
#[tokio::test]
async fn conversations_scope_hides_fires_and_routines_tables() {
let fired = RoutineFiredEvent {
routine: "daily-standup".to_string(),
occurrence: "daily-standup-1".to_string(),
scheduled_at_ms: 1,
fired_at_ms: 2,
..Default::default()
};
let conversations_engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Conversations(vec!["fires-hidden".to_string()]),
vec![PartitionEvents {
partition: "conv-fires-hidden".to_string(),
events: vec![(
0,
Event::trusted(kinds::ROUTINE_FIRED, fired.encode_to_vec()),
)],
}],
ReferenceData {
routines: vec![crate::routine_catalog::RoutineStatusRecord {
name: "daily-standup".to_string(),
uid: "uid-daily-standup".to_string(),
fire_conversation_id: "fire-conv-fixture".to_string(),
ready: true,
phase: None,
message: None,
last_fire_time_ms: None,
next_fire_time_ms: None,
conditions_json: "[]".to_string(),
creator_persona: "persona-1".to_string(),
provenance_conversation_id: "conv-1".to_string(),
schedule_json: r#"{"kind":"cron","expression":"0 9 * * *","timezone":null}"#
.to_string(),
next_fires_json: "[]".to_string(),
suspended: false,
paused_by: None,
paused_at_ms: None,
pause_reason: None,
prompt: "post the morning standup".to_string(),
scope: "private".to_string(),
orphaned: false,
display_name: String::new(),
description: String::new(),
schedule_timezone: "UTC".to_owned(),
}],
..ReferenceData::empty()
},
&[],
QueryLimits::default(),
)
.await
.expect("build conversations engine");
let fires_unqualified = conversations_engine
.execute("SELECT * FROM fires", false)
.await;
assert!(
matches!(fires_unqualified, Err(QueryEngineError::DataFusion(_))),
"fires must not resolve for a non-Fleet scope: {fires_unqualified:?}"
);
let fires_raw_unqualified = conversations_engine
.execute("SELECT * FROM fires_raw", false)
.await;
assert!(
matches!(fires_raw_unqualified, Err(QueryEngineError::DataFusion(_))),
"fires_raw must not resolve for a non-Fleet scope either: {fires_raw_unqualified:?}"
);
let routines_unqualified = conversations_engine
.execute("SELECT * FROM routines", false)
.await;
assert!(
matches!(routines_unqualified, Err(QueryEngineError::DataFusion(_))),
"routines must not resolve for a non-Fleet scope: {routines_unqualified:?} — note \
ReferenceData.routines is never even registered for this scope regardless of the \
(non-empty) value this fixture supplied"
);
let qualified = conversations_engine
.execute("SELECT * FROM datafusion.public.fires", false)
.await;
assert!(
matches!(qualified, Err(QueryEngineError::DataFusion(_))),
"a fully-qualified name must not reach fires either: {qualified:?}"
);
let routines_prompt_unqualified = conversations_engine
.execute("SELECT prompt FROM routines", false)
.await;
assert!(
matches!(
routines_prompt_unqualified,
Err(QueryEngineError::DataFusion(_))
),
"routines.prompt must not resolve for a non-Fleet scope: {routines_prompt_unqualified:?}"
);
let routine_lifecycle_unqualified = conversations_engine
.execute("SELECT * FROM routine_lifecycle", false)
.await;
assert!(
matches!(
routine_lifecycle_unqualified,
Err(QueryEngineError::DataFusion(_))
),
"routine_lifecycle must not resolve for a non-Fleet scope: {routine_lifecycle_unqualified:?}"
);
let routine_lifecycle_raw_unqualified = conversations_engine
.execute("SELECT * FROM routine_lifecycle_raw", false)
.await;
assert!(
matches!(
routine_lifecycle_raw_unqualified,
Err(QueryEngineError::DataFusion(_))
),
"routine_lifecycle_raw must not resolve for a non-Fleet scope either: \
{routine_lifecycle_raw_unqualified:?}"
);
}
#[tokio::test]
async fn owner_scoped_session_cannot_resolve_fires_raw_directly() {
let fired = RoutineFiredEvent {
routine: "daily-standup".to_string(),
occurrence: "daily-standup-1".to_string(),
scheduled_at_ms: 1,
fired_at_ms: 2,
routine_uid: "uid-daily-standup".to_string(),
..Default::default()
};
let reference = ReferenceData {
routines: vec![crate::routine_catalog::RoutineStatusRecord {
name: "daily-standup".to_string(),
uid: "uid-daily-standup".to_string(),
fire_conversation_id: "fire-conv-fixture".to_string(),
ready: true,
phase: None,
message: None,
last_fire_time_ms: None,
next_fire_time_ms: None,
conditions_json: "[]".to_string(),
creator_persona: "persona-owner".to_string(),
provenance_conversation_id: "conv-1".to_string(),
schedule_json: r#"{"kind":"cron","expression":"0 9 * * *","timezone":null}"#
.to_string(),
next_fires_json: "[]".to_string(),
suspended: false,
paused_by: None,
paused_at_ms: None,
pause_reason: None,
prompt: "post the morning standup".to_string(),
scope: "private".to_string(),
orphaned: false,
display_name: String::new(),
description: String::new(),
schedule_timezone: "UTC".to_owned(),
}],
..ReferenceData::empty_except_routines(vec![])
};
let owner_engine = QueryEngine::build_scoped(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-1".to_string()]),
vec![PartitionEvents {
partition: "routine-scheduler".to_string(),
events: vec![(
0,
Event::trusted(kinds::ROUTINE_FIRED, fired.encode_to_vec()),
)],
}],
reference,
&[],
QueryLimits::default(),
true,
)
.await
.expect("build owner-scoped engine");
let fires_owned = owner_engine
.execute("SELECT * FROM fires", false)
.await
.expect("owned fires view resolves");
assert_eq!(fires_owned.batches[0].num_rows(), 1);
let fires_raw_unqualified = owner_engine.execute("SELECT * FROM fires_raw", false).await;
assert!(
matches!(fires_raw_unqualified, Err(QueryEngineError::DataFusion(_))),
"fires_raw must not resolve for an owner-scoped session either: \
{fires_raw_unqualified:?}"
);
let qualified = owner_engine
.execute("SELECT * FROM datafusion.public.fires_raw", false)
.await;
assert!(
matches!(qualified, Err(QueryEngineError::DataFusion(_))),
"a fully-qualified name must not reach fires_raw either: {qualified:?}"
);
}
#[tokio::test]
async fn recreated_routine_name_does_not_leak_the_previous_owners_fires() {
let fire_from_a = RoutineFiredEvent {
routine: "daily-standup".to_string(),
occurrence: "daily-standup-1".to_string(),
scheduled_at_ms: 1,
fired_at_ms: 2,
routine_uid: "uid-A".to_string(),
..Default::default()
};
let fire_from_b = RoutineFiredEvent {
routine: "daily-standup".to_string(),
occurrence: "daily-standup-2".to_string(),
scheduled_at_ms: 3,
fired_at_ms: 4,
routine_uid: "uid-B".to_string(),
..Default::default()
};
let reference = ReferenceData {
routines: vec![crate::routine_catalog::RoutineStatusRecord {
name: "daily-standup".to_string(),
uid: "uid-B".to_string(),
fire_conversation_id: "fire-conv-fixture".to_string(),
ready: true,
phase: None,
message: None,
last_fire_time_ms: None,
next_fire_time_ms: None,
conditions_json: "[]".to_string(),
creator_persona: "persona-b".to_string(),
provenance_conversation_id: "conv-b".to_string(),
schedule_json: r#"{"kind":"cron","expression":"0 9 * * *","timezone":null}"#
.to_string(),
next_fires_json: "[]".to_string(),
suspended: false,
paused_by: None,
paused_at_ms: None,
pause_reason: None,
prompt: "post the morning standup".to_string(),
scope: "private".to_string(),
orphaned: false,
display_name: String::new(),
description: String::new(),
schedule_timezone: "UTC".to_owned(),
}],
..ReferenceData::empty_except_routines(vec![])
};
let owner_engine = QueryEngine::build_scoped(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-b".to_string()]),
vec![PartitionEvents {
partition: "routine-scheduler".to_string(),
events: vec![
(
0,
Event::trusted(kinds::ROUTINE_FIRED, fire_from_a.encode_to_vec()),
),
(
1,
Event::trusted(kinds::ROUTINE_FIRED, fire_from_b.encode_to_vec()),
),
],
}],
reference,
&[],
QueryLimits::default(),
true,
)
.await
.expect("build owner-scoped engine");
let fires_owned = owner_engine
.execute("SELECT occurrence FROM fires", false)
.await
.expect("owned fires view resolves");
assert_eq!(
fires_owned.batches[0].num_rows(),
1,
"only B's own fire should be visible, never A's — same name, different uid"
);
let occurrence = fires_owned.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.unwrap();
assert_eq!(
occurrence.value(0),
"daily-standup-2",
"the visible fire must be B's, never A's stale-uid row"
);
}
#[tokio::test]
async fn empty_uid_fires_and_empty_uid_routines_never_match_each_other() {
let pre_uid_fire = RoutineFiredEvent {
routine: "daily-standup".to_string(),
occurrence: "daily-standup-1".to_string(),
scheduled_at_ms: 1,
fired_at_ms: 2,
routine_uid: String::new(),
..Default::default()
};
let reference = ReferenceData {
routines: vec![crate::routine_catalog::RoutineStatusRecord {
name: "daily-standup".to_string(),
uid: String::new(),
fire_conversation_id: String::new(),
ready: true,
phase: None,
message: None,
last_fire_time_ms: None,
next_fire_time_ms: None,
conditions_json: "[]".to_string(),
creator_persona: "persona-owner".to_string(),
provenance_conversation_id: "conv-1".to_string(),
schedule_json: r#"{"kind":"cron","expression":"0 9 * * *","timezone":null}"#
.to_string(),
next_fires_json: "[]".to_string(),
suspended: false,
paused_by: None,
paused_at_ms: None,
pause_reason: None,
prompt: "post the morning standup".to_string(),
scope: "private".to_string(),
orphaned: false,
display_name: String::new(),
description: String::new(),
schedule_timezone: "UTC".to_owned(),
}],
..ReferenceData::empty_except_routines(vec![])
};
let owner_engine = QueryEngine::build_scoped(
&test_base_state(),
&QueryScope::Conversations(vec!["conv-1".to_string()]),
vec![PartitionEvents {
partition: "routine-scheduler".to_string(),
events: vec![(
0,
Event::trusted(kinds::ROUTINE_FIRED, pre_uid_fire.encode_to_vec()),
)],
}],
reference,
&[],
QueryLimits::default(),
true,
)
.await
.expect("build owner-scoped engine");
let fires_owned = owner_engine
.execute("SELECT * FROM fires", false)
.await
.expect("owned fires view resolves");
let total_rows: usize = fires_owned.batches.iter().map(RecordBatch::num_rows).sum();
assert_eq!(
total_rows, 0,
"an empty-uid fire must stay invisible even against an empty-uid routines row"
);
}
#[tokio::test]
async fn fleet_scope_sees_provenance_pause_and_schedule_columns_on_routines() {
let reference = ReferenceData {
routines: vec![crate::routine_catalog::RoutineStatusRecord {
name: "daily-standup".to_string(),
uid: "uid-daily-standup".to_string(),
fire_conversation_id: "fire-conv-fixture".to_string(),
ready: true,
phase: Some("Ready".to_string()),
message: None,
last_fire_time_ms: None,
next_fire_time_ms: Some(1_784_883_600_000),
conditions_json: "[]".to_string(),
creator_persona: "persona-1".to_string(),
provenance_conversation_id: "conv-1".to_string(),
schedule_json:
r#"{"kind":"cron","expression":"0 9 * * *","timezone":"America/New_York"}"#
.to_string(),
next_fires_json: r"[1784883600000,1784970000000]".to_string(),
suspended: true,
paused_by: Some("persona-1".to_string()),
paused_at_ms: Some(1_784_764_800_000),
pause_reason: Some("rotating out old announcements".to_string()),
prompt: "post the morning standup".to_string(),
scope: "private".to_string(),
orphaned: false,
display_name: String::new(),
description: String::new(),
schedule_timezone: "UTC".to_owned(),
}],
..ReferenceData::empty()
};
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
vec![],
reference,
&[],
QueryLimits::default(),
)
.await
.expect("build");
let result = engine
.execute(
"SELECT creator_persona, provenance_conversation_id, schedule_json, \
next_fires_json, suspended, paused_by, paused_at_ms, pause_reason, prompt \
FROM routines WHERE name = 'daily-standup'",
false,
)
.await
.expect("routines query with spec columns");
assert_eq!(result.batches[0].num_rows(), 1);
let creator_persona = result.batches[0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("creator_persona is Utf8");
assert_eq!(creator_persona.value(0), "persona-1");
let suspended = result.batches[0]
.column(4)
.as_any()
.downcast_ref::<arrow::array::BooleanArray>()
.expect("suspended is Boolean");
assert!(suspended.value(0));
let paused_by = result.batches[0]
.column(5)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("paused_by is Utf8");
assert_eq!(paused_by.value(0), "persona-1");
let prompt = result.batches[0]
.column(8)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.expect("prompt is Utf8");
assert_eq!(prompt.value(0), "post the morning standup");
}
#[tokio::test]
async fn a_bound_parameter_is_a_value_and_never_sql() {
const SQL: &str = "SELECT COUNT(*) AS c FROM events WHERE turn_id = CAST($1 AS VARCHAR)";
let committed_turn = Uuid::from_u128(4242);
let orphaned_turn = Uuid::from_u128(4243);
let partitions = one_partition("conv-param", committed_turn, orphaned_turn);
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let matched = engine
.execute_with_params(SQL, &[&committed_turn.to_string()], false)
.await
.expect("bound query");
assert_eq!(single_count(&matched), 4);
for payload in [
"' OR '1'='1",
"x'; DROP TABLE events; --",
"' UNION SELECT payload FROM events_raw --",
"'; SELECT 1; --",
] {
let result = engine
.execute_with_params(SQL, &[payload], false)
.await
.unwrap_or_else(|e| panic!("`{payload}` must run as a plain value comparison: {e}"));
assert_eq!(
single_count(&result),
0,
"`{payload}` must be compared as literal text"
);
}
let after = engine
.execute_with_params(SQL, &[&committed_turn.to_string()], false)
.await
.expect("bound query");
assert_eq!(single_count(&after), 4);
}
#[tokio::test]
async fn a_surplus_parameter_is_ignored_while_a_missing_one_is_refused() {
const UNFILTERED: &str = "SELECT COUNT(*) AS c FROM events";
const FILTERED: &str = "SELECT COUNT(*) AS c FROM events WHERE turn_id = CAST($1 AS VARCHAR)";
let committed_turn = Uuid::from_u128(4242);
let orphaned_turn = Uuid::from_u128(4243);
let partitions = one_partition("conv-surplus", committed_turn, orphaned_turn);
let engine = QueryEngine::build(
&test_base_state(),
&QueryScope::Fleet,
partitions,
ReferenceData::empty(),
&[],
QueryLimits::default(),
)
.await
.expect("build");
let baseline = engine
.execute_with_params(UNFILTERED, &[], false)
.await
.expect("the unfiltered statement runs");
let all_rows = single_count(&baseline);
assert!(all_rows > 0, "the fixture must have rows to be filtered");
let surplus = engine
.execute_with_params(UNFILTERED, &["never-bound"], false)
.await
.expect("a surplus parameter is silently ignored, not refused");
assert_eq!(
single_count(&surplus),
all_rows,
"a surplus parameter changes nothing: the statement ran unfiltered"
);
let filtered = engine
.execute_with_params(FILTERED, &["never-bound"], false)
.await
.expect("the filtered statement runs");
assert_eq!(
single_count(&filtered),
0,
"with a placeholder to bind, the same value really does filter"
);
let missing = engine.execute_with_params(FILTERED, &[], false).await;
assert!(
matches!(missing, Err(QueryEngineError::DataFusion(_))),
"a placeholder with no value must be refused, not defaulted: {missing:?}"
);
}
fn field_not_found() -> DataFusionError {
DataFusionError::SchemaError(
Box::new(SchemaError::FieldNotFound {
field: Box::new(Column::new(Some("events"), "message")),
valid_fields: vec![
Column::new(Some("events"), "partition"),
Column::new(Some("events"), "turn_id"),
],
}),
Box::new(None),
)
}
#[test]
fn unresolved_column_walks_every_wrapper_the_planner_adds() {
let wrappings: Vec<(&str, DataFusionError)> = vec![
("bare", field_not_found()),
(
"context",
DataFusionError::Context("while planning".to_owned(), Box::new(field_not_found())),
),
(
"diagnostic",
DataFusionError::Diagnostic(
Box::new(datafusion::common::Diagnostic::new_error(
"column not found",
None,
)),
Box::new(field_not_found()),
),
),
(
"context around diagnostic",
DataFusionError::Context(
"while planning".to_owned(),
Box::new(DataFusionError::Diagnostic(
Box::new(datafusion::common::Diagnostic::new_error(
"column not found",
None,
)),
Box::new(field_not_found()),
)),
),
),
(
"shared",
DataFusionError::Shared(std::sync::Arc::new(field_not_found())),
),
(
"collection",
DataFusionError::Collection(vec![field_not_found()]),
),
];
for (label, error) in wrappings {
let unresolved = unresolved_column(&QueryEngineError::DataFusion(error))
.unwrap_or_else(|| panic!("`{label}` must still be recognized as a column failure"));
assert_eq!(unresolved.name, "message", "{label}");
assert_eq!(unresolved.qualifier.as_deref(), Some("events"), "{label}");
assert_eq!(
unresolved.valid_fields,
vec![
(Some("events".to_owned()), "partition".to_owned()),
(Some("events".to_owned()), "turn_id".to_owned()),
],
"{label}"
);
}
}
#[test]
fn unresolved_column_offers_each_column_once() {
let duplicated = DataFusionError::SchemaError(
Box::new(SchemaError::FieldNotFound {
field: Box::new(Column::new(None::<&str>, "nope")),
valid_fields: vec![
Column::new(Some("events"), "turn_id"),
Column::new(Some("events"), "partition"),
Column::new(Some("events"), "position"),
Column::new(Some("events"), "kind"),
Column::new(Some("events"), "turn_id"),
],
}),
Box::new(None),
);
let unresolved = unresolved_column(&QueryEngineError::DataFusion(duplicated))
.expect("a duplicated valid_fields is still a column failure");
assert_eq!(
unresolved.valid_fields,
vec![
(Some("events".to_owned()), "turn_id".to_owned()),
(Some("events".to_owned()), "partition".to_owned()),
(Some("events".to_owned()), "position".to_owned()),
(Some("events".to_owned()), "kind".to_owned()),
],
"first-seen order, each column once"
);
}
#[test]
fn unresolved_column_ignores_every_other_failure() {
let others = vec![
DataFusionError::Execution("the plan failed at runtime".to_owned()),
DataFusionError::Plan("Invalid function 'nope'".to_owned()),
DataFusionError::SchemaError(
Box::new(SchemaError::DuplicateUnqualifiedField {
name: "turn_id".to_owned(),
}),
Box::new(None),
),
DataFusionError::Context(
"while planning".to_owned(),
Box::new(DataFusionError::Execution("still not a column".to_owned())),
),
];
for error in others {
let rendered = error.to_string();
let engine_error = QueryEngineError::DataFusion(error);
assert!(
unresolved_column(&engine_error).is_none(),
"`{rendered}` is not an unresolvable column"
);
assert!(
!is_unresolved_table_error(&engine_error),
"`{rendered}` is not an unresolvable table either"
);
}
let decode = QueryEngineError::Decode {
table: "events_raw",
source: arrow::error::ArrowError::SchemaError("not a DataFusion error".to_owned()),
};
assert!(unresolved_column(&decode).is_none());
assert!(!is_unresolved_table_error(&decode));
}