use std::cell::Cell;
use std::net::{IpAddr, Ipv4Addr};
use std::sync::{Mutex, OnceLock};
use std::sync::atomic::{AtomicU64, Ordering};
use kevy_config::{Config, PeerEntry, ReplicationRole};
use kevy_elect::{
PeerAddr, Transport,
elector::{ElectConfig, ElectJitter, Elector},
message::Role,
};
static ELECT_TRANSPORT: OnceLock<Mutex<Option<Transport>>> = OnceLock::new();
fn slot() -> &'static Mutex<Option<Transport>> {
ELECT_TRANSPORT.get_or_init(|| Mutex::new(None))
}
static SHARD_OFFSETS: OnceLock<Vec<AtomicU64>> = OnceLock::new();
thread_local! {
static MY_SHARD_ID: Cell<Option<usize>> = const { Cell::new(None) };
}
pub(crate) fn install_shard_offsets(nshards: usize) {
let _ = SHARD_OFFSETS.set((0..nshards).map(|_| AtomicU64::new(0)).collect());
}
fn shard_offsets() -> Option<&'static [AtomicU64]> {
SHARD_OFFSETS.get().map(Vec::as_slice)
}
fn aggregate_offset() -> u64 {
shard_offsets()
.map(|slots| {
slots
.iter()
.map(|a| a.load(Ordering::Relaxed))
.fold(0u64, u64::saturating_add)
})
.unwrap_or(0)
}
pub(crate) fn is_configured(cfg: &Config) -> bool {
!cfg.cluster.peers.is_empty() && !cfg.cluster.node_id.is_empty()
}
pub(crate) fn maybe_start(cfg: &Config) {
if !is_configured(cfg) {
return;
}
let listen_port = resolved_elect_port_base(cfg);
let start_role = match cfg.replication.role {
ReplicationRole::Primary => Role::Primary,
_ => Role::Replica,
};
let peer_ids: Vec<String> = cfg
.cluster
.peers
.iter()
.map(|p| p.node_id.clone())
.collect();
let advertised_addr = format!("{}:{}", advertised_host(cfg), cfg.server.port);
let elect_cfg = ElectConfig::default();
let hb_interval = elect_cfg.hb_interval;
let elector = Elector::new(
cfg.cluster.node_id.clone(),
peer_ids,
advertised_addr,
start_role,
elect_cfg,
ElectJitter::System,
);
let self_id = cfg.cluster.node_id.as_str();
let peers: Vec<PeerAddr> = cfg
.cluster
.peers
.iter()
.filter(|p| p.node_id != self_id)
.map(peer_to_addr)
.collect();
let listen = (IpAddr::V4(Ipv4Addr::new(
cfg.server.bind[0],
cfg.server.bind[1],
cfg.server.bind[2],
cfg.server.bind[3],
)), listen_port);
match Transport::spawn(elector, hb_interval, listen, peers) {
Ok(t) => {
*slot().lock().expect("ELECT_TRANSPORT poisoned") = Some(t);
eprintln!(
"kevy: kevy-elect transport up on {}:{} ({} peers, role={})",
cfg.server.bind[0], listen_port,
cfg.cluster.peers.len().saturating_sub(1),
if matches!(start_role, Role::Primary) { "primary" } else { "replica" },
);
}
Err(e) => {
eprintln!("kevy: kevy-elect transport failed to bind {listen_port}: {e}");
}
}
}
pub(crate) fn shutdown() {
if let Ok(mut guard) = slot().lock()
&& let Some(t) = guard.take()
{
t.shutdown();
}
}
pub(crate) fn current_snapshot() -> Option<kevy_elect::ElectorSnapshot> {
slot().lock().ok().and_then(|g| g.as_ref().map(Transport::state_snapshot))
}
pub(crate) fn set_view_offset(offset: u64) {
let shard_id = MY_SHARD_ID.with(|c| {
if let Some(id) = c.get() {
return id;
}
let id = crate::ops::cluster::current_shard_for_elect();
c.set(Some(id));
id
});
if let Some(slots) = shard_offsets()
&& let Some(slot_ref) = slots.get(shard_id)
{
slot_ref.store(offset, Ordering::Relaxed);
}
let agg = aggregate_offset();
if let Ok(guard) = slot().lock()
&& let Some(t) = guard.as_ref()
{
t.set_repl_offset(agg);
}
}
fn resolved_elect_port_base(cfg: &Config) -> u16 {
if cfg.cluster.elect_port_base != 0 {
return cfg.cluster.elect_port_base;
}
cfg.server.port.saturating_add(200)
}
fn advertised_host(cfg: &Config) -> String {
format!(
"{}.{}.{}.{}",
cfg.server.bind[0], cfg.server.bind[1], cfg.server.bind[2], cfg.server.bind[3]
)
}
fn peer_to_addr(p: &PeerEntry) -> PeerAddr {
PeerAddr {
node_id: p.node_id.clone(),
host: p.host.clone(),
port: p.port,
}
}
#[cfg(test)]
mod tests {
use super::*;
fn cfg_with(node_id: &str, peers: &str) -> Config {
let mut c = Config::default();
c.cluster.node_id = node_id.to_string();
c.cluster.peers = PeerEntry::parse_list(peers).unwrap();
c
}
#[test]
fn is_configured_empty_peers_returns_false() {
let cfg = Config::default();
assert!(!is_configured(&cfg));
}
#[test]
fn is_configured_empty_node_id_returns_false() {
let mut cfg = Config::default();
cfg.cluster.peers = PeerEntry::parse_list("a@h:1,b@h:2").unwrap();
assert!(!is_configured(&cfg));
}
#[test]
fn is_configured_both_set_returns_true() {
let cfg = cfg_with("self", "self@127.0.0.1:1,b@127.0.0.1:2");
assert!(is_configured(&cfg));
}
#[test]
fn resolved_elect_port_base_falls_back_when_zero() {
let mut cfg = Config::default();
cfg.server.port = 6004;
cfg.cluster.elect_port_base = 0;
assert_eq!(resolved_elect_port_base(&cfg), 6204);
}
#[test]
fn resolved_elect_port_base_uses_explicit_when_set() {
let mut cfg = Config::default();
cfg.cluster.elect_port_base = 16104;
assert_eq!(resolved_elect_port_base(&cfg), 16104);
}
}