use std::{
collections::BTreeSet,
sync::{
Arc,
atomic::{AtomicBool, AtomicU64, Ordering},
},
};
use parking_lot::Mutex;
use crate::{
aof::aof_address::AofAddress, cluster::i_cluster_session::IClusterSession,
metrics::metrics_item::MetricsItem, types::RespCommand,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ManagerType {
Main,
Object,
Aof,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RoleInfo {
pub role: String,
pub replication_offset: i64,
pub node_id: Option<String>,
}
impl RoleInfo {
pub fn primary(replication_offset: i64) -> Self {
Self {
role: "primary".to_string(),
replication_offset,
node_id: None,
}
}
pub fn replica(replication_offset: i64, node_id: Option<String>) -> Self {
Self {
role: "replica".to_string(),
replication_offset,
node_id,
}
}
}
pub trait IClusterProvider: Send + Sync {
fn create_cluster_session(&self) -> IClusterSession;
fn allow_data_loss(&self) -> bool;
fn flush_config(&self);
fn get_gossip_stats(&self, metrics_disabled: bool) -> Vec<MetricsItem>;
fn get_replication_info(&self) -> Vec<MetricsItem>;
fn get_buffer_pool_stats(&self) -> Vec<MetricsItem>;
fn get_primary_info(&self) -> (AofAddress, Vec<RoleInfo>);
fn get_replica_info(&self) -> RoleInfo;
fn purge_buffer_pool(&self, manager_type: ManagerType);
fn cluster_publish_async(
&self,
cmd: RespCommand,
channel: &[u8],
message: &[u8],
) -> impl Future<Output = ()> + Send;
fn is_primary(&self) -> bool;
fn is_replica(&self) -> bool;
fn is_replica_node(&self, node_id: &str) -> bool;
fn on_checkpoint_initiated(&self, checkpoint_covered_aof_address: &mut AofAddress);
fn recover(&self);
fn reset_gossip_stats(&self);
fn add_new_checkpoint_entry(
&self,
full: bool,
checkpoint_covered_aof_address: AofAddress,
store_checkpoint_token: u128,
object_store_checkpoint_token: u128,
);
fn safe_truncate_aof(&self, truncate_until: &AofAddress);
fn start(&self);
fn update_cluster_auth(&self, cluster_username: Option<&str>, cluster_password: Option<&str>);
fn get_checkpoint_info(&self) -> Vec<MetricsItem>;
fn get_run_id(&self) -> String;
fn prevent_role_change(&self) -> bool;
fn allow_role_change(&self);
}
#[derive(Debug)]
pub struct SingleNodeClusterProvider {
run_id: String,
role_change_blocked: AtomicBool,
local_offset: AtomicU64,
checkpoint_entries: AtomicU64,
last_checkpoint_aof: Mutex<Option<AofAddress>>,
cluster_auth: Mutex<Option<(String, String)>>,
replica_nodes: Mutex<BTreeSet<String>>,
}
impl SingleNodeClusterProvider {
pub fn new(run_id: String) -> Self {
Self {
run_id,
role_change_blocked: AtomicBool::new(false),
local_offset: AtomicU64::new(0),
checkpoint_entries: AtomicU64::new(0),
last_checkpoint_aof: Mutex::new(None),
cluster_auth: Mutex::new(None),
replica_nodes: Mutex::new(BTreeSet::new()),
}
}
pub fn advance_local_offset(&self, delta: u64) {
self.local_offset.fetch_add(delta, Ordering::AcqRel);
}
}
impl Default for SingleNodeClusterProvider {
fn default() -> Self {
Self::new("single-node-run-id".to_string())
}
}
impl IClusterProvider for SingleNodeClusterProvider {
fn create_cluster_session(&self) -> IClusterSession {
IClusterSession::new()
}
fn allow_data_loss(&self) -> bool {
false
}
fn flush_config(&self) {
}
fn get_gossip_stats(&self, metrics_disabled: bool) -> Vec<MetricsItem> {
if metrics_disabled {
return Vec::new();
}
vec![MetricsItem::new(
"NODE_INFO",
format!("myid={},{}", "self", self.run_id),
)]
}
fn get_replication_info(&self) -> Vec<MetricsItem> {
vec![
MetricsItem::new("role", "primary"),
MetricsItem::new(
"master_repl_offset",
self.local_offset.load(Ordering::Acquire).to_string(),
),
]
}
fn get_buffer_pool_stats(&self) -> Vec<MetricsItem> {
Vec::new()
}
fn get_primary_info(&self) -> (AofAddress, Vec<RoleInfo>) {
let offset = AofAddress::default();
(offset, Vec::new())
}
fn get_replica_info(&self) -> RoleInfo {
RoleInfo::primary(self.local_offset.load(Ordering::Acquire) as i64)
}
fn purge_buffer_pool(&self, manager_type: ManagerType) {
let _ = manager_type;
}
async fn cluster_publish_async(&self, _cmd: RespCommand, _channel: &[u8], _message: &[u8]) {
}
fn is_primary(&self) -> bool {
true
}
fn is_replica(&self) -> bool {
false
}
fn is_replica_node(&self, node_id: &str) -> bool {
self.replica_nodes.lock().contains(node_id)
}
fn on_checkpoint_initiated(&self, checkpoint_covered_aof_address: &mut AofAddress) {
*self.last_checkpoint_aof.lock() = Some(*checkpoint_covered_aof_address);
let _ = checkpoint_covered_aof_address;
}
fn recover(&self) {
}
fn reset_gossip_stats(&self) {
}
fn add_new_checkpoint_entry(
&self,
_full: bool,
_checkpoint_covered_aof_address: AofAddress,
_store_checkpoint_token: u128,
_object_store_checkpoint_token: u128,
) {
self.checkpoint_entries.fetch_add(1, Ordering::AcqRel);
}
fn safe_truncate_aof(&self, _truncate_until: &AofAddress) {
}
fn start(&self) {
}
fn update_cluster_auth(&self, cluster_username: Option<&str>, cluster_password: Option<&str>) {
let mut auth = self.cluster_auth.lock();
*auth = match (cluster_username, cluster_password) {
(Some(u), Some(p)) => Some((u.to_string(), p.to_string())),
_ => None,
};
}
fn get_checkpoint_info(&self) -> Vec<MetricsItem> {
vec![
MetricsItem::new(
"CheckpointEntries",
self.checkpoint_entries.load(Ordering::Acquire).to_string(),
),
MetricsItem::new("RunID", self.run_id.clone()),
]
}
fn get_run_id(&self) -> String {
self.run_id.clone()
}
fn prevent_role_change(&self) -> bool {
self
.role_change_blocked
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_ok()
}
fn allow_role_change(&self) {
self.role_change_blocked.store(false, Ordering::Release);
}
}
pub type SharedClusterProvider = Arc<dyn IClusterProvider>;
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn single_node_role_semantics() {
let provider = SingleNodeClusterProvider::default();
assert!(provider.is_primary());
assert!(!provider.is_replica());
assert!(!provider.is_replica_node("anyone"));
let replica_info = provider.get_replica_info();
assert_eq!(replica_info.role, "primary");
let (offset, replicas) = provider.get_primary_info();
assert_eq!(replicas, Vec::new());
assert_eq!(offset, AofAddress::default());
provider.advance_local_offset(128);
assert_eq!(provider.get_replica_info().replication_offset, 128);
}
#[test]
fn role_change_gate_is_pairwise() {
let provider = SingleNodeClusterProvider::default();
assert!(provider.prevent_role_change());
assert!(!provider.prevent_role_change());
provider.allow_role_change();
assert!(provider.prevent_role_change());
provider.allow_role_change();
}
#[test]
fn checkpoint_bookkeeping() {
let provider = SingleNodeClusterProvider::default();
let info = provider.get_checkpoint_info();
assert_eq!(info[0].name, "CheckpointEntries");
assert_eq!(info[0].value, "0");
assert_eq!(info[1].value, "single-node-run-id");
let mut covered = AofAddress::default();
provider.on_checkpoint_initiated(&mut covered);
provider.add_new_checkpoint_entry(true, covered, 1, 2);
provider.add_new_checkpoint_entry(false, covered, 3, 4);
let info = provider.get_checkpoint_info();
assert_eq!(info[0].value, "2");
}
#[test]
fn cluster_auth_atomic_update() {
let provider = SingleNodeClusterProvider::default();
provider.update_cluster_auth(Some("u"), Some("p"));
provider.update_cluster_auth(None, None);
provider.update_cluster_auth(Some("u2"), Some("p2"));
}
#[test]
fn gossip_and_buffers() {
let provider = SingleNodeClusterProvider::default();
assert!(provider.get_gossip_stats(true).is_empty());
assert_eq!(provider.get_gossip_stats(false).len(), 1);
assert!(provider.get_buffer_pool_stats().is_empty());
provider.purge_buffer_pool(ManagerType::Main);
provider.purge_buffer_pool(ManagerType::Object);
provider.purge_buffer_pool(ManagerType::Aof);
let info = provider.get_replication_info();
assert_eq!(info[0].value, "primary");
provider.flush_config();
provider.start();
provider.recover();
provider.reset_gossip_stats();
provider.safe_truncate_aof(&AofAddress::default());
assert_eq!(provider.get_run_id(), "single-node-run-id");
assert!(!provider.allow_data_loss());
}
}