use std::collections::BTreeMap;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex, PoisonError};
use std::time::Duration;
use arrow::array::{
Array, ArrayRef, BooleanArray, FixedSizeBinaryBuilder, StringArray, UInt64Array,
};
use arrow::record_batch::RecordBatch;
use async_trait::async_trait;
use bytes::Bytes;
use datafusion::execution::memory_pool::FairSpillPool;
use datafusion::execution::runtime_env::RuntimeEnv;
use datafusion::execution::runtime_env::RuntimeEnvBuilder;
use datafusion::execution::session_state::SessionStateBuilder;
use futures::StreamExt;
use parquet::arrow::ArrowWriter;
use polyc_projection::family::{CONVERSATION_MESSAGES, CONVERSATION_TURNS, conversation_core};
use polyc_projection_artifact::{
ArtifactReadError, ArtifactReadFuture, ExactArtifactBackend, ExactObjectMetadata,
ExactObjectRange, FleetArtifactAccess, FleetArtifactReader, RealmTopology,
VisibleArtifactAccess, VisibleArtifactReader,
parquet_profile::{canonicalize_writer_output, source_decode_reservation_bytes},
};
use polyc_state::cancel::CancellationToken;
use polyc_state::context::CallContext;
use polyc_state::deadline::{Deadline, MonotonicInstant};
use polyc_state::digest::ContentDigest;
use polyc_state::error::{AmbiguityReason, OutageReach, StateError};
use polyc_state::feed::SourceCheckpoint;
use polyc_state::id::{Audience, NamespaceId, OwnerId, PartitionId};
use polyc_state::immutable::{
AtRestProtection, Classification, ContentReference, Generation as ObjectGeneration,
ObjectDescriptor, Retention,
};
use polyc_state::journal::{
GetJournalSource, JournalAttestation, JournalDirectoryPage, JournalDirectorySnapshot,
JournalDirectorySnapshotId, JournalSourceHead, ListJournalDirectorySnapshot,
ReleaseJournalDirectorySnapshot,
};
use polyc_state::page::PageCompleteness;
use polyc_state::projection::artifact::testing::{FixtureManifestSigner, FixtureManifestTrust};
use polyc_state::projection::artifact::{
AccessRealm, ArtifactManifest, ArtifactSegment, ArtifactTable, ExactObjectRef, ManifestSigner,
ObjectNamespace, SignedArtifactManifest, SourceBounds,
};
use polyc_state::projection::{
FamilyId, ProjectionGeneration, ProjectionHead, ProjectionKey, ProjectionManifest,
ProjectionResolution, PublisherFence, PublisherId, ResolveManifest,
};
use polyc_state::query_audit::memory::MemoryQueryAudit;
use polyc_state::query_audit::{
AuditPhase, BeginOutcome, BeginQueryAudit, QueryAuditRead, QueryAuditWrite, QueryId,
ReadQueryAudit, RequesterId,
};
use polyc_state::receipt::Receipt;
use polyc_state::revision::{CommitRoot, JournalPosition, JournalSource, PartitionIncarnation};
use polyc_state_connect::wire::DeclaredCall;
use sha2::{Digest, Sha256};
use super::super::admission::CoreExecutionAdmissionInput;
use super::super::{
CoreArtifactAuthority, CoreExecutionAdmission, CoreOperationContext, CoreScopeRevalidator,
CurrentCredentialAuthority,
};
use crate::core_resolution::{
CatalogCompiler, CoreAuditContext, CoreCompletionCommand, CoreCompletionContext,
CoreConsistency, CoreExecutionPermit, CoreMetadataAuthority, CoreParameter, CorePlanOutcome,
CorePlanningAuthority, CoreQueryRequest, CoreRequestedBounds, CoreResolutionError,
PreparedCoreQuery, ProjectedCorePolicy, ProjectedCorePolicyInput,
};
use crate::engine::QueryLimits;
use crate::session::QueryScope;
const VISIBLE_NAMESPACE: &str = "conversation-visible";
const FLEET_NAMESPACE: &str = "conversation-fleet";
const MANIFEST_KEY: &str = "conversation/manifests/g1";
const VISIBLE_MESSAGES_KEY: &str = "messages/visible-a-g1.parquet";
const SIGNER: FixtureManifestSigner = FixtureManifestSigner::new(73);
fn lock<T>(mutex: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
mutex.lock().unwrap_or_else(PoisonError::into_inner)
}
fn identity(reference: &ExactObjectRef) -> (String, String, u64) {
(
reference.namespace().as_str().to_owned(),
reference.key().as_str().to_owned(),
reference.generation(),
)
}
#[derive(Debug)]
pub(super) struct TestStore {
objects: Mutex<BTreeMap<(String, String, u64), Vec<u8>>>,
head_defect: Mutex<Option<HeadDefect>>,
pause_ranges: AtomicBool,
heads: AtomicUsize,
ranges: AtomicUsize,
segment_ranges: AtomicUsize,
}
impl TestStore {
fn new() -> Self {
Self {
objects: Mutex::new(BTreeMap::new()),
head_defect: Mutex::new(None),
pause_ranges: AtomicBool::new(false),
heads: AtomicUsize::new(0),
ranges: AtomicUsize::new(0),
segment_ranges: AtomicUsize::new(0),
}
}
fn insert(&self, reference: &ExactObjectRef, bytes: Vec<u8>) {
lock(&self.objects).insert(identity(reference), bytes);
}
pub(super) fn range_calls(&self) -> usize {
self.ranges.load(Ordering::SeqCst)
}
pub(super) fn total_calls(&self) -> usize {
self.heads.load(Ordering::SeqCst) + self.range_calls()
}
pub(super) fn segment_range_calls(&self) -> usize {
self.segment_ranges.load(Ordering::SeqCst)
}
pub(super) fn corrupt_visible_messages(&self) {
let key = (
VISIBLE_NAMESPACE.to_owned(),
VISIBLE_MESSAGES_KEY.to_owned(),
11,
);
lock(&self.objects).get_mut(&key).unwrap()[0] ^= 0xff;
}
pub(super) fn return_wrong_generation(&self) {
*lock(&self.head_defect) = Some(HeadDefect::Generation);
}
pub(super) fn return_wrong_length(&self) {
*lock(&self.head_defect) = Some(HeadDefect::Length);
}
pub(super) fn pause_ranges(&self) {
self.pause_ranges.store(true, Ordering::SeqCst);
}
pub(super) fn add_unreferenced_object(&self) {
let reference = ExactObjectRef::try_new(
ObjectNamespace::try_new(VISIBLE_NAMESPACE).unwrap(),
ContentReference::try_new("messages/unreferenced.parquet").unwrap(),
404,
)
.unwrap();
self.insert(&reference, vec![1, 2, 3]);
}
}
#[derive(Debug, Clone, Copy)]
enum HeadDefect {
Generation,
Length,
}
#[derive(Debug, Clone)]
struct Backend(Arc<TestStore>);
impl ExactArtifactBackend for Backend {
fn declared_protection(&self, _namespace: &ObjectNamespace) -> AtRestProtection {
AtRestProtection::TenantKey
}
fn head_exact(
&self,
reference: &ExactObjectRef,
) -> ArtifactReadFuture<'_, ExactObjectMetadata> {
self.0.heads.fetch_add(1, Ordering::SeqCst);
let key = identity(reference);
let generation = reference.generation();
Box::pin(async move {
let byte_len = {
let objects = lock(&self.0.objects);
let bytes = objects
.get(&key)
.ok_or_else(|| ArtifactReadError::NotFound {
key: key.1.clone(),
generation,
})?;
let byte_len = u64::try_from(bytes.len()).unwrap();
drop(objects);
byte_len
};
let targeted = key.0 == VISIBLE_NAMESPACE && key.1 == VISIBLE_MESSAGES_KEY;
let defect = targeted.then(|| *lock(&self.0.head_defect)).flatten();
let observed_generation = match defect {
Some(HeadDefect::Generation) => generation + 1,
Some(HeadDefect::Length) | None => generation,
};
let observed_len = match defect {
Some(HeadDefect::Length) => byte_len + 1,
Some(HeadDefect::Generation) | None => byte_len,
};
Ok(ExactObjectMetadata::new(observed_generation, observed_len))
})
}
fn read_exact_range(
&self,
reference: &ExactObjectRef,
offset: u64,
len: u64,
) -> ArtifactReadFuture<'_, ExactObjectRange> {
self.0.ranges.fetch_add(1, Ordering::SeqCst);
let key = identity(reference);
if key.1.ends_with(".parquet") {
self.0.segment_ranges.fetch_add(1, Ordering::SeqCst);
}
let generation = reference.generation();
Box::pin(async move {
if self.0.pause_ranges.load(Ordering::SeqCst)
&& key.0 == VISIBLE_NAMESPACE
&& key.1.starts_with("messages/")
{
futures::future::pending::<()>().await;
}
let start = usize::try_from(offset).map_err(|_| ArtifactReadError::Refused {
reason: "fixture offset exceeds usize".to_owned(),
})?;
let end = usize::try_from(offset.saturating_add(len)).map_err(|_| {
ArtifactReadError::Refused {
reason: "fixture range exceeds usize".to_owned(),
}
})?;
let range = {
let objects = lock(&self.0.objects);
let bytes = objects
.get(&key)
.ok_or_else(|| ArtifactReadError::NotFound {
key: key.1.clone(),
generation,
})?;
let range = bytes
.get(start..end)
.ok_or_else(|| ArtifactReadError::Refused {
reason: "fixture received an out-of-range request".to_owned(),
})?
.to_vec();
drop(objects);
range
};
Ok(ExactObjectRange::new(generation, range))
})
}
}
#[derive(Debug)]
pub(super) struct TestState {
source: Mutex<JournalSourceHead>,
manifest: ProjectionManifest,
audit: MemoryQueryAudit,
refuse_audit: Mutex<bool>,
outage: Mutex<bool>,
completion_outage: Mutex<bool>,
completion_response_loss: Mutex<bool>,
completion_attempts: AtomicUsize,
completion_receipt_reads: AtomicUsize,
}
impl TestState {
fn new(source: JournalSource, manifest: ProjectionManifest) -> Self {
Self {
source: Mutex::new(JournalSourceHead::new(source, JournalPosition::new(10))),
manifest,
audit: MemoryQueryAudit::new(),
refuse_audit: Mutex::new(false),
outage: Mutex::new(false),
completion_outage: Mutex::new(false),
completion_response_loss: Mutex::new(false),
completion_attempts: AtomicUsize::new(0),
completion_receipt_reads: AtomicUsize::new(0),
}
}
pub(super) fn refuse_audit(&self) {
*lock(&self.refuse_audit) = true;
}
pub(super) fn recreate_source(&self) {
let partition = lock(&self.source).source().partition().clone();
*lock(&self.source) = JournalSourceHead::new(
JournalSource::new(
partition,
PartitionIncarnation::from_bytes([99; PartitionIncarnation::LEN]),
),
JournalPosition::new(1),
);
}
pub(super) fn set_outage(&self) {
*lock(&self.outage) = true;
}
pub(super) fn arm_completion_response_loss(&self) {
*lock(&self.completion_response_loss) = true;
}
pub(super) fn set_completion_outage(&self) {
*lock(&self.completion_outage) = true;
}
pub(super) fn completion_attempts(&self) -> usize {
self.completion_attempts.load(Ordering::SeqCst)
}
pub(super) fn completion_receipt_reads(&self) -> usize {
self.completion_receipt_reads.load(Ordering::SeqCst)
}
pub(super) fn completion_for(
&self,
query: &str,
) -> Option<polyc_state::query_audit::QueryCompletion> {
self.audit
.audit(
ReadQueryAudit::new(QueryId::new(query), NamespaceId::new("tenant-a")),
&live_context(),
)
.unwrap()
.and_then(|audit| audit.completion().cloned())
}
}
#[async_trait]
impl CoreMetadataAuthority for TestState {
async fn create_directory_snapshot(
&self,
_operation: &CoreOperationContext,
) -> Result<JournalDirectorySnapshot, CoreResolutionError> {
Ok(JournalDirectorySnapshot::new(
JournalDirectorySnapshotId::new("fixture-directory"),
1,
))
}
async fn directory_page(
&self,
_operation: &CoreOperationContext,
request: ListJournalDirectorySnapshot,
) -> Result<JournalDirectoryPage, CoreResolutionError> {
Ok(JournalDirectoryPage::new(
request.snapshot().clone(),
vec![PartitionId::new("conv-a")],
None,
PageCompleteness::Complete,
))
}
async fn release_directory_snapshot(
&self,
_operation: &CoreOperationContext,
_request: ReleaseJournalDirectorySnapshot,
) -> Result<(), CoreResolutionError> {
Ok(())
}
async fn source_head(
&self,
_operation: &CoreOperationContext,
request: GetJournalSource,
) -> Result<Option<JournalSourceHead>, CoreResolutionError> {
if *lock(&self.outage) {
return Err(StateError::Unavailable {
family: polyc_state::journal::family(),
reach: OutageReach::NoDurableEffect,
}
.into());
}
let source = lock(&self.source).clone();
Ok((source.source().partition() == request.partition()).then_some(source))
}
async fn resolve_manifest(
&self,
_operation: &CoreOperationContext,
request: ResolveManifest,
) -> Result<ProjectionResolution, CoreResolutionError> {
let current = if request.key() == self.manifest.key() {
ProjectionHead::Current(Box::new(self.manifest.clone()))
} else {
ProjectionHead::Absent
};
Ok(ProjectionResolution::new(
current,
None,
ProjectionGeneration::new(1),
))
}
async fn begin_audit(
&self,
_operation: &CoreOperationContext,
command: BeginQueryAudit,
) -> Result<BeginOutcome<CoreExecutionPermit>, CoreResolutionError> {
if *lock(&self.refuse_audit) {
return Err(CoreResolutionError::InvalidComposition);
}
let context = CallContext::new(
Deadline::at(MonotonicInstant::from_nanos(u64::MAX)),
CancellationToken::new(),
);
Ok(match self.audit.begin(command, &context)? {
BeginOutcome::Granted(permit) => BeginOutcome::Granted(permit.into()),
BeginOutcome::AlreadyRecorded(receipt) => BeginOutcome::AlreadyRecorded(receipt),
})
}
async fn complete_audit(
&self,
operation: &CoreCompletionContext,
command: &CoreCompletionCommand,
) -> Result<Receipt, CoreResolutionError> {
operation.check()?;
self.completion_attempts.fetch_add(1, Ordering::SeqCst);
if *lock(&self.completion_outage) {
return Err(StateError::Unavailable {
family: polyc_state::query_audit::family(),
reach: OutageReach::NeverDispatched,
}
.into());
}
let CoreCompletionCommand::Local(command) = command else {
return Err(CoreResolutionError::InvalidComposition);
};
let receipt = self
.audit
.complete(command.clone(), operation.local_context()?)?;
if std::mem::take(&mut *lock(&self.completion_response_loss)) {
return Err(StateError::AmbiguousOutcome {
command_id: command.metadata().command_id().clone(),
reason: AmbiguityReason::ResponseLost,
}
.into());
}
Ok(receipt)
}
async fn completion_receipt(
&self,
operation: &CoreCompletionContext,
command: &CoreCompletionCommand,
) -> Result<Option<Receipt>, CoreResolutionError> {
operation.check()?;
self.completion_receipt_reads.fetch_add(1, Ordering::SeqCst);
let CoreCompletionCommand::Local(command) = command else {
return Err(CoreResolutionError::InvalidComposition);
};
let Some(receipt) = self.audit.recorded_receipt(
command.query(),
command.namespace(),
AuditPhase::Completion,
)?
else {
return Ok(None);
};
if !receipt.is_deduplicated() || !receipt.answers(command.metadata()) {
return Err(CoreResolutionError::InvalidComposition);
}
let audit = self
.audit
.audit(
ReadQueryAudit::new(command.query().clone(), command.namespace().clone()),
operation.local_context()?,
)?
.ok_or(CoreResolutionError::InvalidComposition)?;
if audit.completion() != Some(command.completion()) {
return Err(CoreResolutionError::InvalidComposition);
}
Ok(Some(receipt))
}
}
fn live_context() -> CallContext {
CallContext::new(
Deadline::at(MonotonicInstant::from_nanos(u64::MAX)),
CancellationToken::new(),
)
}
#[derive(Debug)]
pub(super) struct TestRevalidator {
current: Mutex<QueryScope>,
}
impl TestRevalidator {
fn new() -> Self {
Self {
current: Mutex::new(QueryScope::Fleet),
}
}
pub(super) fn set(&self, scope: QueryScope) {
*lock(&self.current) = scope;
}
}
#[async_trait]
impl CoreScopeRevalidator for TestRevalidator {
async fn current_scope(
&self,
operation: &CoreOperationContext,
) -> Result<QueryScope, CoreResolutionError> {
operation.check()?;
Ok(lock(&self.current).clone())
}
}
impl super::super::sealed::CredentialProven for TestRevalidator {}
pub(super) struct Harness {
pub(super) store: Arc<TestStore>,
pub(super) state: Arc<TestState>,
pub(super) revalidator: Arc<TestRevalidator>,
visible: Arc<dyn VisibleArtifactAccess>,
fleet: Arc<dyn FleetArtifactAccess>,
trust: Arc<FixtureManifestTrust>,
topology: RealmTopology,
admission: Arc<CoreExecutionAdmission>,
compiler: CatalogCompiler,
planning: CorePlanningAuthority,
runtime: Arc<RuntimeEnv>,
}
impl Harness {
pub(super) fn new() -> Self {
Self::with_policy(ProjectedCorePolicy::default(), 8)
}
pub(super) fn with_release_bounds(result: u64, frame: u64) -> Self {
let policy = ProjectedCorePolicy::try_from(ProjectedCorePolicyInput {
result_release_bytes: result,
response_frame_bytes: frame,
manifest_bytes: 4 * 1024 * 1024,
artifact_file_bytes: 16 * 1024 * 1024,
artifact_range_bytes: 1024 * 1024,
source_decode_bytes: 64 * 1024 * 1024,
})
.unwrap();
Self::with_policy(policy, 8)
}
pub(super) fn with_concurrency(max_concurrent_executions: usize) -> Self {
Self::with_policy(ProjectedCorePolicy::default(), max_concurrent_executions)
}
pub(super) fn with_source_decode_bound(source_decode_bytes: u64) -> Self {
let policy = ProjectedCorePolicy::try_from(ProjectedCorePolicyInput {
result_release_bytes: 32 * 1024 * 1024,
response_frame_bytes: 256 * 1024,
manifest_bytes: 4 * 1024 * 1024,
artifact_file_bytes: 16 * 1024 * 1024,
artifact_range_bytes: 1024 * 1024,
source_decode_bytes,
})
.unwrap();
Self::with_policy(policy, 8)
}
pub(super) fn with_runtime_memory(memory_bytes: usize) -> Self {
Self::with_policy_and_memory(ProjectedCorePolicy::default(), 8, memory_bytes)
}
fn with_policy(policy: ProjectedCorePolicy, max_concurrent_executions: usize) -> Self {
Self::with_policy_and_memory(policy, max_concurrent_executions, 64 * 1024 * 1024)
}
fn with_policy_and_memory(
policy: ProjectedCorePolicy,
max_concurrent_executions: usize,
memory_bytes: usize,
) -> Self {
let store = Arc::new(TestStore::new());
let source = fixture_source();
let record = install_generation(&store, &source);
let state = Arc::new(TestState::new(source, record));
let metadata: Arc<dyn CoreMetadataAuthority> = state.clone();
let planning = CorePlanningAuthority::new(
NamespaceId::new("tenant-a"),
OwnerId::new("projector"),
policy,
metadata,
)
.unwrap();
let visible_namespace = ObjectNamespace::try_new(VISIBLE_NAMESPACE).unwrap();
let fleet_namespace = ObjectNamespace::try_new(FLEET_NAMESPACE).unwrap();
let visible: Arc<dyn VisibleArtifactAccess> = Arc::new(
VisibleArtifactReader::try_new(
Backend(Arc::clone(&store)),
[visible_namespace.clone()],
)
.unwrap(),
);
let fleet: Arc<dyn FleetArtifactAccess> = Arc::new(
FleetArtifactReader::try_new(Backend(Arc::clone(&store)), [fleet_namespace.clone()])
.unwrap(),
);
let topology = RealmTopology::try_new([visible_namespace], [fleet_namespace]).unwrap();
let runtime = RuntimeEnvBuilder::new()
.with_memory_pool(Arc::new(FairSpillPool::new(memory_bytes)))
.build_arc()
.unwrap();
let session = SessionStateBuilder::new()
.with_runtime_env(Arc::clone(&runtime))
.with_default_features()
.build();
Self {
store,
state,
revalidator: Arc::new(TestRevalidator::new()),
visible,
fleet,
trust: Arc::new(FixtureManifestTrust::trusting(&[SIGNER])),
topology,
admission: Arc::new(
CoreExecutionAdmission::try_from(CoreExecutionAdmissionInput {
max_concurrent_executions,
})
.unwrap(),
),
compiler: CatalogCompiler::new(session),
planning,
runtime,
}
}
pub(super) async fn await_completion(
&self,
query: &str,
) -> polyc_state::query_audit::QueryCompletion {
tokio::time::timeout(Duration::from_secs(1), async {
loop {
if let Some(completion) = self.state.completion_for(query) {
break completion;
}
tokio::time::sleep(Duration::from_millis(1)).await;
}
})
.await
.unwrap_or_else(|_| panic!("the guardian settles a terminal for {query}"))
}
pub(super) fn assert_unmatched(&self, query: &str) {
assert!(
self.state.completion_for(query).is_none(),
"{query} must leave its intent unmatched for State's reconciler"
);
}
pub(super) fn reserved_memory(&self) -> usize {
self.runtime.memory_pool.reserved()
}
pub(super) fn visible_message_decode_reservations() -> Vec<u64> {
let scratch = TestStore::new();
artifact_manifest(&scratch, &fixture_source())
.table(CONVERSATION_MESSAGES.as_str())
.unwrap()
.segments()
.iter()
.filter(|segment| segment.realm() == AccessRealm::Visible)
.map(|segment| {
source_decode_reservation_bytes(segment.byte_len(), segment.physical()).unwrap()
})
.collect()
}
pub(super) fn visible_authority(&self) -> CoreArtifactAuthority {
let scope: Arc<dyn CurrentCredentialAuthority> = self.revalidator.clone();
CoreArtifactAuthority::visible(
Arc::clone(&self.visible),
self.trust.clone(),
self.topology.clone(),
scope,
Duration::from_millis(1),
Arc::clone(&self.admission),
)
.unwrap()
}
pub(super) fn fleet_authority(&self) -> CoreArtifactAuthority {
let scope: Arc<dyn CurrentCredentialAuthority> = self.revalidator.clone();
CoreArtifactAuthority::fleet(
Arc::clone(&self.visible),
Arc::clone(&self.fleet),
self.trust.clone(),
self.topology.clone(),
scope,
Duration::from_millis(1),
Arc::clone(&self.admission),
)
.unwrap()
}
pub(super) fn wrong_generation(&self) {
self.store.return_wrong_generation();
}
pub(super) fn wrong_length(&self) {
self.store.return_wrong_length();
}
pub(super) async fn prepare(
&self,
query: &str,
sql: &str,
scope: QueryScope,
row_cap: usize,
) -> PreparedCoreQuery {
self.prepare_with(query, sql, scope, row_cap, Vec::new(), None)
.await
}
pub(super) async fn prepare_with_parameters(
&self,
query: &str,
sql: &str,
scope: QueryScope,
parameters: Vec<CoreParameter>,
) -> PreparedCoreQuery {
self.prepare_with(query, sql, scope, 10, parameters, None)
.await
}
pub(super) async fn prepare_with_timeout(
&self,
query: &str,
sql: &str,
scope: QueryScope,
timeout: Duration,
) -> PreparedCoreQuery {
self.prepare_with(query, sql, scope, 10, Vec::new(), Some(timeout))
.await
}
async fn prepare_with(
&self,
query: &str,
sql: &str,
scope: QueryScope,
row_cap: usize,
parameters: Vec<CoreParameter>,
timeout: Option<Duration>,
) -> PreparedCoreQuery {
match self
.try_prepare_with(query, sql, scope, row_cap, parameters, timeout)
.await
.expect("planning succeeds")
{
CorePlanOutcome::Granted(prepared) => *prepared,
CorePlanOutcome::AlreadyRecorded(_) => panic!("query id is unique"),
}
}
pub(super) async fn try_prepare(
&self,
query: &str,
sql: &str,
scope: QueryScope,
row_cap: usize,
) -> Result<CorePlanOutcome, CoreResolutionError> {
self.try_prepare_with(query, sql, scope, row_cap, Vec::new(), None)
.await
}
async fn try_prepare_with(
&self,
query: &str,
sql: &str,
scope: QueryScope,
row_cap: usize,
parameters: Vec<CoreParameter>,
timeout: Option<Duration>,
) -> Result<CorePlanOutcome, CoreResolutionError> {
let defaults = QueryLimits::default();
let limits = QueryLimits {
row_cap,
timeout: timeout.unwrap_or(defaults.timeout),
..defaults
};
let requested = CoreRequestedBounds::unbounded();
let bounds = self.planning.effective_bounds(&limits, requested)?;
let declared = DeclaredCall::live(Audience::new("state"), limits.timeout);
let audit = CoreAuditContext::from_scoped(
QueryId::new(query),
RequesterId::new("persona:verified-a"),
&declared,
bounds,
);
self.planning
.plan(
&self.compiler,
&limits,
&scope,
false,
audit,
CoreQueryRequest::new(
sql.to_owned(),
parameters,
CoreConsistency::Projected,
requested,
),
)
.await
}
pub(super) async fn collect_text(
&self,
authority: CoreArtifactAuthority,
prepared: PreparedCoreQuery,
) -> Vec<String> {
let mut stream = 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 text = batch
.column(0)
.as_any()
.downcast_ref::<StringArray>()
.unwrap();
rows.extend((0..text.len()).map(|index| text.value(index).to_owned()));
}
rows
}
}
fn fixture_source() -> JournalSource {
JournalSource::new(
PartitionId::new("conv-a"),
PartitionIncarnation::from_bytes([1; PartitionIncarnation::LEN]),
)
}
fn install_generation(store: &TestStore, source: &JournalSource) -> ProjectionManifest {
let manifest = artifact_manifest(store, source);
let (signature, signer_identity) = SIGNER.sign(&manifest.canonical_bytes());
let signed_manifest =
SignedArtifactManifest::try_new(manifest.clone(), signature, signer_identity).unwrap();
let bytes = signed_manifest.envelope_bytes();
let reference = ExactObjectRef::try_new(
ObjectNamespace::try_new(VISIBLE_NAMESPACE).unwrap(),
ContentReference::try_new(MANIFEST_KEY).unwrap(),
20,
)
.unwrap();
store.insert(&reference, bytes.clone());
ProjectionManifest::new(
manifest.key().clone(),
manifest.generation(),
manifest.checkpoint().clone(),
manifest.schema_version(),
manifest.fact_version(),
ObjectDescriptor::new(
manifest.key().object(),
ObjectGeneration::new(manifest.generation().get()),
signed_manifest.envelope_digest(),
OwnerId::new("projector"),
manifest.classification(),
manifest.retention(),
u64::try_from(bytes.len()).unwrap(),
reference.key().clone(),
),
reference,
manifest.publisher().clone(),
manifest.fence().clone(),
)
}
fn artifact_manifest(store: &TestStore, source: &JournalSource) -> ArtifactManifest {
let turns = parquet_turns(source);
let visible_messages_a = parquet_messages(
source,
&[
(2, "visible-one", "user", false),
(3, "visible-two", "assistant", false),
],
);
let visible_messages_b = parquet_messages(source, &[(4, "visible-three", "user", false)]);
let fleet_messages = parquet_messages(source, &[(5, "fleet-secret", "system", true)]);
let turns_segment = segment(
store,
AccessRealm::Visible,
VISIBLE_NAMESPACE,
"turns/visible-g1.parquet",
10,
turns,
1,
None,
);
let visible_segment = segment(
store,
AccessRealm::Visible,
VISIBLE_NAMESPACE,
VISIBLE_MESSAGES_KEY,
11,
visible_messages_a,
2,
Some(SourceBounds::try_new(2, 3).unwrap()),
);
let visible_segment_b = segment(
store,
AccessRealm::Visible,
VISIBLE_NAMESPACE,
"messages/visible-b-g1.parquet",
13,
visible_messages_b,
1,
Some(SourceBounds::try_new(4, 4).unwrap()),
);
let fleet_segment = segment(
store,
AccessRealm::FleetOnly,
FLEET_NAMESPACE,
"messages/fleet-g1.parquet",
12,
fleet_messages,
1,
Some(SourceBounds::try_new(5, 5).unwrap()),
);
let family = conversation_core();
let fingerprint = |table| family.table(table).unwrap().fingerprint();
let key = ProjectionKey::new(
FamilyId::new(family.family_str()),
source.partition().clone(),
);
ArtifactManifest::try_new(
key.clone(),
ObjectNamespace::try_new(VISIBLE_NAMESPACE).unwrap(),
ProjectionGeneration::new(1),
checkpoint(source),
family.versions().schema().get(),
family.versions().fact_model().get(),
family.fingerprint(),
Classification::Confidential,
Retention::For(Duration::from_mins(10)),
PublisherId::new("projector-a"),
PublisherFence::new(key, source.incarnation(), 1),
vec![
ArtifactTable::try_new(
CONVERSATION_TURNS.as_str(),
fingerprint(CONVERSATION_TURNS),
vec![turns_segment],
)
.unwrap(),
ArtifactTable::try_new(
CONVERSATION_MESSAGES.as_str(),
fingerprint(CONVERSATION_MESSAGES),
vec![visible_segment, visible_segment_b, fleet_segment],
)
.unwrap(),
],
)
.unwrap()
}
fn checkpoint(source: &JournalSource) -> SourceCheckpoint {
SourceCheckpoint::try_new(
source.clone(),
JournalPosition::new(1),
JournalPosition::new(10),
10,
JournalAttestation::new(
CommitRoot::from_bytes([7; CommitRoot::LEN]),
11,
vec![8; 64],
vec![9; 32],
),
)
.unwrap()
}
#[allow(clippy::too_many_arguments)]
fn segment(
store: &TestStore,
realm: AccessRealm,
namespace: &str,
key: &str,
generation: u64,
bytes: Vec<u8>,
rows: u64,
bounds: Option<SourceBounds>,
) -> ArtifactSegment {
let reference = ExactObjectRef::try_new(
ObjectNamespace::try_new(namespace).unwrap(),
ContentReference::try_new(key).unwrap(),
generation,
)
.unwrap();
let digest = digest(&bytes);
let byte_len = u64::try_from(bytes.len()).unwrap();
let table = if key.starts_with("turns/") {
CONVERSATION_TURNS
} else {
CONVERSATION_MESSAGES
};
let schema = crate::core_resolution::arrow_schema(conversation_core().table(table).unwrap());
let physical = polyc_projection_artifact::parquet_profile::inspect(
&Bytes::from(bytes.clone()),
schema.as_ref(),
rows,
)
.unwrap();
store.insert(&reference, bytes);
ArtifactSegment::try_new(realm, reference, byte_len, digest, rows, physical, bounds).unwrap()
}
fn digest(bytes: &[u8]) -> ContentDigest {
let mut raw = [0; ContentDigest::LEN];
raw.copy_from_slice(&Sha256::digest(bytes));
ContentDigest::from_bytes(raw)
}
fn parquet_turns(source: &JournalSource) -> Vec<u8> {
let schema = crate::core_resolution::arrow_schema(
conversation_core().table(CONVERSATION_TURNS).unwrap(),
);
let columns: Vec<ArrayRef> = vec![
Arc::new(StringArray::from(vec![source.partition().as_str()])),
Arc::new(incarnations(source, 1)),
Arc::new(StringArray::from(vec!["turn-a"])),
Arc::new(UInt64Array::from(vec![1])),
Arc::new(UInt64Array::from(vec![1])),
Arc::new(UInt64Array::from(vec![6])),
];
write_parquet(&RecordBatch::try_new(schema, columns).unwrap())
}
fn parquet_messages(source: &JournalSource, rows: &[(u64, &str, &str, bool)]) -> Vec<u8> {
let schema = crate::core_resolution::arrow_schema(
conversation_core().table(CONVERSATION_MESSAGES).unwrap(),
);
let columns: Vec<ArrayRef> = vec![
Arc::new(StringArray::from(vec![
source.partition().as_str();
rows.len()
])),
Arc::new(incarnations(source, rows.len())),
Arc::new(UInt64Array::from(
rows.iter().map(|row| row.0).collect::<Vec<_>>(),
)),
Arc::new(StringArray::from(vec!["turn-a"; rows.len()])),
Arc::new(StringArray::from(
rows.iter().map(|row| row.2).collect::<Vec<_>>(),
)),
Arc::new(BooleanArray::from(
rows.iter().map(|row| row.3).collect::<Vec<_>>(),
)),
Arc::new(StringArray::from(
rows.iter().map(|row| row.1).collect::<Vec<_>>(),
)),
Arc::new(StringArray::from(vec!["trusted"; rows.len()])),
];
write_parquet(&RecordBatch::try_new(schema, columns).unwrap())
}
fn incarnations(source: &JournalSource, rows: usize) -> arrow::array::FixedSizeBinaryArray {
let mut builder = FixedSizeBinaryBuilder::with_capacity(rows, 32);
for _ in 0..rows {
builder
.append_value(source.incarnation().as_bytes())
.unwrap();
}
builder.finish()
}
fn write_parquet(batch: &RecordBatch) -> Vec<u8> {
let mut bytes = Vec::new();
let mut writer = ArrowWriter::try_new(
&mut bytes,
batch.schema(),
Some(polyc_projection_artifact::parquet_profile::writer_properties()),
)
.unwrap();
writer.write(batch).unwrap();
writer.close().unwrap();
canonicalize_writer_output(&mut bytes).unwrap();
bytes
}