use std::collections::{BTreeMap, BTreeSet};
use std::sync::{Mutex, PoisonError};
use std::time::Duration;
use super::*;
use polyc_query_credential::session::MemorySources;
use polyc_state::cancel::CancellationToken;
use polyc_state::context::CallContext;
use polyc_state::deadline::{Deadline, MonotonicInstant};
use polyc_state::feed::{SourceCheckpoint, VersionedCheckpoint, VersionedSource};
use polyc_state::id::OperationFamily;
use polyc_state::immutable::{ContentReference, Generation, ObjectDescriptor, Retention};
use polyc_state::journal::{JournalAttestation, JournalDirectorySnapshotId};
use polyc_state::page::PageCompleteness;
use polyc_state::projection::artifact::{ExactObjectRef, ObjectNamespace};
use polyc_state::projection::{ProjectionGeneration, ProjectionHead, PublisherFence, PublisherId};
use polyc_state::query_audit::{QueryAuditWrite, memory::MemoryQueryAudit};
use polyc_state::revision::{CommitRoot, JournalPosition, JournalSource, PartitionIncarnation};
#[derive(Debug, Default)]
struct Counts {
snapshots: usize,
pages: usize,
releases: usize,
source_reads: usize,
resolutions: usize,
audit_writes: usize,
lifecycle: Vec<&'static str>,
operation_contexts: BTreeSet<usize>,
subcall_budgets: Vec<Duration>,
subcall_audiences: BTreeSet<String>,
}
#[derive(Debug)]
struct FakeState {
versioned: polyc_state::versioned::memory::MemoryVersionedState,
persona_memory: polyc_state::persona_memory::journal::MemoryPersonaJournal,
observation: polyc_state::observation::memory::MemoryObservationState,
sources: Mutex<BTreeMap<PartitionId, JournalSourceHead>>,
resolutions: Mutex<BTreeMap<ProjectionKey, ProjectionResolution>>,
snapshot: Mutex<Vec<PartitionId>>,
audit: MemoryQueryAudit,
intents: Mutex<Vec<BeginQueryAudit>>,
counts: Mutex<Counts>,
outage: Mutex<bool>,
refuse_audit: Mutex<bool>,
cross_permit: Mutex<bool>,
settlements: Mutex<usize>,
versioned_namespaces: Mutex<Option<Vec<String>>>,
}
impl FakeState {
fn new(partitions: &[&str]) -> Self {
let owner = OwnerId::new("projector");
let mut sources = BTreeMap::new();
let mut resolutions = BTreeMap::new();
for (ordinal, partition) in partitions.iter().enumerate() {
let partition = PartitionId::new(*partition);
let source = source(partition.clone(), u8::try_from(ordinal + 1).unwrap());
sources.insert(
partition.clone(),
JournalSourceHead::new(source.clone(), JournalPosition::new(20)),
);
for family in [conversation_core(), conversation_execution()] {
let key = ProjectionKey::new(FamilyId::new(family.family_str()), partition.clone());
resolutions.insert(
key,
ProjectionResolution::new(
ProjectionHead::Current(Box::new(manifest_for_family(
family.family_str(),
partition.clone(),
&source,
owner.clone(),
Classification::Confidential,
family.versions().schema().get(),
family.versions().fact_model().get(),
))),
None,
ProjectionGeneration::new(1),
),
);
}
}
Self {
sources: Mutex::new(sources),
resolutions: Mutex::new(resolutions),
snapshot: Mutex::new(
partitions
.iter()
.map(|partition| PartitionId::new(*partition))
.collect(),
),
audit: MemoryQueryAudit::new(),
versioned: polyc_state::versioned::memory::MemoryVersionedState::new(),
persona_memory: polyc_state::persona_memory::journal::MemoryPersonaJournal::new(
polyc_state::id::NamespaceId::new("polychrome"),
),
observation: polyc_state::observation::memory::MemoryObservationState::new(),
intents: Mutex::new(Vec::new()),
counts: Mutex::new(Counts::default()),
outage: Mutex::new(false),
refuse_audit: Mutex::new(false),
cross_permit: Mutex::new(false),
settlements: Mutex::new(0),
versioned_namespaces: Mutex::new(None),
}
}
fn lock<T>(mutex: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
mutex.lock().unwrap_or_else(PoisonError::into_inner)
}
fn set_resolution(&self, partition: &str, resolution: ProjectionResolution) {
self.set_family_resolution(conversation_core(), partition, resolution);
}
fn set_family_resolution(
&self,
family: FamilyEntry,
partition: &str,
resolution: ProjectionResolution,
) {
Self::lock(&self.resolutions).insert(
ProjectionKey::new(
FamilyId::new(family.family_str()),
PartitionId::new(partition),
),
resolution,
);
}
fn set_outage(&self) {
*Self::lock(&self.outage) = true;
}
fn cross_permit(&self) {
*Self::lock(&self.cross_permit) = true;
}
fn settlements(&self) -> usize {
*Self::lock(&self.settlements)
}
fn refuse_audit(&self) {
*Self::lock(&self.refuse_audit) = true;
}
fn source_read_count(&self) -> usize {
Self::lock(&self.counts).source_reads
}
fn resolved_partitions(&self) -> Vec<PartitionId> {
let intents = Self::lock(&self.intents);
intents
.last()
.map(|intent| {
intent
.source()
.pins()
.iter()
.filter_map(|pin| match pin {
SourcePin::Projected(pin) => Some(pin.manifest().key().source().clone()),
SourcePin::Journal(_) | SourcePin::Authoritative(_) => None,
})
.collect()
})
.unwrap_or_default()
}
fn last_source_pins(&self) -> Vec<SourcePin> {
Self::lock(&self.intents)
.last()
.map(|intent| intent.source().pins().to_vec())
.unwrap_or_default()
}
fn source_lifecycle(&self) -> Vec<&'static str> {
Self::lock(&self.counts).lifecycle.clone()
}
fn observe_operation(&self, operation: &CoreOperationContext) {
let address = std::ptr::from_ref(operation).addr();
let declared = operation
.declared()
.expect("the planning seam preflights each metadata call");
let mut counts = Self::lock(&self.counts);
counts.operation_contexts.insert(address);
counts.subcall_budgets.push(declared.budget);
counts
.subcall_audiences
.insert(declared.audience.as_str().to_owned());
}
}
#[async_trait]
impl CoreMetadataAuthority for FakeState {
async fn create_directory_snapshot(
&self,
operation: &CoreOperationContext,
) -> Result<JournalDirectorySnapshot, CoreResolutionError> {
self.observe_operation(operation);
let sequence = {
let mut counts = Self::lock(&self.counts);
counts.snapshots += 1;
counts.snapshots
};
let count = Self::lock(&self.snapshot).len();
Ok(JournalDirectorySnapshot::new(
JournalDirectorySnapshotId::new(format!("snapshot-{sequence}")),
u64::try_from(count).unwrap(),
))
}
async fn directory_page(
&self,
operation: &CoreOperationContext,
request: ListJournalDirectorySnapshot,
) -> Result<JournalDirectoryPage, CoreResolutionError> {
self.observe_operation(operation);
Self::lock(&self.counts).pages += 1;
let all = Self::lock(&self.snapshot);
let partitions = all
.iter()
.filter(|partition| request.start_after().is_none_or(|after| *partition > after))
.take(request.limit() as usize)
.cloned()
.collect::<Vec<_>>();
let consumed = request.start_after().map_or(0, |after| {
all.iter().filter(|partition| *partition <= after).count()
}) + partitions.len();
let truncated = consumed < all.len();
Ok(JournalDirectoryPage::new(
request.snapshot().clone(),
partitions.clone(),
truncated.then(|| partitions.last().unwrap().clone()),
if truncated {
PageCompleteness::Truncated
} else {
PageCompleteness::Complete
},
))
}
async fn release_directory_snapshot(
&self,
operation: &CoreOperationContext,
_request: ReleaseJournalDirectorySnapshot,
) -> Result<(), CoreResolutionError> {
self.observe_operation(operation);
Self::lock(&self.counts).releases += 1;
Ok(())
}
async fn source_head(
&self,
operation: &CoreOperationContext,
request: GetJournalSource,
) -> Result<Option<JournalSourceHead>, CoreResolutionError> {
self.observe_operation(operation);
Self::lock(&self.counts).source_reads += 1;
Self::lock(&self.counts).lifecycle.push("source");
if *Self::lock(&self.outage) {
return Err(StateError::Unavailable {
family: OperationFamily::new("journal"),
reach: polyc_state::error::OutageReach::NoDurableEffect,
}
.into());
}
Ok(Self::lock(&self.sources).get(request.partition()).cloned())
}
async fn resolve_manifest(
&self,
operation: &CoreOperationContext,
request: ResolveManifest,
) -> Result<ProjectionResolution, CoreResolutionError> {
self.observe_operation(operation);
Self::lock(&self.counts).resolutions += 1;
Self::lock(&self.counts).lifecycle.push("projection");
Ok(Self::lock(&self.resolutions)
.get(request.key())
.cloned()
.unwrap_or_else(|| {
ProjectionResolution::new(
ProjectionHead::Absent,
None,
ProjectionGeneration::ORIGIN,
)
}))
}
async fn begin_audit(
&self,
operation: &CoreOperationContext,
command: BeginQueryAudit,
) -> Result<BeginOutcome<CoreExecutionPermit>, CoreResolutionError> {
self.observe_operation(operation);
Self::lock(&self.counts).audit_writes += 1;
Self::lock(&self.counts).lifecycle.push("audit");
Self::lock(&self.intents).push(command.clone());
if *Self::lock(&self.refuse_audit) {
return Err(StateError::Unavailable {
family: polyc_state::query_audit::family(),
reach: polyc_state::error::OutageReach::PossiblyApplied,
}
.into());
}
let command = if *Self::lock(&self.cross_permit) {
BeginQueryAudit::new(
QueryId::new("another-query"),
command.namespace().clone(),
command.requester().clone(),
command.shape(),
command.source().clone(),
command.metadata().digest(),
command.metadata().envelope().clone(),
)
} else {
command
};
match self.audit.begin(command, &live_context())? {
BeginOutcome::Granted(permit) => {
Ok(BeginOutcome::Granted(CoreExecutionPermit::from(permit)))
}
BeginOutcome::AlreadyRecorded(receipt) => Ok(BeginOutcome::AlreadyRecorded(receipt)),
}
}
async fn complete_audit(
&self,
_operation: &CoreCompletionContext,
_command: &CoreCompletionCommand,
) -> Result<Receipt, CoreResolutionError> {
*Self::lock(&self.settlements) += 1;
Err(CoreResolutionError::InvalidComposition)
}
async fn completion_receipt(
&self,
_operation: &CoreCompletionContext,
_command: &CoreCompletionCommand,
) -> Result<Option<Receipt>, CoreResolutionError> {
*Self::lock(&self.settlements) += 1;
Err(CoreResolutionError::InvalidComposition)
}
async fn versioned_source_head(
&self,
_operation: &CoreOperationContext,
scope: &polyc_state::command::CommandScope,
) -> Result<polyc_state::versioned::VersionedSourceHead, CoreResolutionError> {
use polyc_state::versioned::VersionedRead as _;
let context = polyc_state::context::CallContext::new(
polyc_state::deadline::Deadline::after(
polyc_state::deadline::MonotonicInstant::ORIGIN,
std::time::Duration::from_secs(30),
),
polyc_state::cancel::CancellationToken::new(),
);
Ok(self.versioned.source_head(
polyc_state::versioned::GetVersionedSourceHead::new(scope.clone()),
&context,
)?)
}
async fn persona_memory_source_head(
&self,
_operation: &CoreOperationContext,
partition: &polyc_state::persona_memory::journal::MemoryJournalPartition,
) -> Result<polyc_state::persona_memory::journal::PersonaMemorySourceHead, CoreResolutionError>
{
use polyc_state::persona_memory::journal::PersonaMemoryHistory as _;
let context = polyc_state::context::CallContext::new(
polyc_state::deadline::Deadline::after(
polyc_state::deadline::MonotonicInstant::ORIGIN,
std::time::Duration::from_secs(30),
),
polyc_state::cancel::CancellationToken::new(),
);
let page = self.persona_memory.read_history(
partition,
None,
polyc_state::persona_memory::journal::MemoryReplayRange::new(
0,
1,
polyc_state::persona_memory::journal::MAX_PAYLOAD_BYTES,
),
&context,
)?;
Ok(
polyc_state::persona_memory::journal::PersonaMemorySourceHead::new(
page.incarnation(),
page.head(),
),
)
}
async fn persona_memory_directory_page(
&self,
_operation: &CoreOperationContext,
after: Option<&str>,
limit: u32,
) -> Result<polyc_state::persona_memory::journal::MemoryLineagePage, CoreResolutionError> {
use polyc_state::persona_memory::journal::PersonaMemoryHistory as _;
let context = polyc_state::context::CallContext::new(
polyc_state::deadline::Deadline::after(
polyc_state::deadline::MonotonicInstant::ORIGIN,
std::time::Duration::from_secs(30),
),
polyc_state::cancel::CancellationToken::new(),
);
Ok(self.persona_memory.list_lineage(after, limit, &context)?)
}
async fn observed_head(
&self,
_operation: &CoreOperationContext,
collection: &polyc_state::observation::CollectionId,
) -> Result<Option<polyc_state::observation::ObservationHead>, CoreResolutionError> {
use polyc_state::observation::ObservationRead as _;
Ok(self.observation.latest(collection)?)
}
async fn observed_collections(
&self,
_operation: &CoreOperationContext,
kind: polyc_state::observation::CollectionKind,
) -> Result<polyc_state::observation::ObservedCollectionListing, CoreResolutionError> {
use polyc_state::observation::ObservationRead as _;
Ok(self.observation.collections(kind)?)
}
async fn create_versioned_directory_snapshot(
&self,
_operation: &CoreOperationContext,
family: &str,
) -> Result<polyc_state::versioned::VersionedDirectorySnapshot, CoreResolutionError> {
use polyc_state::versioned::VersionedRead as _;
let override_namespaces = Self::lock(&self.versioned_namespaces).clone();
if let Some(namespaces) = override_namespaces {
return Ok(polyc_state::versioned::VersionedDirectorySnapshot::new(
polyc_state::versioned::VersionedDirectorySnapshotId::new("fixture-capture"),
namespaces.len() as u64,
self.versioned.lineage(),
));
}
Ok(self.versioned.create_directory_snapshot(
polyc_state::versioned::CreateVersionedDirectorySnapshot::new(family),
&fake_versioned_context(),
)?)
}
async fn versioned_directory_page(
&self,
_operation: &CoreOperationContext,
request: &polyc_state::versioned::ListVersionedDirectorySnapshot,
) -> Result<polyc_state::versioned::VersionedDirectoryPage, CoreResolutionError> {
use polyc_state::versioned::VersionedRead as _;
let override_namespaces = Self::lock(&self.versioned_namespaces).clone();
if let Some(namespaces) = override_namespaces {
return Ok(polyc_state::versioned::VersionedDirectoryPage::new(
request.snapshot().clone(),
namespaces,
None,
PageCompleteness::Complete,
self.versioned.lineage(),
));
}
Ok(self
.versioned
.directory_page(request.clone(), &fake_versioned_context())?)
}
async fn release_versioned_directory_snapshot(
&self,
_operation: &CoreOperationContext,
snapshot: &polyc_state::versioned::VersionedDirectorySnapshotId,
) -> Result<(), CoreResolutionError> {
use polyc_state::versioned::VersionedRead as _;
if Self::lock(&self.versioned_namespaces).is_some() {
return Ok(());
}
Ok(self.versioned.release_directory_snapshot(
polyc_state::versioned::ReleaseVersionedDirectorySnapshot::new(snapshot.clone()),
&fake_versioned_context(),
)?)
}
}
fn fake_versioned_context() -> polyc_state::context::CallContext {
polyc_state::context::CallContext::new(
polyc_state::deadline::Deadline::after(
polyc_state::deadline::MonotonicInstant::ORIGIN,
std::time::Duration::from_secs(30),
),
polyc_state::cancel::CancellationToken::new(),
)
}
#[test]
fn every_namespace_pins_its_own_versioned_descriptor() {
let lineage = PartitionIncarnation::from_bytes([5; PartitionIncarnation::LEN]);
let namespaces = ["tenant-a", "tenant-b", "polychrome"];
let authority = polyc_projection::family::AuthorityFamily::Credentials;
let pins: Vec<_> = namespaces
.iter()
.map(|namespace| {
super::versioned_pin("credential-lifecycle/v1", authority, namespace, lineage)
})
.collect();
let distinct: std::collections::BTreeSet<&str> =
pins.iter().map(|(key, _)| key.source().as_str()).collect();
assert_eq!(
distinct.len(),
namespaces.len(),
"two tenants shared one pin: {pins:?}"
);
let aggregate = super::state_authority(authority).aggregate();
for ((key, source), namespace) in pins.iter().zip(namespaces) {
assert_eq!(
polyc_state::feed::parse_projection_partition(key.source().as_str()),
Some((aggregate, namespace)),
"a pin must name the pair it was built from"
);
assert_eq!(
key.source(),
source.projection_partition(),
"the key and its evidence must name one partition"
);
}
}
fn versioned_family() -> FamilyEntry {
FamilyEntry::for_test(
"credential-lifecycle/v1",
conversation_core().versions(),
conversation_core().tables(),
polyc_projection::family::SourceKind::Versioned {
family: polyc_projection::family::AuthorityFamily::Credentials,
},
)
}
fn operation() -> CoreOperationContext {
CoreOperationContext::for_test(Duration::from_secs(30))
}
fn versioned_fixture(namespace: &str) -> (Arc<FakeState>, ProjectionKey, VersionedSource) {
let state = Arc::new(FakeState::new(&[]));
*FakeState::lock(&state.versioned_namespaces) = Some(vec![namespace.to_owned()]);
let (key, source) = super::versioned_pin(
"credential-lifecycle/v1",
polyc_projection::family::AuthorityFamily::Credentials,
namespace,
state.versioned.lineage(),
);
(state, key, source)
}
#[allow(clippy::too_many_arguments)]
fn versioned_manifest(
key: &ProjectionKey,
source: &VersionedSource,
owner: OwnerId,
classification: Classification,
schema_version: u32,
fact_version: u32,
) -> ProjectionManifest {
let reference =
ContentReference::try_new(format!("projection/{}/manifest", key.source().as_str()))
.unwrap();
ProjectionManifest::new(
key.clone(),
ProjectionGeneration::new(1),
polyc_state::feed::SourceEvidence::Versioned(
VersionedCheckpoint::try_new(
source.clone(),
JournalPosition::new(20),
ContentDigest::from_bytes([6; ContentDigest::LEN]),
)
.unwrap(),
),
schema_version,
fact_version,
ObjectDescriptor::new(
key.object(),
Generation::new(1),
ContentDigest::from_bytes([5; ContentDigest::LEN]),
owner,
classification,
Retention::For(Duration::from_mins(1)),
128,
reference.clone(),
),
ExactObjectRef::try_new(
ObjectNamespace::try_new("conversation-visible").unwrap(),
reference,
11,
)
.unwrap(),
PublisherId::new("projector-a"),
PublisherFence::new(key.clone(), source.incarnation(), 1),
)
}
fn publish(state: &FakeState, key: &ProjectionKey, manifest: ProjectionManifest) {
FakeState::lock(&state.resolutions).insert(
key.clone(),
ProjectionResolution::new(
ProjectionHead::Current(Box::new(manifest)),
None,
ProjectionGeneration::new(1),
),
);
}
#[derive(Clone, Copy, Debug)]
enum Perturbed {
OtherTenant,
Schema,
FactModel,
Owner,
Classification,
}
impl Perturbed {
const ALL: [Self; 5] = [
Self::OtherTenant,
Self::Schema,
Self::FactModel,
Self::Owner,
Self::Classification,
];
fn manifest(self, key: &ProjectionKey, source: &VersionedSource) -> ProjectionManifest {
let owner = OwnerId::new("projector");
let accepted = versioned_family().versions();
let (schema, fact) = (accepted.schema().get(), accepted.fact_model().get());
match self {
Self::OtherTenant => {
let (other_key, other_source) = super::versioned_pin(
"credential-lifecycle/v1",
polyc_projection::family::AuthorityFamily::Credentials,
"tenant-b",
source.incarnation(),
);
versioned_manifest(
&other_key,
&other_source,
owner,
Classification::Confidential,
schema,
fact,
)
}
Self::Schema => versioned_manifest(
key,
source,
owner,
Classification::Confidential,
schema + 1,
fact,
),
Self::FactModel => versioned_manifest(
key,
source,
owner,
Classification::Confidential,
schema,
fact + 1,
),
Self::Owner => versioned_manifest(
key,
source,
OwnerId::new("somebody-else"),
Classification::Confidential,
schema,
fact,
),
Self::Classification => {
versioned_manifest(key, source, owner, Classification::Public, schema, fact)
}
}
}
}
#[tokio::test]
async fn a_versioned_generation_binding_every_checked_field_is_admitted() {
let family = versioned_family();
let versions = family.versions();
let (state, key, source) = versioned_fixture("tenant-a");
publish(
&state,
&key,
versioned_manifest(
&key,
&source,
OwnerId::new("projector"),
Classification::Confidential,
versions.schema().get(),
versions.fact_model().get(),
),
);
let resolved = authority(state)
.resolve_versioned_family(family, &operation())
.await
.expect("a manifest binding every checked field is admitted");
assert_eq!(resolved.len(), 1, "the tenant's one generation is planned");
assert_eq!(resolved[0].key(), &key);
}
#[tokio::test]
async fn a_versioned_generation_missing_any_checked_field_is_refused() {
let family = versioned_family();
for perturbed in Perturbed::ALL {
let (state, key, source) = versioned_fixture("tenant-a");
publish(&state, &key, perturbed.manifest(&key, &source));
let refusal = authority(state)
.resolve_versioned_family(family, &operation())
.await
.expect_err(&format!("{perturbed:?} must be refused"));
assert!(
matches!(refusal, CoreResolutionError::IncompatibleDescriptor(_)),
"{perturbed:?} was refused for the wrong reason: {refusal:?}"
);
}
}
#[tokio::test]
async fn a_versioned_generation_carrying_journal_evidence_is_refused() {
let family = versioned_family();
let (state, key, source) = versioned_fixture("tenant-a");
let journal = JournalSource::new(
PartitionId::new(key.source().as_str()),
source.incarnation(),
);
publish(
&state,
&key,
manifest_for_family(
"credential-lifecycle/v1",
PartitionId::new(key.source().as_str()),
&journal,
OwnerId::new("projector"),
Classification::Confidential,
1,
1,
),
);
let refusal = authority(state)
.resolve_versioned_family(family, &operation())
.await
.expect_err("journal evidence under a Versioned family must be refused");
assert!(
matches!(refusal, CoreResolutionError::IncompatibleDescriptor(_)),
"refused for the wrong reason: {refusal:?}"
);
}
fn live_context() -> CallContext {
CallContext::new(
Deadline::at(MonotonicInstant::from_nanos(u64::MAX)),
CancellationToken::new(),
)
}
fn source(partition: PartitionId, byte: u8) -> JournalSource {
JournalSource::new(
partition,
PartitionIncarnation::from_bytes([byte; PartitionIncarnation::LEN]),
)
}
fn manifest(
partition: PartitionId,
source: &JournalSource,
owner: OwnerId,
classification: Classification,
schema_version: u32,
fact_version: u32,
) -> ProjectionManifest {
manifest_for_family(
"conversation-core/v1",
partition,
source,
owner,
classification,
schema_version,
fact_version,
)
}
#[allow(clippy::too_many_arguments)]
fn manifest_for_family(
family: &str,
partition: PartitionId,
source: &JournalSource,
owner: OwnerId,
classification: Classification,
schema_version: u32,
fact_version: u32,
) -> ProjectionManifest {
let key = ProjectionKey::new(FamilyId::new(family), partition);
let generation = ProjectionGeneration::new(1);
let reference =
ContentReference::try_new(format!("projection/{}/manifest", key.source().as_str()))
.unwrap();
ProjectionManifest::new(
key.clone(),
generation,
polyc_state::feed::SourceEvidence::Journal(
SourceCheckpoint::try_new(
source.clone(),
JournalPosition::new(10),
JournalPosition::new(20),
20,
JournalAttestation::new(
CommitRoot::from_bytes([7; CommitRoot::LEN]),
22,
vec![8; 64],
vec![9; 32],
),
)
.unwrap(),
),
schema_version,
fact_version,
ObjectDescriptor::new(
key.object(),
Generation::new(1),
ContentDigest::from_bytes([5; ContentDigest::LEN]),
owner,
classification,
Retention::For(Duration::from_mins(1)),
128,
reference.clone(),
),
ExactObjectRef::try_new(
ObjectNamespace::try_new("conversation-visible").unwrap(),
reference,
11,
)
.unwrap(),
PublisherId::new("projector-a"),
PublisherFence::new(key, source.incarnation(), 1),
)
}
fn request(sql: &str) -> CoreQueryRequest {
CoreQueryRequest {
sql: sql.to_owned(),
parameters: Vec::new(),
consistency: CoreConsistency::Projected,
requested_bounds: CoreRequestedBounds::unbounded(),
}
}
fn audit(query: &str) -> CoreAuditContext {
audit_with_declared(
query,
&DeclaredCall::live(Audience::new("state"), Duration::MAX),
)
}
fn audit_with_declared(query: &str, declared: &DeclaredCall) -> CoreAuditContext {
audit_for_bounds(query, declared, CoreRequestedBounds::unbounded())
}
fn audit_for_bounds(
query: &str,
declared: &DeclaredCall,
requested: CoreRequestedBounds,
) -> CoreAuditContext {
audit_for_policy(
query,
declared,
&QueryLimits::default(),
ProjectedCorePolicy::default(),
requested,
)
}
fn audit_for_policy(
query: &str,
declared: &DeclaredCall,
limits: &QueryLimits,
policy: ProjectedCorePolicy,
requested: CoreRequestedBounds,
) -> CoreAuditContext {
CoreAuditContext::from_scoped(
QueryId::new(query),
RequesterId::new("persona:verified-a"),
declared,
EffectiveCoreBounds::mint(limits, policy, requested).unwrap(),
)
}
fn authority(state: Arc<FakeState>) -> CorePlanningAuthority {
authority_with_policy(state, ProjectedCorePolicy::default())
}
fn authority_with_policy(
state: Arc<FakeState>,
policy: ProjectedCorePolicy,
) -> CorePlanningAuthority {
CorePlanningAuthority::new(
NamespaceId::new("tenant-a"),
OwnerId::new("projector"),
policy,
state,
)
.unwrap()
}
#[tokio::test]
async fn schema_only_planning_derives_dependencies_without_a_physical_scan() {
let compiler = CatalogCompiler::new(SessionContext::new().state());
let logical = compiler
.compile(
"SELECT turns.turn_id, messages.text FROM turns JOIN messages USING (turn_id)",
&[],
false,
)
.await
.expect("schema-only plan");
assert_eq!(
logical.dependencies(),
&[CoreTable::Turns, CoreTable::Messages]
);
assert_eq!(compiler.physical_scan_count(), 0);
}
#[tokio::test]
async fn schema_only_planning_closes_mixed_dependencies_without_a_physical_scan() {
let compiler = CatalogCompiler::new(SessionContext::new().state());
let logical = compiler
.compile(
"SELECT messages.text, tool_calls.name \
FROM messages JOIN tool_calls USING (partition, turn_id)",
&[],
false,
)
.await
.expect("mixed schema-only plan");
assert_eq!(
logical.dependencies(),
&[CoreTable::Messages, CoreTable::ToolCalls]
);
assert_eq!(compiler.physical_scan_count(), 0);
}
#[tokio::test]
async fn mixed_sources_are_complete_before_the_durable_intent() {
let state = Arc::new(FakeState::new(&["conv-a"]));
let compiler = CatalogCompiler::new(SessionContext::new().state());
let outcome = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit("q-mixed"),
request(
"SELECT messages.text, tool_calls.name \
FROM messages JOIN tool_calls USING (partition, turn_id)",
),
)
.await
.expect("mixed plan");
assert!(matches!(outcome, CorePlanOutcome::Granted(_)));
let pins = state.last_source_pins();
assert_eq!(pins.len(), 2);
assert!(matches!(pins[0], SourcePin::Projected(_)));
assert!(matches!(pins[1], SourcePin::Projected(_)));
assert_eq!(
state.source_lifecycle(),
["source", "projection", "projection", "audit"]
);
assert_eq!(compiler.physical_scan_count(), 0);
}
#[tokio::test]
async fn execution_only_planning_resolves_only_its_family_manifest() {
let state = Arc::new(FakeState::new(&["conv-a"]));
let compiler = CatalogCompiler::new(SessionContext::new().state());
let outcome = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit("q-execution"),
request("SELECT name FROM tool_calls"),
)
.await
.expect("execution-only plan");
assert!(matches!(outcome, CorePlanOutcome::Granted(_)));
let pins = state.last_source_pins();
assert_eq!(pins.len(), 1);
assert!(matches!(pins[0], SourcePin::Projected(_)));
assert_eq!(state.source_lifecycle(), ["source", "projection", "audit"]);
assert_eq!(compiler.physical_scan_count(), 0);
}
#[tokio::test]
async fn mixed_pin_bound_refuses_before_source_resolution() {
let conversations = (0..17)
.map(|index| format!("c-{index}"))
.collect::<Vec<_>>();
let partitions = conversations
.iter()
.map(|conversation| format!("conv-{conversation}"))
.collect::<Vec<_>>();
let partition_refs = partitions.iter().map(String::as_str).collect::<Vec<_>>();
let state = Arc::new(FakeState::new(&partition_refs));
let compiler = CatalogCompiler::new(SessionContext::new().state());
let error = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations,
memory: MemorySources::default(),
},
false,
audit("q-mixed-bound"),
request(
"SELECT messages.text, tool_calls.name \
FROM messages JOIN tool_calls USING (partition, turn_id)",
),
)
.await
.expect_err("two pins per partition exceed the complete source bound");
assert!(matches!(
error,
CoreResolutionError::State(StateError::BoundsExceeded { requested: 34, .. })
));
assert!(state.source_lifecycle().is_empty());
}
#[tokio::test]
async fn a_missing_core_projection_never_falls_back_to_the_journal_pin() {
let state = Arc::new(FakeState::new(&["conv-a"]));
state.set_resolution(
"conv-a",
ProjectionResolution::new(ProjectionHead::Absent, None, ProjectionGeneration::ORIGIN),
);
let compiler = CatalogCompiler::new(SessionContext::new().state());
let error = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit("q-no-core-fallback"),
request(
"SELECT messages.text, tool_calls.name \
FROM messages JOIN tool_calls USING (partition, turn_id)",
),
)
.await
.expect_err("the projected dependency stays mandatory");
assert!(matches!(error, CoreResolutionError::MissingProjection(_)));
assert_eq!(state.source_lifecycle(), ["source", "projection"]);
assert!(state.last_source_pins().is_empty());
}
#[tokio::test]
async fn schema_only_planning_cannot_see_an_ambient_catalog_table() {
let context = SessionContext::new();
let ambient_scans = Arc::new(AtomicUsize::new(0));
context
.register_table(
"ambient_secret",
Arc::new(SchemaOnlyTable {
schema: Arc::new(Schema::new(vec![Field::new(
"secret",
DataType::Utf8,
false,
)])),
scans: ambient_scans,
}),
)
.unwrap();
let compiler = CatalogCompiler::new(context.state());
let error = compiler
.compile("SELECT secret FROM ambient_secret", &[], false)
.await
.expect_err("the compiler installs a fresh closed catalog");
assert!(matches!(error, CoreResolutionError::DataFusion(_)));
assert_eq!(compiler.physical_scan_count(), 0);
}
#[tokio::test]
async fn sql_predicates_cannot_widen_the_authority_scope() {
let state = Arc::new(FakeState::new(&["conv-a", "conv-b"]));
let compiler = CatalogCompiler::new(SessionContext::new().state());
let outcome = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit("q-scope"),
request("SELECT turn_id FROM turns WHERE partition = 'conv-b'"),
)
.await
.expect("authority scope plans");
assert!(matches!(outcome, CorePlanOutcome::Granted(_)));
assert_eq!(
state.resolved_partitions(),
vec![PartitionId::new("conv-a")]
);
assert_eq!(compiler.physical_scan_count(), 0);
}
#[tokio::test]
async fn exact_intent_replay_returns_evidence_without_a_second_permit() {
let state = Arc::new(FakeState::new(&["conv-a"]));
let compiler = CatalogCompiler::new(SessionContext::new().state());
let planning = authority(Arc::clone(&state));
let scope = QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
};
let first = planning
.plan(
&compiler,
&QueryLimits::default(),
&scope,
false,
audit("q-replay"),
request("SELECT * FROM turns"),
)
.await
.expect("new intent");
let replay = planning
.plan(
&compiler,
&QueryLimits::default(),
&scope,
false,
audit("q-replay"),
request("SELECT * FROM turns"),
)
.await
.expect("exact replay");
assert!(matches!(first, CorePlanOutcome::Granted(_)));
assert!(matches!(replay, CorePlanOutcome::AlreadyRecorded(_)));
assert_eq!(compiler.physical_scan_count(), 0);
}
#[tokio::test]
async fn bounds_and_freshness_refuse_before_source_reads_or_scans() {
let state = Arc::new(FakeState::new(&[]));
let compiler = CatalogCompiler::new(SessionContext::new().state());
let planning = authority(Arc::clone(&state));
let too_many = (0..=MAX_SOURCE_PINS)
.map(|index| format!("c-{index}"))
.collect();
let bound = planning
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: too_many,
memory: MemorySources::default(),
},
false,
audit("q-bound"),
request("SELECT * FROM messages"),
)
.await
.expect_err("the source vector is never truncated");
assert!(matches!(
bound,
CoreResolutionError::State(StateError::BoundsExceeded { .. })
));
let mut fresh = request("SELECT * FROM turns");
fresh.consistency = CoreConsistency::RequireProjectedThrough(JournalPosition::new(41));
let freshness = planning
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit("q-fresh"),
fresh,
)
.await
.expect_err("freshness is recognized before State reads");
assert!(matches!(
freshness,
CoreResolutionError::FreshnessUnsupported { position }
if position == JournalPosition::new(41)
));
assert_eq!(state.source_read_count(), 0);
assert_eq!(compiler.physical_scan_count(), 0);
let duplicate = planning
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned(), "a".to_owned()],
memory: MemorySources::default(),
},
false,
audit("q-duplicate"),
request("SELECT * FROM turns"),
)
.await
.expect_err("one source cannot appear twice");
assert!(matches!(duplicate, CoreResolutionError::DuplicatePartition));
let empty_conversation = planning
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec![String::new()],
memory: MemorySources::default(),
},
false,
audit("q-empty-conversation"),
request("SELECT * FROM turns"),
)
.await
.expect_err("an empty conversation cannot become conv-");
assert!(matches!(
empty_conversation,
CoreResolutionError::EmptyConversationIdentity
));
assert_eq!(state.source_read_count(), 0);
}
#[tokio::test]
async fn parameter_mismatch_refuses_before_state_metadata() {
let state = Arc::new(FakeState::new(&["conv-a"]));
let compiler = CatalogCompiler::new(SessionContext::new().state());
let planning = authority(Arc::clone(&state));
for (index, request) in [
request("SELECT * FROM turns WHERE turn_id = $1"),
request("SELECT * FROM turns"),
]
.into_iter()
.enumerate()
{
let mut request = request;
if index == 1 {
request.parameters.push(CoreParameter::UInt64(1));
}
let error = planning
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit(&format!("q-param-mismatch-{index}")),
request,
)
.await
.expect_err("missing, surplus, and invalid typed parameters fail in planning");
assert!(matches!(error, CoreResolutionError::ParameterMismatch));
}
let mut wrong_type = request("SELECT * FROM turns WHERE source_position = $1");
wrong_type
.parameters
.push(CoreParameter::Utf8("not-a-number".to_owned()));
let error = planning
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit("q-param-type"),
wrong_type,
)
.await
.expect_err("a typed value must satisfy the planned column type");
assert!(matches!(error, CoreResolutionError::DataFusion(_)));
assert_eq!(state.source_read_count(), 0);
assert_eq!(compiler.physical_scan_count(), 0);
}
#[tokio::test]
async fn typed_parameter_values_and_variants_bind_the_audit_shape() {
let compiler = CatalogCompiler::new(SessionContext::new().state());
let cases = [
CoreParameter::Utf8("a".to_owned()),
CoreParameter::Utf8("b".to_owned()),
CoreParameter::UInt64(1),
CoreParameter::Boolean(true),
CoreParameter::Null,
];
let mut shapes = BTreeSet::new();
let mut encodings = BTreeSet::new();
for parameter in cases {
let state = Arc::new(FakeState::new(&["conv-a"]));
let mut parameterized = request("SELECT * FROM turns WHERE $1 = $1");
parameterized.parameters.push(parameter.clone());
let outcome = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit("q-typed"),
parameterized,
)
.await
.expect("closed typed parameter plans");
let CorePlanOutcome::Granted(prepared) = outcome else {
panic!("a distinct query state records a new intent");
};
assert!(
prepared
.compiled
.plan
.get_parameter_names()
.unwrap()
.is_empty(),
"the retained plan owns bound literals, not placeholders or raw parameters"
);
let mut canonical = Vec::new();
push_parameters(&mut canonical, &[parameter]);
encodings.insert(canonical);
shapes.insert(FakeState::lock(&state.intents)[0].shape());
}
assert_eq!(shapes.len(), 5);
assert_eq!(encodings.len(), 5);
assert_eq!(compiler.physical_scan_count(), 0);
}
#[tokio::test]
async fn expired_or_cancelled_operation_stops_before_metadata_reads() {
let declarations = [
DeclaredCall::bounded(Audience::new("state"), Duration::ZERO),
DeclaredCall::live(Audience::new("state"), Duration::MAX).withdrawn(),
];
for (index, declared) in declarations.iter().enumerate() {
let state = Arc::new(FakeState::new(&["conv-a"]));
let compiler = CatalogCompiler::new(SessionContext::new().state());
let error = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit_with_declared(&format!("q-context-{index}"), declared),
request("SELECT * FROM turns"),
)
.await
.expect_err("spent operation context cannot reach State metadata");
assert!(matches!(
error,
CoreResolutionError::State(
StateError::DeadlineExpired { .. } | StateError::Cancelled { .. }
)
));
let counts = {
let counts = FakeState::lock(&state.counts);
(
counts.snapshots,
counts.pages,
counts.releases,
counts.source_reads,
counts.resolutions,
counts.audit_writes,
)
};
assert_eq!(counts, (0, 0, 0, 0, 0, 0));
assert_eq!(compiler.physical_scan_count(), 0);
}
}
#[tokio::test]
async fn effective_bounds_are_stable_clamped_and_permit_owned() {
let default_state = Arc::new(FakeState::new(&["conv-a"]));
let oversized_state = Arc::new(FakeState::new(&["conv-a"]));
let lower_state = Arc::new(FakeState::new(&["conv-a"]));
let compiler = CatalogCompiler::new(SessionContext::new().state());
let declared = DeclaredCall::live(Audience::new("query"), Duration::from_mins(2));
let default_limits = QueryLimits::default();
let mut oversized = CoreRequestedBounds::unbounded();
oversized.timeout = Duration::from_mins(1);
oversized.rows = u64::try_from(default_limits.row_cap).unwrap() * 2;
oversized.result_release_bytes = DEFAULT_CORE_RESULT_RELEASE_BYTES * 2;
oversized.response_frame_bytes = DEFAULT_CORE_RESPONSE_FRAME_BYTES * 2;
let mut lower = CoreRequestedBounds::unbounded();
lower.rows = 100;
let default = authority(Arc::clone(&default_state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit_for_bounds("q-bounds", &declared, CoreRequestedBounds::unbounded()),
request("SELECT * FROM turns"),
)
.await
.unwrap();
let mut oversized_request = request("SELECT * FROM turns");
oversized_request.requested_bounds = oversized;
authority(Arc::clone(&oversized_state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit_for_bounds("q-bounds", &declared, oversized),
oversized_request,
)
.await
.unwrap();
let mut lower_request = request("SELECT * FROM turns");
lower_request.requested_bounds = lower;
authority(Arc::clone(&lower_state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit_for_bounds("q-bounds", &declared, lower),
lower_request,
)
.await
.unwrap();
let default_intent = FakeState::lock(&default_state.intents)[0].clone();
let oversized_intent = FakeState::lock(&oversized_state.intents)[0].clone();
let lower_intent = FakeState::lock(&lower_state.intents)[0].clone();
assert_eq!(default_intent.shape(), oversized_intent.shape());
assert_ne!(default_intent.shape(), lower_intent.shape());
let CorePlanOutcome::Granted(prepared) = default else {
panic!("a new bounded intent grants one permit");
};
assert_eq!(
prepared.bounds.rows,
u64::try_from(default_limits.row_cap).unwrap()
);
assert_eq!(prepared.bounds.timeout, default_limits.timeout);
assert_eq!(
prepared.bounds.result_release_bytes,
DEFAULT_CORE_RESULT_RELEASE_BYTES
);
assert_eq!(
prepared.bounds.artifact_file_bytes,
DEFAULT_CORE_ARTIFACT_FILE_BYTES
);
}
#[tokio::test]
async fn scoped_limits_and_deployment_policy_cannot_be_widened_or_crossed() {
let state = Arc::new(FakeState::new(&["conv-a"]));
let compiler = CatalogCompiler::new(SessionContext::new().state());
let limits = QueryLimits {
timeout: Duration::from_secs(7),
row_cap: 23,
..QueryLimits::default()
};
let policy = ProjectedCorePolicy::try_from(ProjectedCorePolicyInput {
result_release_bytes: 8 * 1024,
response_frame_bytes: 1024,
manifest_bytes: 4096,
artifact_file_bytes: 8192,
artifact_range_bytes: 2048,
source_decode_bytes: 16 * 1024,
})
.unwrap();
let planning = authority_with_policy(Arc::clone(&state), policy);
let effective = planning
.effective_bounds(&limits, CoreRequestedBounds::unbounded())
.unwrap();
assert_eq!(effective.timeout, limits.timeout);
assert_eq!(effective.rows, 23);
assert_eq!(effective.result_release_bytes, 8 * 1024);
assert_eq!(effective.response_frame_bytes, 1024);
assert_eq!(effective.manifest_bytes, 4096);
assert_eq!(effective.artifact_file_bytes, 8192);
assert_eq!(effective.artifact_range_bytes, 2048);
let crossed = planning
.plan(
&compiler,
&limits,
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit_for_policy(
"q-crossed-policy",
&DeclaredCall::live(Audience::new("query"), Duration::MAX),
&limits,
ProjectedCorePolicy::default(),
CoreRequestedBounds::unbounded(),
),
request("SELECT * FROM turns"),
)
.await
.expect_err("bounds minted under another deployment policy fail closed");
assert!(matches!(crossed, CoreResolutionError::InvalidBounds));
let wider_limits = QueryLimits {
timeout: Duration::from_mins(1),
row_cap: 1000,
..QueryLimits::default()
};
let widened = planning
.plan(
&compiler,
&limits,
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit_for_policy(
"q-crossed-limits",
&DeclaredCall::live(Audience::new("query"), Duration::MAX),
&wider_limits,
policy,
CoreRequestedBounds::unbounded(),
),
request("SELECT * FROM turns"),
)
.await
.expect_err("bounds minted above the scoped query limits fail closed");
assert!(matches!(widened, CoreResolutionError::InvalidBounds));
assert_eq!(state.source_read_count(), 0);
assert_eq!(compiler.physical_scan_count(), 0);
}
#[tokio::test]
async fn audit_retry_digest_does_not_depend_on_remaining_time() {
let state = Arc::new(FakeState::new(&["conv-a"]));
let compiler = CatalogCompiler::new(SessionContext::new().state());
let planning = authority(Arc::clone(&state));
let scope = QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
};
let first_call = DeclaredCall::live(Audience::new("query"), Duration::from_mins(2));
let retry_call = DeclaredCall::live(Audience::new("query"), Duration::from_secs(20));
let first = planning
.plan(
&compiler,
&QueryLimits::default(),
&scope,
false,
audit_with_declared("q-time-retry", &first_call),
request("SELECT * FROM messages"),
)
.await
.unwrap();
let retry = planning
.plan(
&compiler,
&QueryLimits::default(),
&scope,
false,
audit_with_declared("q-time-retry", &retry_call),
request("SELECT * FROM messages"),
)
.await
.unwrap();
assert!(matches!(first, CorePlanOutcome::Granted(_)));
assert!(matches!(retry, CorePlanOutcome::AlreadyRecorded(_)));
let evidence = {
let intents = FakeState::lock(&state.intents);
(
intents[0].shape(),
intents[1].shape(),
intents[0].metadata().digest(),
intents[1].metadata().digest(),
)
};
assert_eq!(evidence.0, evidence.1);
assert_eq!(evidence.2, evidence.3);
}
#[tokio::test]
async fn unsorted_authority_scope_produces_the_same_source_and_shape() {
let first_state = Arc::new(FakeState::new(&["conv-a", "conv-b"]));
let second_state = Arc::new(FakeState::new(&["conv-a", "conv-b"]));
let explain_enabled_state = Arc::new(FakeState::new(&["conv-a", "conv-b"]));
let compiler = CatalogCompiler::new(SessionContext::new().state());
authority(Arc::clone(&first_state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["b".to_owned(), "a".to_owned()],
memory: MemorySources::default(),
},
false,
audit("q-deterministic"),
request("SELECT * FROM messages"),
)
.await
.unwrap();
authority(Arc::clone(&second_state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned(), "b".to_owned()],
memory: MemorySources::default(),
},
false,
audit("q-deterministic"),
request("SELECT * FROM messages"),
)
.await
.unwrap();
authority(Arc::clone(&explain_enabled_state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned(), "b".to_owned()],
memory: MemorySources::default(),
},
true,
audit("q-deterministic"),
request("SELECT * FROM messages"),
)
.await
.unwrap();
let first = FakeState::lock(&first_state.intents)[0].clone();
let second = FakeState::lock(&second_state.intents)[0].clone();
let explain_enabled = FakeState::lock(&explain_enabled_state.intents)[0].clone();
assert_eq!(first.shape(), second.shape());
assert_eq!(first.source(), second.source());
assert_eq!(first.canonical_bytes(), second.canonical_bytes());
assert_ne!(first.shape(), explain_enabled.shape());
}
#[tokio::test]
async fn fleet_uses_and_releases_one_anchored_directory_snapshot() {
let state = Arc::new(FakeState::new(&[
"conv-b",
"conv-",
"conv-a",
"routine-scheduler",
]));
sort_snapshot(&state);
let compiler = CatalogCompiler::new(SessionContext::new().state());
let outcome = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Fleet,
true,
audit("q-fleet"),
request("SELECT * FROM messages"),
)
.await
.expect("anchored Fleet plan");
let CorePlanOutcome::Granted(prepared) = outcome else {
panic!("new fleet intent grants one permit");
};
assert_eq!(
prepared.partitions,
vec![PartitionId::new("conv-a"), PartitionId::new("conv-b")]
);
let counts = {
let counts = FakeState::lock(&state.counts);
(
counts.snapshots,
counts.pages,
counts.releases,
counts.operation_contexts.len(),
counts.subcall_budgets.clone(),
counts.subcall_audiences.clone(),
)
};
assert_eq!((counts.0, counts.1, counts.2), (1, 1, 1));
assert_eq!(counts.3, 1, "cleanup reuses the anchored operation context");
assert!(counts.4.windows(2).all(|pair| pair[1] <= pair[0]));
assert_eq!(
counts.5,
BTreeSet::from([polyc_state_connect::STATE_AUDIENCE.to_owned()])
);
}
#[tokio::test]
async fn fleet_retry_ignores_ephemeral_directory_snapshot_identity() {
let state = Arc::new(FakeState::new(&["conv-a"]));
let compiler = CatalogCompiler::new(SessionContext::new().state());
let planning = authority(Arc::clone(&state));
let first = planning
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Fleet,
false,
audit("q-fleet-retry"),
request("SELECT * FROM turns"),
)
.await
.expect("first snapshot records the intent");
let retry = planning
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Fleet,
false,
audit("q-fleet-retry"),
request("SELECT * FROM turns"),
)
.await
.expect("a later snapshot id does not change the request shape");
assert!(matches!(first, CorePlanOutcome::Granted(_)));
assert!(matches!(retry, CorePlanOutcome::AlreadyRecorded(_)));
let counts = {
let counts = FakeState::lock(&state.counts);
(counts.snapshots, counts.pages, counts.releases)
};
assert_eq!(counts, (2, 2, 2));
assert_eq!(compiler.physical_scan_count(), 0);
}
fn sort_snapshot(state: &FakeState) {
FakeState::lock(&state.snapshot).sort();
}
#[tokio::test]
async fn descriptor_and_state_failures_stop_before_audit_or_physical_scan() {
let accepted = conversation_core().versions();
let (good_schema, good_fact) = (accepted.schema().get(), accepted.fact_model().get());
let cases = [
(
OwnerId::new("crossed-owner"),
Classification::Confidential,
good_schema,
good_fact,
),
(
OwnerId::new("projector"),
Classification::Internal,
good_schema,
good_fact,
),
(
OwnerId::new("projector"),
Classification::Confidential,
good_schema + 1,
good_fact,
),
(
OwnerId::new("projector"),
Classification::Confidential,
good_schema,
good_fact + 1,
),
];
for (index, (owner, classification, schema, fact)) in cases.into_iter().enumerate() {
let state = Arc::new(FakeState::new(&["conv-a"]));
let current_source = FakeState::lock(&state.sources)
.get(&PartitionId::new("conv-a"))
.unwrap()
.source()
.clone();
let broken = manifest(
PartitionId::new("conv-a"),
¤t_source,
owner,
classification,
schema,
fact,
);
state.set_resolution(
"conv-a",
ProjectionResolution::new(
ProjectionHead::Current(Box::new(broken)),
None,
ProjectionGeneration::new(1),
),
);
let compiler = CatalogCompiler::new(SessionContext::new().state());
let error = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit(&format!("q-broken-{index}")),
request("SELECT * FROM turns"),
)
.await
.expect_err("incompatible descriptor fails closed");
assert!(matches!(
error,
CoreResolutionError::IncompatibleDescriptor(_)
));
assert_eq!(FakeState::lock(&state.counts).audit_writes, 0);
assert_eq!(compiler.physical_scan_count(), 0);
}
}
#[tokio::test]
async fn another_projection_family_cannot_satisfy_conversation_core() {
let state = Arc::new(FakeState::new(&["conv-a"]));
let current_source = FakeState::lock(&state.sources)
.get(&PartitionId::new("conv-a"))
.unwrap()
.source()
.clone();
let wrong_family = manifest_for_family(
"conversation-search/v1",
PartitionId::new("conv-a"),
¤t_source,
OwnerId::new("projector"),
Classification::Confidential,
1,
1,
);
state.set_resolution(
"conv-a",
ProjectionResolution::new(
ProjectionHead::Current(Box::new(wrong_family)),
None,
ProjectionGeneration::new(1),
),
);
let compiler = CatalogCompiler::new(SessionContext::new().state());
let error = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit("q-wrong-family"),
request("SELECT * FROM turns"),
)
.await
.expect_err("another projection family cannot satisfy core");
assert!(matches!(
error,
CoreResolutionError::IncompatibleDescriptor(_)
));
assert_eq!(FakeState::lock(&state.counts).audit_writes, 0);
}
#[tokio::test]
async fn source_recreation_and_crossed_source_responses_fail_closed() {
let state = Arc::new(FakeState::new(&["conv-a"]));
let recreated = source(PartitionId::new("conv-a"), 9);
let broken = manifest(
PartitionId::new("conv-a"),
&recreated,
OwnerId::new("projector"),
Classification::Confidential,
1,
1,
);
state.set_resolution(
"conv-a",
ProjectionResolution::new(
ProjectionHead::Current(Box::new(broken)),
None,
ProjectionGeneration::new(1),
),
);
let compiler = CatalogCompiler::new(SessionContext::new().state());
let recreated_error = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit("q-recreated"),
request("SELECT * FROM turns"),
)
.await
.expect_err("a descriptor for another incarnation is not current");
assert!(matches!(
recreated_error,
CoreResolutionError::IncompatibleDescriptor(_)
));
assert_eq!(FakeState::lock(&state.counts).audit_writes, 0);
let crossed = Arc::new(FakeState::new(&["conv-a"]));
FakeState::lock(&crossed.sources).insert(
PartitionId::new("conv-a"),
JournalSourceHead::new(
source(PartitionId::new("conv-b"), 1),
JournalPosition::new(20),
),
);
let crossed_error = authority(Arc::clone(&crossed))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit("q-crossed-source"),
request("SELECT * FROM turns"),
)
.await
.expect_err("source responses bind the requested partition");
assert!(matches!(
crossed_error,
CoreResolutionError::SourceMismatch(_)
));
assert_eq!(FakeState::lock(&crossed.counts).resolutions, 0);
}
#[tokio::test]
async fn a_superseded_or_missing_administrator_audit_source_fails_closed() {
let partition = PartitionId::new(polyc_projection::family::ADMIN_AUDIT_PARTITION);
let state = Arc::new(FakeState::new(&["conv-a"]));
FakeState::lock(&state.sources).insert(
partition.clone(),
JournalSourceHead::new(source(partition.clone(), 5), JournalPosition::new(20)),
);
state.set_family_resolution(
polyc_projection::family::administrator_audit(),
polyc_projection::family::ADMIN_AUDIT_PARTITION,
ProjectionResolution::new(
ProjectionHead::Superseded {
generation: ProjectionGeneration::new(1),
source: Box::new(polyc_state::feed::ProjectionSource::Journal(source(
partition.clone(),
6,
))),
},
None,
ProjectionGeneration::new(1),
),
);
let compiler = CatalogCompiler::new(SessionContext::new().state());
let superseded = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Fleet,
false,
audit("q-admin-superseded"),
request("SELECT position FROM admin_model_changes"),
)
.await
.expect_err("a superseded administrator-audit source must not resolve");
assert!(
matches!(superseded, CoreResolutionError::Superseded(_)),
"{superseded:?}"
);
let absent = Arc::new(FakeState::new(&["conv-a"]));
let missing = authority(Arc::clone(&absent))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Fleet,
false,
audit("q-admin-missing"),
request("SELECT position FROM admin_model_changes"),
)
.await
.expect_err("an administrator-audit partition with no head must not resolve");
assert!(
matches!(missing, CoreResolutionError::MissingSource(_)),
"{missing:?}"
);
}
#[tokio::test]
async fn a_superseded_or_missing_query_audit_source_fails_closed() {
let state = Arc::new(FakeState::new(&["conv-a"]));
state.set_family_resolution(
polyc_projection::family::query_audit(),
polyc_projection::family::QUERY_AUDIT_SOURCE,
ProjectionResolution::new(
ProjectionHead::Superseded {
generation: ProjectionGeneration::new(1),
source: Box::new(polyc_state::feed::ProjectionSource::QueryAudit(
polyc_state::query_audit::AuditSource::new(PartitionIncarnation::from_bytes(
[200; PartitionIncarnation::LEN],
)),
)),
},
None,
ProjectionGeneration::new(1),
),
);
let compiler = CatalogCompiler::new(SessionContext::new().state());
let superseded = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Fleet,
false,
audit("q-query-audit-superseded"),
request("SELECT position FROM query_audit_intents"),
)
.await
.expect_err("a superseded query-audit source must not resolve");
assert!(
matches!(superseded, CoreResolutionError::Superseded(_)),
"{superseded:?}"
);
let absent = Arc::new(FakeState::new(&["conv-a"]));
let missing = authority(Arc::clone(&absent))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Fleet,
false,
audit("q-query-audit-missing"),
request("SELECT position FROM query_audit_intents"),
)
.await
.expect_err("a query-audit family with no published generation must not resolve");
assert!(
matches!(missing, CoreResolutionError::MissingProjection(_)),
"{missing:?}"
);
}
#[tokio::test]
async fn a_superseded_or_missing_persona_memory_source_fails_closed() {
let partition_name = "persona-p-mem";
let partition =
polyc_state::persona_memory::journal::MemoryJournalPartition::parse(partition_name)
.expect("valid persona-memory partition name");
let scope = memory_scope(Some("p"), &[]);
let state = Arc::new(FakeState::new(&["conv-a"]));
state.set_family_resolution(
polyc_projection::family::persona_memory(),
partition_name,
ProjectionResolution::new(
ProjectionHead::Superseded {
generation: ProjectionGeneration::new(1),
source: Box::new(polyc_state::feed::ProjectionSource::PersonaMemory(
polyc_state::persona_memory::journal::PersonaMemorySource::new(
partition.clone(),
PartitionIncarnation::from_bytes([200; PartitionIncarnation::LEN]),
),
)),
},
None,
ProjectionGeneration::new(1),
),
);
let compiler = CatalogCompiler::new(SessionContext::new().state());
let superseded = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&scope,
false,
audit("q-memory-superseded"),
request("SELECT fact_id FROM memory_facts"),
)
.await
.expect_err("a superseded persona-memory source must not resolve");
assert!(
matches!(superseded, CoreResolutionError::Superseded(_)),
"{superseded:?}"
);
let absent = Arc::new(FakeState::new(&["conv-a"]));
let missing = authority(Arc::clone(&absent))
.plan(
&compiler,
&QueryLimits::default(),
&scope,
false,
audit("q-memory-missing"),
request("SELECT fact_id FROM memory_facts"),
)
.await
.expect_err("a persona-memory partition with no published generation must not resolve");
assert!(
matches!(missing, CoreResolutionError::MissingProjection(_)),
"{missing:?}"
);
}
#[test]
fn each_realm_gate_refuses_the_administrator_table_on_its_own() {
assert!(!CoreTable::AdminModelChanges.visible_in(CoreRealm::Visible));
assert!(CoreTable::AdminModelChanges.visible_in(CoreRealm::Fleet));
assert!(!family_readable_in(
CoreRealm::Visible,
polyc_projection::family::administrator_audit()
));
assert!(family_readable_in(
CoreRealm::Fleet,
polyc_projection::family::administrator_audit()
));
assert_eq!(
CoreTable::AdminModelChanges.physical_schema().realm_shape(),
polyc_projection::family::RealmShape::FleetOnly
);
}
#[tokio::test]
async fn the_administrator_audit_family_resolves_by_declaration() {
let partition = PartitionId::new(polyc_projection::family::ADMIN_AUDIT_PARTITION);
let state = Arc::new(FakeState::new(&["conv-a"]));
let admin_source = source(partition.clone(), 5);
FakeState::lock(&state.sources).insert(
partition.clone(),
JournalSourceHead::new(admin_source.clone(), JournalPosition::new(20)),
);
state.set_family_resolution(
polyc_projection::family::administrator_audit(),
polyc_projection::family::ADMIN_AUDIT_PARTITION,
ProjectionResolution::new(
ProjectionHead::Current(Box::new(manifest_for_family(
polyc_projection::family::ADMINISTRATOR_AUDIT,
partition.clone(),
&admin_source,
OwnerId::new("projector"),
Classification::Confidential,
polyc_projection::family::ADMINISTRATOR_AUDIT_SCHEMA.get(),
polyc_projection::family::ADMINISTRATOR_AUDIT_FACT_MODEL.get(),
))),
None,
ProjectionGeneration::new(1),
),
);
let compiler = CatalogCompiler::new(SessionContext::new().state());
let planned = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Fleet,
false,
audit("q-admin-fleet"),
request("SELECT position FROM admin_model_changes"),
)
.await
.expect("a Fleet session resolves the declared partition");
let CorePlanOutcome::Granted(prepared) = planned else {
panic!("a new intent grants one permit");
};
assert_eq!(
prepared.manifests.len(),
1,
"one declaration, one key, one generation"
);
assert_eq!(prepared.manifests[0].key().source(), &partition);
let refused = authority(state)
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit("q-admin-visible"),
request("SELECT position FROM admin_model_changes"),
)
.await
.expect_err("a conversation-scoped session cannot resolve it");
assert!(
matches!(
refused,
CoreResolutionError::TableOutsideRealm | CoreResolutionError::FamilyOutsideRealm
),
"{refused:?}"
);
}
#[tokio::test]
async fn missing_superseded_outage_and_audit_refusal_fail_closed() {
for head in [
ProjectionHead::Absent,
ProjectionHead::Superseded {
generation: ProjectionGeneration::new(1),
source: Box::new(polyc_state::feed::ProjectionSource::Journal(source(
PartitionId::new("conv-a"),
9,
))),
},
] {
let state = Arc::new(FakeState::new(&["conv-a"]));
state.set_resolution(
"conv-a",
ProjectionResolution::new(head, None, ProjectionGeneration::new(1)),
);
let compiler = CatalogCompiler::new(SessionContext::new().state());
let error = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit("q-head"),
request("SELECT * FROM turns"),
)
.await
.expect_err("missing current projection fails closed");
assert!(matches!(
error,
CoreResolutionError::MissingProjection(_) | CoreResolutionError::Superseded(_)
));
}
let outage = Arc::new(FakeState::new(&["conv-a"]));
outage.set_outage();
let compiler = CatalogCompiler::new(SessionContext::new().state());
assert!(
authority(Arc::clone(&outage))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default()
},
false,
audit("q-outage"),
request("SELECT * FROM turns"),
)
.await
.is_err()
);
let refused = Arc::new(FakeState::new(&["conv-a"]));
refused.refuse_audit();
assert!(
authority(Arc::clone(&refused))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default()
},
false,
audit("q-refused"),
request("SELECT * FROM turns"),
)
.await
.is_err()
);
assert_eq!(compiler.physical_scan_count(), 0);
}
#[tokio::test]
async fn a_crossed_permit_settles_nothing_and_refuses() {
let state = Arc::new(FakeState::new(&["conv-a"]));
state.cross_permit();
let compiler = CatalogCompiler::new(SessionContext::new().state());
let error = authority(Arc::clone(&state))
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations: vec!["a".to_owned()],
memory: MemorySources::default(),
},
false,
audit("q-crossed-permit"),
request("SELECT text FROM messages"),
)
.await
.expect_err("a crossed permit is refused");
assert!(matches!(error, CoreResolutionError::CrossedPermit));
tokio::time::sleep(Duration::from_millis(50)).await;
assert_eq!(
state.settlements(),
0,
"an abandoned guardian never presents a completion for another trail"
);
}
#[test]
fn the_signer_table_is_fleet_only_and_the_handoff_table_is_not() {
assert!(CoreTable::Handoffs.visible_in(CoreRealm::Visible));
assert!(CoreTable::Handoffs.visible_in(CoreRealm::Fleet));
assert!(!CoreTable::HandoffSigners.visible_in(CoreRealm::Visible));
assert!(CoreTable::HandoffSigners.visible_in(CoreRealm::Fleet));
}
#[test]
fn no_key_material_reaches_the_public_handoff_schema() {
let names: Vec<String> = CoreTable::Handoffs
.public_schema()
.fields()
.iter()
.map(|field| field.name().clone())
.collect();
assert!(
!names.iter().any(|name| name.contains("signed_by")),
"the public handoff schema must carry no signing key: {names:?}"
);
assert!(names.iter().any(|name| name == "signature_status"));
}
#[test]
fn both_delegation_tables_resolve_from_their_public_names() {
for table in [CoreTable::Handoffs, CoreTable::HandoffSigners] {
assert_eq!(CoreTable::from_name(table.name()).unwrap(), table);
assert_eq!(
table.family().family_str(),
polyc_projection::family::CONVERSATION_DELEGATION
);
}
}
#[test]
fn every_resolvable_table_belongs_to_a_registered_family() {
let registered: BTreeSet<&str> = PROJECTED_FAMILIES
.iter()
.map(polyc_projection::family::FamilyEntry::family_str)
.collect();
for table in [
CoreTable::Turns,
CoreTable::Messages,
CoreTable::Usage,
CoreTable::ModelCall,
CoreTable::ToolCalls,
CoreTable::TurnFailed,
CoreTable::Summary,
CoreTable::Handoffs,
CoreTable::HandoffSigners,
] {
assert!(
registered.contains(table.family().family_str()),
"`{}` resolves to the unregistered family `{}`",
table.name(),
table.family().family_str()
);
}
}
#[test]
fn every_registered_family_table_resolves_by_name() {
for family in sql_servable_families() {
for schema in family.tables() {
let name = schema.table().as_str();
let resolved = CoreTable::from_name(name)
.unwrap_or_else(|_| panic!("`{name}` is registered but does not resolve"));
assert_eq!(resolved.table(), schema.table());
assert_eq!(resolved.family().family_str(), family.family_str());
}
}
}
#[test]
fn every_public_view_fits_the_surface_and_bounds() {
for table in CoreTable::ALL {
if let Some(sql) = table.public_view_sql() {
assert_eq!(
check_statement_allowed(sql, false)
.unwrap_or_else(|err| panic!("{table:?} must be allowed, got {err:?}")),
AllowedStatement::Query
);
}
}
}
#[test]
fn every_table_round_trips_its_name() {
for table in CoreTable::ALL {
let resolved = CoreTable::from_name(table.name())
.unwrap_or_else(|_| panic!("`{}` must resolve by name", table.name()));
assert_eq!(
resolved,
table,
"`{}` resolved to another table",
table.name()
);
}
}
#[test]
fn every_variant_is_listed_in_all() {
let mut declared: Vec<String> = sql_servable_families()
.flat_map(|family| family.tables().iter().map(|id| id.table().to_string()))
.collect();
declared.sort();
let mut listed: Vec<String> = CoreTable::ALL
.iter()
.map(|table| table.name().to_string())
.collect();
listed.sort();
assert_eq!(
listed, declared,
"CoreTable::ALL and the closed families disagree — a table missing here is a table \
no ALL-driven test checks"
);
}
#[test]
fn search_index_is_registered_but_not_sql_servable() {
assert!(
PROJECTED_FAMILIES
.iter()
.any(|family| family.family_str() == polyc_projection::family::SEARCH_INDEX),
"search-index/v1 must be registered for the projector to publish it"
);
assert!(
sql_servable_families()
.all(|family| family.family_str() != polyc_projection::family::SEARCH_INDEX),
"search-index/v1 must not be SQL-servable until a reader ships"
);
assert!(
CoreTable::from_name("postings").is_err(),
"no CoreTable may resolve search-index/v1's postings table yet"
);
assert!(
CoreTable::from_name("coverage").is_err(),
"no CoreTable may resolve search-index/v1's coverage table yet"
);
}
#[test]
fn no_fleet_only_table_is_visible() {
use polyc_projection::family::RealmShape;
for table in CoreTable::ALL {
let fleet_only = table.physical_schema().realm_shape() == RealmShape::FleetOnly;
assert_eq!(
!table.visible_in(CoreRealm::Visible),
fleet_only,
"`{}` disagrees with its family about the visible realm",
table.name()
);
assert!(
table.visible_in(CoreRealm::Fleet),
"`{}` must resolve for a Fleet reader",
table.name()
);
}
}
#[tokio::test]
async fn catalog_listing_agrees_with_planning() {
let scopes = [
(
"visible, conversation-scoped, no memory owner",
QueryScope::Conversations {
conversations: vec!["conv-a".to_owned()],
memory: MemorySources::default(),
},
),
(
"visible, persona-scoped, memory owner",
QueryScope::Conversations {
conversations: vec!["conv-a".to_owned()],
memory: MemorySources {
owner: Some("persona-a".to_owned()),
participants: Vec::new(),
},
},
),
("fleet", QueryScope::Fleet),
(
"fleet, one conversation",
QueryScope::FleetConversation {
conversation: "conv-a".to_owned(),
},
),
];
let compiler = CatalogCompiler::new(SessionContext::new().state());
for (label, scope) in &scopes {
let realm = CoreRealm::from_scope(scope);
let listed: BTreeSet<CoreTable> = catalog_tables(realm, scope).into_iter().collect();
let state = Arc::new(FakeState::new(&["conv-a"]));
let authority = authority(state);
for table in CoreTable::ALL {
let outcome = authority
.plan(
&compiler,
&QueryLimits::default(),
scope,
false,
audit(&format!("q-catalog-agreement-{label}-{}", table.name())),
request(&format!("SELECT * FROM {} LIMIT 0", table.name())),
)
.await;
let visibility_refused = matches!(
outcome,
Err(CoreResolutionError::TableOutsideRealm
| CoreResolutionError::FamilyOutsideRealm
| CoreResolutionError::TableOutsideAudience)
);
assert_eq!(
listed.contains(&table),
!visibility_refused,
"{label}: catalog and planning disagree on `{}` (outcome: {outcome:?})",
table.name()
);
}
}
}
#[test]
fn attribution_provenance_is_fleet_only_absolutely() {
use polyc_projection::family::RealmShape;
assert_eq!(
CoreTable::AttributionProvenance
.physical_schema()
.realm_shape(),
RealmShape::FleetOnly,
"a Visible-realm composition must never resolve attribution_provenance"
);
assert!(!CoreTable::AttributionProvenance.visible_in(CoreRealm::Visible));
assert!(CoreTable::AttributionProvenance.visible_in(CoreRealm::Fleet));
}
#[test]
fn only_a_conversation_scoped_family_is_readable_in_the_visible_realm() {
let journal = polyc_projection::family::conversation_core();
let versioned = polyc_projection::family::credential_lifecycle();
assert!(super::family_readable_in(super::CoreRealm::Fleet, journal));
assert!(super::family_readable_in(
super::CoreRealm::Fleet,
versioned
));
assert!(super::family_readable_in(
super::CoreRealm::Visible,
journal
));
assert!(
!super::family_readable_in(super::CoreRealm::Visible, versioned),
"a visible session must not read a family it cannot be bounded to"
);
}
#[test]
fn persona_memory_is_the_second_family_readable_in_the_visible_realm() {
let persona_memory = polyc_projection::family::persona_memory();
assert!(super::family_readable_in(
super::CoreRealm::Visible,
persona_memory
));
assert!(super::family_readable_in(
super::CoreRealm::Fleet,
persona_memory
));
}
fn memory_scope(owner: Option<&str>, participants: &[&str]) -> QueryScope {
QueryScope::Conversations {
conversations: Vec::new(),
memory: MemorySources {
owner: owner.map(str::to_owned),
participants: participants.iter().map(|id| (*id).to_owned()).collect(),
},
}
}
#[test]
fn an_owner_table_with_a_verified_owner_and_participants_is_admitted() {
let scope = memory_scope(Some("p-owner"), &["p-other"]);
super::check_persona_memory_admission(
&[CoreTable::MemoryFacts],
&scope,
CoreRealm::Visible,
false,
)
.expect("the execution registry narrows owner-table files to the verified owner");
}
#[test]
fn an_owner_table_with_no_owner_is_refused() {
let scope = memory_scope(None, &[]);
let error = super::check_persona_memory_admission(
&[CoreTable::MemoryFacts],
&scope,
CoreRealm::Visible,
false,
)
.unwrap_err();
assert!(matches!(error, CoreResolutionError::TableOutsideAudience));
}
#[test]
fn an_owner_table_with_only_its_owner_is_admitted() {
let scope = memory_scope(Some("p-owner"), &[]);
super::check_persona_memory_admission(
&[CoreTable::MemoryFacts],
&scope,
CoreRealm::Visible,
false,
)
.expect("an owner reading only their own owner-audience table is admitted");
}
#[test]
fn a_portable_table_admits_participants() {
let scope = memory_scope(Some("p-owner"), &["p-other"]);
super::check_persona_memory_admission(
&[CoreTable::MemoryPortableFacts],
&scope,
CoreRealm::Visible,
false,
)
.expect("a portable table is AnyAuthorized: participants are exactly the point");
}
#[test]
fn persona_memory_cannot_mix_with_a_journal_directory_family() {
let scope = memory_scope(Some("p-owner"), &[]);
let error = super::check_persona_memory_admission(
&[CoreTable::MemoryPortableFacts, CoreTable::Turns],
&scope,
CoreRealm::Visible,
false,
)
.unwrap_err();
assert!(matches!(error, CoreResolutionError::FamilyOutsideRealm));
}
#[test]
fn composite_trace_memory_may_pin_its_trace_family() {
let scope = memory_scope(Some("p-owner"), &["p-other"]);
super::check_persona_memory_admission(
&[CoreTable::MemoryPortableFacts, CoreTable::TraceTurns],
&scope,
CoreRealm::Visible,
true,
)
.expect("the fixed composite statement pins both authorities");
}
#[test]
fn fleet_scope_is_exempt_from_both_persona_memory_admission_rules() {
super::check_persona_memory_admission(
&[CoreTable::MemoryFacts, CoreTable::Turns],
&QueryScope::Fleet,
CoreRealm::Fleet,
false,
)
.expect("Fleet reads every family and every audience whole");
}
const GET_MY_PAYMENTS_SQL: &str = "SELECT partition, position, turn_id, direction, reference, \
amount_base_units, asset, recipient, method, tool_call_id, payer_kind, timestamp_unix \
FROM payments WHERE subject = $1 \
UNION ALL \
SELECT partition, position, turn_id, direction, reference, amount_base_units, asset, \
recipient, method, tool_call_id, payer_kind, timestamp_unix FROM outbound_payments \
WHERE subject = $1 \
ORDER BY timestamp_unix DESC, partition DESC, position DESC LIMIT 500";
const GET_MY_REFUSALS_SQL: &str = "SELECT partition, position, turn_id, reason, reason_detail, \
merchant_host, requested_base_units, permitted_base_units, tool_call_id, timestamp_unix \
FROM refusals WHERE subject = $1 \
ORDER BY timestamp_unix DESC, partition DESC, position DESC LIMIT 500";
fn state_with_financial(conversations: &[String]) -> Arc<FakeState> {
let partitions: Vec<String> = conversations
.iter()
.map(|id| format!("conv-{id}"))
.collect();
let refs: Vec<&str> = partitions.iter().map(String::as_str).collect();
let state = Arc::new(FakeState::new(&refs));
let family = conversation_financial();
let versions = family.versions();
for (ordinal, partition) in partitions.iter().enumerate() {
let partition = PartitionId::new(partition);
let journal = source(partition.clone(), u8::try_from(ordinal + 1).unwrap());
let key = ProjectionKey::new(FamilyId::new(family.family_str()), partition.clone());
publish(
&state,
&key,
manifest_for_family(
family.family_str(),
partition,
&journal,
OwnerId::new("projector"),
Classification::Confidential,
versions.schema().get(),
versions.fact_model().get(),
),
);
}
state
}
async fn plan_my_statement(
conversations: Vec<String>,
sql: &str,
label: &str,
) -> Result<CorePlanOutcome, CoreResolutionError> {
let state = state_with_financial(&conversations);
let compiler = CatalogCompiler::new(SessionContext::new().state());
let mut request = request(sql);
request
.parameters
.push(CoreParameter::Utf8("persona-a".to_owned()));
authority(state)
.plan(
&compiler,
&QueryLimits::default(),
&QueryScope::Conversations {
conversations,
memory: MemorySources::default(),
},
false,
audit(&format!("q-479-{label}")),
request,
)
.await
}
#[tokio::test]
async fn a_persona_wide_statement_past_the_source_pin_bound_is_refused() {
for (label, sql) in [
("payments", GET_MY_PAYMENTS_SQL),
("refusals", GET_MY_REFUSALS_SQL),
] {
let over = usize::try_from(MAX_SOURCE_PINS).unwrap() + 1;
let conversations: Vec<String> = (0..over).map(|index| format!("c-{index}")).collect();
let error = plan_my_statement(conversations, sql, label)
.await
.expect_err("a persona-wide statement past the source-pin bound must refuse");
assert!(
matches!(
error,
CoreResolutionError::State(StateError::BoundsExceeded {
bound: BoundKind::CommandRecords,
limit: 32,
requested: 33,
})
),
"{label} refused for the wrong reason: {error:?}"
);
}
}
#[tokio::test]
async fn a_persona_wide_statement_at_the_source_pin_bound_is_granted() {
for (label, sql) in [
("payments", GET_MY_PAYMENTS_SQL),
("refusals", GET_MY_REFUSALS_SQL),
] {
let conversations: Vec<String> = (0..MAX_SOURCE_PINS)
.map(|index| format!("c-{index}"))
.collect();
let outcome = plan_my_statement(conversations, sql, label)
.await
.expect("a persona-wide statement at the source-pin bound must plan");
assert!(
matches!(outcome, CorePlanOutcome::Granted(_)),
"{label} planned to an unexpected outcome"
);
}
}