use gnitz_foundation::env::env_num;
use rustc_hash::FxHashMap;
use gnitz_zset::algebra::index_entries;
use gnitz_zset::schema::key::{key_range_between_cuts, KeyCut};
use gnitz_zset::schema::{index_spec_and_schema, KeySpec, SchemaDescriptor};
pub use crate::storage::Cut;
use crate::storage::{RecoverySource, StoreBudgets, Table};
use gnitz_wire::{PkColList, PkKeys, ViewProps};
use gnitz_zset::repr::{Batch, PkSetGather, ReadCursor, StorageError, StoredRow};
use gnitz_zset::schema::{Placement, Slot};
mod build;
mod circuit_state;
mod delta;
mod dirs;
mod disk_usage;
mod ingest;
mod repartition;
mod store;
mod store_lifecycle;
mod unique_pk;
pub use circuit_state::{CircuitState, StateIdx, StateLayout};
pub(crate) use delta::{delta_round, delta_round_prefix};
use dirs::ensure_dir;
pub use dirs::{lock_data_dir, relation_dir, relations_dir, ChildAddr, ChildKind, DirLock};
pub use disk_usage::{disk_usage, DiskUsage};
use store::Store;
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum IndexClaim {
ForeignKey,
Index { id: u64, unique: bool },
}
pub struct SecondaryIndex {
store: Store,
cols: PkColList,
key_spec: KeySpec,
claims: Vec<IndexClaim>,
covers_pk: bool,
}
impl SecondaryIndex {
pub fn cols(&self) -> PkColList {
self.cols
}
pub fn schema(&self) -> SchemaDescriptor {
*self.store.schema()
}
pub fn key_spec(&self) -> KeySpec {
self.key_spec
}
pub fn is_unique(&self) -> bool {
self.claims
.iter()
.any(|c| matches!(c, IndexClaim::Index { unique: true, .. }))
}
pub fn claims(&self) -> &[IndexClaim] {
&self.claims
}
pub fn cursor(&self) -> ReadCursor {
self.store.held().open_cursor()
}
pub(crate) fn gather(&self, spans: PkKeys) -> PkSetGather {
self.store.held().gather(spans, Cut::Now)
}
pub(crate) fn cursor_over(&self, r: &gnitz_wire::KeyRange) -> ReadCursor {
let t = self.store.held();
t.range_cursor(self.key_spec.range_keys(t.schema().pk_stride(), r))
}
pub(crate) fn project_and_ingest(&mut self, source: &Batch) -> Result<bool, StorageError> {
let table = self.store.held_mut();
let projected = index_entries(source, &self.key_spec, table.schema());
if projected.is_empty() {
return Ok(false);
}
table.ingest_owned_batch(projected).map(|()| true)
}
pub fn resumed(&self) -> bool {
self.store.held().resumed_from_checkpoint()
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum RelationKind {
SystemCatalog,
BaseTable,
View(ViewProps),
Stream,
}
impl RelationKind {
#[inline]
pub fn class(self) -> gnitz_wire::RelClass {
match self {
RelationKind::Stream => gnitz_wire::RelClass::Stream,
RelationKind::View(props) => props.into(),
RelationKind::BaseTable | RelationKind::SystemCatalog => gnitz_wire::RelClass::Table,
}
}
#[inline]
pub fn noun(self) -> &'static str {
self.class().noun()
}
#[inline]
pub fn is_base_table(self) -> bool {
matches!(self, RelationKind::BaseTable)
}
#[inline]
pub fn is_ingestion_point(self) -> bool {
matches!(self, RelationKind::BaseTable | RelationKind::Stream)
}
#[inline]
pub fn is_view(self) -> bool {
matches!(self, RelationKind::View(_))
}
fn admits_index(self) -> bool {
match self {
RelationKind::BaseTable | RelationKind::View(ViewProps::Plain | ViewProps::Fed { .. }) => true,
RelationKind::View(ViewProps::Bounded { .. }) | RelationKind::Stream | RelationKind::SystemCatalog => false,
}
}
#[inline]
pub fn has_delta_feed(self) -> bool {
matches!(self, RelationKind::View(ViewProps::Fed { .. }))
}
#[inline]
pub fn is_bounded(self) -> bool {
matches!(self, RelationKind::View(ViewProps::Bounded { .. }))
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum Residency {
Master,
Origin,
Worker,
}
impl Residency {
#[inline]
pub(crate) fn owns_stores(self) -> bool {
!matches!(self, Residency::Master)
}
}
pub struct Relation {
id: u64,
store: Store,
unsealed: Option<Batch>,
delta: Option<Box<Table>>,
indexes: Vec<SecondaryIndex>,
kind: RelationKind,
placement: Placement,
pk_repeats: bool,
}
impl Relation {
pub(crate) fn delta(&self) -> Option<&Table> {
self.delta.as_deref()
}
pub fn schema(&self) -> SchemaDescriptor {
*self.store.schema()
}
pub fn id(&self) -> u64 {
self.id
}
pub fn kind(&self) -> RelationKind {
self.kind
}
pub fn placement(&self) -> Placement {
self.placement
}
pub fn pk_repeats(&self) -> bool {
self.pk_repeats
}
pub fn unique_indexes_to_check(&self) -> impl Iterator<Item = &SecondaryIndex> + '_ {
self.indexes.iter().filter(|ic| ic.is_unique() && !ic.covers_pk)
}
pub fn cursor(&self) -> ReadCursor {
self.table().open_cursor()
}
pub fn gather(&self, keys: PkKeys, cut: Cut) -> PkSetGather {
self.table().gather(keys, cut)
}
pub fn cursor_for_keys(&self, keys: &Batch, cut: Cut) -> ReadCursor {
self.table().cursor_for_keys(keys, cut)
}
pub fn cursor_between(&self, first: &[u8], last: &[u8], cut: Cut) -> ReadCursor {
self.table().cursor_between(first, last, cut)
}
pub fn for_each_positive_with_prefix(&self, prefix: &[u8], f: impl FnMut(&ReadCursor)) {
let table = self.table();
let band = key_range_between_cuts(
KeyCut::min_of(prefix),
KeyCut::above(prefix),
table.schema().pk_stride(),
);
table.range_cursor(band).for_each_positive_while(|_| true, f);
}
pub fn full_scan(&self) -> std::rc::Rc<Batch> {
self.table().full_scan()
}
pub fn live_row_at(&self, key: &[u8]) -> (i64, Option<StoredRow>) {
self.table().live_row_at(key)
}
pub fn indexes(&self) -> &[SecondaryIndex] {
&self.indexes
}
pub fn index_on(&self, cols: &[u32]) -> Option<&SecondaryIndex> {
self.indexes.iter().find(|ix| ix.cols.as_slice() == cols)
}
pub(crate) fn has_pending(&self) -> bool {
match &self.store {
Store::Held(t) => t.has_pending(),
Store::Absent(_) => self.unsealed.is_some(),
}
}
pub(crate) fn table(&self) -> &Table {
self.store.held()
}
}
#[derive(Clone, Copy)]
pub struct RelationSpec {
pub id: u64,
pub kind: RelationKind,
pub schema: SchemaDescriptor,
pub placement: Placement,
pub pk_repeats: bool,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub struct StoreConfig {
pub ram_tier_bytes: usize,
pub scan_chunk_rows: usize,
pub adhoc_group_cap: usize,
pub key_spans_spill_bytes: usize,
}
impl Default for StoreConfig {
fn default() -> Self {
StoreConfig {
ram_tier_bytes: crate::storage::DEFAULT_RAM_TIER_BYTES,
scan_chunk_rows: 65_536,
adhoc_group_cap: 65_536,
key_spans_spill_bytes: 128 << 20,
}
}
}
impl StoreConfig {
pub fn from_env(prefix: &str) -> Self {
let d = StoreConfig::default();
StoreConfig {
ram_tier_bytes: env_num(&format!("{prefix}RAM_TIER_BYTES"), d.ram_tier_bytes),
scan_chunk_rows: env_num(&format!("{prefix}SCAN_CHUNK_ROWS"), d.scan_chunk_rows),
adhoc_group_cap: env_num(&format!("{prefix}ADHOC_GROUP_CAP"), d.adhoc_group_cap),
key_spans_spill_bytes: env_num(&format!("{prefix}KEY_SPANS_SPILL_BYTES"), d.key_spans_spill_bytes),
}
}
}
pub struct RelationRegistry {
pub(crate) tables: FxHashMap<u64, Relation>,
pub(crate) base_dir: String,
pub(crate) slot: Slot,
pub(crate) config: StoreConfig,
pub(crate) residency: Residency,
}
impl RelationRegistry {
pub fn new(base_dir: &str, slot: Slot, config: StoreConfig) -> Self {
RelationRegistry {
tables: FxHashMap::default(),
base_dir: base_dir.to_string(),
slot,
config: StoreConfig {
scan_chunk_rows: config.scan_chunk_rows.max(1),
..config
},
residency: Residency::Origin,
}
}
pub fn master(base_dir: &str, worker_count: u32, config: StoreConfig) -> Self {
RelationRegistry {
residency: Residency::Master,
..Self::new(base_dir, Slot::new(0, worker_count), config)
}
}
pub fn unregister(&mut self, id: u64) {
self.tables.remove(&id);
}
pub fn add_index(&mut self, owner: u64, claim: IndexClaim, cols: PkColList) -> Result<(), String> {
self.enter_index(owner, claim, cols, None)
}
pub fn reopen_index(
&mut self,
owner: u64,
claim: IndexClaim,
cols: PkColList,
generation: u64,
) -> Result<(), String> {
self.enter_index(owner, claim, cols, Some(generation))
}
fn enter_index(
&mut self,
owner: u64,
claim: IndexClaim,
cols: PkColList,
resume_at: Option<u64>,
) -> Result<(), String> {
let unique = matches!(claim, IndexClaim::Index { unique: true, .. });
let owner_schema = self.index_owner(owner, unique)?.schema();
let entry = self.tables.get_mut(&owner).expect("resolved above");
if let Some(ix) = entry.indexes.iter_mut().find(|ix| ix.cols == cols) {
debug_assert!(!ix.claims.contains(&claim), "{claim:?} claimed twice");
ix.claims.push(claim);
return Ok(());
}
let (key_spec, index_schema) = index_spec_and_schema(cols.as_slice(), &owner_schema)?;
let mut ix = SecondaryIndex {
cols,
store: Store::Absent(Box::new(index_schema)),
key_spec,
claims: vec![claim],
covers_pk: owner_schema.covers_pk(cols.as_slice()),
};
if self.residency.owns_stores() {
ix.store = Store::Held(Box::new(self.open_child(
owner,
ChildKind::Index(ix.cols),
index_schema,
RecoverySource::Rederive { resume_at },
self.store_budgets(),
)?));
if !ix.resumed() {
let owner_store = &self.tables[&owner].store;
ingest::fill_indexes(owner_store, self.config.scan_chunk_rows, owner, &mut [&mut ix])?;
}
}
self.tables.get_mut(&owner).expect("resolved above").indexes.push(ix);
Ok(())
}
pub fn release_index(&mut self, owner: u64, index_id: u64) {
let Some(entry) = self.tables.get_mut(&owner) else {
return;
};
for ix in &mut entry.indexes {
ix.claims
.retain(|c| !matches!(*c, IndexClaim::Index { id, .. } if id == index_id));
}
entry.indexes.retain(|ix| !ix.claims.is_empty());
}
pub fn swap_schema(&mut self, id: u64, schema: SchemaDescriptor) -> Result<(), String> {
let entry = self.relation_mut_or_err(id)?;
if !schema.is_trailing_append_of(entry.store.schema()) {
return Err(format!(
"ALTER on table {id}: the new descriptor is not a trailing append of the stored one"
));
}
entry
.store
.swap_schema(schema)
.map_err(|e| format!("ALTER on table {id}: rebinding shards: {e}"))
}
pub fn has_id(&self, id: u64) -> bool {
self.tables.contains_key(&id)
}
pub fn base_dir(&self) -> &str {
&self.base_dir
}
pub fn slot(&self) -> Slot {
self.slot
}
pub fn child_dir(&self, id: u64, kind: ChildKind<'_>) -> String {
ChildAddr { kind, slot: self.slot }.dir(&relation_dir(&self.base_dir, id))
}
pub fn scan_chunk_rows(&self) -> usize {
self.config.scan_chunk_rows
}
pub fn set_scan_chunk_rows(&mut self, rows: usize) {
self.config.scan_chunk_rows = rows.max(1);
}
pub fn relation(&self, id: u64) -> Option<&Relation> {
self.tables.get(&id)
}
pub fn relation_or_err(&self, id: u64) -> Result<&Relation, String> {
self.relation(id).ok_or_else(|| Self::unregistered(id))
}
pub fn index_owner(&self, owner: u64, unique: bool) -> Result<&Relation, String> {
let e = self.relation_or_err(owner)?;
if !e.kind.admits_index() {
return Err(format!(
"Index: owner {owner} is a {}; only a base table or a view without a capacity can be indexed",
e.kind.noun()
));
}
if unique && e.kind.is_view() {
return Err(format!(
"Index: view {owner} cannot carry a UNIQUE index: its circuit cannot refuse a duplicate"
));
}
if e.pk_repeats {
return Err(format!(
"Index: owner {owner} repeats its primary key, so an index entry does not name one row"
));
}
Ok(e)
}
pub(crate) fn relation_mut_or_err(&mut self, id: u64) -> Result<&mut Relation, String> {
self.tables.get_mut(&id).ok_or_else(|| Self::unregistered(id))
}
fn unregistered(id: u64) -> String {
format!("relation {id} is not registered")
}
pub fn any_delta_feed(&self) -> bool {
self.tables.values().any(|r| r.kind.has_delta_feed())
}
pub fn view_ids(&self) -> impl Iterator<Item = u64> + '_ {
self.tables.iter().filter(|(_, e)| e.kind.is_view()).map(|(&id, _)| id)
}
pub(crate) fn store_budgets(&self) -> StoreBudgets {
StoreBudgets::new(self.config.ram_tier_bytes)
}
}
#[cfg(test)]
#[path = "tests/relation.rs"]
mod tests;