use std::collections::HashMap;
use std::sync::RwLock;
use crate::datatypes::keyfun::KeyFun;
use crate::replication::ReplicationStrategy;
#[derive(Clone, Debug, Default, Eq, PartialEq)]
pub struct BucketProps {
pub keyfun: Option<KeyFun>,
pub strategy: Option<ReplicationStrategy>,
pub n_val: Option<u8>,
pub custom_keyfun_module: Option<String>,
pub allow_mult: Option<bool>,
pub precommit_module: Option<String>,
pub postcommit_module: Option<String>,
pub ttl_seconds: Option<u64>,
pub r: Option<u32>,
pub w: Option<u32>,
pub pr: Option<u32>,
pub pw: Option<u32>,
pub dw: Option<u32>,
}
impl BucketProps {
#[must_use]
pub fn effective_keyfun_with(&self, default: KeyFun) -> KeyFun {
match self.keyfun.clone() {
Some(KeyFun::Custom(id)) => {
let module = self.custom_keyfun_module.clone().unwrap_or_else(|| {
if id.is_empty() {
String::new()
} else {
id
}
});
KeyFun::Custom(module)
}
Some(other) => other,
None => default,
}
}
#[must_use]
pub fn effective_strategy_with(&self, default: ReplicationStrategy) -> ReplicationStrategy {
self.strategy.unwrap_or(default)
}
#[must_use]
pub fn effective_n_val_with(&self, default: u8) -> u8 {
self.n_val.unwrap_or(default)
}
#[must_use]
pub fn effective_ttl_seconds(&self) -> u64 {
self.ttl_seconds.unwrap_or(0)
}
#[must_use]
pub fn effective_allow_mult(&self) -> bool {
self.allow_mult.unwrap_or(false)
}
#[must_use]
pub fn precommit_module(&self) -> Option<&str> {
self.precommit_module.as_deref()
}
#[must_use]
pub fn postcommit_module(&self) -> Option<&str> {
self.postcommit_module.as_deref()
}
#[must_use]
pub fn effective_r(&self, n_val: u8, request: Option<u32>) -> u32 {
crate::quorum::resolve(request, n_val, self.r)
}
#[must_use]
pub fn effective_w(&self, n_val: u8, request: Option<u32>) -> u32 {
crate::quorum::resolve(request, n_val, self.w)
}
#[must_use]
pub fn effective_pr(&self, n_val: u8, request: Option<u32>) -> u32 {
match request.or(self.pr) {
None => 0,
some => crate::quorum::resolve(some, n_val, self.pr),
}
}
#[must_use]
pub fn effective_pw(&self, n_val: u8, request: Option<u32>) -> u32 {
match request.or(self.pw) {
None => 0,
some => crate::quorum::resolve(some, n_val, self.pw),
}
}
#[must_use]
pub fn effective_dw(&self, n_val: u8, request: Option<u32>) -> u32 {
crate::quorum::resolve(request, n_val, self.dw)
}
#[must_use]
pub fn effective_keyfun(&self) -> KeyFun {
self.effective_keyfun_with(KeyFun::default())
}
#[must_use]
pub fn effective_strategy(&self) -> ReplicationStrategy {
self.effective_strategy_with(ReplicationStrategy::default())
}
#[must_use]
pub fn effective_n_val(&self) -> u8 {
self.effective_n_val_with(3)
}
}
#[derive(Debug)]
pub struct BucketPropsRegistry {
inner: RwLock<RegistryInner>,
}
#[derive(Debug)]
struct RegistryInner {
by_bucket: HashMap<(Vec<u8>, Vec<u8>), BucketProps>,
default_keyfun: KeyFun,
default_strategy: ReplicationStrategy,
default_n_val: u8,
}
impl BucketPropsRegistry {
#[must_use]
pub fn new() -> Self {
Self {
inner: RwLock::new(RegistryInner {
by_bucket: HashMap::new(),
default_keyfun: KeyFun::Std,
default_strategy: ReplicationStrategy::Topology,
default_n_val: 3,
}),
}
}
#[must_use]
pub fn new_riak_defaults() -> Self {
Self {
inner: RwLock::new(RegistryInner {
by_bucket: HashMap::new(),
default_keyfun: KeyFun::Std,
default_strategy: ReplicationStrategy::Successors,
default_n_val: 3,
}),
}
}
pub fn set(&self, bucket_type: &[u8], bucket: &[u8], props: BucketProps) {
let key = (Self::norm_type(bucket_type), bucket.to_vec());
let mut inner = self.inner.write().expect("registry rwlock poisoned");
inner.by_bucket.insert(key, props);
}
#[must_use]
pub fn resolve(&self, bucket_type: &[u8], bucket: &[u8]) -> BucketProps {
let key = (Self::norm_type(bucket_type), bucket.to_vec());
let inner = self.inner.read().expect("registry rwlock poisoned");
let mut p = inner.by_bucket.get(&key).cloned().unwrap_or_default();
if p.keyfun.is_none() {
p.keyfun = Some(inner.default_keyfun.clone());
}
if p.strategy.is_none() {
p.strategy = Some(inner.default_strategy);
}
if p.n_val.is_none() {
p.n_val = Some(inner.default_n_val);
}
p
}
#[must_use]
pub fn defaults(&self) -> BucketProps {
let inner = self.inner.read().expect("registry rwlock poisoned");
BucketProps {
keyfun: Some(inner.default_keyfun.clone()),
strategy: Some(inner.default_strategy),
n_val: Some(inner.default_n_val),
custom_keyfun_module: None,
allow_mult: None,
precommit_module: None,
postcommit_module: None,
ttl_seconds: None,
r: None,
w: None,
pr: None,
pw: None,
dw: None,
}
}
pub fn set_default_strategy(&self, strategy: ReplicationStrategy) {
let mut inner = self.inner.write().expect("registry rwlock poisoned");
inner.default_strategy = strategy;
}
fn norm_type(bucket_type: &[u8]) -> Vec<u8> {
if bucket_type.is_empty() {
b"default".to_vec()
} else {
bucket_type.to_vec()
}
}
}
impl Default for BucketPropsRegistry {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn fresh_registry_returns_mode_defaults() {
let reg = BucketPropsRegistry::new();
let p = reg.resolve(b"default", b"users");
assert_eq!(p.effective_keyfun(), KeyFun::Std);
assert_eq!(p.effective_strategy(), ReplicationStrategy::Topology);
assert_eq!(p.effective_n_val(), 3);
}
#[test]
fn riak_default_swaps_strategy_to_successors() {
let reg = BucketPropsRegistry::new_riak_defaults();
let p = reg.resolve(b"default", b"users");
assert_eq!(p.effective_strategy(), ReplicationStrategy::Successors);
}
#[test]
fn override_takes_precedence() {
let reg = BucketPropsRegistry::new_riak_defaults();
reg.set(
b"default",
b"users",
BucketProps {
keyfun: Some(KeyFun::BucketOnly),
strategy: Some(ReplicationStrategy::Topology),
n_val: Some(5),
..BucketProps::default()
},
);
let p = reg.resolve(b"default", b"users");
assert_eq!(p.effective_keyfun(), KeyFun::BucketOnly);
assert_eq!(p.effective_strategy(), ReplicationStrategy::Topology);
assert_eq!(p.effective_n_val(), 5);
}
#[test]
fn ttl_defaults_to_zero_and_round_trips() {
let reg = BucketPropsRegistry::new_riak_defaults();
assert_eq!(reg.resolve(b"default", b"c").effective_ttl_seconds(), 0);
reg.set(
b"default",
b"cache",
BucketProps {
ttl_seconds: Some(3600),
..BucketProps::default()
},
);
assert_eq!(
reg.resolve(b"default", b"cache").effective_ttl_seconds(),
3600
);
assert_eq!(
BucketProps {
ttl_seconds: Some(0),
..BucketProps::default()
}
.effective_ttl_seconds(),
0
);
}
#[test]
fn quorum_resolves_request_over_bucket_default_over_quorum() {
use crate::quorum::{QUORUM_ALL, QUORUM_ONE};
let p = BucketProps::default();
assert_eq!(p.effective_r(3, None), 2);
assert_eq!(p.effective_w(3, None), 2);
assert_eq!(p.effective_r(3, Some(QUORUM_ALL)), 3);
assert_eq!(p.effective_w(3, Some(QUORUM_ONE)), 1);
let all = BucketProps {
r: Some(QUORUM_ALL),
w: Some(QUORUM_ONE),
..BucketProps::default()
};
assert_eq!(all.effective_r(3, None), 3);
assert_eq!(all.effective_w(3, None), 1);
assert_eq!(all.effective_r(3, Some(QUORUM_ONE)), 1);
assert_eq!(p.effective_pr(3, None), 0);
assert_eq!(p.effective_pw(3, None), 0);
assert_eq!(p.effective_pr(3, Some(QUORUM_ALL)), 3);
assert_eq!(p.effective_dw(3, None), 2);
assert_eq!(p.effective_dw(3, Some(QUORUM_ALL)), 3);
let dw_all = BucketProps {
dw: Some(QUORUM_ALL),
..BucketProps::default()
};
assert_eq!(dw_all.effective_dw(3, None), 3);
assert_eq!(dw_all.effective_dw(3, Some(QUORUM_ONE)), 1);
}
#[test]
fn empty_bucket_type_normalises_to_default() {
let reg = BucketPropsRegistry::new();
reg.set(
b"",
b"users",
BucketProps {
keyfun: Some(KeyFun::BucketOnly),
..BucketProps::default()
},
);
assert_eq!(
reg.resolve(b"", b"users").effective_keyfun(),
KeyFun::BucketOnly
);
assert_eq!(
reg.resolve(b"default", b"users").effective_keyfun(),
KeyFun::BucketOnly
);
}
#[test]
fn missing_buckets_fall_through_to_defaults() {
let reg = BucketPropsRegistry::new_riak_defaults();
let p = reg.resolve(b"default", b"never-set");
assert_eq!(p.effective_strategy(), ReplicationStrategy::Successors);
}
#[test]
fn custom_keyfun_takes_module_id_from_field() {
let p = BucketProps {
keyfun: Some(KeyFun::Custom(String::new())),
custom_keyfun_module: Some("reverse".to_string()),
..BucketProps::default()
};
assert_eq!(p.effective_keyfun(), KeyFun::Custom("reverse".to_string()));
let p2 = BucketProps {
keyfun: Some(KeyFun::Custom("embedded".to_string())),
..BucketProps::default()
};
assert_eq!(
p2.effective_keyfun(),
KeyFun::Custom("embedded".to_string())
);
let p3 = BucketProps {
keyfun: Some(KeyFun::Custom(String::new())),
..BucketProps::default()
};
assert_eq!(p3.effective_keyfun(), KeyFun::Custom(String::new()));
}
#[test]
fn set_default_strategy_changes_unconfigured_buckets_only() {
let reg = BucketPropsRegistry::new_riak_defaults();
reg.set(
b"default",
b"explicit",
BucketProps {
strategy: Some(ReplicationStrategy::Topology),
..BucketProps::default()
},
);
reg.set_default_strategy(ReplicationStrategy::Topology);
assert_eq!(
reg.resolve(b"default", b"never-set").effective_strategy(),
ReplicationStrategy::Topology,
);
assert_eq!(
reg.resolve(b"default", b"explicit").effective_strategy(),
ReplicationStrategy::Topology,
);
}
}