use std::collections::BTreeSet;
use std::fs;
use super::dirs::{cluster_children, subdir_names};
use super::{ChildAddr, ChildKind};
use crate::storage::{flush_barrier, fsync_dir, link_store, retire_store, RecoverySource, StoreBudgets, Table};
use gnitz_zset::algebra::ScatterPlan;
use gnitz_zset::repr::{from_runs, Batch, StorageError};
use gnitz_zset::schema::{Placement, SchemaDescriptor, Slot};
fn relay_source(rel_dir: &str, launched: u32) -> Result<Option<u32>, String> {
let mut foreign = BTreeSet::new();
let names = subdir_names(rel_dir).map_err(|e| format!("repartition: list {rel_dir}: {e}"))?;
for name in names {
match ChildAddr::parse(&name) {
Some(ChildAddr { kind: ChildKind::Rows, slot }) if slot.of != launched => {
foreign.insert(slot.of);
}
Some(_) => {}
None => {
return Err(format!(
"{rel_dir}/{name} is in no child-directory grammar this build knows; \
it may hold rows written by an incompatible layout. Refusing to boot — \
delete the relation directory deliberately if the data is expendable."
))
}
}
}
let complete = |of: u32| -> Result<bool, String> {
for m in cluster_children(of).map(|c| c.manifest(rel_dir)) {
if !fs::exists(&m).map_err(|e| format!("repartition: {m}: {e}"))? {
return Ok(false);
}
}
Ok(true)
};
if foreign.is_empty() || complete(launched)? {
return Ok(None);
}
for of in foreign {
if complete(of)? {
return Ok(Some(of));
}
}
Ok(None)
}
pub(super) fn repartition_relation(
rel_dir: &str,
schema: &SchemaDescriptor,
placement: Placement,
launched: u32,
ram_tier_bytes: usize,
chunk_rows: usize,
) -> Result<(), String> {
let Some(source) = relay_source(rel_dir, launched)? else {
return Ok(());
};
gnitz_info!("repartition: {} from {} to {} worker(s)", rel_dir, source, launched);
relay(rel_dir, schema, placement, source, launched, ram_tier_bytes, chunk_rows)
.map_err(|e| format!("repartition {rel_dir} to {launched} worker(s): {e}"))
}
fn relay(
rel_dir: &str,
schema: &SchemaDescriptor,
placement: Placement,
source: u32,
launched: u32,
ram_tier_bytes: usize,
chunk_rows: usize,
) -> Result<(), StorageError> {
remove_set(rel_dir, launched)?;
if placement.is_replicated() {
link_targets(rel_dir, source, launched)?;
} else {
rewrite_targets(rel_dir, schema, placement, source, launched, ram_tier_bytes, chunk_rows)?;
}
fsync_dir(rel_dir)?;
remove_set(rel_dir, source)
}
fn remove_set(rel_dir: &str, of: u32) -> Result<(), StorageError> {
cluster_children(of).try_for_each(|c| retire_store(&c.dir(rel_dir)))
}
fn link_targets(rel_dir: &str, source: u32, launched: u32) -> Result<(), StorageError> {
let source_dir = ChildAddr {
kind: ChildKind::Rows,
slot: Slot::new(0, source),
}
.dir(rel_dir);
cluster_children(launched).try_for_each(|target| link_store(&source_dir, &target.dir(rel_dir)))
}
fn rewrite_targets(
rel_dir: &str,
schema: &SchemaDescriptor,
placement: Placement,
source: u32,
launched: u32,
ram_tier_bytes: usize,
chunk_rows: usize,
) -> Result<(), StorageError> {
let budgets = StoreBudgets::new(ram_tier_bytes);
let open = |of: u32| {
cluster_children(of)
.map(|c| Table::new(&c.dir(rel_dir), *schema, RecoverySource::SalReplay, budgets))
.collect::<Result<Vec<Table>, _>>()
};
let sources = open(source)?;
for t in &sources {
t.verify_shards()?;
}
let mut cursor = from_runs(sources.iter().flat_map(Table::runs), *schema, 0);
let mut targets = open(launched)?;
let mut buffers: Vec<Batch> = targets.iter().map(|_| Batch::empty_with_schema(schema)).collect();
let plan = ScatterPlan::native(placement);
let mut rows: Vec<Vec<u32>> = Vec::new();
while let Some(chunk) = cursor.drain_chunk(chunk_rows) {
let slots = plan.route(&chunk, &mut rows, launched as usize);
for ((target, buffer), idx) in targets.iter_mut().zip(&mut buffers).zip(slots.iter()) {
if idx.is_empty() {
continue;
}
buffer.append_above(chunk.ascending_subset(idx));
if buffer.total_bytes() >= ram_tier_bytes {
write_run(target, buffer)?;
}
}
}
for (target, buffer) in targets.iter_mut().zip(&mut buffers) {
if !buffer.is_empty() {
write_run(target, buffer)?;
}
}
flush_barrier(targets.iter_mut(), 0)
}
fn write_run(target: &mut Table, run: &mut Batch) -> Result<(), StorageError> {
target.append_terminal_run(run)?;
run.clear();
Ok(())
}
#[cfg(test)]
#[path = "tests/repartition.rs"]
mod tests;