use std::sync::OnceLock;
use kevy_config::Config;
use kevy_scope::{MigrationState, MigrationTable, OwnershipTable, Routing, Scope};
static OWNERSHIP: OnceLock<Option<OwnershipTable>> = OnceLock::new();
static PEER_ADDRS: OnceLock<Vec<(String, String)>> = OnceLock::new();
static MIGRATIONS: OnceLock<MigrationTable> = OnceLock::new();
fn migrations() -> &'static MigrationTable {
MIGRATIONS.get_or_init(MigrationTable::new)
}
#[allow(dead_code)] pub(crate) fn migration_start(
prefix: Vec<u8>,
from: String,
to: String,
) -> Result<(), kevy_scope::MigrationError> {
migrations().start(prefix, from, to)
}
#[allow(dead_code)] pub(crate) fn migration_commit(prefix: &[u8]) -> Option<MigrationState> {
migrations().commit(prefix)
}
pub(crate) fn migration_abort(prefix: &[u8]) -> Option<MigrationState> {
migrations().abort(prefix)
}
thread_local! {
static INGESTING_PREFIX: std::cell::RefCell<Option<Vec<u8>>> =
const { std::cell::RefCell::new(None) };
}
pub(crate) struct IngestGuard {
_priv: (),
}
impl IngestGuard {
pub(crate) fn enter(prefix: Vec<u8>) -> Self {
INGESTING_PREFIX.with(|c| *c.borrow_mut() = Some(prefix));
Self { _priv: () }
}
}
impl Drop for IngestGuard {
fn drop(&mut self) {
INGESTING_PREFIX.with(|c| *c.borrow_mut() = None);
}
}
fn is_ingesting_match(key: &[u8]) -> bool {
INGESTING_PREFIX.with(|c| c.borrow().as_ref().is_some_and(|p| key.starts_with(p)))
}
pub(crate) fn self_node_id() -> Option<&'static str> {
SELF_NODE_ID.get().and_then(Option::as_deref)
}
pub(crate) fn install(cfg: &Config) -> Result<(), String> {
if !cfg.cluster.scopes.is_empty() {
let scopes: Vec<Scope> = cfg
.cluster
.scopes
.iter()
.map(|e| {
let s = Scope::new(e.prefix.clone(), e.writer.clone());
match e.fallback.as_ref() {
Some(fb) => s.with_fallback(fb.clone()),
None => s,
}
})
.collect();
let table = OwnershipTable::new(scopes).map_err(|e| e.to_string())?;
for s in table.scopes_without_fallback() {
let prefix_lossy = String::from_utf8_lossy(s.prefix());
eprintln!(
"kevy: WARN scope {prefix_lossy:?} has no fallback declared — writes for this scope fail if writer {:?} is DOWN",
s.writer(),
);
}
let _ = OWNERSHIP.set(Some(table));
} else {
let _ = OWNERSHIP.set(None);
}
let peers: Vec<(String, String)> = cfg
.cluster
.peers
.iter()
.map(|p| {
let reported = p.client_port.unwrap_or(p.port);
(p.node_id.clone(), format!("{}:{}", p.host, reported))
})
.collect();
let _ = PEER_ADDRS.set(peers);
Ok(())
}
#[inline]
pub(crate) fn is_active() -> bool {
matches!(OWNERSHIP.get(), Some(Some(_)))
}
static SELF_NODE_ID: OnceLock<Option<String>> = OnceLock::new();
pub(crate) fn install_self_id(cfg: &Config) {
let id = if cfg.cluster.node_id.is_empty() {
None
} else {
Some(cfg.cluster.node_id.clone())
};
let _ = SELF_NODE_ID.set(id);
}
pub(crate) fn route_write(key: &[u8]) -> Option<WriteRedirect> {
if is_ingesting_match(key) {
return None;
}
if let Some(m) = migrations().match_migrating(key) {
return Some(WriteRedirect::Quiesced {
to_addr: resolve_addr(&m.to),
});
}
if let Some(m) = migrations().match_migrated(key) {
return Some(WriteRedirect::Misdirected(resolve_addr(&m.to)));
}
let table = OWNERSHIP.get()?.as_ref()?;
let self_id = self_node_id()?;
let routing = match crate::elect_integration::current_snapshot() {
Some(snap) => table.route_with_fallback_state(key, self_id, |id| {
snap.down_peers.iter().any(|d| d == id)
}),
None => table.route(key, self_id),
};
match routing {
Routing::Owned | Routing::Unknown => None,
Routing::Misdirected { target } => Some(WriteRedirect::Misdirected(resolve_addr(target))),
}
}
pub(crate) enum WriteRedirect {
Misdirected(String),
Quiesced { to_addr: String },
}
fn resolve_addr(node_id: &str) -> String {
let Some(peers) = PEER_ADDRS.get() else {
return node_id.to_string();
};
peers
.iter()
.find(|(id, _)| id == node_id)
.map(|(_, addr)| addr.clone())
.unwrap_or_else(|| node_id.to_string())
}
pub(crate) fn encode_misdirected(out: &mut Vec<u8>, target: &str) {
out.extend_from_slice(b"-MISDIRECTED writer is ");
out.extend_from_slice(target.as_bytes());
out.extend_from_slice(b"\r\n");
}
pub(crate) fn encode_quiesced(out: &mut Vec<u8>, target: &str) {
out.extend_from_slice(b"-QUIESCED migrating to ");
out.extend_from_slice(target.as_bytes());
out.extend_from_slice(b"\r\n");
}
#[allow(dead_code)] pub(crate) fn peer_addr(node_id: &str) -> Option<String> {
let peers = PEER_ADDRS.get()?;
peers
.iter()
.find(|(id, _)| id == node_id)
.map(|(_, addr)| addr.clone())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn encode_misdirected_wire_shape() {
let mut out = Vec::new();
encode_misdirected(&mut out, "10.0.0.1:6004");
assert_eq!(out, b"-MISDIRECTED writer is 10.0.0.1:6004\r\n");
}
#[test]
fn resolve_addr_falls_back_to_node_id_when_no_peers() {
let v = resolve_addr("kevy-some-unmapped-id");
assert_eq!(v, "kevy-some-unmapped-id");
}
}