use std::io::{Read, Write};
use std::net::{Shutdown, TcpListener, TcpStream};
use std::sync::atomic::Ordering::Relaxed;
use std::sync::atomic::{AtomicBool, AtomicU64};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use yo_common::lock::Lock;
use yo_common::{Code, Error};
use super::super::Server;
use super::super::pubsub::{self, Kind};
use super::{
FLAG_FAIL, FLAG_HANDSHAKE, FLAG_MASTER, FLAG_MEET, FLAG_MIGRATE_TO, FLAG_MYSELF, FLAG_NOADDR,
FLAG_NOFAILOVER, FLAG_PFAIL, FLAG_SLAVE, Map, Node, SLOTS, new_id,
};
use crate::reply::Out;
const SIG: &[u8; 4] = b"RCmb";
const PROTO_VER: u16 = 1;
const NAME_LEN: usize = 40;
const IP_LEN: usize = 46;
const BITMAP_LEN: usize = SLOTS / 8;
const HDR_LEN: usize = 2256;
const GOSSIP_LEN: usize = 104;
const MAX_PACKET: usize = 64 * 1024 * 1024;
const O_TOTLEN: usize = 4;
const O_VER: usize = 8;
const O_PORT: usize = 10;
const O_TYPE: usize = 12;
const O_COUNT: usize = 14;
const O_CURRENT_EPOCH: usize = 16;
const O_CONFIG_EPOCH: usize = 24;
const O_OFFSET: usize = 32;
const O_SENDER: usize = 40;
const O_SLOTS: usize = 80;
const O_SLAVEOF: usize = 2128;
const O_MYIP: usize = 2168;
const O_EXTENSIONS: usize = 2214;
const O_PPORT: usize = 2246;
const O_CPORT: usize = 2248;
const O_FLAGS: usize = 2250;
const O_STATE: usize = 2252;
const O_MFLAGS: usize = 2253;
const G_NAME: usize = 0;
const G_PING: usize = 40;
const G_PONG: usize = 44;
const G_IP: usize = 48;
const G_PORT: usize = 94;
const G_CPORT: usize = 96;
const G_FLAGS: usize = 98;
const T_PING: u16 = 0;
const T_PONG: u16 = 1;
const T_MEET: u16 = 2;
const T_FAIL: u16 = 3;
const T_PUBLISH: u16 = 4;
const T_AUTH_REQUEST: u16 = 5;
const T_AUTH_ACK: u16 = 6;
const T_UPDATE: u16 = 7;
const T_MFSTART: u16 = 8;
const T_MODULE: u16 = 9;
const T_PUBLISHSHARD: u16 = 10;
const MF_PAUSED: u8 = 1;
const MF_FORCEACK: u8 = 2;
const MF_EXT_DATA: u8 = 4;
const X_HOSTNAME: u16 = 0;
const X_HUMAN_NAME: u16 = 1;
const X_FORGOTTEN: u16 = 2;
const X_SHARDID: u16 = 3;
const X_SECRET: u16 = 4;
const STATE_OK: u8 = 0;
const STATE_FAIL: u8 = 1;
const BLACKLIST_MS: u64 = 60_000;
const CRON_MS: u64 = 100;
const PING_MS: u64 = 1000;
const NODE_TIMEOUT_MS: u64 = 15_000;
const CONNECT_TIMEOUT: Duration = Duration::from_millis(500);
const ELECTION_DELAY_MS: u64 = 500;
const RANK_DELAY_MS: u64 = 1000;
const ELECTION_TIMEOUT_MS: u64 = if NODE_TIMEOUT_MS * 2 > 2000 {
NODE_TIMEOUT_MS * 2
} else {
2000
};
const STALE_DATA_MS: u64 = 10_000 + NODE_TIMEOUT_MS * 10;
const MANUAL_TIMEOUT_MS: u64 = 5000;
const MANUAL_PAUSE_MULT: u64 = 2;
const OFFSET_UNKNOWN: u64 = u64::MAX;
#[derive(Default)]
pub(super) struct Vote {
given: AtomicU64,
at: AtomicU64,
count: AtomicU64,
epoch: AtomicU64,
rank: AtomicU64,
sent: AtomicBool,
}
impl Vote {
pub(super) fn given(&self) -> u64 {
self.given.load(Relaxed)
}
pub(super) fn reload(&self, epoch: u64) {
self.given.store(epoch, Relaxed);
}
}
pub(super) struct Manual {
end: AtomicU64,
can_start: AtomicBool,
offset: AtomicU64,
held: AtomicU64,
}
impl Default for Manual {
fn default() -> Manual {
Manual {
end: AtomicU64::new(0),
can_start: AtomicBool::new(false),
offset: AtomicU64::new(OFFSET_UNKNOWN),
held: AtomicU64::new(0),
}
}
}
fn be16(p: &[u8], at: usize) -> u16 {
p.get(at..at + 2)
.map_or(0, |b| u16::from_be_bytes([b[0], b[1]]))
}
fn be32(p: &[u8], at: usize) -> u32 {
p.get(at..at + 4)
.map_or(0, |b| u32::from_be_bytes([b[0], b[1], b[2], b[3]]))
}
fn be64(p: &[u8], at: usize) -> u64 {
p.get(at..at + 8).map_or(0, |b| {
u64::from_be_bytes([b[0], b[1], b[2], b[3], b[4], b[5], b[6], b[7]])
})
}
fn put16(p: &mut [u8], at: usize, v: u16) {
p[at..at + 2].copy_from_slice(&v.to_be_bytes());
}
fn put64(p: &mut [u8], at: usize, v: u64) {
p[at..at + 8].copy_from_slice(&v.to_be_bytes());
}
fn text(p: &[u8], at: usize, len: usize) -> String {
let Some(raw) = p.get(at..at + len) else {
return String::new();
};
let end = raw.iter().position(|b| *b == 0).unwrap_or(len);
String::from_utf8_lossy(&raw[..end]).into_owned()
}
fn node_id(p: &[u8], at: usize) -> Option<String> {
let raw = p.get(at..at + NAME_LEN)?;
if raw.iter().all(|b| *b == 0) {
return None;
}
if !raw.iter().all(u8::is_ascii_hexdigit) {
return None;
}
Some(String::from_utf8_lossy(raw).into_owned())
}
fn bit(bitmap: &[u8], slot: usize) -> bool {
bitmap
.get(slot / 8)
.is_some_and(|b| b & (1 << (slot % 8)) != 0)
}
fn set_bit(bitmap: &mut [u8], slot: usize) {
bitmap[slot / 8] |= 1 << (slot % 8);
}
fn pick(n: usize) -> u16 {
let mut raw = [0u8; 8];
yo_common::entropy::fill(&mut raw);
(u64::from_be_bytes(raw) % n as u64) as u16
}
pub(super) struct Wire {
pub(super) inbound: bool,
pub(super) node: Lock<String>,
pub(super) created: u64,
peer: String,
local: String,
sock: Mutex<TcpStream>,
dead: AtomicBool,
sent: AtomicU64,
}
impl Wire {
fn new(sock: TcpStream, inbound: bool, node: &str, now: u64) -> Option<Arc<Wire>> {
let peer = sock.peer_addr().ok()?.ip().to_string();
let local = sock.local_addr().ok()?.ip().to_string();
Some(Arc::new(Wire {
inbound,
node: Lock::new(String::from(node)),
created: now,
peer,
local,
sock: Mutex::new(sock),
dead: AtomicBool::new(false),
sent: AtomicU64::new(0),
}))
}
fn named(&self) -> String {
let held = self.node.lock();
yo_alloc::allow(|| held.clone())
}
fn name(&self, id: &str) {
let mut held = self.node.lock();
yo_alloc::allow(|| {
held.clear();
held.push_str(id);
});
}
fn send(&self, packet: &[u8]) {
if self.dead.load(Relaxed) {
return;
}
let Ok(mut sock) = self.sock.lock() else {
self.dead.store(true, Relaxed);
return;
};
if sock.write_all(packet).is_err() {
self.dead.store(true, Relaxed);
let _ = sock.shutdown(Shutdown::Both);
return;
}
self.sent.store(packet.len() as u64, Relaxed);
}
fn kill(&self) {
self.dead.store(true, Relaxed);
if let Ok(sock) = self.sock.lock() {
let _ = sock.shutdown(Shutdown::Both);
}
}
}
#[derive(Default)]
pub(super) struct Bus {
links: Lock<Vec<Arc<Wire>>>,
blacklist: Lock<Vec<(String, u64)>>,
pub(super) secret: Lock<String>,
on: AtomicBool,
dirty: AtomicBool,
}
impl Bus {
fn add(&self, wire: &Arc<Wire>) {
let mut links = self.links.lock();
yo_alloc::allow(|| {
links.retain(|held| !held.dead.load(Relaxed));
links.push(Arc::clone(wire));
});
}
fn outbound(&self, id: &str) -> Option<Arc<Wire>> {
let links = self.links.lock();
links
.iter()
.find(|held| !held.inbound && !held.dead.load(Relaxed) && *held.node.lock() == *id)
.map(Arc::clone)
}
fn all(&self) -> Vec<Arc<Wire>> {
let links = self.links.lock();
yo_alloc::allow(|| {
links
.iter()
.filter(|held| !held.dead.load(Relaxed))
.map(Arc::clone)
.collect()
})
}
fn cut(&self, id: &str) {
let mut links = self.links.lock();
links.retain(|held| {
let theirs = held.node.lock();
if *theirs != *id {
return true;
}
drop(theirs);
held.kill();
false
});
}
fn blacklisted(&self, id: &str, now: u64) -> bool {
let mut list = self.blacklist.lock();
list.retain(|(_, until)| *until > now);
list.iter().any(|(held, _)| held == id)
}
fn blacklist(&self, id: &str, now: u64) {
let mut list = self.blacklist.lock();
yo_alloc::allow(|| {
list.retain(|(held, until)| held != id && *until > now);
list.push((String::from(id), now + BLACKLIST_MS));
});
}
}
impl Server {
pub fn start_cluster_bus(self: &Arc<Server>) -> Result<(), String> {
if !self.cluster_enabled() || self.cluster.bus.on.swap(true, Relaxed) {
return Ok(());
}
let (port, id) = {
let map = self.cluster.map.lock();
(
map.nodes[0].bus,
yo_alloc::allow(|| map.nodes[0].id.clone()),
)
};
let door = TcpListener::bind(("0.0.0.0", port))
.map_err(|e| format!("cluster bus port {port} could not be bound: {e}"))?;
let _ = id;
let accepting = Arc::clone(self);
spawn("yo-bus-accept", move || accept(&accepting, &door));
let ticking = Arc::clone(self);
spawn("yo-bus-cron", move || cron(&ticking));
Ok(())
}
}
fn spawn(name: &str, body: impl FnOnce() + Send + 'static) {
yo_alloc::allow(|| {
let _ = std::thread::Builder::new()
.name(String::from(name))
.spawn(body);
});
}
fn accept(server: &Arc<Server>, door: &TcpListener) {
loop {
let Ok((sock, _)) = door.accept() else {
std::thread::sleep(Duration::from_millis(CRON_MS));
continue;
};
let _ = sock.set_nodelay(true);
let now = server.now_ms();
let Some(wire) = Wire::new(sock, true, "", now) else {
continue;
};
server.cluster.bus.add(&wire);
let reading = Arc::clone(server);
spawn("yo-bus-link", move || pump(&reading, &wire));
}
}
fn unreachable(server: &Arc<Server>, id: &str) {
let mut map = server.cluster.map.lock();
if let Some(at) = map.find(id.as_bytes())
&& map.nodes[usize::from(at)].ping_sent == 0
{
map.nodes[usize::from(at)].ping_sent = server.now_ms();
}
}
fn dial(server: &Arc<Server>, id: &str, host: &str, bus: u16, meet: bool) -> bool {
let Ok(addrs) = std::net::ToSocketAddrs::to_socket_addrs(&(host, bus)) else {
unreachable(server, id);
return false;
};
let mut sock = None;
for addr in addrs {
if let Ok(open) = TcpStream::connect_timeout(&addr, CONNECT_TIMEOUT) {
sock = Some(open);
break;
}
}
let Some(sock) = sock else {
unreachable(server, id);
return false;
};
let _ = sock.set_nodelay(true);
let now = server.now_ms();
let Some(wire) = Wire::new(sock, false, id, now) else {
unreachable(server, id);
return false;
};
server.cluster.bus.add(&wire);
let packet = {
let mut map = server.cluster.map.lock();
let at = map.find(id.as_bytes());
if let Some(at) = at {
map.nodes[usize::from(at)].linked = true;
map.nodes[usize::from(at)].flags &= !FLAG_MEET;
if !meet && map.nodes[usize::from(at)].ping_sent == 0 {
map.nodes[usize::from(at)].ping_sent = now;
}
}
ping(server, &map, if meet { T_MEET } else { T_PING }, at)
};
wire.send(&packet);
let reading = Arc::clone(server);
spawn("yo-bus-link", move || pump(&reading, &wire));
true
}
fn pump(server: &Arc<Server>, wire: &Arc<Wire>) {
let sock = {
let Ok(held) = wire.sock.lock() else {
return;
};
held.try_clone()
};
let Ok(mut sock) = sock else {
wire.kill();
return;
};
let mut head = [0u8; 8];
let mut buf: Vec<u8> = yo_alloc::allow(|| Vec::with_capacity(HDR_LEN * 2));
loop {
if wire.dead.load(Relaxed) || sock.read_exact(&mut head).is_err() {
break;
}
if &head[0..4] != SIG {
break;
}
let total = u32::from_be_bytes([head[4], head[5], head[6], head[7]]) as usize;
if !(16..=MAX_PACKET).contains(&total) {
break;
}
yo_alloc::allow(|| {
buf.clear();
buf.resize(total, 0);
});
buf[..8].copy_from_slice(&head);
if sock.read_exact(&mut buf[8..]).is_err() {
break;
}
if !process(server, wire, &buf) {
break;
}
}
wire.kill();
let name = wire.named();
if !name.is_empty() && !wire.inbound {
let mut map = server.cluster.map.lock();
if let Some(at) = map.find(name.as_bytes()) {
map.nodes[usize::from(at)].linked = false;
}
}
}
fn header(server: &Server, map: &Map, kind: u16) -> Vec<u8> {
let mut p = yo_alloc::allow(|| vec![0u8; HDR_LEN]);
p[0..4].copy_from_slice(SIG);
put16(&mut p, O_VER, PROTO_VER);
put16(&mut p, O_TYPE, kind);
put16(&mut p, O_PORT, map.nodes[0].port);
put16(&mut p, O_CPORT, map.nodes[0].bus);
put16(&mut p, O_PPORT, 0);
put16(&mut p, O_FLAGS, map.nodes[0].flags);
let speaking = map.nodes[0]
.master
.filter(|_| map.nodes[0].flags & FLAG_SLAVE != 0)
.map_or(0, usize::from);
put64(&mut p, O_CURRENT_EPOCH, server.cluster.epoch.load(Relaxed));
put64(&mut p, O_CONFIG_EPOCH, map.nodes[speaking].epoch);
put64(&mut p, O_OFFSET, server.repl_offset());
p[O_SENDER..O_SENDER + NAME_LEN].copy_from_slice(map.nodes[0].id.as_bytes());
for slot in 0..SLOTS {
if map.owner[slot] == Some(speaking as u16) {
set_bit(&mut p[O_SLOTS..O_SLOTS + BITMAP_LEN], slot);
}
}
if let Some(master) = map.nodes[0].master {
let id = map.nodes[usize::from(master)].id.as_bytes();
p[O_SLAVEOF..O_SLAVEOF + NAME_LEN].copy_from_slice(id);
}
let _ = O_MYIP;
p[O_STATE] = if server.cluster_up() {
STATE_OK
} else {
STATE_FAIL
};
p[O_MFLAGS] = MF_EXT_DATA;
if map.nodes[0].is_master() && server.cluster.manual.end.load(Relaxed) != 0 {
p[O_MFLAGS] |= MF_PAUSED;
}
p
}
fn seal(p: &mut [u8]) {
let total = p.len() as u32;
p[O_TOTLEN..O_TOTLEN + 4].copy_from_slice(&total.to_be_bytes());
}
fn ping(server: &Server, map: &Map, kind: u16, to: Option<u16>) -> Vec<u8> {
let mut p = header(server, map, kind);
let now = server.now_ms();
let mut want = (map.nodes.len() / 10).max(3);
if map.nodes.len() >= 2 {
want = want.min(map.nodes.len() - 2);
} else {
want = 0;
}
let mut count = 0u16;
let mut sent: Vec<u16> = yo_alloc::allow(|| Vec::with_capacity(want + 4));
let mut tries = want * 3;
while sent.len() < want && tries > 0 {
tries -= 1;
let at = pick(map.nodes.len());
if at == 0 || Some(at) == to || sent.contains(&at) {
continue;
}
let node = &map.nodes[usize::from(at)];
if node.flags & (FLAG_HANDSHAKE | FLAG_NOADDR) != 0 {
continue;
}
if !node.linked && map.runs(at).is_empty() {
continue;
}
sent.push(at);
}
for at in 1..map.nodes.len() as u16 {
if map.nodes[usize::from(at)].flags & FLAG_PFAIL != 0 && !sent.contains(&at) {
sent.push(at);
}
}
for at in sent {
let node = &map.nodes[usize::from(at)];
let mut entry = [0u8; GOSSIP_LEN];
entry[G_NAME..G_NAME + NAME_LEN].copy_from_slice(node.id.as_bytes());
entry[G_PING..G_PING + 4].copy_from_slice(&((node.ping_sent / 1000) as u32).to_be_bytes());
entry[G_PONG..G_PONG + 4].copy_from_slice(&((node.pong_recv / 1000) as u32).to_be_bytes());
let host = node.host.as_bytes();
let take = host.len().min(IP_LEN - 1);
entry[G_IP..G_IP + take].copy_from_slice(&host[..take]);
entry[G_PORT..G_PORT + 2].copy_from_slice(&node.port.to_be_bytes());
entry[G_CPORT..G_CPORT + 2].copy_from_slice(&node.bus.to_be_bytes());
entry[G_FLAGS..G_FLAGS + 2].copy_from_slice(&node.flags.to_be_bytes());
yo_alloc::allow(|| p.extend_from_slice(&entry));
count += 1;
}
let _ = now;
put16(&mut p, O_COUNT, count);
let mut exts = 0u16;
push_ext(&mut p, X_SHARDID, map.nodes[0].shard.as_bytes());
exts += 1;
let secret = server.cluster.bus.secret.lock();
if secret.len() == NAME_LEN {
push_ext(&mut p, X_SECRET, secret.as_bytes());
exts += 1;
}
drop(secret);
put16(&mut p, O_EXTENSIONS, exts);
seal(&mut p);
p
}
fn push_ext(p: &mut Vec<u8>, kind: u16, data: &[u8]) {
let len = 8 + data.len().div_ceil(8) * 8;
yo_alloc::allow(|| {
p.extend_from_slice(&(len as u32).to_be_bytes());
p.extend_from_slice(&kind.to_be_bytes());
p.extend_from_slice(&[0, 0]);
p.extend_from_slice(data);
p.resize(p.len() + (len - 8 - data.len()), 0);
});
}
fn fail_packet(server: &Server, map: &Map, about: &str) -> Vec<u8> {
let mut p = header(server, map, T_FAIL);
yo_alloc::allow(|| p.extend_from_slice(about.as_bytes()));
seal(&mut p);
p
}
fn publish_packet(server: &Server, map: &Map, shard: bool, channel: &[u8], body: &[u8]) -> Vec<u8> {
let kind = if shard { T_PUBLISHSHARD } else { T_PUBLISH };
let mut p = header(server, map, kind);
yo_alloc::allow(|| {
p.extend_from_slice(&(channel.len() as u32).to_be_bytes());
p.extend_from_slice(&(body.len() as u32).to_be_bytes());
p.extend_from_slice(channel);
p.extend_from_slice(body);
});
seal(&mut p);
p
}
fn update_packet(server: &Server, map: &Map, about: u16) -> Vec<u8> {
let mut p = header(server, map, T_UPDATE);
let node = &map.nodes[usize::from(about)];
let mut body = [0u8; 8 + NAME_LEN + BITMAP_LEN];
body[0..8].copy_from_slice(&node.epoch.to_be_bytes());
body[8..8 + NAME_LEN].copy_from_slice(node.id.as_bytes());
for slot in 0..SLOTS {
if map.owner[slot] == Some(about) {
set_bit(&mut body[8 + NAME_LEN..], slot);
}
}
yo_alloc::allow(|| p.extend_from_slice(&body));
seal(&mut p);
p
}
#[derive(Default)]
struct Todo {
reply: Vec<Vec<u8>>,
shout: Vec<Vec<u8>>,
deliver: Option<(bool, Vec<u8>, Vec<u8>)>,
save: bool,
recount: bool,
lost: Vec<u16>,
demoted: bool,
close: bool,
}
fn process(server: &Arc<Server>, wire: &Arc<Wire>, p: &[u8]) -> bool {
if be16(p, O_VER) != PROTO_VER {
return true;
}
let kind = be16(p, O_TYPE);
let Some(explen) = expected(p, kind) else {
return true;
};
if explen != p.len() {
return true;
}
let now = server.now_ms();
let mut todo = Todo::default();
{
let mut map = server.cluster.map.lock();
digest(server, wire, p, kind, now, &mut map, &mut todo);
}
for packet in &todo.reply {
wire.send(packet);
}
if !todo.shout.is_empty() {
for link in server.cluster.bus.all() {
for packet in &todo.shout {
link.send(packet);
}
}
}
if let Some((shard, channel, body)) = todo.deliver
&& server.anyone_subscribed()
{
let kind = if shard { Kind::Shard } else { Kind::Channel };
pubsub::deliver(server, kind, &channel, &body);
}
if todo.save {
server.cluster.bus.dirty.store(true, Relaxed);
}
if todo.recount {
server.recount_coverage();
}
server.asm_slots_moved(&todo.lost, todo.demoted);
!todo.close
}
fn expected(p: &[u8], kind: u16) -> Option<usize> {
match kind {
T_PING | T_PONG | T_MEET => {
let count = usize::from(be16(p, O_COUNT));
let mut len = HDR_LEN.checked_add(count.checked_mul(GOSSIP_LEN)?)?;
if p.get(O_MFLAGS).is_some_and(|f| f & MF_EXT_DATA != 0) {
let mut left = be16(p, O_EXTENSIONS);
let mut at = len;
while left > 0 {
left -= 1;
let extlen = be32(p, at) as usize;
if extlen < 8 || !extlen.is_multiple_of(8) || p.len().checked_sub(len)? < extlen
{
return None;
}
len += extlen;
at += extlen;
}
}
Some(len)
}
T_FAIL => Some(HDR_LEN + NAME_LEN),
T_PUBLISH | T_PUBLISHSHARD => {
let channel = be32(p, HDR_LEN) as usize;
let body = be32(p, HDR_LEN + 4) as usize;
HDR_LEN
.checked_add(8)?
.checked_add(channel)?
.checked_add(body)
}
T_AUTH_REQUEST | T_AUTH_ACK | T_MFSTART => Some(HDR_LEN),
T_UPDATE => Some(HDR_LEN + 8 + NAME_LEN + BITMAP_LEN),
_ => Some(p.len()),
}
}
#[allow(clippy::too_many_lines)]
fn digest(
server: &Arc<Server>,
wire: &Arc<Wire>,
p: &[u8],
kind: u16,
now: u64,
map: &mut Map,
todo: &mut Todo,
) {
let flags = be16(p, O_FLAGS);
let claimed = node_id(p, O_SENDER);
let linked = wire.named();
let mut sender = None;
if !linked.is_empty()
&& let Some(at) = map.find(linked.as_bytes())
&& map.nodes[usize::from(at)].flags & FLAG_HANDSHAKE == 0
{
sender = Some(at);
}
if sender.is_none()
&& let Some(id) = claimed.as_deref()
{
sender = map.find(id.as_bytes());
if sender.is_some() && linked.is_empty() {
wire.name(id);
}
}
if let Some(at) = sender {
let node = &mut map.nodes[usize::from(at)];
if p.get(O_MFLAGS).is_some_and(|f| f & MF_EXT_DATA != 0) {
node.flags |= super::FLAG_EXTENSIONS;
}
node.data_recv = now;
}
let sender_epoch = be64(p, O_CONFIG_EPOCH);
if let Some(at) = sender
&& map.nodes[usize::from(at)].flags & FLAG_HANDSHAKE == 0
{
let theirs = be64(p, O_CURRENT_EPOCH);
server.cluster.epoch.fetch_max(theirs, Relaxed);
let node = &mut map.nodes[usize::from(at)];
if sender_epoch > node.epoch {
node.epoch = sender_epoch;
todo.save = true;
}
node.offset = be64(p, O_OFFSET);
let manual = &server.cluster.manual;
if manual.end.load(Relaxed) != 0
&& manual.offset.load(Relaxed) == OFFSET_UNKNOWN
&& map.nodes[0].flags & FLAG_SLAVE != 0
&& map.nodes[0].master == Some(at)
&& p.get(O_MFLAGS).is_some_and(|f| f & MF_PAUSED != 0)
{
manual.offset.store(be64(p, O_OFFSET), Relaxed);
}
}
if kind == T_PING || kind == T_MEET {
if (kind == T_MEET || map.nodes[0].host.is_empty()) && map.nodes[0].host != wire.local {
yo_alloc::allow(|| map.nodes[0].host.clone_from(&wire.local));
todo.save = true;
}
if sender.is_none() && kind == T_MEET {
let host = {
let announced = text(p, O_MYIP, IP_LEN);
if announced.is_empty() {
yo_alloc::allow(|| wire.peer.clone())
} else {
announced
}
};
let id = yo_alloc::allow(|| String::from_utf8_lossy(&new_id()).into_owned());
let node = yo_alloc::allow(|| {
Node::new(
id,
host,
be16(p, O_PORT),
be16(p, O_CPORT),
FLAG_HANDSHAKE,
now,
)
});
yo_alloc::allow(|| map.nodes.push(node));
todo.save = true;
gossip(server, p, now, map, todo);
}
todo.reply.push(ping(server, map, T_PONG, sender));
}
match kind {
T_PING | T_PONG | T_MEET => {}
T_FAIL => {
if sender.is_some()
&& let Some(id) = node_id(p, HDR_LEN)
&& let Some(at) = map.find(id.as_bytes())
&& map.nodes[usize::from(at)].flags & (FLAG_FAIL | FLAG_MYSELF) == 0
{
let node = &mut map.nodes[usize::from(at)];
node.flags |= FLAG_FAIL;
node.flags &= !FLAG_PFAIL;
node.fail_time = now;
todo.save = true;
}
return;
}
T_PUBLISH | T_PUBLISHSHARD => {
if sender.is_none() {
todo.close = true;
return;
}
let channel = be32(p, HDR_LEN) as usize;
let body = be32(p, HDR_LEN + 4) as usize;
let at = HDR_LEN + 8;
todo.deliver = yo_alloc::allow(|| {
Some((
kind == T_PUBLISHSHARD,
p[at..at + channel].to_vec(),
p[at + channel..at + channel + body].to_vec(),
))
});
return;
}
T_UPDATE => {
if sender.is_none() {
todo.close = true;
return;
}
let epoch = be64(p, HDR_LEN);
let Some(id) = node_id(p, HDR_LEN + 8) else {
return;
};
let Some(about) = map.find(id.as_bytes()) else {
return;
};
if epoch <= map.nodes[usize::from(about)].epoch {
return;
}
map.nodes[usize::from(about)].epoch = epoch;
map.nodes[usize::from(about)].flags &= !FLAG_SLAVE;
map.nodes[usize::from(about)].flags |= FLAG_MASTER;
map.nodes[usize::from(about)].master = None;
claim_slots(
server,
map,
about,
epoch,
&p[HDR_LEN + 8 + NAME_LEN..],
todo,
);
todo.save = true;
return;
}
T_AUTH_REQUEST => {
if let Some(at) = sender
&& vote_if_needed(server, map, at, p, now)
{
let mut packet = header(server, map, T_AUTH_ACK);
seal(&mut packet);
yo_alloc::allow(|| todo.reply.push(packet));
todo.save = true;
}
return;
}
T_AUTH_ACK => {
let vote = &server.cluster.vote;
if let Some(at) = sender
&& map.nodes[usize::from(at)].is_master()
&& !map.runs(at).is_empty()
&& be64(p, O_CURRENT_EPOCH) >= vote.epoch.load(Relaxed)
{
vote.count.fetch_add(1, Relaxed);
}
return;
}
T_MFSTART => {
if let Some(at) = sender {
stand_down(server, map, at, now, todo);
}
return;
}
T_MODULE => return,
_ => return,
}
if !wire.inbound {
let held = map.find(linked.as_bytes());
if let Some(at) = held
&& map.nodes[usize::from(at)].flags & FLAG_HANDSHAKE != 0
{
match sender {
Some(known) => {
let host = yo_alloc::allow(|| wire.peer.clone());
let node = &mut map.nodes[usize::from(known)];
if node.host != host {
node.host = host;
node.port = be16(p, O_PORT);
node.bus = be16(p, O_CPORT);
}
map.forget(at);
todo.save = true;
todo.close = true;
return;
}
None => {
let Some(id) = claimed.clone() else {
return;
};
wire.name(&id);
let node = &mut map.nodes[usize::from(at)];
yo_alloc::allow(|| node.id = id);
node.flags &= !(FLAG_HANDSHAKE | FLAG_MEET);
node.flags |= flags & (FLAG_MASTER | FLAG_SLAVE);
node.pong_recv = now;
node.ping_sent = 0;
sender = Some(at);
todo.save = true;
}
}
} else if let Some(at) = held
&& Some(map.nodes[usize::from(at)].id.as_str()) != claimed.as_deref()
{
let node = &mut map.nodes[usize::from(at)];
node.flags |= FLAG_NOADDR;
node.host.clear();
node.port = 0;
node.bus = 0;
node.linked = false;
todo.save = true;
todo.close = true;
return;
}
}
let Some(at) = sender else {
return;
};
let node = &mut map.nodes[usize::from(at)];
node.flags &= !FLAG_NOFAILOVER;
node.flags |= flags & FLAG_NOFAILOVER;
if kind == T_PING && !wire.inbound {
let host = yo_alloc::allow(|| wire.peer.clone());
if node.host != host {
node.host = host;
node.port = be16(p, O_PORT);
node.bus = be16(p, O_CPORT);
todo.save = true;
}
}
if !wire.inbound && kind == T_PONG {
node.pong_recv = now;
node.ping_sent = 0;
if node.flags & FLAG_PFAIL != 0 {
node.flags &= !FLAG_PFAIL;
todo.save = true;
} else if node.flags & FLAG_FAIL != 0 {
clear_failure(map, at, now);
todo.save = true;
}
}
let follows = node_id(p, O_SLAVEOF);
match follows {
None => {
if map.nodes[usize::from(at)].flags & FLAG_SLAVE != 0 {
let node = &mut map.nodes[usize::from(at)];
node.flags &= !FLAG_SLAVE;
node.flags |= FLAG_MASTER;
node.master = None;
todo.save = true;
}
}
Some(id) => {
let master = map.find(id.as_bytes());
if map.nodes[usize::from(at)].is_master() {
let same_shard = master.is_some_and(|m| {
map.nodes[usize::from(m)].shard == map.nodes[usize::from(at)].shard
});
if same_shard && sender_epoch >= map.nodes[usize::from(at)].epoch {
let m = master.expect("same shard means there is one");
for slot in 0..SLOTS {
if map.owner[slot] == Some(at) {
map.owner[slot] = Some(m);
}
}
let promoted = &mut map.nodes[usize::from(m)];
promoted.flags &= !FLAG_SLAVE;
promoted.flags |= FLAG_MASTER;
promoted.master = None;
promoted.epoch = sender_epoch;
} else if !same_shard {
for slot in 0..SLOTS {
if map.owner[slot] == Some(at) {
map.owner[slot] = None;
}
}
}
let node = &mut map.nodes[usize::from(at)];
node.flags &= !(FLAG_MASTER | FLAG_MIGRATE_TO);
node.flags |= FLAG_SLAVE;
todo.save = true;
}
if let Some(m) = master
&& map.nodes[usize::from(at)].master != Some(m)
&& m != at
{
map.nodes[usize::from(at)].master = Some(m);
let shard = yo_alloc::allow(|| map.nodes[usize::from(m)].shard.clone());
yo_alloc::allow(|| map.nodes[usize::from(at)].shard = shard);
todo.save = true;
}
}
}
let speaking = if map.nodes[usize::from(at)].is_master() {
Some(at)
} else {
map.nodes[usize::from(at)].master
};
let claim = &p[O_SLOTS..O_SLOTS + BITMAP_LEN];
let dirty = speaking
.is_some_and(|m| (0..SLOTS).any(|slot| bit(claim, slot) != (map.owner[slot] == Some(m))));
if dirty && map.nodes[usize::from(at)].is_master() {
claim_slots(server, map, at, sender_epoch, claim, todo);
}
if dirty {
for slot in 0..SLOTS {
if !bit(claim, slot) {
continue;
}
let Some(owner) = map.owner[slot] else {
continue;
};
if owner == at {
continue;
}
if map.nodes[usize::from(owner)].epoch > sender_epoch {
todo.reply.push(update_packet(server, map, owner));
break;
}
}
}
if map.nodes[0].is_master()
&& map.nodes[usize::from(at)].is_master()
&& sender_epoch == map.nodes[0].epoch
&& map.nodes[0].id < map.nodes[usize::from(at)].id
{
let next = server.cluster.epoch.fetch_add(1, Relaxed) + 1;
map.nodes[0].epoch = next;
todo.save = true;
}
gossip(server, p, now, map, todo);
extensions(server, p, now, map, at, todo);
}
fn gossip(server: &Arc<Server>, p: &[u8], now: u64, map: &mut Map, todo: &mut Todo) {
let count = usize::from(be16(p, O_COUNT));
let sender_master = node_id(p, O_SENDER)
.and_then(|id| map.find(id.as_bytes()))
.is_none_or(|at| map.nodes[usize::from(at)].is_master());
for entry in 0..count {
let base = HDR_LEN + entry * GOSSIP_LEN;
let Some(id) = node_id(p, base + G_NAME) else {
continue;
};
let flags = be16(p, base + G_FLAGS);
let host = text(p, base + G_IP, IP_LEN);
let port = be16(p, base + G_PORT);
let bus = be16(p, base + G_CPORT);
if let Some(at) = map.find(id.as_bytes()) {
if at == 0 {
continue;
}
if sender_master {
let sender = node_id(p, O_SENDER).unwrap_or_default();
if flags & (FLAG_FAIL | FLAG_PFAIL) != 0 {
report(map, at, &sender, now);
if mark_failing(map, at, now) {
let about = yo_alloc::allow(|| map.nodes[usize::from(at)].id.clone());
todo.shout.push(fail_packet(server, map, &about));
todo.save = true;
}
} else {
unreport(map, at, &sender);
}
}
let node = &mut map.nodes[usize::from(at)];
if node.flags & (FLAG_FAIL | FLAG_PFAIL) == 0
&& node.ping_sent == 0
&& node.reports.is_empty()
{
let heard = u64::from(be32(p, base + G_PONG)) * 1000;
if heard > node.pong_recv && heard <= now + 500 {
node.pong_recv = heard;
}
} else if node.down()
&& flags & (FLAG_FAIL | FLAG_PFAIL | FLAG_NOADDR | FLAG_HANDSHAKE) == 0
&& !host.is_empty()
&& (node.host != host || node.port != port)
{
yo_alloc::allow(|| node.host = host);
node.port = port;
node.bus = bus;
node.linked = false;
node.flags &= !FLAG_NOADDR;
todo.save = true;
}
continue;
}
if flags & FLAG_NOADDR != 0 {
continue;
}
if server.cluster.bus.blacklisted(&id, now) {
continue;
}
let flags = flags & !(FLAG_MYSELF | FLAG_HANDSHAKE | FLAG_MEET);
let node = yo_alloc::allow(|| Node::new(id, host, port, bus, flags, now));
yo_alloc::allow(|| map.nodes.push(node));
todo.save = true;
}
}
fn extensions(server: &Arc<Server>, p: &[u8], now: u64, map: &mut Map, at: u16, todo: &mut Todo) {
if !p.get(O_MFLAGS).is_some_and(|f| f & MF_EXT_DATA != 0) {
return;
}
let mut left = be16(p, O_EXTENSIONS);
let mut base = HDR_LEN + usize::from(be16(p, O_COUNT)) * GOSSIP_LEN;
while left > 0 && base + 8 <= p.len() {
left -= 1;
let extlen = be32(p, base) as usize;
if extlen < 8 || base + extlen > p.len() {
return;
}
let kind = be16(p, base + 4);
let data = &p[base + 8..base + extlen];
match kind {
X_SHARDID => {
let id = String::from_utf8_lossy(&data[..NAME_LEN.min(data.len())]);
if id.len() == NAME_LEN && map.nodes[usize::from(at)].shard != id {
yo_alloc::allow(|| map.nodes[usize::from(at)].shard = id.into_owned());
todo.save = true;
}
}
X_SECRET => {
let theirs = String::from_utf8_lossy(&data[..NAME_LEN.min(data.len())]);
let mut mine = server.cluster.bus.secret.lock();
if theirs.len() == NAME_LEN && **mine > *theirs {
yo_alloc::allow(|| *mine = theirs.into_owned());
}
}
X_FORGOTTEN => {
if data.len() < NAME_LEN + 8 {
return;
}
let Some(id) = node_id(data, 0) else {
return;
};
let ttl = be64(data, NAME_LEN);
if map.nodes[0].id == id
|| map.nodes[0].master.map(usize::from)
== map.find(id.as_bytes()).map(usize::from)
{
base += extlen;
continue;
}
server.cluster.bus.blacklist(&id, now + ttl);
if let Some(gone) = map.find(id.as_bytes())
&& gone != 0
{
server.cluster.bus.cut(&id);
map.forget(gone);
todo.save = true;
}
}
X_HOSTNAME | X_HUMAN_NAME => {}
_ => {}
}
base += extlen;
}
}
fn report(map: &mut Map, at: u16, from: &str, now: u64) {
if from.is_empty() {
return;
}
let node = &mut map.nodes[usize::from(at)];
if let Some(held) = node.reports.iter_mut().find(|(who, _)| who == from) {
held.1 = now;
return;
}
yo_alloc::allow(|| node.reports.push((String::from(from), now)));
}
fn unreport(map: &mut Map, at: u16, from: &str) {
map.nodes[usize::from(at)]
.reports
.retain(|(who, _)| who != from);
}
fn mark_failing(map: &mut Map, at: u16, now: u64) -> bool {
let cutoff = now.saturating_sub(NODE_TIMEOUT_MS * 2);
map.nodes[usize::from(at)]
.reports
.retain(|(_, when)| *when > cutoff);
if map.nodes[usize::from(at)].flags & (FLAG_PFAIL | FLAG_FAIL) == 0 {
return false;
}
if map.nodes[usize::from(at)].flags & FLAG_FAIL != 0 {
return false;
}
let needed = map.voters() / 2 + 1;
let mut votes = map.nodes[usize::from(at)].reports.len();
if map.nodes[0].is_master() && !map.runs(0).is_empty() {
votes += 1;
}
if votes < needed {
return false;
}
let node = &mut map.nodes[usize::from(at)];
node.flags &= !FLAG_PFAIL;
node.flags |= FLAG_FAIL;
node.fail_time = now;
true
}
fn clear_failure(map: &mut Map, at: u16, now: u64) {
let node = &map.nodes[usize::from(at)];
let replica = !node.is_master();
let empty = map.runs(at).is_empty();
let stale = now.saturating_sub(node.fail_time) > NODE_TIMEOUT_MS * 2;
if replica || empty || stale {
let node = &mut map.nodes[usize::from(at)];
node.flags &= !FLAG_FAIL;
node.fail_time = 0;
}
}
fn claim_slots(
server: &Arc<Server>,
map: &mut Map,
owner: u16,
epoch: u64,
claim: &[u8],
todo: &mut Todo,
) {
let ours = if map.nodes[0].is_master() {
Some(0)
} else {
map.nodes[0].master
};
let mut lost = 0usize;
for slot in 0..SLOTS {
if !bit(claim, slot) {
continue;
}
if map.owner[slot] == Some(owner) || map.importing[slot].is_some() {
continue;
}
let held = map.owner[slot];
let newer = held.is_none_or(|at| map.nodes[usize::from(at)].epoch <= epoch);
if !newer {
continue;
}
if held == ours {
lost += 1;
}
if held == Some(0) {
yo_alloc::allow(|| todo.lost.push(slot as u16));
}
map.owner[slot] = Some(owner);
map.migrating[slot] = None;
todo.save = true;
}
if lost > 0 && ours.is_some_and(|at| map.runs(at).is_empty()) && owner != 0 {
todo.demoted = true;
manual_reset(server);
map.nodes[0].flags &= !FLAG_MASTER;
map.nodes[0].flags |= FLAG_SLAVE;
map.nodes[0].master = Some(owner);
let shard = yo_alloc::allow(|| map.nodes[usize::from(owner)].shard.clone());
yo_alloc::allow(|| map.nodes[0].shard = shard);
let (host, port) = {
let node = &map.nodes[usize::from(owner)];
(yo_alloc::allow(|| node.host.clone()), node.port)
};
let following = Arc::clone(server);
spawn("yo-bus-follow", move || {
following.follow_master(&host, port);
});
}
todo.recount = true;
}
fn vote_if_needed(server: &Arc<Server>, map: &mut Map, at: u16, p: &[u8], now: u64) -> bool {
if map.nodes[0].flags & FLAG_SLAVE != 0 || map.runs(0).is_empty() {
return false;
}
let epoch = server.cluster.epoch.load(Relaxed);
if be64(p, O_CURRENT_EPOCH) < epoch {
return false;
}
if server.cluster.vote.given.load(Relaxed) == epoch {
return false;
}
let asking = &map.nodes[usize::from(at)];
let Some(master) = asking.master.filter(|_| asking.flags & FLAG_SLAVE != 0) else {
return false;
};
let forced = p.get(O_MFLAGS).is_some_and(|f| f & MF_FORCEACK != 0);
if map.nodes[usize::from(master)].flags & FLAG_FAIL == 0 && !forced {
return false;
}
if now.saturating_sub(map.nodes[usize::from(master)].voted_time) < NODE_TIMEOUT_MS * 2 {
return false;
}
let claimed = be64(p, O_CONFIG_EPOCH);
for slot in 0..SLOTS {
if !bit(&p[O_SLOTS..O_SLOTS + BITMAP_LEN], slot) {
continue;
}
if map.owner[slot].is_some_and(|held| map.nodes[usize::from(held)].epoch > claimed) {
return false;
}
}
server.cluster.vote.given.store(epoch, Relaxed);
map.nodes[usize::from(master)].voted_time = now;
true
}
fn rank_of(map: &Map, master: u16, mine: u64) -> u64 {
map.nodes
.iter()
.skip(1)
.filter(|node| {
node.master == Some(master)
&& node.flags & FLAG_SLAVE != 0
&& node.flags & FLAG_NOFAILOVER == 0
&& node.offset > mine
})
.count() as u64
}
enum Step {
Idle,
Announce,
Ask(Vec<u8>),
Won,
}
fn decide(server: &Arc<Server>, map: &mut Map, now: u64) -> Step {
let vote = &server.cluster.vote;
let manual = &server.cluster.manual;
let asked = manual.end.load(Relaxed) != 0;
let ready = asked && manual.can_start.load(Relaxed);
let me = &map.nodes[0];
if me.flags & FLAG_SLAVE == 0 || (me.flags & FLAG_NOFAILOVER != 0 && !ready) {
return Step::Idle;
}
let Some(master) = me.master else {
return Step::Idle;
};
if map.nodes[usize::from(master)].flags & FLAG_FAIL == 0 && !ready {
return Step::Idle;
}
if map.runs(master).is_empty() {
return Step::Idle;
}
if !ready && server.master_silence(now).saturating_sub(NODE_TIMEOUT_MS) > STALE_DATA_MS {
return Step::Idle;
}
let since = now as i64 - vote.at.load(Relaxed) as i64;
if since > (ELECTION_TIMEOUT_MS * 2) as i64 {
let rank = if asked {
0
} else {
rank_of(map, master, server.repl_offset())
};
let at = if asked {
now
} else {
now + ELECTION_DELAY_MS
+ u64::from(pick(ELECTION_DELAY_MS as usize))
+ rank * RANK_DELAY_MS
};
vote.at.store(at, Relaxed);
vote.rank.store(rank, Relaxed);
vote.count.store(0, Relaxed);
vote.sent.store(false, Relaxed);
return Step::Announce;
}
if !vote.sent.load(Relaxed) {
if !asked {
let rank = rank_of(map, master, server.repl_offset());
let was = vote.rank.load(Relaxed);
if rank > was {
vote.at.fetch_add((rank - was) * RANK_DELAY_MS, Relaxed);
vote.rank.store(rank, Relaxed);
}
}
if now < vote.at.load(Relaxed) {
return Step::Idle;
}
}
if since > ELECTION_TIMEOUT_MS as i64 {
return Step::Idle;
}
if !vote.sent.load(Relaxed) {
let epoch = server.cluster.epoch.fetch_add(1, Relaxed) + 1;
vote.epoch.store(epoch, Relaxed);
vote.sent.store(true, Relaxed);
let mut packet = header(server, map, T_AUTH_REQUEST);
if asked {
packet[O_MFLAGS] |= MF_FORCEACK;
}
seal(&mut packet);
return Step::Ask(packet);
}
if vote.count.load(Relaxed) < (map.size() / 2 + 1) as u64 {
return Step::Idle;
}
let epoch = vote.epoch.load(Relaxed);
replace_master(map, master, epoch);
Step::Won
}
fn replace_master(map: &mut Map, master: u16, epoch: u64) {
if map.nodes[0].epoch < epoch {
map.nodes[0].epoch = epoch;
}
map.nodes[0].flags &= !FLAG_SLAVE;
map.nodes[0].flags |= FLAG_MASTER;
map.nodes[0].master = None;
for slot in 0..SLOTS {
if map.owner[slot] == Some(master) {
map.owner[slot] = Some(0);
}
}
}
fn won(server: &Arc<Server>) {
server.stop_following();
manual_reset(server);
server.recount_coverage();
server.cluster.bus.dirty.store(true, Relaxed);
let _ = super::save(server);
server.cluster_broadcast_pong();
}
fn failover(server: &Arc<Server>, now: u64) {
let step = {
let mut map = server.cluster.map.lock();
decide(server, &mut map, now)
};
match step {
Step::Idle => {}
Step::Announce => server.cluster_broadcast_pong(),
Step::Ask(packet) => {
for link in server.cluster.bus.all() {
link.send(&packet);
}
server.cluster.bus.dirty.store(true, Relaxed);
}
Step::Won => won(server),
}
}
fn manual_reset(server: &Server) {
let manual = &server.cluster.manual;
let held = manual.held.swap(0, Relaxed);
if held != 0 {
server.lift(held, false);
}
manual.end.store(0, Relaxed);
manual.can_start.store(false, Relaxed);
manual.offset.store(OFFSET_UNKNOWN, Relaxed);
}
fn stand_down(server: &Arc<Server>, map: &Map, at: u16, now: u64, todo: &mut Todo) {
let asking = &map.nodes[usize::from(at)];
if asking.master != Some(0) || asking.flags & FLAG_SLAVE == 0 || !map.nodes[0].is_master() {
return;
}
server.cluster.asm.cancel(None, now as i64);
manual_reset(server);
let manual = &server.cluster.manual;
manual.end.store(now + MANUAL_TIMEOUT_MS, Relaxed);
let until = now + MANUAL_TIMEOUT_MS * MANUAL_PAUSE_MULT;
manual.held.store(until, Relaxed);
server.pause(until, false);
yo_alloc::allow(|| todo.reply.push(ping(server, map, T_PING, Some(at))));
}
fn manual_cron(server: &Arc<Server>, now: u64) {
let manual = &server.cluster.manual;
let end = manual.end.load(Relaxed);
if end == 0 {
return;
}
if end < now {
manual_reset(server);
return;
}
if manual.can_start.load(Relaxed) {
return;
}
let offset = manual.offset.load(Relaxed);
if offset != OFFSET_UNKNOWN && offset == server.repl_offset() {
manual.can_start.store(true, Relaxed);
}
}
pub(super) fn manual_failover(
server: &Arc<Server>,
force: bool,
takeover: bool,
) -> Result<(), Error> {
let now = server.now_ms();
let mut ask: Option<(String, Vec<u8>)> = None;
let mut promoted = false;
{
let mut map = server.cluster.map.lock();
if map.nodes[0].is_master() {
return Err(Error::new(
Code::Invalid,
"You should send CLUSTER FAILOVER to a replica",
));
}
let Some(master) = map.nodes[0].master else {
return Err(Error::new(
Code::Invalid,
"I'm a replica but my master is unknown to me",
));
};
let node = &map.nodes[usize::from(master)];
if !force && (node.flags & FLAG_FAIL != 0 || !node.linked) {
return Err(Error::new(
Code::Invalid,
"Master is down or failed, please use CLUSTER FAILOVER FORCE",
));
}
manual_reset(server);
server
.cluster
.manual
.end
.store(now + MANUAL_TIMEOUT_MS, Relaxed);
if takeover {
super::bump_without_consensus(server, &mut map);
let epoch = map.nodes[0].epoch;
replace_master(&mut map, master, epoch);
promoted = true;
} else if force {
server.cluster.manual.can_start.store(true, Relaxed);
} else {
let mut packet = header(server, &map, T_MFSTART);
seal(&mut packet);
let id = yo_alloc::allow(|| map.nodes[usize::from(master)].id.clone());
ask = Some((id, packet));
}
}
if promoted {
won(server);
}
if let Some((id, packet)) = ask
&& let Some(link) = server.cluster.bus.outbound(&id)
{
link.send(&packet);
}
Ok(())
}
fn cron(server: &Arc<Server>) {
let mut tick = 0u64;
loop {
std::thread::sleep(Duration::from_millis(CRON_MS));
tick += 1;
let now = server.now_ms();
let mut dial_list: Vec<(String, String, u16, bool)> = Vec::new();
let mut ping_list: Vec<(String, Vec<u8>)> = Vec::new();
let mut shout: Vec<Vec<u8>> = Vec::new();
let mut follow: Option<(String, u16)> = None;
let mut save = false;
let paused = server.cluster.manual.held.load(Relaxed) != 0;
{
let mut map = server.cluster.map.lock();
let count = map.nodes.len();
for at in 1..count as u16 {
let node = &map.nodes[usize::from(at)];
if node.flags & FLAG_HANDSHAKE != 0
&& now.saturating_sub(node.data_recv) > NODE_TIMEOUT_MS.max(1000)
{
map.forget(at);
save = true;
break;
}
if node.flags & FLAG_NOADDR != 0 || node.host.is_empty() {
continue;
}
if !node.linked {
let meet = node.flags & FLAG_MEET != 0;
yo_alloc::allow(|| {
dial_list.push((node.id.clone(), node.host.clone(), node.bus, meet));
});
continue;
}
let quiet = now.saturating_sub(node.pong_recv);
let due = node.ping_sent == 0 && quiet > PING_MS;
let waiting = paused && node.master == Some(0);
if due
|| waiting
|| (tick.is_multiple_of(10) && oldest_of_five(&map, now) == Some(at))
{
let packet = ping(server, &map, T_PING, Some(at));
let id = yo_alloc::allow(|| map.nodes[usize::from(at)].id.clone());
map.nodes[usize::from(at)].ping_sent = now;
ping_list.push((id, packet));
}
}
for at in 1..map.nodes.len() as u16 {
let node = &map.nodes[usize::from(at)];
if node.flags & (FLAG_HANDSHAKE | FLAG_FAIL | FLAG_PFAIL) != 0 {
continue;
}
let waiting = if node.ping_sent == 0 {
0
} else {
now.saturating_sub(node.ping_sent)
};
let quiet = now.saturating_sub(node.data_recv);
if waiting.min(quiet) > NODE_TIMEOUT_MS {
map.nodes[usize::from(at)].flags |= FLAG_PFAIL;
save = true;
}
if mark_failing(&mut map, at, now) {
let about = yo_alloc::allow(|| map.nodes[usize::from(at)].id.clone());
shout.push(fail_packet(server, &map, &about));
save = true;
}
}
if map.nodes[0].flags & FLAG_SLAVE != 0
&& !server.following()
&& let Some(master) = map.nodes[0].master
&& let Some(node) = map.nodes.get(usize::from(master))
&& node.flags & FLAG_NOADDR == 0
&& !node.host.is_empty()
{
follow = yo_alloc::allow(|| Some((node.host.clone(), node.port)));
}
}
if let Some((host, port)) = follow {
server.follow_master(&host, port);
}
for (id, host, bus, meet) in dial_list {
if !dial(server, &id, &host, bus, meet) {
continue;
}
let mut map = server.cluster.map.lock();
if let Some(at) = map.find(id.as_bytes()) {
map.nodes[usize::from(at)].linked = true;
}
}
for (id, packet) in ping_list {
if let Some(link) = server.cluster.bus.outbound(&id) {
link.send(&packet);
} else {
let mut map = server.cluster.map.lock();
if let Some(at) = map.find(id.as_bytes()) {
map.nodes[usize::from(at)].linked = false;
}
}
}
if !shout.is_empty() {
for link in server.cluster.bus.all() {
for packet in &shout {
link.send(packet);
}
}
}
if save {
server.cluster.bus.dirty.store(true, Relaxed);
}
server.recount_coverage();
manual_cron(server, now);
failover(server, now);
server.asm_cron();
server.asm_relax();
if tick.is_multiple_of(10) && server.cluster.bus.dirty.swap(false, Relaxed) {
let _ = super::save(server);
}
}
}
fn oldest_of_five(map: &Map, now: u64) -> Option<u16> {
if map.nodes.len() < 2 {
return None;
}
let mut best: Option<(u16, u64)> = None;
for _ in 0..5 {
let at = pick(map.nodes.len());
if at == 0 {
continue;
}
let node = &map.nodes[usize::from(at)];
if node.ping_sent != 0 || node.flags & (FLAG_HANDSHAKE | FLAG_NOADDR) != 0 {
continue;
}
let quiet = now.saturating_sub(node.pong_recv);
if best.is_none_or(|(_, held)| quiet > held) {
best = Some((at, quiet));
}
}
best.map(|(at, _)| at)
}
impl Server {
pub(super) fn cluster_meet(&self, host: &str, port: u16, bus: u16) {
let now = self.now_ms();
let mut map = self.cluster.map.lock();
let known = map
.nodes
.iter()
.any(|node| node.host == host && node.port == port);
if known {
return;
}
let id = yo_alloc::allow(|| String::from_utf8_lossy(&new_id()).into_owned());
let node = yo_alloc::allow(|| {
Node::new(
id,
yo_alloc::allow(|| String::from(host)),
port,
bus,
FLAG_HANDSHAKE | FLAG_MEET,
now,
)
});
yo_alloc::allow(|| map.nodes.push(node));
}
pub(super) fn cluster_blacklisted(&self, id: &str) -> bool {
self.cluster.bus.blacklisted(id, self.now_ms())
}
pub(super) fn cluster_forget(&self, at: u16) {
let now = self.now_ms();
let id = {
let mut map = self.cluster.map.lock();
let id = yo_alloc::allow(|| map.nodes[usize::from(at)].id.clone());
map.forget(at);
id
};
self.cluster.bus.blacklist(&id, now);
self.cluster.bus.cut(&id);
let packet = {
let map = self.cluster.map.lock();
let mut p = ping(self, &map, T_PING, None);
let mut body = [0u8; NAME_LEN + 8];
body[..NAME_LEN].copy_from_slice(id.as_bytes());
body[NAME_LEN..].copy_from_slice(&BLACKLIST_MS.to_be_bytes());
push_ext(&mut p, X_FORGOTTEN, &body);
let exts = be16(&p, O_EXTENSIONS) + 1;
put16(&mut p, O_EXTENSIONS, exts);
seal(&mut p);
p
};
for link in self.cluster.bus.all() {
link.send(&packet);
}
self.cluster.bus.dirty.store(true, Relaxed);
}
pub(super) fn cluster_replicate(self: &Arc<Server>, at: u16) {
manual_reset(self);
let (host, port) = {
let mut map = self.cluster.map.lock();
for slot in 0..SLOTS {
if map.owner[slot] == Some(0) {
map.owner[slot] = None;
}
}
map.nodes[0].flags &= !(FLAG_MASTER | FLAG_MIGRATE_TO);
map.nodes[0].flags |= FLAG_SLAVE;
map.nodes[0].master = Some(at);
let shard = yo_alloc::allow(|| map.nodes[usize::from(at)].shard.clone());
yo_alloc::allow(|| map.nodes[0].shard = shard);
let node = &map.nodes[usize::from(at)];
(yo_alloc::allow(|| node.host.clone()), node.port)
};
self.recount_coverage();
self.cluster.bus.dirty.store(true, Relaxed);
if !host.is_empty() {
self.follow_master(&host, port);
}
}
pub(crate) fn cluster_publish(&self, shard: bool, channel: &[u8], body: &[u8]) {
if !self.cluster_enabled() || !self.cluster.bus.on.load(Relaxed) {
return;
}
let packet = {
let map = self.cluster.map.lock();
publish_packet(self, &map, shard, channel, body)
};
for link in self.cluster.bus.all() {
link.send(&packet);
}
}
pub(super) fn cluster_broadcast_pong(&self) {
if !self.cluster_enabled() || !self.cluster.bus.on.load(Relaxed) {
return;
}
let packet = {
let map = self.cluster.map.lock();
ping(self, &map, T_PONG, None)
};
for link in self.cluster.bus.all() {
link.send(&packet);
}
}
pub(super) fn cluster_links(&self, out: &mut Out) {
let links = self.cluster.bus.all();
let at = out.len();
let mut n = 0;
for link in links {
let node = link.named();
if node.is_empty() {
continue;
}
out.map(6);
out.bulk(b"direction");
out.bulk(if link.inbound {
b"from".as_slice()
} else {
b"to".as_slice()
});
out.bulk(b"node");
out.bulk(node.as_bytes());
out.bulk(b"create-time");
out.int(link.created as i64);
out.bulk(b"events");
out.bulk(b"r");
out.bulk(b"send-buffer-allocated");
out.int(link.sent.load(Relaxed) as i64);
out.bulk(b"send-buffer-used");
out.int(0);
n += 1;
}
out.close_array(at, n);
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_bitmap_round_trips_through_the_reference_bit_order() {
let mut bitmap = [0u8; BITMAP_LEN];
for slot in [0usize, 1, 7, 8, 1234, 5061, 12182, SLOTS - 1] {
set_bit(&mut bitmap, slot);
}
for slot in 0..SLOTS {
let want = matches!(slot, 0 | 1 | 7 | 8 | 1234 | 5061 | 12182) || slot == SLOTS - 1;
assert_eq!(bit(&bitmap, slot), want, "slot {slot}");
}
assert_eq!(bitmap[0], 0b1000_0011);
assert_eq!(bitmap[1], 0b0000_0001);
}
#[test]
fn an_extension_is_padded_to_the_eight_byte_boundary() {
let mut p = Vec::new();
push_ext(&mut p, X_SHARDID, &[b'a'; 40]);
assert_eq!(p.len(), 48);
assert_eq!(be32(&p, 0), 48);
assert_eq!(be16(&p, 4), X_SHARDID);
let mut q = Vec::new();
push_ext(&mut q, X_HOSTNAME, b"node1.example\0");
assert_eq!(q.len(), 8 + 16);
assert_eq!(be32(&q, 0), 24);
}
#[test]
fn a_packet_of_the_wrong_length_is_refused() {
let mut p = vec![0u8; HDR_LEN];
p[0..4].copy_from_slice(SIG);
put16(&mut p, O_VER, PROTO_VER);
put16(&mut p, O_TYPE, T_PING);
put16(&mut p, O_COUNT, 0);
p[O_MFLAGS] = 0;
seal(&mut p);
assert_eq!(expected(&p, T_PING), Some(HDR_LEN));
put16(&mut p, O_COUNT, 1);
assert_eq!(expected(&p, T_PING), Some(HDR_LEN + GOSSIP_LEN));
assert_ne!(expected(&p, T_PING), Some(p.len()));
}
const T0: u64 = 1_700_000_000_000;
fn shard() -> Arc<Server> {
let mut server = Server::new();
server.enable_cluster("", 7000);
let server = Arc::new(server);
let master = server.cluster_pretend_node("1".repeat(40).as_str(), "10.0.0.1", 7001);
let other = server.cluster_pretend_node("2".repeat(40).as_str(), "10.0.0.2", 7002);
let third = server.cluster_pretend_node("3".repeat(40).as_str(), "10.0.0.3", 7003);
server.cluster_pretend_follower(master);
{
let mut map = server.cluster.map.lock();
for slot in 0..SLOTS {
map.owner[slot] = Some(match slot {
0..=5460 => master,
5461..=10922 => other,
_ => third,
});
}
}
server.recount_coverage();
server
}
fn asking(epoch: u64, config: u64, forced: bool) -> Vec<u8> {
let mut p = vec![0u8; HDR_LEN];
put64(&mut p, O_CURRENT_EPOCH, epoch);
put64(&mut p, O_CONFIG_EPOCH, config);
for slot in 0..=5460 {
set_bit(&mut p[O_SLOTS..O_SLOTS + BITMAP_LEN], slot);
}
p[O_MFLAGS] = if forced { MF_FORCEACK } else { 0 };
p
}
#[test]
fn a_node_that_cannot_be_reached_is_treated_as_pinged() {
let server = shard();
let id = "1".repeat(40);
assert_eq!(server.cluster.map.lock().nodes[1].ping_sent, 0);
unreachable(&server, &id);
let first = server.cluster.map.lock().nodes[1].ping_sent;
assert!(first > 0, "the clock is running");
unreachable(&server, &id);
assert_eq!(server.cluster.map.lock().nodes[1].ping_sent, first);
unreachable(&server, &"9".repeat(40));
}
#[test]
fn a_vote_is_given_once_and_only_when_every_condition_holds() {
let server = shard();
{
let mut map = server.cluster.map.lock();
map.nodes[0].flags &= !FLAG_SLAVE;
map.nodes[0].flags |= FLAG_MASTER;
map.nodes[0].master = None;
for slot in 10923..SLOTS {
map.owner[slot] = Some(0);
}
let mut replica = Node::new(
"4".repeat(40),
"10.0.0.4".to_owned(),
7004,
7004 + 10000,
FLAG_SLAVE,
0,
);
replica.master = Some(1);
map.nodes.push(replica);
}
let at = 4u16;
let ask = |server: &Arc<Server>, p: &[u8], now: u64| {
let mut map = server.cluster.map.lock();
vote_if_needed(server, &mut map, at, p, now)
};
server.cluster.epoch.store(1, Relaxed);
assert!(!ask(&server, &asking(1, 0, false), T0));
assert!(ask(&server, &asking(1, 0, true), T0));
assert!(!ask(&server, &asking(1, 0, true), T0));
server.cluster.epoch.store(2, Relaxed);
{
let mut map = server.cluster.map.lock();
map.nodes[1].flags |= FLAG_FAIL;
}
assert!(!ask(&server, &asking(2, 0, false), T0 + NODE_TIMEOUT_MS));
let later = T0 + NODE_TIMEOUT_MS * 2 + 1;
assert!(!ask(&server, &asking(1, 0, false), later));
{
let mut map = server.cluster.map.lock();
map.nodes[4].flags &= !FLAG_SLAVE;
}
assert!(!ask(&server, &asking(2, 0, false), later));
{
let mut map = server.cluster.map.lock();
map.nodes[4].flags |= FLAG_SLAVE;
}
{
let mut map = server.cluster.map.lock();
map.owner[3000] = Some(2);
map.nodes[2].epoch = 7;
}
assert!(!ask(&server, &asking(2, 6, false), later));
assert!(ask(&server, &asking(2, 7, false), later));
assert_eq!(server.cluster.vote.given(), 2);
server.cluster.epoch.store(3, Relaxed);
{
let mut map = server.cluster.map.lock();
for slot in 0..SLOTS {
if map.owner[slot] == Some(0) {
map.owner[slot] = Some(3);
}
}
}
assert!(!ask(
&server,
&asking(3, 7, false),
later + NODE_TIMEOUT_MS * 3
));
}
#[test]
fn an_election_waits_then_asks_then_wins() {
let server = shard();
let step = |server: &Arc<Server>, now: u64| {
let mut map = server.cluster.map.lock();
decide(server, &mut map, now)
};
assert!(matches!(step(&server, T0), Step::Idle));
{
let mut map = server.cluster.map.lock();
map.nodes[1].flags |= FLAG_FAIL;
}
assert!(matches!(step(&server, T0), Step::Announce));
let at = server.cluster.vote.at.load(Relaxed);
assert!(
(T0 + 500..=T0 + 1000).contains(&at),
"half a second plus up to half a second more, got {at}"
);
assert!(matches!(step(&server, at - 1), Step::Idle));
{
let mut map = server.cluster.map.lock();
let mut ahead = Node::new(
"5".repeat(40),
"10.0.0.5".to_owned(),
7005,
7005 + 10000,
FLAG_SLAVE,
0,
);
ahead.master = Some(1);
ahead.offset = 900;
map.nodes.push(ahead);
}
assert!(matches!(step(&server, at), Step::Idle));
assert_eq!(server.cluster.vote.at.load(Relaxed), at + RANK_DELAY_MS);
{
let mut map = server.cluster.map.lock();
map.nodes[4].offset = 0;
}
assert!(matches!(step(&server, at), Step::Idle));
assert_eq!(server.cluster.vote.at.load(Relaxed), at + RANK_DELAY_MS);
let now = at + RANK_DELAY_MS;
let Step::Ask(packet) = step(&server, now) else {
panic!("the turn has come");
};
assert_eq!(be64(&packet, O_CURRENT_EPOCH), 1);
assert_eq!(server.cluster.vote.epoch.load(Relaxed), 1);
assert!(bit(&packet[O_SLOTS..O_SLOTS + BITMAP_LEN], 5460));
assert!(!bit(&packet[O_SLOTS..O_SLOTS + BITMAP_LEN], 5461));
assert!(matches!(step(&server, now), Step::Idle), "asked already");
server.cluster.vote.count.store(1, Relaxed);
assert!(matches!(step(&server, now), Step::Idle));
server.cluster.vote.count.store(2, Relaxed);
assert!(matches!(step(&server, now), Step::Won));
let map = server.cluster.map.lock();
assert!(map.nodes[0].is_master());
assert_eq!(map.nodes[0].master, None);
assert_eq!(map.nodes[0].epoch, 1, "the epoch it stood under");
assert_eq!(map.owner[0], Some(0));
assert_eq!(map.owner[5460], Some(0));
assert_eq!(map.owner[5461], Some(2), "somebody else's slots are theirs");
}
#[test]
fn a_replica_that_should_not_stand_does_not() {
let refused = |now: u64, set: fn(&Arc<Server>)| {
let server = shard();
{
let mut map = server.cluster.map.lock();
map.nodes[1].flags |= FLAG_FAIL;
}
set(&server);
let mut map = server.cluster.map.lock();
matches!(decide(&server, &mut map, now), Step::Idle)
};
assert!(refused(T0, |server| {
let mut map = server.cluster.map.lock();
map.nodes[0].flags &= !FLAG_SLAVE;
map.nodes[0].flags |= FLAG_MASTER;
}));
assert!(refused(T0, |server| {
let mut map = server.cluster.map.lock();
map.nodes[0].flags |= FLAG_NOFAILOVER;
}));
assert!(refused(T0, |server| {
let mut map = server.cluster.map.lock();
for slot in 0..=5460 {
map.owner[slot] = Some(2);
}
}));
let stale = T0 + NODE_TIMEOUT_MS + STALE_DATA_MS + 1;
assert!(refused(stale, |server| {
server.pretend_following("10.0.0.1", 7001, false);
server.pretend_master_down_at(T0);
}));
assert!(!refused(stale - 1000, |server| {
server.pretend_following("10.0.0.1", 7001, false);
server.pretend_master_down_at(T0);
}));
}
fn standing() -> Arc<Server> {
let server = shard();
{
let mut map = server.cluster.map.lock();
map.nodes[0].flags &= !FLAG_SLAVE;
map.nodes[0].flags |= FLAG_MASTER;
map.nodes[0].master = None;
map.nodes[1].flags &= !FLAG_MASTER;
map.nodes[1].flags |= FLAG_SLAVE;
map.nodes[1].master = Some(0);
for slot in 0..=5460 {
map.owner[slot] = Some(0);
}
}
server.recount_coverage();
server
}
#[test]
fn a_master_asked_to_stand_down_stops_writing_and_says_where() {
let server = standing();
let mut todo = Todo::default();
{
let map = server.cluster.map.lock();
stand_down(&server, &map, 2, T0, &mut todo);
}
assert!(todo.reply.is_empty(), "only a replica of this node may ask");
assert_eq!(server.pause_ends(), 0);
{
let map = server.cluster.map.lock();
stand_down(&server, &map, 1, T0, &mut todo);
}
let reply = todo
.reply
.first()
.expect("answered at once, not by the cron");
assert_eq!(be16(reply, O_TYPE), T_PING);
assert!(
reply[O_MFLAGS] & MF_PAUSED != 0,
"and the answer says the writes have stopped"
);
assert_eq!(be64(reply, O_OFFSET), server.repl_offset());
assert_eq!(
server.pause_ends(),
T0 + MANUAL_TIMEOUT_MS * MANUAL_PAUSE_MULT
);
assert_eq!(
server.paused(T0),
Some(false),
"the writes and not the reads"
);
manual_cron(&server, T0 + MANUAL_TIMEOUT_MS + 1);
assert_eq!(server.cluster.manual.end.load(Relaxed), 0);
assert_eq!(server.pause_ends(), 0);
}
#[test]
fn a_manual_failover_waits_for_the_offsets_to_meet() {
let server = shard();
let step = |now: u64| {
let mut map = server.cluster.map.lock();
decide(&server, &mut map, now)
};
let manual = &server.cluster.manual;
manual.end.store(T0 + MANUAL_TIMEOUT_MS, Relaxed);
assert!(matches!(step(T0), Step::Idle));
manual_cron(&server, T0);
assert!(!manual.can_start.load(Relaxed), "no offset from the master");
manual.offset.store(server.repl_offset() + 1, Relaxed);
manual_cron(&server, T0);
assert!(!manual.can_start.load(Relaxed), "still behind the master");
assert!(matches!(step(T0), Step::Idle));
manual.offset.store(server.repl_offset(), Relaxed);
manual_cron(&server, T0);
assert!(manual.can_start.load(Relaxed));
assert!(matches!(step(T0), Step::Announce));
assert_eq!(server.cluster.vote.at.load(Relaxed), T0);
assert_eq!(server.cluster.vote.rank.load(Relaxed), 0);
let Step::Ask(packet) = step(T0) else {
panic!("nothing to wait for on this path");
};
assert!(
packet[O_MFLAGS] & MF_FORCEACK != 0,
"the master is up and is in on it, so say so"
);
server.cluster.vote.count.store(2, Relaxed);
assert!(matches!(step(T0), Step::Won));
let map = server.cluster.map.lock();
assert!(map.nodes[0].is_master());
assert_eq!(map.owner[5460], Some(0));
assert_eq!(
map.owner[5461],
Some(2),
"and nothing that was not its master's"
);
}
#[test]
fn cluster_failover_refuses_where_the_reference_refuses() {
let ends = |e: Error, want: &str| {
let text = e.to_string();
assert!(text.ends_with(want), "wanted {want}, got {text}");
};
let server = standing();
ends(
manual_failover(&server, false, false).unwrap_err(),
"You should send CLUSTER FAILOVER to a replica",
);
{
let mut map = server.cluster.map.lock();
map.nodes[0].flags &= !FLAG_MASTER;
map.nodes[0].flags |= FLAG_SLAVE;
map.nodes[0].master = None;
}
ends(
manual_failover(&server, false, false).unwrap_err(),
"I'm a replica but my master is unknown to me",
);
let server = shard();
server.cluster.map.lock().nodes[1].linked = false;
ends(
manual_failover(&server, false, false).unwrap_err(),
"Master is down or failed, please use CLUSTER FAILOVER FORCE",
);
server.cluster.map.lock().nodes[1].linked = true;
server.cluster.map.lock().nodes[1].flags |= FLAG_FAIL;
ends(
manual_failover(&server, false, false).unwrap_err(),
"Master is down or failed, please use CLUSTER FAILOVER FORCE",
);
server.cluster.map.lock().nodes[1].flags &= !FLAG_FAIL;
manual_failover(&server, false, false).unwrap();
assert_ne!(server.cluster.manual.end.load(Relaxed), 0);
assert!(
!server.cluster.manual.can_start.load(Relaxed),
"the plain form waits for the master to answer"
);
assert!(server.cluster.map.lock().nodes[0].flags & FLAG_SLAVE != 0);
let server = shard();
manual_failover(&server, true, false).unwrap();
assert!(server.cluster.manual.can_start.load(Relaxed));
assert!(server.cluster.map.lock().nodes[0].flags & FLAG_SLAVE != 0);
let server = shard();
manual_failover(&server, true, true).unwrap();
assert_eq!(
server.cluster.manual.end.load(Relaxed),
0,
"nothing left to wait for"
);
let map = server.cluster.map.lock();
assert!(map.nodes[0].is_master());
assert_eq!(map.nodes[0].epoch, 1, "under an epoch it made up");
assert_eq!(map.owner[0], Some(0));
assert_eq!(map.owner[5460], Some(0));
assert_eq!(map.owner[5461], Some(2));
}
#[test]
fn a_node_id_has_to_be_forty_hex_characters() {
let mut p = vec![0u8; NAME_LEN * 3];
assert_eq!(node_id(&p, 0), None, "all zeros is the reference's no node");
p[0..NAME_LEN].copy_from_slice(&[b'a'; NAME_LEN]);
assert_eq!(node_id(&p, 0).as_deref(), Some("a".repeat(40).as_str()));
p[5] = b'z';
assert_eq!(node_id(&p, 0), None);
assert_eq!(node_id(&p, NAME_LEN * 3), None, "off the end is not an id");
}
}