#![allow(clippy::format_push_string)]
use kevy_config::Config;
use kevy_resp::{ArgvView, encode_array_len, encode_bulk, encode_integer, encode_simple_string};
use kevy_store::Store;
use super::wrong_args;
use crate::state::Ctx;
fn node_id(i: usize) -> String {
format!("{:040x}", i + 1)
}
fn advertised_ip(cfg: &Config) -> String {
let [a, b, c, d] = cfg.server.bind;
if [a, b, c, d] == [0, 0, 0, 0] { "127.0.0.1".into() } else { format!("{a}.{b}.{c}.{d}") }
}
pub(crate) fn cmd_cluster<A: ArgvView + ?Sized>(
ctx: &Ctx<'_>,
store: &mut Store,
args: &A,
out: &mut Vec<u8>,
) {
let cfg = &ctx.state.config();
let sub = match args.get(1) {
Some(s) => s.to_ascii_uppercase(),
None => return wrong_args(out, "cluster"),
};
let n = cfg.server.threads.max(1);
let enabled = cfg.cluster.enabled;
match sub.as_slice() {
b"INFO" => {
let known = if enabled { cfg.cluster.peers.len().max(1) } else { 1 };
let body = format!(
"cluster_enabled:{}\r\ncluster_state:ok\r\n\
cluster_slots_assigned:16384\r\ncluster_slots_ok:16384\r\n\
cluster_slots_pfail:0\r\ncluster_slots_fail:0\r\n\
cluster_known_nodes:{}\r\ncluster_size:{}\r\n\
cluster_current_epoch:0\r\ncluster_my_epoch:0\r\n",
u8::from(enabled),
known,
if enabled { n } else { 1 },
);
encode_bulk(out, body.as_bytes());
}
b"NODES" if enabled => {
encode_bulk(out, nodes_text(cfg, n, live_role(ctx), ctx.shard.shard_id()).as_bytes());
}
b"NODES" => {
let role = live_role(ctx);
let body = format!(
"0000000000000000000000000000000000000000 :0@0 myself,{role} - 0 0 0 connected 0-16383\r\n",
);
encode_bulk(out, body.as_bytes());
}
b"SLOTS" if enabled => encode_slots(cfg, n, out),
b"SHARDS" if enabled => encode_shards(cfg, n, out),
b"SLOTS" | b"SHARDS" => encode_array_len(out, 0),
b"MYID" if enabled => encode_bulk(out, node_id(ctx.shard.shard_id()).as_bytes()),
b"MYID" => encode_bulk(out, b"0000000000000000000000000000000000000000"),
b"KEYSLOT" => match args.get(2) {
Some(key) => encode_integer(out, i64::from(kevy_hash::key_hash_slot(key))),
None => wrong_args(out, "cluster|keyslot"),
},
b"COUNTKEYSINSLOT" => {
let Some(slot) = args
.get(2)
.and_then(|s| std::str::from_utf8(s).ok())
.and_then(|s| s.parse::<u16>().ok())
.filter(|&s| s < 16384)
else {
return wrong_args(out, "cluster|countkeysinslot");
};
let mut count = 0i64;
store.snapshot_each(|key, _, _| {
if kevy_hash::key_hash_slot(key) == slot {
count += 1;
}
});
encode_integer(out, count);
}
_ => encode_simple_string(out, "OK"),
}
}
fn for_each_node(cfg: &Config, n: usize, mut f: impl FnMut(usize, &str, i64, u16, u16)) {
let ip = advertised_ip(cfg);
let base = i64::from(crate::cluster_port_base(cfg));
for i in 0..n {
let (start, end) = kevy_rt::shard_slot_range(i, n);
f(i, &ip, base + i as i64, start, end);
}
}
fn live_role(ctx: &Ctx<'_>) -> &'static str {
if ctx.state.replication.current_upstream().is_some() { "slave" } else { "master" }
}
fn nodes_text(cfg: &Config, n: usize, role: &str, me: usize) -> String {
let mut body = String::new();
for_each_node(cfg, n, |i, ip, port, start, end| {
let flags = if i == me { format!("myself,{role}") } else { role.to_string() };
body.push_str(&format!(
"{} {ip}:{port}@{port} {flags} - 0 0 {} connected {start}-{end}\r\n",
node_id(i),
i + 1,
));
});
body
}
fn encode_slots(cfg: &Config, n: usize, out: &mut Vec<u8>) {
encode_array_len(out, n as i64);
for_each_node(cfg, n, |i, ip, port, start, end| {
encode_array_len(out, 3);
encode_integer(out, i64::from(start));
encode_integer(out, i64::from(end));
encode_array_len(out, 4);
encode_bulk(out, ip.as_bytes());
encode_integer(out, port);
encode_bulk(out, node_id(i).as_bytes());
encode_array_len(out, 0);
});
}
fn encode_shards(cfg: &Config, n: usize, out: &mut Vec<u8>) {
encode_array_len(out, n as i64);
for_each_node(cfg, n, |i, ip, port, start, end| {
encode_array_len(out, 4); encode_bulk(out, b"slots");
encode_array_len(out, 2);
encode_integer(out, i64::from(start));
encode_integer(out, i64::from(end));
encode_bulk(out, b"nodes");
encode_array_len(out, 1);
encode_array_len(out, 12); encode_bulk(out, b"id");
encode_bulk(out, node_id(i).as_bytes());
encode_bulk(out, b"port");
encode_integer(out, port);
encode_bulk(out, b"ip");
encode_bulk(out, ip.as_bytes());
encode_bulk(out, b"endpoint");
encode_bulk(out, ip.as_bytes());
encode_bulk(out, b"role");
encode_bulk(out, b"master");
encode_bulk(out, b"health");
encode_bulk(out, b"online");
});
}