use crate::{
db::{
commit::CommitSchemaFingerprint,
data::{CanonicalSlotReader, StructuralRowContract},
index::{
IndexEntryValue, IndexId, IndexKey, IndexKeyKind, IndexState, IndexStore,
IndexStoreVisit, RawIndexStoreKey,
},
key_taxonomy::PrimaryKeyValue,
predicate::{PredicateProgram, normalize, parse_sql_predicate},
schema::{
AcceptedCatalogIdentity, PersistedSchemaSnapshot, SchemaExpressionIndexRebuildTarget,
SchemaFieldPathIndexRebuildTarget, SchemaVersion,
accepted_schema_cache_fingerprint_for_persisted_snapshot,
mutation::SchemaMutationRequest,
},
},
error::InternalError,
types::EntityTag,
};
use std::{collections::BTreeSet, mem::size_of};
const MAX_SOURCE_ROWS: usize = 65_536;
const MAX_SOURCE_ROW_BYTES: usize = 256 * 1024 * 1024;
const MAX_PROJECTION_ENTRIES: usize = 131_072;
const MAX_DELETION_KEYS: usize = 65_536;
const MAX_STAGED_RAW_BYTES: usize = 256 * 1024 * 1024;
const MAX_PROJECTION_WORK_UNITS: usize = 262_144;
#[derive(Clone, Copy)]
pub(in crate::db) struct SchemaUserIndexDomainRow<'a> {
primary_key_value: PrimaryKeyValue,
slots: &'a dyn CanonicalSlotReader,
encoded_row_bytes: usize,
}
impl<'a> SchemaUserIndexDomainRow<'a> {
#[must_use]
pub(in crate::db) fn new(
primary_key_value: impl Into<PrimaryKeyValue>,
slots: &'a dyn CanonicalSlotReader,
encoded_row_bytes: usize,
) -> Self {
Self {
primary_key_value: primary_key_value.into(),
slots,
encoded_row_bytes,
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub(in crate::db) struct StagedUserIndexDomainEntry {
key: RawIndexStoreKey,
value: IndexEntryValue,
}
impl StagedUserIndexDomainEntry {
#[must_use]
pub(in crate::db) const fn key(&self) -> &RawIndexStoreKey {
&self.key
}
#[cfg(test)]
#[must_use]
pub(in crate::db) const fn value(&self) -> &IndexEntryValue {
&self.value
}
pub(in crate::db) fn into_parts(self) -> (RawIndexStoreKey, IndexEntryValue) {
(self.key, self.value)
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(in crate::db) struct StagedUserIndexDomainUsage {
source_rows: usize,
source_row_bytes: usize,
accepted_before_entries: usize,
accepted_after_entries: usize,
projection_entries: usize,
deletion_keys: usize,
staged_raw_bytes: usize,
projection_work_units: usize,
}
impl StagedUserIndexDomainUsage {
#[cfg(any(test, feature = "sql"))]
#[must_use]
pub(in crate::db) const fn source_rows(self) -> usize {
self.source_rows
}
#[cfg(any(test, feature = "sql"))]
#[must_use]
pub(in crate::db) const fn accepted_before_entries(self) -> usize {
self.accepted_before_entries
}
#[cfg(any(test, feature = "sql"))]
#[must_use]
pub(in crate::db) const fn accepted_after_entries(self) -> usize {
self.accepted_after_entries
}
#[cfg(test)]
#[must_use]
pub(in crate::db) const fn staged_raw_bytes(self) -> usize {
self.staged_raw_bytes
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub(in crate::db) struct StagedUserIndexDomainReplacement {
store_path: &'static str,
entity_tag: EntityTag,
accepted_before_identity: AcceptedCatalogIdentity,
accepted_after_version: SchemaVersion,
accepted_after_fingerprint: CommitSchemaFingerprint,
deletion_keys: Vec<RawIndexStoreKey>,
final_entries: Vec<StagedUserIndexDomainEntry>,
usage: StagedUserIndexDomainUsage,
}
impl StagedUserIndexDomainReplacement {
#[must_use]
pub(in crate::db) const fn store_path(&self) -> &'static str {
self.store_path
}
#[must_use]
pub(in crate::db) const fn entity_tag(&self) -> EntityTag {
self.entity_tag
}
#[must_use]
pub(in crate::db) const fn accepted_before_identity(&self) -> AcceptedCatalogIdentity {
self.accepted_before_identity
}
#[must_use]
pub(in crate::db) const fn accepted_after_version(&self) -> SchemaVersion {
self.accepted_after_version
}
#[must_use]
pub(in crate::db) const fn accepted_after_fingerprint(&self) -> CommitSchemaFingerprint {
self.accepted_after_fingerprint
}
#[cfg(test)]
#[must_use]
pub(in crate::db) const fn deletion_keys(&self) -> &[RawIndexStoreKey] {
self.deletion_keys.as_slice()
}
#[cfg(test)]
#[must_use]
pub(in crate::db) const fn final_entries(&self) -> &[StagedUserIndexDomainEntry] {
self.final_entries.as_slice()
}
#[cfg(any(test, feature = "sql"))]
#[must_use]
pub(in crate::db) const fn usage(&self) -> StagedUserIndexDomainUsage {
self.usage
}
pub(in crate::db) fn into_apply_parts(
self,
) -> (Vec<RawIndexStoreKey>, Vec<StagedUserIndexDomainEntry>) {
(self.deletion_keys, self.final_entries)
}
}
pub(in crate::db) struct StagedUserIndexDomainReplacementBuilder {
store_path: &'static str,
entity_tag: EntityTag,
accepted_before_identity: AcceptedCatalogIdentity,
accepted_after_version: SchemaVersion,
accepted_after_fingerprint: CommitSchemaFingerprint,
before_projection: PreparedUserIndexProjection,
after_projection: PreparedUserIndexProjection,
expected_before: Vec<StagedUserIndexDomainEntry>,
final_entries: Vec<StagedUserIndexDomainEntry>,
budget: StagedUserIndexDomainBudget,
}
impl StagedUserIndexDomainReplacementBuilder {
pub(in crate::db) fn new(
accepted_before_identity: AcceptedCatalogIdentity,
accepted_before: &PersistedSchemaSnapshot,
accepted_after: &PersistedSchemaSnapshot,
predicate_row_contract: Option<&StructuralRowContract>,
index_store: &IndexStore,
) -> Result<Self, StagedUserIndexDomainError> {
validate_stage_authority(
accepted_before_identity,
accepted_before,
accepted_after,
predicate_row_contract,
index_store,
)?;
let entity_tag = accepted_before_identity.entity_tag();
let before_projection = PreparedUserIndexProjection::from_snapshot(
entity_tag,
accepted_before,
predicate_row_contract,
)?;
let after_projection = PreparedUserIndexProjection::from_snapshot(
entity_tag,
accepted_after,
predicate_row_contract,
)?;
Ok(Self {
store_path: accepted_before_identity.store_path(),
entity_tag,
accepted_before_identity,
accepted_after_version: accepted_after.version(),
accepted_after_fingerprint: accepted_schema_cache_fingerprint_for_persisted_snapshot(
accepted_after,
)
.map_err(StagedUserIndexDomainError::Fingerprint)?,
before_projection,
after_projection,
expected_before: Vec::new(),
final_entries: Vec::new(),
budget: StagedUserIndexDomainBudget::standard(),
})
}
pub(in crate::db) fn observe_row(
&mut self,
row: &SchemaUserIndexDomainRow<'_>,
) -> Result<(), StagedUserIndexDomainError> {
self.budget.consume_source_row(row.encoded_row_bytes)?;
self.before_projection.derive_row(
self.entity_tag,
row,
&mut self.expected_before,
&mut self.budget,
)?;
self.after_projection.derive_row(
self.entity_tag,
row,
&mut self.final_entries,
&mut self.budget,
)
}
pub(in crate::db) fn finish(
mut self,
index_store: &IndexStore,
) -> Result<StagedUserIndexDomainReplacement, StagedUserIndexDomainError> {
if index_store.state() != IndexState::Ready {
return Err(StagedUserIndexDomainError::IndexStoreNotReady);
}
validate_projection(
&mut self.expected_before,
&self.before_projection.unique_index_ids,
ProjectionAuthority::AcceptedBefore,
self.accepted_before_identity.entity_path(),
)?;
validate_projection(
&mut self.final_entries,
&self.after_projection.unique_index_ids,
ProjectionAuthority::CandidateAfter,
self.accepted_before_identity.entity_path(),
)?;
let observed_before =
observe_current_user_index_domain(index_store, self.entity_tag, &mut self.budget)?;
if observed_before != self.expected_before {
return Err(StagedUserIndexDomainError::CurrentDomainMismatch);
}
let deletion_keys = observed_before
.into_iter()
.map(|entry| entry.key)
.collect::<Vec<_>>();
validate_insertion_collisions(index_store, &deletion_keys, &self.final_entries)?;
self.budget.finish_sort_workspace(
self.expected_before.len(),
self.final_entries.len(),
deletion_keys.len(),
)?;
self.budget
.record_projection_counts(self.expected_before.len(), self.final_entries.len());
Ok(StagedUserIndexDomainReplacement {
store_path: self.store_path,
entity_tag: self.entity_tag,
accepted_before_identity: self.accepted_before_identity,
accepted_after_version: self.accepted_after_version,
accepted_after_fingerprint: self.accepted_after_fingerprint,
deletion_keys,
final_entries: self.final_entries,
usage: self.budget.usage(),
})
}
}
pub(in crate::db) enum StagedUserIndexDomainError {
AcceptedAfterEntityMismatch,
AcceptedBeforeIdentityMismatch,
AcceptedIndexStoreMismatch,
CurrentDomainMismatch,
DeletionKeyLimitExceeded,
DuplicateRawKey,
DuplicateUniqueKey,
CandidateUniqueConflict { entity_path: &'static str },
Fingerprint(InternalError),
InsertionCollision,
IndexStoreNotReady,
KeyDerivation(InternalError),
KeyEncode,
MissingPredicateRowContract,
PhysicalKeyDecode,
PredicateEvaluation(InternalError),
PredicateParse,
ProjectionEntryLimitExceeded,
ProjectionWorkLimitExceeded,
RowContractEntityMismatch,
SourceRowBytesLimitExceeded,
SourceRowLimitExceeded,
StagedRawBytesLimitExceeded,
UnsupportedAcceptedIndex,
}
impl StagedUserIndexDomainError {
pub(in crate::db) fn into_internal_error(self) -> InternalError {
match self {
Self::Fingerprint(error)
| Self::KeyDerivation(error)
| Self::PredicateEvaluation(error) => error,
Self::CandidateUniqueConflict { entity_path } => {
InternalError::index_violation(entity_path, &[])
}
Self::CurrentDomainMismatch
| Self::DuplicateRawKey
| Self::DuplicateUniqueKey
| Self::InsertionCollision
| Self::PhysicalKeyDecode => InternalError::store_corruption(),
Self::AcceptedAfterEntityMismatch
| Self::AcceptedBeforeIdentityMismatch
| Self::AcceptedIndexStoreMismatch
| Self::RowContractEntityMismatch => InternalError::store_invariant(),
Self::DeletionKeyLimitExceeded
| Self::IndexStoreNotReady
| Self::KeyEncode
| Self::MissingPredicateRowContract
| Self::PredicateParse
| Self::ProjectionEntryLimitExceeded
| Self::ProjectionWorkLimitExceeded
| Self::SourceRowBytesLimitExceeded
| Self::SourceRowLimitExceeded
| Self::StagedRawBytesLimitExceeded
| Self::UnsupportedAcceptedIndex => InternalError::store_unsupported(),
}
}
}
enum PreparedUserIndexTarget {
FieldPath(SchemaFieldPathIndexRebuildTarget),
Expression(SchemaExpressionIndexRebuildTarget),
}
struct PreparedUserIndex {
target: PreparedUserIndexTarget,
predicate: Option<PredicateProgram>,
}
impl PreparedUserIndex {
fn derive_key(
&self,
entity_tag: EntityTag,
row: &SchemaUserIndexDomainRow<'_>,
) -> Result<Option<IndexKey>, StagedUserIndexDomainError> {
if let Some(predicate) = self.predicate.as_ref()
&& !predicate
.eval_with_structural_slot_reader(row.slots)
.map_err(StagedUserIndexDomainError::PredicateEvaluation)?
{
return Ok(None);
}
match &self.target {
PreparedUserIndexTarget::FieldPath(target) => {
IndexKey::new_from_slots_with_field_path_rebuild_target(
entity_tag,
row.primary_key_value,
target,
row.slots,
)
}
PreparedUserIndexTarget::Expression(target) => {
IndexKey::new_from_slots_with_expression_rebuild_target(
entity_tag,
row.primary_key_value,
target,
row.slots,
)
}
}
.map_err(StagedUserIndexDomainError::KeyDerivation)
}
}
#[derive(Clone, Copy)]
enum ProjectionAuthority {
AcceptedBefore,
CandidateAfter,
}
struct PreparedUserIndexProjection {
indexes: Vec<PreparedUserIndex>,
unique_index_ids: BTreeSet<IndexId>,
}
impl PreparedUserIndexProjection {
fn from_snapshot(
entity_tag: EntityTag,
snapshot: &PersistedSchemaSnapshot,
predicate_row_contract: Option<&StructuralRowContract>,
) -> Result<Self, StagedUserIndexDomainError> {
let mut indexes = Vec::with_capacity(snapshot.indexes().len());
let mut unique_index_ids = BTreeSet::new();
for index in snapshot.indexes() {
let request = if index.key().is_field_path_only() {
SchemaMutationRequest::from_accepted_field_path_index(index)
} else {
SchemaMutationRequest::from_accepted_expression_index(index)
}
.map_err(|_| StagedUserIndexDomainError::UnsupportedAcceptedIndex)?;
let target = match request {
SchemaMutationRequest::AddFieldPathIndex { target } => {
PreparedUserIndexTarget::FieldPath(target)
}
SchemaMutationRequest::AddExpressionIndex { target } => {
PreparedUserIndexTarget::Expression(target)
}
SchemaMutationRequest::ExactMatch | SchemaMutationRequest::AppendOnlyFields(_) => {
return Err(StagedUserIndexDomainError::UnsupportedAcceptedIndex);
}
};
let predicate = index
.predicate_sql()
.map(|sql| {
let row_contract = predicate_row_contract
.ok_or(StagedUserIndexDomainError::MissingPredicateRowContract)?;
parse_sql_predicate(sql)
.map(|predicate| {
PredicateProgram::compile_with_row_contract(
row_contract,
&normalize(&predicate),
)
})
.map_err(|_| StagedUserIndexDomainError::PredicateParse)
})
.transpose()?;
if index.unique() {
unique_index_ids.insert(IndexId::new(entity_tag, index.ordinal()));
}
indexes.push(PreparedUserIndex { target, predicate });
}
Ok(Self {
indexes,
unique_index_ids,
})
}
fn derive_row(
&self,
entity_tag: EntityTag,
row: &SchemaUserIndexDomainRow<'_>,
entries: &mut Vec<StagedUserIndexDomainEntry>,
budget: &mut StagedUserIndexDomainBudget,
) -> Result<(), StagedUserIndexDomainError> {
for index in &self.indexes {
budget.consume_projection_work()?;
let Some(key) = index.derive_key(entity_tag, row)? else {
continue;
};
let key = key
.to_raw()
.map_err(|_| StagedUserIndexDomainError::KeyEncode)?;
let value = IndexEntryValue::presence();
budget.consume_projection_entry(key.as_bytes().len())?;
entries.push(StagedUserIndexDomainEntry { key, value });
}
Ok(())
}
}
fn validate_stage_authority(
accepted_before_identity: AcceptedCatalogIdentity,
accepted_before: &PersistedSchemaSnapshot,
accepted_after: &PersistedSchemaSnapshot,
predicate_row_contract: Option<&StructuralRowContract>,
index_store: &IndexStore,
) -> Result<(), StagedUserIndexDomainError> {
let accepted_before_fingerprint =
accepted_schema_cache_fingerprint_for_persisted_snapshot(accepted_before)
.map_err(StagedUserIndexDomainError::Fingerprint)?;
let entity_path_matches =
accepted_before_identity.entity_path() == accepted_before.entity_path();
let schema_version_matches =
accepted_before_identity.accepted_schema_version() == accepted_before.version();
let schema_fingerprint_matches =
accepted_before_identity.accepted_schema_fingerprint() == accepted_before_fingerprint;
if !(entity_path_matches && schema_version_matches && schema_fingerprint_matches) {
return Err(StagedUserIndexDomainError::AcceptedBeforeIdentityMismatch);
}
if accepted_after.entity_path() != accepted_before.entity_path() {
return Err(StagedUserIndexDomainError::AcceptedAfterEntityMismatch);
}
let store_path = accepted_before_identity.store_path();
if accepted_before
.indexes()
.iter()
.chain(accepted_after.indexes())
.any(|index| index.store() != store_path)
{
return Err(StagedUserIndexDomainError::AcceptedIndexStoreMismatch);
}
if predicate_row_contract.is_some_and(|row_contract| {
row_contract.entity_path() != accepted_before_identity.entity_path()
}) {
return Err(StagedUserIndexDomainError::RowContractEntityMismatch);
}
if index_store.state() != IndexState::Ready {
return Err(StagedUserIndexDomainError::IndexStoreNotReady);
}
Ok(())
}
fn validate_projection(
entries: &mut [StagedUserIndexDomainEntry],
unique_index_ids: &BTreeSet<IndexId>,
authority: ProjectionAuthority,
entity_path: &'static str,
) -> Result<(), StagedUserIndexDomainError> {
entries.sort_unstable_by(|left, right| left.key.cmp(&right.key));
for pair in entries.windows(2) {
if pair[0].key == pair[1].key {
return Err(StagedUserIndexDomainError::DuplicateRawKey);
}
let left = IndexKey::try_from_raw(&pair[0].key)
.map_err(|_| StagedUserIndexDomainError::PhysicalKeyDecode)?;
let right = IndexKey::try_from_raw(&pair[1].key)
.map_err(|_| StagedUserIndexDomainError::PhysicalKeyDecode)?;
if left.index_id() == right.index_id()
&& unique_index_ids.contains(left.index_id())
&& left.has_same_components(&right)
{
return Err(match authority {
ProjectionAuthority::AcceptedBefore => {
StagedUserIndexDomainError::DuplicateUniqueKey
}
ProjectionAuthority::CandidateAfter => {
StagedUserIndexDomainError::CandidateUniqueConflict { entity_path }
}
});
}
}
Ok(())
}
fn observe_current_user_index_domain(
index_store: &IndexStore,
entity_tag: EntityTag,
budget: &mut StagedUserIndexDomainBudget,
) -> Result<Vec<StagedUserIndexDomainEntry>, StagedUserIndexDomainError> {
let mut entries = Vec::new();
let result = index_store.visit_entries(|raw_key, value| {
budget.consume_projection_work()?;
let key = IndexKey::try_from_raw(raw_key)
.map_err(|_| StagedUserIndexDomainError::PhysicalKeyDecode)?;
if key.key_kind() == IndexKeyKind::User && key.index_id().entity_tag() == entity_tag {
budget.consume_deletion_key(raw_key.as_bytes().len())?;
entries.push(StagedUserIndexDomainEntry {
key: raw_key.clone(),
value: value.clone(),
});
}
Ok(IndexStoreVisit::Continue)
});
result?;
Ok(entries)
}
fn validate_insertion_collisions(
index_store: &IndexStore,
deletion_keys: &[RawIndexStoreKey],
final_entries: &[StagedUserIndexDomainEntry],
) -> Result<(), StagedUserIndexDomainError> {
for entry in final_entries {
if index_store.get(entry.key()).is_some()
&& deletion_keys.binary_search(entry.key()).is_err()
{
return Err(StagedUserIndexDomainError::InsertionCollision);
}
}
Ok(())
}
struct StagedUserIndexDomainBudget {
usage: StagedUserIndexDomainUsage,
}
impl StagedUserIndexDomainBudget {
const fn standard() -> Self {
Self {
usage: StagedUserIndexDomainUsage {
source_rows: 0,
source_row_bytes: 0,
accepted_before_entries: 0,
accepted_after_entries: 0,
projection_entries: 0,
deletion_keys: 0,
staged_raw_bytes: 0,
projection_work_units: 0,
},
}
}
fn consume_source_row(
&mut self,
encoded_row_bytes: usize,
) -> Result<(), StagedUserIndexDomainError> {
self.usage.source_rows = self
.usage
.source_rows
.checked_add(1)
.ok_or(StagedUserIndexDomainError::SourceRowLimitExceeded)?;
if self.usage.source_rows > MAX_SOURCE_ROWS {
return Err(StagedUserIndexDomainError::SourceRowLimitExceeded);
}
self.usage.source_row_bytes = self
.usage
.source_row_bytes
.checked_add(encoded_row_bytes)
.ok_or(StagedUserIndexDomainError::SourceRowBytesLimitExceeded)?;
if self.usage.source_row_bytes > MAX_SOURCE_ROW_BYTES {
return Err(StagedUserIndexDomainError::SourceRowBytesLimitExceeded);
}
Ok(())
}
fn consume_projection_work(&mut self) -> Result<(), StagedUserIndexDomainError> {
self.usage.projection_work_units = self
.usage
.projection_work_units
.checked_add(1)
.ok_or(StagedUserIndexDomainError::ProjectionWorkLimitExceeded)?;
if self.usage.projection_work_units > MAX_PROJECTION_WORK_UNITS {
return Err(StagedUserIndexDomainError::ProjectionWorkLimitExceeded);
}
Ok(())
}
fn consume_projection_entry(
&mut self,
key_bytes: usize,
) -> Result<(), StagedUserIndexDomainError> {
self.usage.projection_entries = self
.usage
.projection_entries
.checked_add(1)
.ok_or(StagedUserIndexDomainError::ProjectionEntryLimitExceeded)?;
if self.usage.projection_entries > MAX_PROJECTION_ENTRIES {
return Err(StagedUserIndexDomainError::ProjectionEntryLimitExceeded);
}
self.consume_staged_bytes(
key_bytes
.checked_add(1)
.and_then(|bytes| bytes.checked_add(size_of::<StagedUserIndexDomainEntry>()))
.ok_or(StagedUserIndexDomainError::StagedRawBytesLimitExceeded)?,
)
}
fn consume_deletion_key(&mut self, key_bytes: usize) -> Result<(), StagedUserIndexDomainError> {
self.usage.deletion_keys = self
.usage
.deletion_keys
.checked_add(1)
.ok_or(StagedUserIndexDomainError::DeletionKeyLimitExceeded)?;
if self.usage.deletion_keys > MAX_DELETION_KEYS {
return Err(StagedUserIndexDomainError::DeletionKeyLimitExceeded);
}
self.consume_staged_bytes(
key_bytes
.checked_add(1)
.and_then(|bytes| bytes.checked_add(size_of::<StagedUserIndexDomainEntry>()))
.ok_or(StagedUserIndexDomainError::StagedRawBytesLimitExceeded)?,
)
}
fn finish_sort_workspace(
&mut self,
before_entries: usize,
after_entries: usize,
deletion_keys: usize,
) -> Result<(), StagedUserIndexDomainError> {
let entry_workspace = before_entries
.checked_add(after_entries)
.and_then(|count| count.checked_mul(size_of::<StagedUserIndexDomainEntry>()))
.ok_or(StagedUserIndexDomainError::StagedRawBytesLimitExceeded)?;
let deletion_workspace = deletion_keys
.checked_mul(size_of::<RawIndexStoreKey>())
.ok_or(StagedUserIndexDomainError::StagedRawBytesLimitExceeded)?;
self.consume_staged_bytes(
entry_workspace
.checked_add(deletion_workspace)
.ok_or(StagedUserIndexDomainError::StagedRawBytesLimitExceeded)?,
)
}
fn consume_staged_bytes(&mut self, bytes: usize) -> Result<(), StagedUserIndexDomainError> {
self.usage.staged_raw_bytes = self
.usage
.staged_raw_bytes
.checked_add(bytes)
.ok_or(StagedUserIndexDomainError::StagedRawBytesLimitExceeded)?;
if self.usage.staged_raw_bytes > MAX_STAGED_RAW_BYTES {
return Err(StagedUserIndexDomainError::StagedRawBytesLimitExceeded);
}
Ok(())
}
const fn record_projection_counts(
&mut self,
accepted_before_entries: usize,
accepted_after_entries: usize,
) {
self.usage.accepted_before_entries = accepted_before_entries;
self.usage.accepted_after_entries = accepted_after_entries;
}
const fn usage(&self) -> StagedUserIndexDomainUsage {
self.usage
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn standard_budget_fixes_and_admits_the_design_measurement_point() {
let design_rows = std::hint::black_box(2_048usize);
let design_indexes = std::hint::black_box(8usize);
let design_domain_entries = design_rows * design_indexes;
let design_projection_entries = design_domain_entries * 2;
let design_work_units = design_domain_entries * 3;
assert_eq!(MAX_SOURCE_ROWS, 65_536);
assert_eq!(MAX_SOURCE_ROW_BYTES, 256 * 1024 * 1024);
assert_eq!(MAX_PROJECTION_ENTRIES, 131_072);
assert_eq!(MAX_DELETION_KEYS, 65_536);
assert_eq!(MAX_STAGED_RAW_BYTES, 256 * 1024 * 1024);
assert_eq!(MAX_PROJECTION_WORK_UNITS, 262_144);
assert!(MAX_SOURCE_ROWS >= design_rows);
assert!(MAX_PROJECTION_ENTRIES >= design_projection_entries);
assert!(MAX_DELETION_KEYS >= design_domain_entries);
assert!(MAX_PROJECTION_WORK_UNITS >= design_work_units);
}
#[test]
fn standard_budget_rejects_each_aggregate_dimension_at_its_boundary() {
let mut row_count = StagedUserIndexDomainBudget::standard();
for _ in 0..MAX_SOURCE_ROWS {
assert!(row_count.consume_source_row(0).is_ok());
}
assert!(matches!(
row_count.consume_source_row(0),
Err(StagedUserIndexDomainError::SourceRowLimitExceeded),
));
let mut row_bytes = StagedUserIndexDomainBudget::standard();
assert!(matches!(
row_bytes.consume_source_row(MAX_SOURCE_ROW_BYTES + 1),
Err(StagedUserIndexDomainError::SourceRowBytesLimitExceeded),
));
let mut entry_count = StagedUserIndexDomainBudget::standard();
for _ in 0..MAX_PROJECTION_ENTRIES {
assert!(entry_count.consume_projection_entry(0).is_ok());
}
assert!(matches!(
entry_count.consume_projection_entry(0),
Err(StagedUserIndexDomainError::ProjectionEntryLimitExceeded),
));
let mut deletion_count = StagedUserIndexDomainBudget::standard();
for _ in 0..MAX_DELETION_KEYS {
assert!(deletion_count.consume_deletion_key(0).is_ok());
}
assert!(matches!(
deletion_count.consume_deletion_key(0),
Err(StagedUserIndexDomainError::DeletionKeyLimitExceeded),
));
let mut work = StagedUserIndexDomainBudget::standard();
for _ in 0..MAX_PROJECTION_WORK_UNITS {
assert!(work.consume_projection_work().is_ok());
}
assert!(matches!(
work.consume_projection_work(),
Err(StagedUserIndexDomainError::ProjectionWorkLimitExceeded),
));
let mut staged_bytes = StagedUserIndexDomainBudget::standard();
assert!(matches!(
staged_bytes.consume_staged_bytes(MAX_STAGED_RAW_BYTES + 1),
Err(StagedUserIndexDomainError::StagedRawBytesLimitExceeded),
));
}
}