use gnitz_wire::{widen_pk_be, KeyRange, NARROW_PK_MAX_BYTES};
use crate::schema::{key, ColumnTable, SchemaDescriptor};
#[inline(always)]
fn bucket(h: u64, num_workers: usize) -> usize {
debug_assert!(num_workers >= 1, "worker routing: num_workers must be >= 1");
((h as u128 * num_workers as u128) >> 64) as usize
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub struct Slot {
pub rank: u32,
pub of: u32,
}
impl Slot {
pub const SOLO: Slot = Slot { rank: 0, of: 1 };
pub fn new(rank: u32, of: u32) -> Slot {
assert!(rank < of, "slot {rank} of {of}");
Slot { rank, of }
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum Placement {
Replicated,
Local,
Keyed { dist_stride: u8 },
}
impl Placement {
pub const REPLICA_OWNER: u32 = 0;
pub fn full_pk(schema: &SchemaDescriptor) -> Placement {
Placement::Keyed { dist_stride: schema.pk_stride() as u8 }
}
pub fn keyed(schema: &SchemaDescriptor, prefix_cols: usize) -> Placement {
let pk = schema.pk_cols();
assert!(
(1..=pk.len()).contains(&prefix_cols),
"Placement::keyed: prefix {prefix_cols} of a {}-column PK",
pk.len()
);
let dist_stride = pk[..prefix_cols]
.iter()
.map(|&c| schema.columns[c as usize].size())
.sum();
Placement::Keyed { dist_stride }
}
#[inline]
pub const fn counts_on(self, rank: u32) -> bool {
!self.is_replicated() || rank == Self::REPLICA_OWNER
}
#[inline]
pub const fn is_key_routed(self) -> bool {
matches!(self, Placement::Keyed { .. })
}
#[inline]
pub const fn is_replicated(self) -> bool {
matches!(self, Placement::Replicated)
}
pub fn confined_worker(self, schema: &SchemaDescriptor, range: &KeyRange, num_workers: usize) -> Option<usize> {
if !range.walks_pk(schema.pk_cols()) {
return None;
}
let Some((start, end)) = schema.pk_range_keys(range) else {
return Some(0);
};
let Placement::Keyed { dist_stride } = self else {
return None;
};
key::range_shares_prefix(&start, end.as_ref(), dist_stride as usize)
.then(|| worker_for_pk_bytes(&start.pk_bytes()[..dist_stride as usize], num_workers))
}
pub fn owner(self, pk: &[u8], num_workers: usize) -> Option<usize> {
match self {
Placement::Keyed { dist_stride } => Some(worker_for_pk_bytes(&pk[..dist_stride as usize], num_workers)),
_ => None,
}
}
}
#[inline(always)]
pub(crate) fn worker_for_key(pk: u128, num_workers: usize) -> usize {
let lo = pk as u64;
let hi = (pk >> 64) as u64;
bucket(
lo.wrapping_mul(0x9e3779b97f4a7c15_u64) ^ hi.wrapping_mul(0x6c62272e07bb0142_u64),
num_workers,
)
}
#[inline(always)]
pub(crate) fn worker_for_pk_bytes(bytes: &[u8], num_workers: usize) -> usize {
if bytes.len() <= NARROW_PK_MAX_BYTES {
worker_for_key(widen_pk_be(bytes), num_workers)
} else {
worker_for_wide_pk(bytes, num_workers)
}
}
#[inline(never)]
fn worker_for_wide_pk(bytes: &[u8], num_workers: usize) -> usize {
bucket(gnitz_wire::checksum(bytes), num_workers)
}
pub fn ground_owner(num_workers: usize) -> usize {
worker_for_key(gnitz_wire::global_group_key(), num_workers)
}
#[cfg(test)]
#[path = "tests/route.rs"]
mod tests;