use kevy_config::Config;
use kevy_scope::{MigrationState, MigrationTable, OwnershipTable, Routing, Scope};
use super::{RuntimeState, ShardCtx};
#[derive(Debug)]
pub(crate) struct ScopeState {
ownership: Option<OwnershipTable>,
peer_addrs: Vec<(String, String)>,
self_node_id: Option<String>,
migrations: MigrationTable,
}
impl ScopeState {
pub(crate) fn from_config(cfg: &Config) -> Result<Self, String> {
let ownership = build_ownership(cfg)?;
let peer_addrs: 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 self_node_id =
if cfg.cluster.node_id.is_empty() { None } else { Some(cfg.cluster.node_id.clone()) };
Ok(Self { ownership, peer_addrs, self_node_id, migrations: MigrationTable::new() })
}
#[inline]
pub(crate) fn is_active(&self) -> bool {
self.ownership.is_some()
}
pub(crate) fn self_node_id(&self) -> Option<&str> {
self.self_node_id.as_deref()
}
pub(crate) fn peer_addr(&self, node_id: &str) -> Option<String> {
self.peer_addrs.iter().find(|(id, _)| id == node_id).map(|(_, addr)| addr.clone())
}
pub(crate) fn migration_start(
&self,
prefix: Vec<u8>,
from: String,
to: String,
) -> Result<(), kevy_scope::MigrationError> {
self.migrations.start(prefix, from, to)
}
pub(crate) fn migration_commit(&self, prefix: &[u8]) -> Option<MigrationState> {
self.migrations.commit(prefix)
}
pub(crate) fn migration_abort(&self, prefix: &[u8]) -> Option<MigrationState> {
self.migrations.abort(prefix)
}
fn resolve_addr(&self, node_id: &str) -> String {
self.peer_addr(node_id).unwrap_or_else(|| node_id.to_string())
}
}
pub(crate) enum WriteRedirect {
Misdirected(String),
Quiesced { to_addr: String },
}
impl RuntimeState {
pub(crate) fn route_write(&self, key: &[u8], shard: &ShardCtx) -> Option<WriteRedirect> {
if shard.ingesting_matches(key) {
return None;
}
let scope = &self.scope;
if let Some(m) = scope.migrations.match_migrating(key) {
return Some(WriteRedirect::Quiesced { to_addr: scope.resolve_addr(&m.to) });
}
if let Some(m) = scope.migrations.match_migrated(key) {
return Some(WriteRedirect::Misdirected(scope.resolve_addr(&m.to)));
}
let table = scope.ownership.as_ref()?;
let self_id = scope.self_node_id()?;
let routing = match self.election.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(scope.resolve_addr(target)))
}
}
}
}
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");
}
fn build_ownership(cfg: &Config) -> Result<Option<OwnershipTable>, String> {
if cfg.cluster.scopes.is_empty() {
return Ok(None);
}
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(),
);
}
Ok(Some(table))
}
#[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_unmapped() {
let scope = ScopeState::from_config(&Config::default()).unwrap();
assert_eq!(scope.resolve_addr("kevy-some-unmapped-id"), "kevy-some-unmapped-id");
}
#[test]
fn from_config_default_is_inert() {
let scope = ScopeState::from_config(&Config::default()).unwrap();
assert!(!scope.is_active());
assert!(scope.self_node_id().is_none());
}
}