use super::dirs::remove_children;
use super::*;
use gnitz_foundation::fault::Seam;
static TABLE_CREATE_DELAY: Seam = Seam::new("GNITZ_INJECT_TABLE_CREATE_DELAY_MS");
impl RelationRegistry {
pub fn register(&mut self, spec: RelationSpec) -> Result<(), String> {
let stores = self.build_relation_store(spec, None)?;
self.enter(spec, stores);
Ok(())
}
pub fn reopen_view(&mut self, spec: RelationSpec, generation: u64) -> Result<(), String> {
let stores = self.build_relation_store(spec, Some(generation))?;
self.enter(spec, stores);
Ok(())
}
fn enter(&mut self, spec: RelationSpec, (store, delta): (Store, Option<Box<Table>>)) {
let RelationSpec { id, kind, placement, pk_repeats, .. } = spec;
let prev = self.tables.insert(
id,
Relation {
id,
store,
unsealed: None,
delta,
indexes: Vec::new(),
kind,
placement,
pk_repeats,
},
);
debug_assert!(prev.is_none(), "relation {id} registered twice");
}
pub(super) fn build_relation_store(
&self,
spec: RelationSpec,
resume_at: Option<u64>,
) -> Result<(Store, Option<Box<Table>>), String> {
let RelationSpec { id, kind, schema, placement, .. } = spec;
let absent = || Ok((Store::Absent(Box::new(schema)), None));
let (recovery, budgets, feed) = match kind {
RelationKind::Stream => return absent(),
RelationKind::SystemCatalog => {
let dir = relation_dir(&self.base_dir, id);
let table = Table::new(&dir, schema, RecoverySource::SalReplay, self.store_budgets())
.map_err(|e| format!("open store '{dir}': {e}"))?;
return Ok((Store::Held(Box::new(table)), None));
}
_ if !self.residency.owns_stores() => {
ensure_dir(&relation_dir(&self.base_dir, id))?;
return absent();
}
RelationKind::BaseTable => {
if let Some(ms) = TABLE_CREATE_DELAY.count() {
std::thread::sleep(std::time::Duration::from_millis(ms));
}
(RecoverySource::SalReplay, self.store_budgets(), None)
}
RelationKind::View(p) => {
if resume_at.is_none() {
let slot = self.slot;
remove_children(&relation_dir(&self.base_dir, id), |c| {
matches!(c.kind, ChildKind::Scratch(_)) && c.slot == slot
})
.map_err(|e| format!("remove the operator traces of view {id}: {e}"))?;
}
(
RecoverySource::Rederive { resume_at },
self.store_budgets().bounded(p.capacity_bytes()),
p.delta_bytes(),
)
}
};
let rows = self.open_child_as(id, ChildKind::Rows, schema, recovery, budgets)?;
let delta = match feed {
Some(budget) if placement.counts_on(self.slot.rank) => {
let delta_schema = super::delta::make_delta_schema(&schema)
.ok_or_else(|| format!("view {id} has too many columns to carry a delta feed"))?;
let table = self.open_child(
id,
ChildKind::Delta,
delta_schema,
RecoverySource::Rederive { resume_at: None },
self.store_budgets().delta(budget),
)?;
Some(Box::new(table))
}
_ => None,
};
Ok((Store::Held(Box::new(rows)), delta))
}
pub(super) fn open_child_as(
&self,
id: u64,
kind: ChildKind<'_>,
schema: SchemaDescriptor,
recovery: RecoverySource,
budgets: StoreBudgets,
) -> Result<Table, String> {
let table = self.open_child(id, kind, schema, recovery, budgets)?;
if table.recovery_source() != recovery {
return Err(format!(
"relation {id}: {kind:?} store did not resume from its manifest"
));
}
Ok(table)
}
pub(super) fn open_child(
&self,
id: u64,
kind: ChildKind<'_>,
schema: SchemaDescriptor,
recovery: RecoverySource,
budgets: StoreBudgets,
) -> Result<Table, String> {
let dir = self.child_dir(id, kind);
Table::new(&dir, schema, recovery, budgets).map_err(|e| format!("open store '{dir}': {e}"))
}
}