#![cfg(test)]
use std::collections::BTreeMap;
use std::sync::Arc;
use std::time::Duration;
use arrow::array::{Array, BooleanArray, Int64Array, StringArray, StringViewArray, UInt64Array};
use datafusion::catalog::TableProvider;
use datafusion::datasource::empty::EmptyTable;
use datafusion::error::DataFusionError;
use datafusion::execution::context::SessionContext;
use futures::StreamExt;
use polyc_state::query_audit::{ErrorClass, QueryOutcome, Truncation};
use super::admission::CoreExecutionAdmissionInput;
use super::{
CoreExecutionAdmission, CoreExecutionError, manifest_visible_to_table, manifests_match_plan,
request_context,
};
use crate::core_resolution::{CoreParameter, CorePlanOutcome, CoreTable, arrow_schema};
use polyc_query_credential::session::{MemorySources, QueryScope};
mod support;
use support::{HANG_GUARD, Harness};
const DEFAULT_ARTIFACT_RANGE_BYTES: u64 = 4 * 1024 * 1024;
#[test]
fn an_empty_authorized_scope_needs_no_manifest_for_each_declared_family() {
assert!(manifests_match_plan(
&[],
&[],
&[CoreTable::Usage],
&QueryScope::Fleet,
));
assert!(manifests_match_plan(
&[],
&[],
&[CoreTable::Turns, CoreTable::Usage],
&QueryScope::Fleet,
));
}
#[test]
fn a_participant_plan_registers_no_owner_table_file_from_a_foreign_partition() {
use polyc_state::feed::ProjectionSource;
use polyc_state::persona_memory::journal::{MemoryJournalPartition, PersonaMemorySource};
use polyc_state::revision::PartitionIncarnation;
let source = |partition: &str| {
ProjectionSource::PersonaMemory(PersonaMemorySource::new(
MemoryJournalPartition::parse(partition).expect("memory partition"),
PartitionIncarnation::from_bytes([9; PartitionIncarnation::LEN]),
))
};
let scope = QueryScope::Conversations {
conversations: vec!["conversation-a".to_owned()],
memory: MemorySources {
owner: Some("owner".to_owned()),
participants: vec!["peer".to_owned()],
},
};
let owner = source("persona-owner-mem");
let peer = source("persona-peer-mem");
assert!(manifest_visible_to_table(
CoreTable::MemoryFacts,
&owner,
&scope
));
assert!(
!manifest_visible_to_table(CoreTable::MemoryFacts, &peer, &scope),
"registering every dependency must not open a peer's owner-table file"
);
assert!(manifest_visible_to_table(
CoreTable::MemoryPortableFacts,
&peer,
&scope
));
}
#[test]
fn a_versioned_family_binds_every_namespace_its_authority_holds() {
let authority = polyc_state::versioned::authority::AuthorityFamily::Credentials;
let keyed = |namespace: &str| {
let scope = polyc_state::versioned::authority::scope(
&polyc_state::id::NamespaceId::new(namespace),
authority,
);
let source = polyc_state::feed::VersionedSource::new(
scope,
polyc_state::revision::PartitionIncarnation::from_bytes([3; 32]),
);
support::versioned_manifest(source.projection_partition().clone())
};
let manifests = vec![keyed("tenant-a"), keyed("tenant-b"), keyed("polychrome")];
assert!(
manifests_match_plan(
&manifests,
&[],
&[CoreTable::CredentialLifecycle],
&QueryScope::Fleet,
),
"every namespace the authority holds is one this plan reads"
);
for foreign in [
polyc_state::id::PartitionId::new("9:not-credentials/8:tenant-a"),
polyc_state::id::PartitionId::new("conv-a"),
] {
assert!(
!manifests_match_plan(
&[support::versioned_manifest(foreign.clone())],
&[],
&[CoreTable::CredentialLifecycle],
&QueryScope::Fleet,
),
"{foreign:?} is not a key this family's authority produces"
);
}
}
#[test]
fn a_fixed_source_family_with_no_generation_is_refused() {
assert!(
!manifests_match_plan(
&[],
&[],
&[CoreTable::AdminModelChanges],
&QueryScope::Fleet,
),
"an unpublished administrative trail must not bind as an empty one"
);
}
#[test]
fn aggregate_execution_admission_refuses_zero() {
let error = CoreExecutionAdmission::try_from(CoreExecutionAdmissionInput {
max_concurrent_executions: 0,
})
.unwrap_err();
assert!(matches!(error, CoreExecutionError::InvalidComposition(_)));
}
#[test]
fn aggregate_execution_admission_refuses_more_than_the_permit_ceiling() {
let error = CoreExecutionAdmission::try_from(CoreExecutionAdmissionInput {
max_concurrent_executions: usize::MAX,
})
.unwrap_err();
assert!(matches!(error, CoreExecutionError::InvalidComposition(_)));
}
#[tokio::test]
#[allow(
clippy::significant_drop_tightening,
reason = "the bound query must remain alive to prove bind-time I/O separately from execution I/O"
)]
async fn signed_exact_parquet_streams_only_visible_rows() {
let harness = Harness::new();
let prepared = harness
.prepare(
"visible-exact",
"SELECT text AS first, role, text AS duplicate FROM messages ORDER BY position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let authority = harness.visible_authority();
let bound = authority.bind(prepared).await.expect("exact files bind");
let reads_after_binding = harness.store.range_calls();
assert!(
reads_after_binding > 0,
"binding proves complete files and footers"
);
let mut stream = bound.execute();
let first = stream
.next()
.await
.expect("one batch")
.expect("verified batch");
assert_eq!(
harness.store.range_calls(),
reads_after_binding,
"DataFusion scans zero-copy slices of the retained verified image",
);
assert_eq!(first.schema().field(0).name(), "first");
assert_eq!(first.schema().field(1).name(), "role");
assert_eq!(first.schema().field(2).name(), "duplicate");
let left = first
.column(0)
.as_any()
.downcast_ref::<StringArray>()
.unwrap();
let duplicate = first
.column(2)
.as_any()
.downcast_ref::<StringArray>()
.unwrap();
let values = (0..left.len())
.map(|index| left.value(index).to_owned())
.collect::<Vec<_>>();
assert_eq!(values, ["visible-one", "visible-two", "visible-three"]);
assert_eq!(left, duplicate);
assert!(stream.next().await.is_none());
}
#[tokio::test]
async fn execution_history_keeps_the_legacy_public_shape_over_exact_files() {
let harness = Harness::new();
let prepared = harness
.prepare(
"execution-public-shape",
"SELECT * FROM tool_calls ORDER BY position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut stream = harness
.visible_authority()
.bind(prepared)
.await
.expect("execution files bind")
.execute();
let batch = stream.next().await.unwrap().unwrap();
assert_eq!(
batch
.schema()
.fields()
.iter()
.map(|field| field.name().as_str())
.collect::<Vec<_>>(),
[
"partition",
"position",
"turn_id",
"tool_call_id",
"block_type",
"name",
"arguments",
"result",
"first_party",
"internal_only",
"trust",
]
);
assert_eq!(
batch.num_rows(),
2,
"visible scope never opens the Fleet file"
);
let arguments = batch.column(6);
let result = batch.column(7);
let first_party = batch
.column(8)
.as_any()
.downcast_ref::<BooleanArray>()
.unwrap();
assert!(!arguments.is_null(0));
assert!(result.is_null(0));
assert!(first_party.is_null(0));
assert!(arguments.is_null(1));
assert!(!result.is_null(1));
assert!(first_party.value(1));
assert!(stream.next().await.is_none());
}
#[tokio::test]
#[allow(
clippy::too_many_lines,
reason = "one failure-sensitive fixture keeps the five public row sets and the signed excision source together"
)]
async fn prepared_journal_rows_match_every_projected_execution_view() {
use buffa::Message as _;
use polyc_crypto::approval::{ApprovalSigner, EXCISION_SCOPE_SOURCE_ONLY, excision_payload};
use polyc_eventlog_model::Event;
use polyc_projector::artifact::{EncodingBounds, RowSource, encode_execution_generation};
use polyc_proto::kinds;
use polyc_proto::proto::polychrome::{
agent::v1::{
Content, FunctionCallContent, FunctionResultContent, Message, ToolCallContent,
ToolResultContent, content, function_result_content, tool_call_content,
tool_result_content,
},
events::v1::{ModelCallEvent, SummaryEvent, TurnFailedEvent, UsageEvent},
harness::v1::TurnFailureKind,
};
fn tagged(base: &str, turn: uuid::Uuid) -> String {
kinds::tagged(base, &turn)
}
fn tool_call(id: &str) -> Vec<u8> {
let arguments = serde_json::from_value::<buffa_types::google::protobuf::Struct>(
serde_json::json!({"q": id}),
)
.map(buffa::MessageField::some)
.unwrap();
Message {
role: "assistant".to_owned(),
content: buffa::MessageField::some(Content {
r#type: Some(content::Type::ToolCall(Box::new(ToolCallContent {
id: id.to_owned(),
r#type: Some(tool_call_content::Type::FunctionCall(Box::new(
FunctionCallContent {
name: "search".to_owned(),
arguments,
..Default::default()
},
))),
..Default::default()
}))),
..Default::default()
}),
..Default::default()
}
.encode_to_vec()
}
fn tool_result(id: &str) -> Vec<u8> {
let response = Some(function_result_content::Result::Response(Box::new(
serde_json::from_value::<buffa_types::google::protobuf::Struct>(
serde_json::json!({"ok": true}),
)
.unwrap(),
)));
Message {
role: "tool".to_owned(),
content: buffa::MessageField::some(Content {
r#type: Some(content::Type::ToolResult(Box::new(ToolResultContent {
call_id: id.to_owned(),
first_party: true,
r#type: Some(tool_result_content::Type::FunctionResult(Box::new(
FunctionResultContent {
name: "search".to_owned(),
result: response,
..Default::default()
},
))),
..Default::default()
}))),
..Default::default()
}),
..Default::default()
}
.encode_to_vec()
}
let kept = uuid::Uuid::parse_str("01950000-0000-7000-8000-00000000e501").unwrap();
let excised = uuid::Uuid::parse_str("01950000-0000-7000-8000-00000000e502").unwrap();
let summary = uuid::Uuid::parse_str("01950000-0000-7000-8000-00000000e503").unwrap();
let signer = ApprovalSigner::from_key_bytes(&[23; 32]).unwrap();
let (marker, _, _) = excision_payload(
"a",
EXCISION_SCOPE_SOURCE_ONLY,
&[11],
"persona-a",
"parity proof",
&signer,
);
let mut events = vec![
(1, Event::new(tagged(kinds::TURN_START, kept), Vec::new())),
(
2,
Event::new(
tagged(kinds::USAGE, kept),
UsageEvent {
input_tokens: 7,
output_tokens: 3,
..Default::default()
}
.encode_to_vec(),
),
),
(
3,
Event::new(
tagged(kinds::MODEL_CALL, kept),
ModelCallEvent {
provider: "provider-a".to_owned(),
model: "model-a".to_owned(),
captured_clock_unix_ms: 1_234,
..Default::default()
}
.encode_to_vec(),
),
),
(
4,
Event::new(tagged(kinds::OUTPUT_MSG, kept), tool_call("kept")),
),
(
5,
Event::new(tagged(kinds::OUTPUT_MSG, kept), tool_result("kept")),
),
(
6,
Event::new(
tagged(kinds::TURN_FAILED, kept),
TurnFailedEvent {
kind: TurnFailureKind::Timeout.into(),
message: "provider timed out".to_owned(),
..Default::default()
}
.encode_to_vec(),
),
),
(
7,
Event::new(tagged(kinds::TURN_COMPLETE, kept), Vec::new()),
),
(
8,
Event::new(
tagged(kinds::SUMMARY, summary),
SummaryEvent {
text: "fleet summary".to_owned(),
covers_through_position: 7,
..Default::default()
}
.encode_to_vec(),
),
),
(
9,
Event::new(tagged(kinds::TURN_START, excised), Vec::new()),
),
(
10,
Event::new(tagged(kinds::OUTPUT_MSG, excised), tool_call("removed")),
),
(
11,
Event::new(tagged(kinds::OUTPUT_MSG, excised), tool_result("removed")),
),
(
12,
Event::new(tagged(kinds::TURN_COMPLETE, excised), Vec::new()),
),
(13, Event::new(kinds::TAINT_EXCISION, marker)),
];
let (prepared, facts) =
polyc_facts::prepare_and_fold_conversation_execution(&mut events, "conv-a").unwrap();
assert!(prepared.excised.contains_key(&10) && prepared.excised.contains_key(&11));
let segments = encode_execution_generation(
polyc_projection::family::conversation_execution(),
&RowSource::new("conv-a".to_owned(), [1; 32]),
&facts,
EncodingBounds::new(128, 128, 4 * 1024 * 1024, 32 * 1024 * 1024),
)
.unwrap();
let harness = Harness::with_encoded_execution(segments);
for (index, (sql, expected, fleet)) in [
(
"SELECT CAST(input_tokens AS VARCHAR) || ':' || CAST(output_tokens AS VARCHAR) FROM usage",
&["7:3"][..],
false,
),
(
"SELECT provider || ':' || model || ':' || CAST(captured_clock_unix_ms AS VARCHAR) FROM model_call",
&["provider-a:model-a:1234"][..],
false,
),
(
"SELECT tool_call_id || ':' || block_type FROM tool_calls ORDER BY position",
&["kept:call", "kept:result"][..],
false,
),
(
"SELECT failure_kind || ':' || message FROM turn_failed",
&["timeout:provider timed out"][..],
false,
),
(
"SELECT turn_id || ':' || text || ':' || CAST(covers_through_position AS VARCHAR) FROM summary",
&["01950000-0000-7000-8000-00000000e503:fleet summary:7"][..],
true,
),
]
.into_iter()
.enumerate()
{
let scope = if fleet {
QueryScope::Fleet
} else {
QueryScope::Conversations { conversations: vec!["a".to_owned()], memory: MemorySources::default() }
};
let prepared = harness
.prepare(&format!("execution-parity-{index}"), sql, scope, 20)
.await;
let authority = if fleet {
harness.fleet_authority()
} else {
harness.visible_authority()
};
let rows = harness.collect_text(authority, prepared).await;
assert_eq!(
rows.iter().map(String::as_str).collect::<Vec<_>>(),
expected
);
}
}
#[tokio::test]
async fn mixed_family_plans_read_each_exact_descriptor_once() {
let harness = Harness::new();
let prepared = harness
.prepare(
"mixed-exact-families",
"SELECT m.text, u.input_tokens FROM messages m \
JOIN usage u USING (partition, turn_id) ORDER BY m.position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut stream = harness
.visible_authority()
.bind(prepared)
.await
.expect("both family manifests bind")
.execute();
let batch = stream.next().await.unwrap().unwrap();
assert_eq!(batch.num_rows(), 3);
let usage = batch
.column(1)
.as_any()
.downcast_ref::<UInt64Array>()
.unwrap();
assert!((0..usage.len()).all(|row| usage.value(row) == 7));
assert!(stream.next().await.is_none());
}
#[tokio::test]
async fn a_visible_session_cannot_read_credential_lifecycle_history() {
let harness = Harness::new();
for table in [
"credential_lifecycle",
"credential_key_lifecycle",
"credential_lifecycle_refusals",
] {
let refusal = harness
.try_prepare(
"visible-credentials",
&format!("SELECT position FROM {table}"),
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await
.unwrap_err();
assert!(
matches!(
refusal,
crate::core_resolution::CoreResolutionError::TableOutsideRealm
| crate::core_resolution::CoreResolutionError::FamilyOutsideRealm
),
"{table} was not refused for a visible session: {refusal:?}"
);
}
}
#[tokio::test]
async fn a_visible_session_cannot_read_the_routine_lifecycle_or_the_persona_directory() {
let harness = Harness::new();
for table in [
"routine_lifecycle",
"routine_setup",
"fires",
"persona_profiles",
"persona_identities",
"persona_participations",
"persona_visibility",
"persona_refusals",
] {
let refusal = harness
.try_prepare(
"visible-fleet-only-families",
&format!("SELECT position FROM {table}"),
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await
.unwrap_err();
assert!(
matches!(
refusal,
crate::core_resolution::CoreResolutionError::TableOutsideRealm
| crate::core_resolution::CoreResolutionError::FamilyOutsideRealm
),
"{table} was not refused for a visible session: {refusal:?}"
);
}
}
fn administrator_audit_segments() -> Vec<polyc_projector::artifact::EncodedSegment> {
use polyc_facts::{AdminModelChangeFact, AdminSignatureStatus, AdministratorAuditFacts};
let change = |position: u64, status: AdminSignatureStatus, model: &str| AdminModelChangeFact {
position,
signature_status: status,
principal: if status == AdminSignatureStatus::Malformed {
String::new()
} else {
"admin:root".to_owned()
},
previous_provider: "prov-a".to_owned(),
previous_model: "model-a".to_owned(),
new_provider: "prov-b".to_owned(),
new_model: model.to_owned(),
changed_at_ms: 1_750_000_000_000,
};
let facts = AdministratorAuditFacts {
model_changes: vec![
change(1, AdminSignatureStatus::VerifiedCurrent, "model-current"),
change(2, AdminSignatureStatus::VerifiedRetired, "model-retired"),
change(3, AdminSignatureStatus::Malformed, ""),
],
};
polyc_projector::artifact::encode_administrator_audit_generation(
polyc_projection::family::administrator_audit(),
&polyc_projector::artifact::RowSource::new(
polyc_projection::family::ADMIN_AUDIT_PARTITION.to_owned(),
[2; 32],
),
&facts,
polyc_projector::artifact::EncodingBounds::default(),
)
.unwrap()
}
#[tokio::test]
async fn only_a_fleet_session_can_read_the_administrator_history() {
let harness = Harness::with_encoded_administrator_audit(administrator_audit_segments());
let refusal = harness
.try_prepare(
"visible-administrator-audit",
"SELECT principal FROM admin_model_changes",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await
.expect_err("a conversation-scoped session cannot name the administrator history");
assert!(
matches!(
refusal,
crate::core_resolution::CoreResolutionError::TableOutsideRealm
| crate::core_resolution::CoreResolutionError::FamilyOutsideRealm
),
"{refusal:?}"
);
let prepared = harness
.prepare(
"fleet-administrator-audit",
"SELECT principal FROM admin_model_changes ORDER BY position",
QueryScope::Fleet,
10,
)
.await;
let rows = harness
.collect_text(harness.fleet_authority(), prepared)
.await;
assert_eq!(
rows,
["admin:root", "admin:root", ""],
"the Fleet session reads every row, including the malformed record's empty payload"
);
}
#[tokio::test]
async fn a_session_scoped_to_the_partition_still_cannot_read_it() {
let harness = Harness::with_encoded_administrator_audit(administrator_audit_segments());
let refusal = harness
.try_prepare(
"partition-scoped-administrator-audit",
"SELECT principal FROM admin_model_changes",
QueryScope::Conversations {
conversations: vec![polyc_projection::family::ADMIN_AUDIT_PARTITION.to_owned()],
memory: MemorySources::default(),
},
10,
)
.await
.expect_err("naming the partition must not admit the family");
assert!(
matches!(
refusal,
crate::core_resolution::CoreResolutionError::TableOutsideRealm
| crate::core_resolution::CoreResolutionError::FamilyOutsideRealm
),
"{refusal:?}"
);
}
#[tokio::test]
async fn the_administrator_history_carries_every_verdict_it_published() {
let harness = Harness::with_encoded_administrator_audit(administrator_audit_segments());
let prepared = harness
.prepare(
"fleet-administrator-audit-status",
"SELECT signature_status FROM admin_model_changes ORDER BY position",
QueryScope::Fleet,
10,
)
.await;
let rows = harness
.collect_text(harness.fleet_authority(), prepared)
.await;
assert_eq!(
rows,
[
polyc_facts::AdminSignatureStatus::VerifiedCurrent.as_str(),
polyc_facts::AdminSignatureStatus::VerifiedRetired.as_str(),
polyc_facts::AdminSignatureStatus::Malformed.as_str(),
]
);
}
fn query_audit_segments() -> Vec<polyc_projector::artifact::EncodedSegment> {
use polyc_facts::{
QueryAuditCompletionFact, QueryAuditFacts, QueryAuditIntentFact, QueryAuditSourcePinFact,
};
let facts = QueryAuditFacts {
intents: vec![QueryAuditIntentFact {
position: 1,
namespace: "tenant-a".to_owned(),
query_id: "q-1".to_owned(),
requester: "requester-1".to_owned(),
shape_digest: [1; 32],
recorded_at_nanos: 1_000,
source_pin_count: 1,
command_id: "begin-1".to_owned(),
command_digest: [2; 32],
}],
completions: vec![QueryAuditCompletionFact {
position: 2,
namespace: "tenant-a".to_owned(),
query_id: "q-1".to_owned(),
intent_position: 1,
outcome: "succeeded",
error_class: None,
duration_nanos: 500,
rows: 3,
truncation: "complete",
truncated_at: None,
command_id: "complete-1".to_owned(),
command_digest: [3; 32],
}],
pins: vec![QueryAuditSourcePinFact {
position: 1,
namespace: "tenant-a".to_owned(),
query_id: "q-1".to_owned(),
pin_index: 0,
pin_kind: "authoritative",
family: None,
source_partition: None,
pin_source_incarnation: None,
projection_generation: None,
evidence_kind: None,
evidence_position: None,
evidence_journal_position: None,
schema_version: None,
fact_version: None,
artifact_digest: None,
anchor_head: None,
revision: Some(7),
}],
};
polyc_projector::artifact::encode_query_audit_generation(
polyc_projection::family::query_audit(),
&polyc_projector::artifact::RowSource::new(
polyc_projection::family::QUERY_AUDIT_SOURCE.to_owned(),
[42; 32],
),
&facts,
polyc_projector::artifact::EncodingBounds::default(),
)
.unwrap()
}
#[tokio::test]
async fn only_a_fleet_session_can_read_the_query_audit_history() {
let harness = Harness::with_encoded_query_audit(query_audit_segments());
for table in [
"query_audit_intents",
"query_audit_completions",
"query_audit_source_pins",
] {
let refusal = harness
.try_prepare(
&format!("visible-{table}"),
&format!("SELECT namespace FROM {table}"),
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await
.expect_err("a conversation-scoped session cannot name the query-audit history");
assert!(
matches!(
refusal,
crate::core_resolution::CoreResolutionError::TableOutsideRealm
| crate::core_resolution::CoreResolutionError::FamilyOutsideRealm
),
"{table}: {refusal:?}"
);
}
let prepared = harness
.prepare(
"fleet-query-audit-intents",
"SELECT query_id FROM query_audit_intents ORDER BY position",
QueryScope::Fleet,
10,
)
.await;
let rows = harness
.collect_text(harness.fleet_authority(), prepared)
.await;
assert_eq!(rows, ["q-1"]);
}
#[tokio::test]
async fn naming_the_query_audit_source_does_not_admit_it() {
let harness = Harness::with_encoded_query_audit(query_audit_segments());
let refusal = harness
.try_prepare(
"partition-scoped-query-audit",
"SELECT query_id FROM query_audit_intents",
QueryScope::Conversations {
conversations: vec![polyc_projection::family::QUERY_AUDIT_SOURCE.to_owned()],
memory: MemorySources::default(),
},
10,
)
.await
.expect_err("naming the source string must not admit the family");
assert!(matches!(
refusal,
crate::core_resolution::CoreResolutionError::TableOutsideRealm
| crate::core_resolution::CoreResolutionError::FamilyOutsideRealm
));
}
#[tokio::test]
async fn the_query_audit_tables_join_on_position() {
let harness = Harness::with_encoded_query_audit(query_audit_segments());
let prepared = harness
.prepare(
"fleet-query-audit-join",
"SELECT c.outcome FROM query_audit_completions c \
JOIN query_audit_intents i ON c.intent_position = i.position \
WHERE i.query_id = 'q-1'",
QueryScope::Fleet,
10,
)
.await;
let rows = harness
.collect_text(harness.fleet_authority(), prepared)
.await;
assert_eq!(rows, ["succeeded"]);
}
#[tokio::test]
async fn a_query_over_the_query_audit_history_never_reads_its_own_row() {
let harness = Harness::with_encoded_query_audit(query_audit_segments());
let prepared = harness
.prepare(
"self-audit-check",
"SELECT query_id FROM query_audit_intents ORDER BY position",
QueryScope::Fleet,
10,
)
.await;
let rows = harness
.collect_text(harness.fleet_authority(), prepared)
.await;
assert_eq!(
rows,
["q-1"],
"the query's own id cannot appear in a generation resolved before its intent existed"
);
let audit = harness.state.audit_for("self-audit-check");
let pinned = audit
.intent()
.source()
.pins()
.iter()
.find_map(|pin| match pin {
polyc_state::query_audit::SourcePin::Projected(projected)
if projected.manifest().key().source().as_str()
== polyc_projection::family::QUERY_AUDIT_SOURCE =>
{
match projected.manifest().evidence() {
polyc_state::feed::SourceEvidence::QueryAudit(checkpoint) => Some(checkpoint),
_ => None,
}
}
_ => None,
});
assert!(
pinned.is_some(),
"the plan must pin exactly one query-audit generation, naming the frontier it stood on"
);
let order = harness.state.call_order();
let resolved_at = order.iter().position(|call| *call == "resolve_manifest");
let began_at = order.iter().position(|call| *call == "begin_audit");
assert!(
matches!((resolved_at, began_at), (Some(resolved), Some(began)) if resolved < began),
"resolve_manifest must precede begin_audit: {order:?}"
);
}
#[tokio::test]
async fn summary_is_fleet_only_in_planning_and_in_physical_release() {
let harness = Harness::new();
let refusal = harness
.try_prepare(
"visible-summary",
"SELECT text FROM summary",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await
.expect_err("visible authority cannot name the Fleet-only table");
assert!(matches!(
refusal,
crate::core_resolution::CoreResolutionError::TableOutsideRealm
));
let prepared = harness
.prepare(
"fleet-summary",
"SELECT turn_id, text FROM summary",
QueryScope::Fleet,
10,
)
.await;
let rows = harness
.collect_text(harness.fleet_authority(), prepared)
.await;
assert_eq!(rows, ["summary-a"]);
}
#[tokio::test]
async fn successful_eof_waits_for_durable_completion() {
let harness = Harness::new();
let prepared = harness
.prepare(
"terminal-success",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut stream = harness
.visible_authority()
.bind(prepared)
.await
.unwrap()
.execute();
let mut rows = 0;
while rows < 3 {
rows += stream.next().await.unwrap().unwrap().num_rows();
}
assert!(
harness.state.completion_for("terminal-success").is_none(),
"a released batch is not a successful terminal"
);
assert!(stream.next().await.is_none());
let completion = harness
.state
.completion_for("terminal-success")
.expect("EOF follows a durable completion");
assert_eq!(completion.outcome(), QueryOutcome::Succeeded);
assert_eq!(harness.state.completion_attempts(), 1);
}
#[tokio::test]
async fn ambiguous_completion_settles_the_exact_receipt() {
let harness = Harness::new();
harness.state.arm_completion_response_loss();
let prepared = harness
.prepare(
"terminal-ambiguous",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let rows = harness
.collect_text(harness.visible_authority(), prepared)
.await;
assert_eq!(rows.len(), 3);
assert_eq!(
harness.state.completion_attempts(),
1,
"receipt settlement must not construct or send a second command"
);
assert_eq!(harness.state.completion_receipt_reads(), 1);
assert_eq!(
harness
.state
.completion_for("terminal-ambiguous")
.unwrap()
.outcome(),
QueryOutcome::Succeeded
);
}
#[tokio::test]
async fn dropped_stream_records_cancellation_and_never_success() {
let harness = Harness::new();
let prepared = harness
.prepare(
"terminal-drop",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut stream = harness
.visible_authority()
.bind(prepared)
.await
.unwrap()
.execute();
let delivered = stream.next().await.unwrap().unwrap().num_rows();
assert!(delivered > 0, "the consumer received rows before dropping");
drop(stream);
let completion = harness.await_completion("terminal-drop").await;
assert_eq!(
completion.outcome(),
QueryOutcome::Failed(ErrorClass::Cancelled)
);
assert_eq!(
completion.rows(),
polyc_state::query_audit::RowCount::new(u64::try_from(delivered).unwrap()),
"a cancelled terminal reports what the consumer received, never zero"
);
}
#[tokio::test]
async fn drop_while_eof_is_pending_selects_cancellation_before_dispatch() {
let harness = Harness::new();
let prepared = harness
.prepare(
"terminal-pending-drop",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let dispatch = prepared.pause_completion_dispatch();
let mut stream = harness
.visible_authority()
.bind(prepared)
.await
.unwrap()
.execute();
let mut rows = 0;
while rows < 3 {
rows += stream.next().await.unwrap().unwrap().num_rows();
}
let mut eof = Box::pin(stream.next());
tokio::select! {
() = dispatch.wait() => {}
result = &mut eof => panic!("EOF escaped before terminal dispatch: {result:?}"),
}
drop(eof);
drop(stream);
dispatch.resume();
let completion = harness.await_completion("terminal-pending-drop").await;
assert_eq!(
completion.outcome(),
QueryOutcome::Failed(ErrorClass::Cancelled)
);
}
#[tokio::test]
async fn completion_outage_exposes_no_successful_eof() {
let harness = Harness::new();
harness.state.set_completion_outage();
let prepared = harness
.prepare(
"terminal-outage",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut stream = harness
.visible_authority()
.bind(prepared)
.await
.unwrap()
.execute();
let mut rows = 0;
while rows < 3 {
rows += stream.next().await.unwrap().unwrap().num_rows();
}
let error = stream.next().await.unwrap().unwrap_err();
assert!(matches!(
error,
CoreExecutionError::AuditCompletionUnavailable
));
assert!(harness.state.completion_for("terminal-outage").is_none());
assert_eq!(harness.state.completion_attempts(), 3);
}
#[tokio::test]
async fn exact_provider_executes_duplicate_and_zero_column_projections() {
let harness = Harness::new();
let prepared = harness
.prepare(
"direct-provider-projection",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let bound = harness.visible_authority().bind(prepared).await.unwrap();
let provider = bound.provider(CoreTable::Messages);
let context = SessionContext::new();
let reads_after_binding = harness.store.range_calls();
let projection = vec![6, 4, 6];
let plan = provider
.scan(&context.state(), Some(&projection), &[], None)
.await
.unwrap();
assert_eq!(harness.store.range_calls(), reads_after_binding);
let batches = datafusion::physical_plan::collect(plan, context.task_ctx())
.await
.unwrap();
assert_eq!(
batches
.iter()
.map(arrow::record_batch::RecordBatch::num_rows)
.sum::<usize>(),
3
);
assert!(batches.iter().all(|batch| batch.num_columns() == 3));
for batch in &batches {
assert_eq!(batch.column(0), batch.column(2));
}
let empty = Vec::new();
let plan = provider
.scan(&context.state(), Some(&empty), &[], None)
.await
.unwrap();
let batches = datafusion::physical_plan::collect(plan, context.task_ctx())
.await
.unwrap();
assert_eq!(
batches
.iter()
.map(arrow::record_batch::RecordBatch::num_rows)
.sum::<usize>(),
3
);
assert!(batches.iter().all(|batch| batch.num_columns() == 0));
drop(bound);
}
#[tokio::test]
async fn zero_column_scan_preserves_row_count() {
let harness = Harness::new();
let prepared = harness
.prepare(
"zero-column",
"SELECT COUNT(*) AS held FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut stream = harness
.visible_authority()
.bind(prepared)
.await
.unwrap()
.execute();
let batch = stream.next().await.unwrap().unwrap();
let count = batch
.column(0)
.as_any()
.downcast_ref::<Int64Array>()
.unwrap();
assert_eq!(count.value(0), 3);
}
#[tokio::test]
async fn unsupported_filter_is_evaluated_above_the_exact_scan() {
let harness = Harness::new();
let prepared = harness
.prepare(
"filter-above-scan",
"SELECT text FROM messages WHERE position > 2 ORDER BY position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let rows = harness
.collect_text(harness.visible_authority(), prepared)
.await;
assert_eq!(rows, ["visible-two", "visible-three"]);
}
#[tokio::test]
async fn fleet_combines_disjoint_realms_while_visible_cannot() {
let harness = Harness::new();
let visible = harness
.prepare(
"realm-visible",
"SELECT text FROM messages ORDER BY position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let fleet = harness
.prepare(
"realm-fleet",
"SELECT text FROM messages ORDER BY position",
QueryScope::Fleet,
10,
)
.await;
let visible_rows = harness
.collect_text(harness.visible_authority(), visible)
.await;
let fleet_rows = harness.collect_text(harness.fleet_authority(), fleet).await;
assert_eq!(
visible_rows,
["visible-one", "visible-two", "visible-three"]
);
assert_eq!(
fleet_rows,
[
"visible-one",
"visible-two",
"visible-three",
"fleet-secret"
]
);
}
#[tokio::test]
async fn unreferenced_objects_and_scope_widening_add_no_files() {
let harness = Harness::new();
harness.store.add_unreferenced_object();
let prepared = harness
.prepare(
"unreferenced",
"SELECT text FROM messages ORDER BY position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let rows = harness
.collect_text(harness.visible_authority(), prepared)
.await;
assert_eq!(rows, ["visible-one", "visible-two", "visible-three"]);
}
#[tokio::test]
async fn permit_refusal_performs_no_artifact_io() {
let harness = Harness::new();
harness.state.refuse_audit();
let outcome = harness
.try_prepare(
"refused",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
assert!(outcome.is_err());
assert_eq!(harness.store.total_calls(), 0);
}
#[tokio::test]
async fn exact_digest_failure_refuses_before_a_row() {
let harness = Harness::new();
let prepared = harness
.prepare(
"corrupt",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
harness.store.corrupt_visible_messages();
let error = harness
.visible_authority()
.bind(prepared)
.await
.err()
.expect("complete-file proof catches corruption");
assert!(matches!(error, CoreExecutionError::Artifact(_)));
assert_eq!(
harness.state.completion_for("corrupt").unwrap().outcome(),
QueryOutcome::Failed(ErrorClass::Internal),
"stored bytes that disagree with their signed digest are corruption, not caller syntax"
);
assert_eq!(harness.reserved_memory(), 0);
}
#[tokio::test]
async fn dropping_a_bound_query_before_execute_records_cancellation() {
let harness = Harness::new();
let prepared = harness
.prepare(
"bound-drop",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let bound = harness.visible_authority().bind(prepared).await.unwrap();
drop(bound);
let completion = harness.await_completion("bound-drop").await;
assert_eq!(
completion.outcome(),
QueryOutcome::Failed(ErrorClass::Cancelled)
);
}
#[tokio::test]
async fn exact_generation_and_length_disagreement_refuse_before_a_row() {
for (query, defect) in [
(
"wrong-generation",
Harness::wrong_generation as fn(&Harness),
),
("wrong-length", Harness::wrong_length as fn(&Harness)),
] {
let harness = Harness::new();
let prepared = harness
.prepare(
query,
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
defect(&harness);
let error = harness
.visible_authority()
.bind(prepared)
.await
.err()
.expect("exact metadata disagreement is closed");
assert!(matches!(error, CoreExecutionError::Artifact(_)));
}
}
#[tokio::test]
async fn first_release_revalidation_fails_closed_on_scope_narrowing() {
let harness = Harness::new();
let prepared = harness
.prepare(
"narrowed",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut stream = harness
.visible_authority()
.bind(prepared)
.await
.unwrap()
.execute();
harness.revalidator.set(QueryScope::Conversations {
conversations: Vec::new(),
memory: MemorySources::default(),
});
let error = stream.next().await.unwrap().unwrap_err();
assert!(matches!(error, CoreExecutionError::AuthorityNarrowed));
assert_eq!(
harness.state.completion_for("narrowed").unwrap().outcome(),
QueryOutcome::Failed(ErrorClass::Denied)
);
assert!(stream.next().await.is_none());
}
#[tokio::test]
async fn first_release_revalidation_fails_closed_on_source_recreation() {
let harness = Harness::new();
let prepared = harness
.prepare(
"recreated",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut stream = harness
.visible_authority()
.bind(prepared)
.await
.unwrap()
.execute();
harness.state.recreate_source();
let error = stream.next().await.unwrap().unwrap_err();
assert!(matches!(error, CoreExecutionError::SourceChanged(_)));
assert_eq!(
harness.state.completion_for("recreated").unwrap().outcome(),
QueryOutcome::Failed(ErrorClass::Unavailable)
);
}
#[tokio::test]
async fn first_release_revalidation_fails_closed_on_state_outage() {
let harness = Harness::new();
let prepared = harness
.prepare(
"revalidation-outage",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut stream = harness
.visible_authority()
.bind(prepared)
.await
.unwrap()
.execute();
harness.state.set_outage();
let error = stream.next().await.unwrap().unwrap_err();
assert!(matches!(error, CoreExecutionError::Resolution(_)));
assert_eq!(
harness
.state
.completion_for("revalidation-outage")
.unwrap()
.outcome(),
QueryOutcome::Failed(ErrorClass::Unavailable)
);
}
#[tokio::test]
async fn periodic_revalidation_stops_subsequent_batches() {
for defect in ["narrow", "recreate", "outage"] {
let harness = Harness::new();
let prepared = harness
.prepare(
&format!("periodic-{defect}"),
"SELECT text FROM messages",
QueryScope::Fleet,
10,
)
.await;
let mut stream = harness
.fleet_authority()
.bind(prepared)
.await
.unwrap()
.execute();
let first = stream.next().await.unwrap().unwrap();
assert!(first.num_rows() > 0);
tokio::time::sleep(Duration::from_millis(2)).await;
match defect {
"narrow" => harness.revalidator.set(QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
}),
"recreate" => harness.state.recreate_source(),
"outage" => harness.state.set_outage(),
_ => unreachable!(),
}
let error = stream.next().await.unwrap().unwrap_err();
assert!(matches!(
(defect, error),
("narrow", CoreExecutionError::AuthorityNarrowed)
| ("recreate", CoreExecutionError::SourceChanged(_))
| ("outage", CoreExecutionError::Resolution(_))
));
}
}
#[tokio::test]
async fn row_cap_releases_exactly_the_cap() {
let harness = Harness::new();
for (query, cap, expected) in [
("cap-minus-one", 2, 2),
("cap", 3, 3),
("cap-plus-one", 4, 3),
] {
let prepared = harness
.prepare(
query,
"SELECT text FROM messages ORDER BY position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
cap,
)
.await;
let rows = harness
.collect_text(harness.visible_authority(), prepared)
.await;
assert_eq!(rows.len(), expected);
let completion = harness.state.completion_for(query).unwrap();
assert_eq!(
completion.truncation(),
if cap < 3 {
Truncation::TruncatedAt(u64::try_from(expected).unwrap())
} else {
Truncation::Complete
}
);
}
}
#[tokio::test]
async fn release_byte_bound_refuses_before_release() {
let harness = Harness::with_release_bounds(1, 1);
let prepared = harness
.prepare(
"release-byte-bound",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut stream = harness
.visible_authority()
.bind(prepared)
.await
.unwrap()
.execute();
let error = stream.next().await.unwrap().unwrap_err();
assert!(matches!(error, CoreExecutionError::ReleaseBound { .. }));
}
#[tokio::test]
#[allow(
clippy::significant_drop_tightening,
reason = "the test observes the bound query's live pool reservation before consuming it"
)]
async fn retained_scan_performs_no_backend_reads_and_drop_releases_memory() {
let harness = Harness::new();
let prepared = harness
.prepare(
"drop-cancels",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let bound = harness.visible_authority().bind(prepared).await.unwrap();
let calls = harness.store.range_calls();
assert!(bound.source_decode_reservation.size() > 0);
assert!(harness.reserved_memory() > 0);
let mut stream = bound.execute();
assert!(stream.next().await.unwrap().is_ok());
assert_eq!(harness.store.range_calls(), calls);
drop(stream);
tokio::time::timeout(HANG_GUARD, async {
while harness.reserved_memory() != 0 {
tokio::task::yield_now().await;
}
})
.await
.expect("the cancelled producer releases its source reservation");
}
#[tokio::test]
async fn source_decode_limit_refuses_before_segment_body_reads() {
let harness = Harness::with_source_decode_bound(1);
let prepared = harness
.prepare(
"decode-limit",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let error = harness
.visible_authority()
.bind(prepared)
.await
.err()
.expect("the authenticated envelope exceeds the deployment limit");
assert!(matches!(
error,
CoreExecutionError::SourceDecodeBound { .. }
));
assert_eq!(harness.store.segment_range_calls(), 0);
assert_eq!(harness.reserved_memory(), 0);
}
#[tokio::test]
async fn aggregate_source_decode_reservation_uses_the_shared_runtime_pool() {
let reservations = Harness::visible_message_decode_reservations();
assert!(reservations.len() >= 2);
let aggregate = reservations.iter().copied().sum::<u64>() + DEFAULT_ARTIFACT_RANGE_BYTES + 1;
let pool_limit = usize::try_from(aggregate - 1).unwrap();
assert!(reservations.iter().all(|value| *value < pool_limit as u64));
let harness = Harness::with_runtime_memory(pool_limit);
let prepared = harness
.prepare(
"aggregate-decode-pool",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let error = harness
.visible_authority()
.bind(prepared)
.await
.err()
.expect("the summed reservation exceeds the shared pool");
assert!(matches!(
error,
CoreExecutionError::DataFusion(DataFusionError::ResourcesExhausted(_))
));
assert_eq!(harness.store.segment_range_calls(), 0);
assert_eq!(harness.reserved_memory(), 0);
}
#[tokio::test]
#[allow(
clippy::significant_drop_tightening,
reason = "execute consumes the bound query; an earlier drop would withdraw it before the deadline can fire"
)]
async fn original_deadline_stops_retained_decode_before_release() {
let harness = Harness::new();
let prepared = harness
.prepare(
"range-deadline",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut bound = harness.visible_authority().bind(prepared).await.unwrap();
bound.restart_deadline_for_test(Duration::from_millis(100));
let mut stream = bound.execute();
tokio::time::sleep(Duration::from_millis(110)).await;
let error = stream.next().await.unwrap().unwrap_err();
assert!(matches!(error, CoreExecutionError::Deadline));
assert_eq!(
harness
.state
.completion_for("range-deadline")
.unwrap()
.outcome(),
QueryOutcome::Failed(ErrorClass::Deadline),
"terminal settlement has a fresh server-owned budget"
);
}
#[tokio::test]
async fn expired_deadline_refuses_before_artifact_io() {
let harness = Harness::new();
let prepared = harness
.prepare(
"bind-deadline",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await
.with_expired_operation_for_test();
let error = harness
.visible_authority()
.bind(prepared)
.await
.err()
.expect("expired operation cannot enter artifacts");
assert!(matches!(error, CoreExecutionError::Deadline));
assert_eq!(harness.store.total_calls(), 0);
}
#[tokio::test]
async fn original_deadline_stops_a_pending_bind_time_artifact_read() {
let harness = Harness::new();
let prepared = harness
.prepare(
"bind-range-deadline",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await
.with_fresh_operation_for_test(Duration::from_millis(500));
harness.store.pause_ranges();
let error = harness
.visible_authority()
.bind(prepared)
.await
.err()
.expect("the original deadline stops stalled artifact admission");
assert!(matches!(error, CoreExecutionError::Deadline));
assert!(harness.store.range_calls() > 0);
assert_eq!(harness.reserved_memory(), 0);
}
#[tokio::test]
async fn aggregate_execution_admission_queues_before_artifact_reads_and_releases_on_drop() {
let harness = Harness::with_concurrency(1);
let first = harness
.prepare(
"admission-first",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let first = harness.visible_authority().bind(first).await.unwrap();
let calls_while_held = harness.store.total_calls();
let second = harness
.prepare(
"admission-second",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await
.with_fresh_operation_for_test(Duration::from_millis(100));
let error = harness
.visible_authority()
.bind(second)
.await
.err()
.expect("the sole aggregate slot remains owned");
assert!(matches!(error, CoreExecutionError::Deadline));
assert_eq!(harness.store.total_calls(), calls_while_held);
drop(first);
let third = harness
.prepare(
"admission-third",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let rebound = harness.visible_authority().bind(third).await.unwrap();
drop(rebound);
assert!(harness.store.total_calls() > calls_while_held);
}
#[tokio::test]
async fn typed_parameter_plan_executes_without_replacement_sql() {
let harness = Harness::new();
let prepared = harness
.prepare_with_parameters(
"typed-real-plan",
"SELECT text FROM messages WHERE position > $1 ORDER BY position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
vec![CoreParameter::UInt64(2)],
)
.await;
let rows = harness
.collect_text(harness.visible_authority(), prepared)
.await;
assert_eq!(rows, ["visible-two", "visible-three"]);
}
#[tokio::test]
async fn fresh_request_catalogs_cannot_observe_each_other() {
let family = polyc_projection::family::conversation_core();
let turns: Arc<dyn TableProvider> = Arc::new(EmptyTable::new(arrow_schema(
family
.table(polyc_projection::family::CONVERSATION_TURNS)
.unwrap(),
)));
let messages: Arc<dyn TableProvider> = Arc::new(EmptyTable::new(arrow_schema(
family
.table(polyc_projection::family::CONVERSATION_MESSAGES)
.unwrap(),
)));
let base = SessionContext::new().state();
let (first, _) = request_context(&base, &BTreeMap::from([(CoreTable::Turns, turns)]))
.await
.unwrap();
let (second, _) = request_context(&base, &BTreeMap::from([(CoreTable::Messages, messages)]))
.await
.unwrap();
assert!(first.table_exist("turns").unwrap());
assert!(!first.table_exist("messages").unwrap());
assert!(second.table_exist("messages").unwrap());
assert!(!second.table_exist("turns").unwrap());
}
#[tokio::test]
async fn prepared_realm_cannot_be_replaced_after_audit() {
let harness = Harness::new();
let prepared = harness
.prepare(
"crossed-realm",
"SELECT text FROM messages",
QueryScope::Fleet,
10,
)
.await;
let calls = harness.store.total_calls();
let error = harness
.visible_authority()
.bind(prepared)
.await
.err()
.expect("authority realm is retained in prepared state");
assert!(matches!(error, CoreExecutionError::RealmMismatch));
assert_eq!(harness.store.total_calls(), calls);
}
#[tokio::test]
async fn exact_audit_replay_never_yields_a_second_execution_capability() {
let harness = Harness::new();
let scope = QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
};
let first = harness
.try_prepare("replay", "SELECT text FROM messages", scope.clone(), 10)
.await
.unwrap();
let second = harness
.try_prepare("replay", "SELECT text FROM messages", scope, 10)
.await
.unwrap();
assert!(matches!(first, CorePlanOutcome::Granted(_)));
assert!(matches!(second, CorePlanOutcome::AlreadyRecorded(_)));
assert_eq!(harness.store.total_calls(), 0);
}
#[tokio::test]
#[allow(
clippy::significant_drop_tightening,
reason = "execute consumes the bound query; an earlier drop would withdraw it before the deadline can fire"
)]
async fn an_idle_consumer_settles_the_deadline_and_releases_its_slot() {
let harness = Harness::new();
let prepared = harness
.prepare(
"idle-consumer",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut bound = harness.visible_authority().bind(prepared).await.unwrap();
bound.restart_deadline_for_test(Duration::from_millis(60));
let stream = bound.execute();
let completion = harness.await_completion("idle-consumer").await;
assert_eq!(
completion.outcome(),
QueryOutcome::Failed(ErrorClass::Deadline)
);
assert_eq!(
completion.rows(),
polyc_state::query_audit::RowCount::new(0)
);
drop(stream);
tokio::time::timeout(HANG_GUARD, async {
while harness.reserved_memory() != 0 {
tokio::task::yield_now().await;
}
})
.await
.expect("the expired producer releases its source reservation");
}
#[tokio::test]
async fn a_descheduled_producer_still_records_the_exact_delivered_rows() {
let harness = Harness::new();
let prepared = harness
.prepare(
"descheduled-report",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut stream = {
let mut bound = harness.visible_authority().bind(prepared).await.unwrap();
bound.delay_report_for_test(Duration::from_millis(200));
bound.execute()
};
let delivered = stream.next().await.unwrap().unwrap().num_rows();
assert!(
delivered > 0,
"the consumer received rows before withdrawing"
);
drop(stream);
let completion = harness.await_completion("descheduled-report").await;
assert_eq!(
completion.outcome(),
QueryOutcome::Failed(ErrorClass::Cancelled)
);
assert_eq!(
completion.rows(),
polyc_state::query_audit::RowCount::new(u64::try_from(delivered).unwrap()),
"a withdrawn query reports what the consumer received, whatever the scheduler did"
);
}
#[tokio::test]
#[allow(
clippy::significant_drop_tightening,
reason = "execute consumes the bound query; an earlier drop would withdraw it before the deadline can fire"
)]
async fn a_measured_deadline_survives_a_withdrawal_at_the_dispatch_boundary() {
let harness = Harness::new();
let prepared = harness
.prepare(
"deadline-vs-withdrawal",
"SELECT text FROM messages",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let dispatch = prepared.pause_completion_dispatch();
let mut bound = harness.visible_authority().bind(prepared).await.unwrap();
bound.restart_deadline_for_test(Duration::from_millis(60));
let stream = bound.execute();
tokio::time::timeout(HANG_GUARD, dispatch.wait())
.await
.expect("the guardian reaches the dispatch boundary");
drop(stream);
dispatch.resume();
let completion = harness.await_completion("deadline-vs-withdrawal").await;
assert_eq!(
completion.outcome(),
QueryOutcome::Failed(ErrorClass::Deadline),
"a withdrawal after a measured deadline never rewrites its class"
);
}
#[tokio::test]
#[allow(
clippy::significant_drop_tightening,
reason = "the bound query outlives the stream; dropping it early records cancellation, not the success this case is about"
)]
async fn a_successful_stream_releases_one_terminal_and_then_ends() {
let harness = Harness::new();
let prepared = harness
.prepare(
"one-terminal",
"SELECT text FROM messages ORDER BY position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let bound = harness
.visible_authority()
.bind(prepared)
.await
.expect("exact files bind");
let mut stream =
crate::core_service::ProjectedResultStream::start(bound.execute(), 64 * 1024, u64::MAX)
.expect("a representable frame ceiling");
let mut kinds = Vec::new();
let mut released_rows = 0_u64;
let mut released_bytes = 0_u64;
let mut terminal = None;
for _ in 0..32 {
match stream.next_frame().await {
Some(frame) => kinds.push(match frame.expect("a produced frame") {
polyc_query_model::ResultFrame::Schema(_) => "schema",
polyc_query_model::ResultFrame::Data(data) => {
released_rows += data.rows();
released_bytes += data.arrow_ipc().len() as u64;
"data"
}
polyc_query_model::ResultFrame::Terminal(frame) => {
terminal = Some(frame);
"terminal"
}
}),
None => break,
}
}
assert_eq!(
kinds.iter().filter(|kind| **kind == "terminal").count(),
1,
"exactly one terminal: {kinds:?}"
);
assert_eq!(
kinds.last(),
Some(&"terminal"),
"the terminal is last: {kinds:?}"
);
assert!(
stream.next_frame().await.is_none(),
"the stream stays ended after its terminal"
);
let terminal = terminal.expect("the stream released a terminal");
assert!(released_rows > 0, "this case releases rows to count");
assert_eq!(
terminal.rows(),
released_rows,
"terminal rows are released rows"
);
assert_eq!(
terminal.result_bytes(),
released_bytes,
"terminal bytes are the bytes released across data frames"
);
}
#[tokio::test]
#[allow(
clippy::significant_drop_tightening,
reason = "the bound query outlives the refused stream; dropping it early records a withdrawal this case says nothing about"
)]
async fn a_frame_ceiling_below_the_schema_is_the_callers_bound() {
let harness = Harness::new();
let prepared = harness
.prepare(
"ceiling-below-schema",
"SELECT text FROM messages ORDER BY position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let bound = harness
.visible_authority()
.bind(prepared)
.await
.expect("exact files bind");
let Err(refusal) =
crate::core_service::ProjectedResultStream::start(bound.execute(), 1, u64::MAX)
else {
panic!("a one-byte ceiling cannot carry a schema");
};
let crate::core_service::QueryServiceError::Encode(encode) = refusal else {
panic!("a ceiling refusal is an encode refusal: {refusal:?}");
};
assert!(
matches!(
encode,
crate::core_service::FrameEncodeError::SchemaTooLarge { .. }
),
"the schema is measured against the caller's ceiling: {encode:?}"
);
assert_eq!(
encode.class(),
polyc_query_model::ErrorClass::Bounds,
"a ceiling the caller chose is the caller's bound"
);
let completion = harness.await_completion("ceiling-below-schema").await;
assert_eq!(
completion.outcome(),
polyc_state::query_audit::QueryOutcome::Failed(
polyc_state::query_audit::ErrorClass::Bounds
),
"the durable record says the bound refused it, not that the caller withdrew"
);
}
#[tokio::test]
#[allow(
clippy::significant_drop_tightening,
reason = "each bound query must outlive the stream it produced; dropping one early records a withdrawal these cases exist to rule out"
)]
async fn a_row_over_the_ceiling_ends_with_a_bounds_terminal() {
let harness = Harness::new();
let measured = harness
.prepare(
"ceiling-measure",
"SELECT text FROM messages ORDER BY position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let measured = harness
.visible_authority()
.bind(measured)
.await
.expect("exact files bind");
let mut measuring =
crate::core_service::ProjectedResultStream::start(measured.execute(), 64 * 1024, u64::MAX)
.expect("a representable frame ceiling");
let polyc_query_model::ResultFrame::Schema(schema) = measuring
.next_frame()
.await
.expect("a first frame")
.expect("a produced frame")
else {
panic!("the first frame is the schema");
};
let schema_bytes = schema.arrow_ipc().len() as u64;
drop(measuring);
let prepared = harness
.prepare(
"row-over-ceiling",
"SELECT text FROM messages ORDER BY position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let bound = harness
.visible_authority()
.bind(prepared)
.await
.expect("exact files bind");
let mut stream =
crate::core_service::ProjectedResultStream::start(bound.execute(), schema_bytes, u64::MAX)
.expect("the schema fits its own size");
let mut kinds = Vec::new();
let mut terminal = None;
for _ in 0..32 {
match stream.next_frame().await {
Some(frame) => kinds.push(match frame.expect("no frame is an error") {
polyc_query_model::ResultFrame::Schema(_) => "schema",
polyc_query_model::ResultFrame::Data(_) => "data",
polyc_query_model::ResultFrame::Terminal(frame) => {
terminal = Some(frame);
"terminal"
}
}),
None => break,
}
}
assert_eq!(
kinds,
vec!["schema", "terminal"],
"no data frame fits, and the stream still ends with one terminal"
);
let terminal = terminal.expect("the stream released a terminal");
assert_eq!(
terminal.outcome(),
polyc_query_model::QueryOutcome::Failed(polyc_query_model::ErrorClass::Bounds),
"the caller's ceiling is the caller's bound"
);
assert_eq!(terminal.rows(), 0, "no row was released");
drop(stream);
let completion = harness.await_completion("row-over-ceiling").await;
assert_eq!(
completion.outcome(),
polyc_state::query_audit::QueryOutcome::Failed(
polyc_state::query_audit::ErrorClass::Bounds
),
"the durable record says the bound refused it, not that the caller withdrew"
);
}
#[tokio::test]
#[allow(
clippy::significant_drop_tightening,
reason = "execute consumes the bound query; an earlier drop would withdraw it before the deadline can fire"
)]
async fn a_late_framing_failure_takes_the_producers_selected_terminal() {
let harness = Harness::new();
let measured = harness
.prepare(
"late-frame-measure",
"SELECT text FROM messages ORDER BY position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut measuring = crate::core_service::ProjectedResultStream::start(
harness
.visible_authority()
.bind(measured)
.await
.expect("exact files bind")
.execute(),
64 * 1024,
u64::MAX,
)
.expect("a representable frame ceiling");
let polyc_query_model::ResultFrame::Schema(schema) = measuring
.next_frame()
.await
.expect("a first frame")
.expect("a produced frame")
else {
panic!("the first frame is the schema");
};
let schema_bytes = schema.arrow_ipc().len() as u64;
drop(measuring);
let prepared = harness
.prepare(
"late-framing-failure",
"SELECT text FROM messages ORDER BY position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut bound = harness
.visible_authority()
.bind(prepared)
.await
.expect("exact files bind");
bound.restart_deadline_for_test(Duration::from_millis(500));
let mut stream =
crate::core_service::ProjectedResultStream::start(bound.execute(), schema_bytes, u64::MAX)
.expect("the schema fits its own size");
assert!(matches!(
stream.next_frame().await,
Some(Ok(polyc_query_model::ResultFrame::Schema(_)))
));
stream.request_buffered_batch();
tokio::time::timeout(HANG_GUARD, async {
while stream.buffered_frames() == 0 {
tokio::task::yield_now().await;
}
})
.await
.expect("the producer buffered one batch");
tokio::time::timeout(HANG_GUARD, async {
while !stream.terminal_selected() {
tokio::task::yield_now().await;
}
})
.await
.expect("the producer selected its deadline");
let terminal = match stream.next_frame().await.expect("a terminal follows") {
Ok(polyc_query_model::ResultFrame::Terminal(frame)) => frame,
Ok(polyc_query_model::ResultFrame::Data(data)) => {
panic!("a terminal query released {} rows", data.rows())
}
Ok(polyc_query_model::ResultFrame::Schema(_)) => {
panic!("the schema is released exactly once")
}
Err(error) => panic!("no frame is an error: {error}"),
};
assert_eq!(
terminal.outcome(),
polyc_query_model::QueryOutcome::Failed(polyc_query_model::ErrorClass::Deadline),
"the caller takes the terminal that already became the record"
);
assert_eq!(terminal.rows(), 0, "the buffered batch reached no caller");
drop(stream);
let completion = harness.await_completion("late-framing-failure").await;
assert_one_result(&terminal, &completion);
}
fn assert_one_result(
terminal: &polyc_query_model::TerminalFrame,
completion: &polyc_state::query_audit::QueryCompletion,
) {
let durable_outcome = match terminal.outcome() {
polyc_query_model::QueryOutcome::Succeeded => {
polyc_state::query_audit::QueryOutcome::Succeeded
}
polyc_query_model::QueryOutcome::Failed(class) => {
polyc_state::query_audit::QueryOutcome::Failed(crate::core_evidence::durable_class_of(
class,
))
}
};
assert_eq!(completion.outcome(), durable_outcome, "same outcome");
assert_eq!(
completion.rows(),
polyc_state::query_audit::RowCount::new(terminal.rows()),
"same rows"
);
let durable_truncation = match terminal.truncation() {
polyc_query_model::Truncation::Complete => polyc_state::query_audit::Truncation::Complete,
polyc_query_model::Truncation::TruncatedAt(rows) => {
polyc_state::query_audit::Truncation::TruncatedAt(rows)
}
};
assert_eq!(
completion.truncation(),
durable_truncation,
"same truncation"
);
assert_eq!(
terminal.source(),
&crate::core_evidence::evidence_of(completion.source())
.expect("the source vector is representable"),
"same premises"
);
}
async fn byte_ceiling_stream(
harness: &Harness,
query: &str,
frame_ceiling: u64,
release_ceiling: u64,
) -> crate::core_service::ProjectedResultStream {
let prepared = harness
.prepare(
query,
"SELECT m.text FROM messages m, messages n ORDER BY m.position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
64,
)
.await;
let bound = harness
.visible_authority()
.bind(prepared)
.await
.expect("exact files bind");
crate::core_service::ProjectedResultStream::start(
bound.execute(),
frame_ceiling,
release_ceiling,
)
.expect("a schema under both ceilings")
}
async fn drain_frames(
stream: &mut crate::core_service::ProjectedResultStream,
) -> (Vec<(u64, u64)>, polyc_query_model::TerminalFrame, u64) {
let mut frames = Vec::new();
let mut terminal = None;
for _ in 0..64 {
match stream.next_frame().await {
Some(frame) => match frame.expect("no frame is an error") {
polyc_query_model::ResultFrame::Schema(_) => {}
polyc_query_model::ResultFrame::Data(data) => {
frames.push((data.rows(), data.arrow_ipc().len() as u64));
}
polyc_query_model::ResultFrame::Terminal(frame) => terminal = Some(frame),
},
None => break,
}
assert!(
stream.frames_built() <= u64::try_from(frames.len()).unwrap() + 1,
"the encoder builds at most one frame beyond the released ones"
);
}
let built = stream.frames_built();
(
frames,
terminal.expect("the stream released a terminal"),
built,
)
}
#[tokio::test]
#[allow(
clippy::significant_drop_tightening,
reason = "each bound query outlives the stream it produced; dropping one early records a withdrawal this case exists to rule out"
)]
async fn release_stops_at_the_callers_encoded_byte_ceiling() {
let harness = Harness::new();
let mut whole =
byte_ceiling_stream(&harness, "byte-ceiling-measure", 64 * 1024, u64::MAX).await;
let (frames, terminal, _) = drain_frames(&mut whole).await;
assert_eq!(
frames.len(),
1,
"the whole result fits one frame: {frames:?}"
);
let (all_rows, all_bytes) = frames[0];
assert!(all_rows > 1, "this case needs a result that can be split");
assert_eq!(
terminal.truncation(),
polyc_query_model::Truncation::Complete,
"the unbounded run is not truncated"
);
drop(whole);
let ceiling = all_bytes - 1;
let mut bounded = byte_ceiling_stream(&harness, "encoded-byte-ceiling", ceiling, ceiling).await;
let (frames, terminal, built) = drain_frames(&mut bounded).await;
let released_rows: u64 = frames.iter().map(|(rows, _)| rows).sum();
let released_bytes: u64 = frames.iter().map(|(_, bytes)| bytes).sum();
assert!(!frames.is_empty(), "the caller receives what fits");
assert!(
released_bytes <= ceiling,
"released {released_bytes} bytes, over the caller's {ceiling} byte ceiling"
);
assert!(
released_rows < all_rows,
"the ceiling actually bound: {released_rows} of {all_rows} rows"
);
assert_eq!(
terminal.result_bytes(),
released_bytes,
"the terminal reports the bytes this stream released"
);
assert_eq!(terminal.rows(), released_rows, "and the rows");
assert_eq!(
terminal.outcome(),
polyc_query_model::QueryOutcome::Succeeded,
"a caller's own ceiling is not a failure"
);
assert_eq!(
terminal.truncation(),
polyc_query_model::Truncation::TruncatedAt(released_rows),
"stopping at the caller's byte ceiling is a truncation"
);
assert_eq!(
built,
u64::try_from(frames.len()).unwrap() + 1,
"the encoder builds the released frames and the one that did not fit"
);
drop(bounded);
let completion = harness.await_completion("encoded-byte-ceiling").await;
assert_ne!(
completion.outcome(),
polyc_state::query_audit::QueryOutcome::Failed(
polyc_state::query_audit::ErrorClass::Cancelled
),
"the caller withdrew nothing; it read exactly the result it asked for"
);
assert_eq!(
completion.outcome(),
polyc_state::query_audit::QueryOutcome::Succeeded,
"the durable outcome is the terminal's outcome"
);
assert_eq!(
completion.rows(),
polyc_state::query_audit::RowCount::new(released_rows),
"the durable row count is the rows the caller received"
);
assert_one_result(&terminal, &completion);
}
#[tokio::test]
async fn an_abandoned_stream_records_only_the_rows_the_caller_received() {
let harness = Harness::new();
let mut whole = byte_ceiling_stream(&harness, "abandon-measure", 64 * 1024, u64::MAX).await;
let (frames, _terminal, _built) = drain_frames(&mut whole).await;
let (all_rows, all_bytes) = frames[0];
assert!(all_rows > 1, "this case needs a result that can be split");
drop(whole);
let mut abandoned =
byte_ceiling_stream(&harness, "abandon-midway", all_bytes - 1, u64::MAX).await;
let mut released_rows = 0;
for _ in 0..2 {
match abandoned.next_frame().await.expect("a frame") {
Ok(polyc_query_model::ResultFrame::Data(data)) => released_rows += data.rows(),
Ok(_schema) => {}
Err(error) => panic!("no frame is an error: {error}"),
}
}
assert!(released_rows > 0, "the caller received a frame");
assert!(
released_rows < all_rows,
"and stopped with rows still to come: {released_rows} of {all_rows}"
);
drop(abandoned);
let completion = harness.await_completion("abandon-midway").await;
assert_eq!(
completion.outcome(),
polyc_state::query_audit::QueryOutcome::Failed(
polyc_state::query_audit::ErrorClass::Cancelled
),
"abandoning a stream is still a withdrawal"
);
assert_eq!(
completion.rows(),
polyc_state::query_audit::RowCount::new(released_rows),
"a withdrawal reports the rows the caller received, not the rows the \
batch carried"
);
assert_eq!(
completion.truncation(),
polyc_state::query_audit::Truncation::Complete,
"abandoning is not stopping at a ceiling; the outcome carries that, \
not the truncation"
);
}
#[tokio::test]
#[allow(
clippy::significant_drop_tightening,
reason = "the bound query outlives the stream; the case turns on what the producer does while the consumer still holds a batch"
)]
async fn framing_stops_when_the_query_settles_under_the_caller() {
let harness = Harness::new();
let mut whole = byte_ceiling_stream(&harness, "settle-measure", 64 * 1024, u64::MAX).await;
let (frames, _terminal, _built) = drain_frames(&mut whole).await;
let (all_rows, all_bytes) = frames[0];
assert!(all_rows > 1, "this case needs a result that can be split");
drop(whole);
let prepared = harness
.prepare(
"settle-under-caller",
"SELECT m.text FROM messages m, messages n ORDER BY m.position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut bound = harness
.visible_authority()
.bind(prepared)
.await
.expect("exact files bind");
bound.restart_deadline_for_test(Duration::from_millis(500));
let mut stream =
crate::core_service::ProjectedResultStream::start(bound.execute(), all_bytes - 1, u64::MAX)
.expect("a schema under both ceilings");
let mut released_rows = 0;
for _ in 0..2 {
match stream.next_frame().await.expect("a frame") {
Ok(polyc_query_model::ResultFrame::Data(data)) => released_rows += data.rows(),
Ok(_schema) => {}
Err(error) => panic!("no frame is an error: {error}"),
}
}
assert!(released_rows > 0, "the caller received a frame");
assert!(released_rows < all_rows, "with rows still unframed");
let built_before = stream.frames_built();
tokio::time::timeout(HANG_GUARD, async {
while !stream.terminal_selected() {
tokio::task::yield_now().await;
}
})
.await
.expect("the producer selected its deadline");
let terminal = loop {
match stream.next_frame().await.expect("a terminal follows") {
Ok(polyc_query_model::ResultFrame::Terminal(frame)) => break frame,
Ok(polyc_query_model::ResultFrame::Data(data)) => {
panic!("a settled query released {} more rows", data.rows())
}
Ok(polyc_query_model::ResultFrame::Schema(_)) => {}
Err(error) => panic!("no frame is an error: {error}"),
}
};
assert_eq!(
stream.frames_built(),
built_before + 1,
"one frame is built and refused admission; the rest of the batch is not framed"
);
assert_eq!(
terminal.outcome(),
polyc_query_model::QueryOutcome::Failed(polyc_query_model::ErrorClass::Deadline),
"the producer's measured deadline is the terminal"
);
assert_eq!(
terminal.rows(),
released_rows,
"and reports what was released"
);
drop(stream);
let completion = harness.await_completion("settle-under-caller").await;
assert_eq!(
completion.outcome(),
polyc_state::query_audit::QueryOutcome::Failed(
polyc_state::query_audit::ErrorClass::Deadline
),
"a measured deadline outranks the withdrawal that follows it"
);
assert_eq!(
completion.rows(),
polyc_state::query_audit::RowCount::new(released_rows),
"and the record counts exactly the rows the caller received"
);
assert_one_result(&terminal, &completion);
}
fn hex_lower(bytes: &[u8]) -> String {
use std::fmt::Write as _;
bytes.iter().fold(String::new(), |mut out, byte| {
let _ = write!(out, "{byte:02x}");
out
})
}
#[tokio::test]
#[allow(
clippy::too_many_lines,
reason = "one failure-sensitive fixture keeps both public row sets and the realm split together"
)]
async fn folded_journal_rows_match_every_projected_delegation_view() {
use buffa::Message as _;
use polyc_crypto::signing_role::{HandoffSigner, RoleTrustSet};
use polyc_eventlog_model::Event;
use polyc_projector::artifact::{EncodingBounds, RowSource, encode_delegation_generation};
use polyc_proto::kinds;
use polyc_proto::proto::polychrome::handoff::v1::{Handoff, HandoffDenied};
let current = HandoffSigner::from_seed(0x00de_1e6a);
let retired = HandoffSigner::from_seed(0x00de_1e6b);
let stranger = HandoffSigner::from_seed(0x00de_1e6c);
let trust = RoleTrustSet::<polyc_crypto::signing_role::HandoffRole>::from_public_keys(vec![
current.public_key_bytes(),
retired.public_key_bytes(),
])
.expect("a two-key handoff trust set");
let turn = uuid::Uuid::from_u128(0x5150);
let mut spawn = Handoff {
child_conversation_id: "conv-child".to_owned(),
child_agent_id: "researcher".to_owned(),
carried_count: 3,
reason: "delegate the search".to_owned(),
..Default::default()
};
polyc_crypto::handoff::sign_handoff_into(¤t, &mut spawn);
let mut denied = HandoffDenied {
parent_conversation_id: "conv-a".to_owned(),
parent_agent_id: "assistant".to_owned(),
child_agent_id: "banned".to_owned(),
reason: "delegate the payment".to_owned(),
denial_reason: "that agent is not in the allowlist".to_owned(),
allowed: vec!["coding".to_owned(), "research".to_owned()],
..Default::default()
};
polyc_crypto::handoff::sign_handoff_denied_into(&retired, &mut denied);
let mut foreign = Handoff {
child_conversation_id: "conv-foreign".to_owned(),
child_agent_id: "outsider".to_owned(),
carried_count: 0,
reason: "delegate from an untrusted key".to_owned(),
..Default::default()
};
polyc_crypto::handoff::sign_handoff_into(&stranger, &mut foreign);
let events = vec![
(
7,
Event::new(kinds::tagged(kinds::HANDOFF, &turn), spawn.encode_to_vec()),
),
(
8,
Event::new(
kinds::tagged(kinds::HANDOFF_DENIED, &turn),
denied.encode_to_vec(),
),
),
(
9,
Event::new(kinds::HANDOFF.to_owned(), foreign.encode_to_vec()),
),
];
let facts = polyc_facts::fold_conversation_delegation(&events, &trust).expect("the fold");
assert_eq!(
facts.handoffs.len(),
3,
"no committed-turn gate applies here"
);
assert_eq!(facts.signers.len(), 3);
let segments = encode_delegation_generation(
polyc_projection::family::conversation_delegation(),
&RowSource::new("conv-a".to_owned(), [1; 32]),
&facts,
EncodingBounds::default(),
)
.expect("the delegation generation encodes");
let harness = Harness::with_encoded_delegation(segments);
let prepared = harness
.prepare(
"delegation-parity-visible",
"SELECT partition || '|' || CAST(position AS VARCHAR) \
|| '|' || COALESCE(turn_id, 'NULL') \
|| '|' || phase \
|| '|' || COALESCE(child_conversation_id, 'NULL') \
|| '|' || COALESCE(child_agent_id, 'NULL') \
|| '|' || COALESCE(CAST(carried_count AS VARCHAR), 'NULL') \
|| '|' || COALESCE(reason, 'NULL') \
|| '|' || COALESCE(parent_agent_id, 'NULL') \
|| '|' || COALESCE(denial_reason, 'NULL') \
|| '|' || COALESCE(allowed, 'NULL') \
|| '|' || signature_status FROM handoffs ORDER BY position",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
20,
)
.await;
let rows = harness
.collect_text(harness.visible_authority(), prepared)
.await;
assert_eq!(
rows,
vec![
format!(
"conv-a|7|{turn}|handoff|conv-child|researcher|3|\
delegate the search|NULL|NULL|NULL|verified"
),
format!(
"conv-a|8|{turn}|handoff_denied|NULL|banned|NULL|\
delegate the payment|assistant|that agent is not in the allowlist|\
[\"coding\",\"research\"]|verified"
),
"conv-a|9|NULL|handoff|conv-foreign|outsider|0|\
delegate from an untrusted key|NULL|NULL|NULL|untrusted"
.to_owned(),
],
"every projected column must equal the fold, verdicts included"
);
let refusal = harness
.try_prepare(
"delegation-parity-visible-signers",
"SELECT signed_by FROM handoff_signers",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
20,
)
.await
.expect_err("a visible authority cannot name the Fleet-only signer table");
assert!(matches!(
refusal,
crate::core_resolution::CoreResolutionError::TableOutsideRealm
));
let prepared = harness
.prepare(
"delegation-parity-fleet-signers",
"SELECT CAST(position AS VARCHAR) \
|| '|' || encode(arrow_cast(signed_by, 'Binary'), 'hex') \
|| '|' || signer_key_id \
FROM handoff_signers ORDER BY position",
QueryScope::Fleet,
20,
)
.await;
let signers = harness
.collect_text(harness.fleet_authority(), prepared)
.await;
assert_eq!(
signers,
facts
.signers
.iter()
.map(|signer| format!(
"{}|{}|{}",
signer.position,
hex_lower(&signer.signed_by),
signer.signer_key_id
))
.collect::<Vec<_>>(),
"the Fleet rows must equal the signer facts the fold derived, key included"
);
assert_eq!(signers.len(), 3);
let prepared = harness
.prepare(
"delegation-parity-no-key",
"SELECT phase FROM handoffs",
QueryScope::Fleet,
20,
)
.await;
let schema = crate::core_resolution::CoreTable::Handoffs.public_schema();
assert!(
!schema
.fields()
.iter()
.any(|field| field.name() == "signed_by"),
"the handoff table must carry no signing key in any realm"
);
drop(
harness
.collect_text(harness.fleet_authority(), prepared)
.await,
);
}
#[tokio::test]
async fn json_expressions_cross_the_wire_as_text() {
let harness = Harness::new();
let prepared = harness
.prepare(
"json-e2e",
"SELECT json_get_str(arguments, 'q') AS str_form, \
arguments ->> 'q' AS operator_text, \
arguments -> 'q' AS operator_union, \
json_get(arguments, 'q') AS union_form \
FROM tool_calls WHERE block_type = 'call'",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await;
let mut stream = harness
.visible_authority()
.bind(prepared)
.await
.unwrap()
.execute();
let batch = stream.next().await.unwrap().unwrap();
assert!(stream.next().await.is_none(), "one released batch");
assert_eq!(batch.num_rows(), 1);
let schema = batch.schema();
for field in schema.fields() {
assert!(
!matches!(field.data_type(), arrow::datatypes::DataType::Union(..)),
"no union column crosses the boundary: {field:?}"
);
}
for name in ["operator_union", "union_form"] {
let field = schema.field_with_name(name).expect("column exists");
assert_eq!(
field
.metadata()
.get("ARROW:extension:name")
.map(String::as_str),
Some("arrow.json"),
"{name} is marked JSON text"
);
}
let text_at = |name: &str| -> String {
let column = batch.column(schema.index_of(name).unwrap());
column
.as_any()
.downcast_ref::<StringArray>()
.map(|array| array.value(0).to_owned())
.or_else(|| {
column
.as_any()
.downcast_ref::<StringViewArray>()
.map(|array| array.value(0).to_owned())
})
.unwrap_or_else(|| panic!("{name} is text: {:?}", column.data_type()))
};
assert_eq!(text_at("str_form"), "a");
assert_eq!(text_at("operator_text"), "a");
assert_eq!(text_at("operator_union"), "\"a\"");
assert_eq!(text_at("union_form"), "\"a\"");
let (mut encoder, _schema_frame) =
crate::core_service::ipc::FrameEncoder::start(Arc::clone(&schema), 1 << 20)
.expect("the schema frames");
encoder.begin(batch);
let frame = encoder
.next()
.expect("the batch frames")
.expect("one data frame");
assert!(encoder.next().unwrap().is_none());
let mut reader =
arrow::ipc::reader::StreamReader::try_new(std::io::Cursor::new(frame.arrow_ipc()), None)
.expect("the frame decodes");
let decoded = reader.next().unwrap().unwrap();
assert!(reader.next().is_none(), "one batch per frame");
let mut bytes = Vec::new();
{
let mut writer = arrow::json::ArrayWriter::new(&mut bytes);
writer.write_batches(&[&decoded]).unwrap();
writer.finish().unwrap();
}
let objects: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
let row = &objects[0];
assert_eq!(row["str_form"], "a");
assert_eq!(row["operator_text"], "a");
assert_eq!(row["operator_union"], "\"a\"");
assert_eq!(row["union_form"], "\"a\"");
}
#[tokio::test]
async fn a_json_expression_over_a_fleet_table_is_refused_by_the_realm() {
let harness = Harness::new();
let refusal = harness
.try_prepare(
"json-fleet-boundary",
"SELECT json_get_str(text, 'q') FROM summary",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
10,
)
.await
.expect_err("the table boundary refuses, not the function");
assert!(matches!(
refusal,
crate::core_resolution::CoreResolutionError::TableOutsideRealm
));
}
const TRACE_WALLET: &str = "0x1111111111111111111111111111111111111111";
const TRACE_PAYER: &str = "0x4444444444444444444444444444444444444444";
fn signed_payment_trace_segments() -> Vec<polyc_projector::artifact::EncodedSegment> {
use polyc_crypto::approval::{
ApprovalSigner, ReceiptPayload, WalletLinkLifecyclePayload, receipt_payload,
wallet_link_lifecycle_payload,
};
use polyc_crypto::signing_role::{HandoffRole, RoleSigner, RoleTrustSet, SubagentRole};
use polyc_eventlog_model::Event;
use polyc_projector::artifact::{EncodingBounds, RowSource, encode_trace_generation};
use polyc_proto::kinds;
let signer = ApprovalSigner::from_seed(290);
let (link, _, _) = wallet_link_lifecycle_payload(
&WalletLinkLifecyclePayload {
kind: kinds::WALLET_LINK_LIFECYCLE,
transition: "linked",
subject: "persona-linker",
wallet_address: TRACE_WALLET,
currency: "0x2222222222222222222222222222222222222222",
chain_id: "42431",
limit_base_units: "5000000",
limit_human: "5",
period_secs: "86400",
expiry_unix: "1780086400",
recipients: "0x3333333333333333333333333333333333333333",
conversation_id: "conv-wallet-link",
timestamp: "1780000000",
},
&signer,
);
let (receipt, _, _) = receipt_payload(
&ReceiptPayload {
kind: kinds::OUTBOUND_PAYMENT_RECEIPT,
reference: "0xtx-290",
amount: "0.05",
currency: "USD",
recipient: "0x5555555555555555555555555555555555555555",
method: "tempo",
timestamp: "2026-09-04T00:00:00Z",
tool_call_id: "call-pay",
approval_pos: "3",
approved_args_hash: "abc",
subject: "persona-payer",
payer_kind: "linked_wallet",
paying_account: TRACE_PAYER,
},
&signer,
);
let mut events = vec![
(1, Event::new(kinds::WALLET_LINK_LIFECYCLE.to_owned(), link)),
(
2,
Event::new(kinds::OUTBOUND_PAYMENT_RECEIPT.to_owned(), receipt),
),
];
let trust = polyc_facts::TraceTrust {
approval: vec![signer.public_key_bytes()],
handoff: RoleTrustSet::from_public_keys(vec![
RoleSigner::<HandoffRole>::from_seed(1).public_key_bytes(),
])
.expect("a handoff trust set"),
subagent: RoleTrustSet::from_public_keys(vec![
RoleSigner::<SubagentRole>::from_seed(2).public_key_bytes(),
])
.expect("a subagent trust set"),
question: vec![vec![9_u8; 32]],
};
let prepared = polyc_facts::prepare_conversation_core(&mut events, "conv-a");
let facts = polyc_facts::fold_conversation_trace(&events, &prepared.excised, &trust)
.expect("the trace fold");
encode_trace_generation(
polyc_projection::family::conversation_trace(),
&RowSource::new("conv-a".to_owned(), [1; 32]),
&facts,
EncodingBounds::default(),
)
.expect("the trace generation encodes")
}
async fn participant_rows(sql: &str) -> Vec<String> {
let harness = Harness::with_encoded_trace(
polyc_projection::family::conversation_trace(),
signed_payment_trace_segments(),
);
let prepared = harness
.prepare(
"trace-participant",
sql,
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
20,
)
.await;
harness
.collect_text(harness.visible_authority(), prepared)
.await
}
#[tokio::test]
async fn a_participant_cannot_name_the_participant_addresses() {
let harness = Harness::with_encoded_trace(
polyc_projection::family::conversation_trace(),
signed_payment_trace_segments(),
);
let refusal = harness
.try_prepare(
"trace-participant-addresses",
"SELECT address FROM trace_participant_addresses",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
20,
)
.await
.expect_err("a participant must not read the participant addresses");
assert!(
matches!(
refusal,
crate::core_resolution::CoreResolutionError::TableOutsideRealm
),
"refused for the wrong reason: {refusal:?}"
);
}
#[tokio::test]
async fn a_participant_reads_every_other_receipt_and_wallet_field() {
assert_eq!(
participant_rows(
"SELECT direction || '|' || reference || '|' || amount || '|' || currency \
|| '|' || recipient || '|' || method || '|' || tool_call_id \
|| '|' || CAST(approval_pos AS VARCHAR) || '|' || approved_args_hash \
|| '|' || subject || '|' || kind || '|' || payer_kind \
FROM trace_payment_receipts",
)
.await,
vec![
"outbound|0xtx-290|0.05|USD|0x5555555555555555555555555555555555555555|tempo|\
call-pay|3|abc|persona-payer|outbound_payment_receipt|linked_wallet"
.to_owned()
],
"a participant read keeps every receipt field"
);
assert_eq!(
participant_rows(
"SELECT transition || '|' || subject || '|' || currency \
|| '|' || CAST(chain_id AS VARCHAR) || '|' || CAST(limit_base_units AS VARCHAR) \
|| '|' || limit_human || '|' || CAST(period_secs AS VARCHAR) \
|| '|' || CAST(expiry_unix AS VARCHAR) || '|' || conversation_id \
|| '|' || CAST(timestamp AS VARCHAR) FROM trace_wallet_links",
)
.await,
vec![
"linked|persona-linker|0x2222222222222222222222222222222222222222|42431|5000000|5|\
86400|1780086400|conv-wallet-link|1780000000"
.to_owned()
],
"a participant read keeps every wallet-link field"
);
assert_eq!(
participant_rows("SELECT text_value FROM trace_lists WHERE list_kind = 'recipients'").await,
vec!["0x3333333333333333333333333333333333333333".to_owned()],
"the permitted recipients are counterparties and stay readable"
);
}
#[tokio::test]
async fn a_trace_generation_published_under_an_older_schema_is_refused() {
use polyc_projection::family::{CONVERSATION_TRACE, FamilyEntry, conversation_trace};
use polyc_projection::version::{SchemaVersion, VersionSet};
let shipped = conversation_trace();
let older = FamilyEntry::for_test(
CONVERSATION_TRACE,
VersionSet::new(SchemaVersion::new(1), shipped.versions().fact_model()),
shipped.tables(),
shipped.source(),
);
let harness = Harness::with_encoded_trace(older, signed_payment_trace_segments());
let refusal = harness
.try_prepare(
"trace-older-schema-generation",
"SELECT payer_kind FROM trace_payment_receipts",
QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
20,
)
.await
.expect_err("a generation published under schema version 1 must be refused");
assert!(
matches!(
refusal,
crate::core_resolution::CoreResolutionError::IncompatibleDescriptor(_)
),
"refused for the wrong reason: {refusal:?}"
);
}
async fn fleet_address_rows(
harness: &Harness,
scope: QueryScope,
) -> Vec<(u64, u64, String, String)> {
let prepared = harness
.prepare_with_parameters(
"trace-addresses",
polyc_query_model::statements::COMPOSITE_TRACE_ADDRESSES_SQL,
scope,
vec![CoreParameter::Utf8("conv-a".to_owned())],
)
.await;
let mut stream = harness
.fleet_authority()
.bind(prepared)
.await
.expect("binding succeeds")
.execute();
let mut rows = Vec::new();
while let Some(batch) = stream.next().await {
let batch = batch.expect("released batch");
let number = |name: &str| {
batch
.column_by_name(name)
.and_then(|column| column.as_any().downcast_ref::<UInt64Array>().cloned())
.unwrap_or_else(|| panic!("the statement returns a `{name}` number"))
};
let text = |name: &str| {
batch
.column_by_name(name)
.and_then(|column| column.as_any().downcast_ref::<StringArray>().cloned())
.unwrap_or_else(|| panic!("the statement returns a `{name}` text"))
};
let (position, ordinal) = (number("position"), number("ordinal"));
let (kind, address) = (text("step_kind"), text("address"));
for index in 0..batch.num_rows() {
rows.push((
position.value(index),
ordinal.value(index),
kind.value(index).to_owned(),
address.value(index).to_owned(),
));
}
}
rows
}
#[tokio::test]
async fn an_admin_fleet_read_returns_each_participant_address_exactly_once() {
let harness = Harness::with_encoded_trace(
polyc_projection::family::conversation_trace(),
signed_payment_trace_segments(),
);
let rows = fleet_address_rows(
&harness,
QueryScope::FleetConversation {
conversation: "a".to_owned(),
},
)
.await;
assert_eq!(
rows,
vec![
(1, 0, "wallet_link".to_owned(), TRACE_WALLET.to_owned()),
(2, 0, "payment_receipt".to_owned(), TRACE_PAYER.to_owned()),
]
);
}
#[tokio::test]
async fn an_admin_fleet_read_of_a_trace_with_no_address_step_returns_no_rows() {
use polyc_projector::artifact::{EncodingBounds, RowSource, encode_trace_generation};
let segments = encode_trace_generation(
polyc_projection::family::conversation_trace(),
&RowSource::new("conv-a".to_owned(), [1; 32]),
&polyc_facts::ConversationTraceFacts::default(),
EncodingBounds::default(),
)
.expect("a trace generation with no step encodes");
let harness =
Harness::with_encoded_trace(polyc_projection::family::conversation_trace(), segments);
let rows = fleet_address_rows(
&harness,
QueryScope::FleetConversation {
conversation: "a".to_owned(),
},
)
.await;
assert_eq!(rows, Vec::new());
}
#[tokio::test]
async fn a_conversation_bounded_fleet_read_ignores_other_conversations() {
let harness = Harness::with_encoded_trace(
polyc_projection::family::conversation_trace(),
signed_payment_trace_segments(),
);
harness.state.add_bare_conversation("conv-b");
let refusal = harness
.try_prepare(
"trace-addresses-fleet-wide",
"SELECT address FROM trace_participant_addresses",
QueryScope::Fleet,
20,
)
.await
.expect_err("a Fleet-wide read pins the conversation with no generation");
assert!(
matches!(
refusal,
crate::core_resolution::CoreResolutionError::MissingProjection(_)
),
"refused for the wrong reason: {refusal:?}"
);
let rows = fleet_address_rows(
&harness,
QueryScope::FleetConversation {
conversation: "a".to_owned(),
},
)
.await;
assert_eq!(
rows.len(),
2,
"the bounded read returns its own rows: {rows:?}"
);
}
#[tokio::test]
async fn a_conversation_bounded_fleet_read_refuses_persona_memory() {
let harness = Harness::with_encoded_trace(
polyc_projection::family::conversation_trace(),
signed_payment_trace_segments(),
);
let refusal = harness
.try_prepare(
"memory-under-fleet-conversation",
"SELECT fact_id FROM memory_portable_facts",
QueryScope::FleetConversation {
conversation: "a".to_owned(),
},
20,
)
.await
.expect_err("a one-conversation Fleet scope has no memory to read");
assert!(
matches!(
refusal,
crate::core_resolution::CoreResolutionError::FamilyOutsideRealm
),
"refused for the wrong reason: {refusal:?}"
);
}
#[tokio::test]
async fn a_conversation_bounded_fleet_read_refuses_every_other_family() {
let harness = Harness::with_encoded_trace(
polyc_projection::family::conversation_trace(),
signed_payment_trace_segments(),
);
for (query, sql) in [
("fleet-conversation-core", "SELECT turn_id FROM turns"),
(
"fleet-conversation-security",
"SELECT position FROM approval_details",
),
(
"fleet-conversation-admin",
"SELECT change_kind FROM credential_lifecycle",
),
] {
let refusal = harness
.try_prepare(
query,
sql,
QueryScope::FleetConversation {
conversation: "a".to_owned(),
},
20,
)
.await
.expect_err("the one-conversation Fleet scope reads the trace family alone");
assert!(
matches!(
refusal,
crate::core_resolution::CoreResolutionError::FamilyOutsideRealm
),
"`{sql}` refused for the wrong reason: {refusal:?}"
);
}
}