pub use crate::kbucket::K_VALUE;
use super::*;
#[derive(Debug, Clone)]
pub struct PendingNode<TVal> {
node: Node<IdBytes, TVal>,
status: NodeStatus,
replace: Instant,
}
#[derive(PartialEq, Eq, Debug, Copy, Clone)]
pub enum NodeStatus {
Connected,
Disconnected,
}
impl<TVal> PendingNode<TVal> {
pub fn status(&self) -> NodeStatus {
self.status
}
pub fn value_mut(&mut self) -> &mut TVal {
&mut self.node.value
}
pub fn is_ready(&self) -> bool {
Instant::now() >= self.replace
}
pub fn into_node(self) -> Node<IdBytes, TVal> {
self.node
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Node<TKey, TVal> {
pub key: TKey,
pub value: TVal,
}
#[derive(Copy, Clone, Debug, PartialEq, Eq, PartialOrd, Ord)]
pub struct Position(usize);
#[derive(Debug, Clone)]
pub struct KBucket<TVal> {
nodes: ArrayVec<[Node<IdBytes, TVal>; K_VALUE.get()]>,
first_connected_pos: Option<usize>,
pending: Option<PendingNode<TVal>>,
pending_timeout: Duration,
}
#[must_use]
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum InsertResult {
Inserted,
Pending {
disconnected: IdBytes,
},
Full,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AppliedPending<TVal> {
pub inserted: Node<IdBytes, TVal>,
pub evicted: Option<Node<IdBytes, TVal>>,
}
impl<TVal> KBucket<TVal>
where
TVal: Clone,
{
pub fn new(pending_timeout: Duration) -> Self {
KBucket {
nodes: ArrayVec::new(),
first_connected_pos: None,
pending: None,
pending_timeout,
}
}
pub fn pending(&self) -> Option<&PendingNode<TVal>> {
self.pending.as_ref()
}
pub fn pending_mut(&mut self) -> Option<&mut PendingNode<TVal>> {
self.pending.as_mut()
}
pub fn as_pending(&self, key: &IdBytes) -> Option<&PendingNode<TVal>> {
self.pending()
.filter(|p| p.node.key.as_ref() == key.as_ref())
}
#[expect(unused)] pub fn get(&self, key: &IdBytes) -> Option<&Node<IdBytes, TVal>> {
self.position(key).map(|p| &self.nodes[p.0])
}
pub fn iter(&self) -> impl Iterator<Item = (&Node<IdBytes, TVal>, NodeStatus)> {
self.nodes
.iter()
.enumerate()
.map(move |(p, n)| (n, self.status(Position(p))))
}
pub fn apply_pending(&mut self) -> Option<AppliedPending<TVal>> {
if let Some(pending) = self.pending.take() {
if pending.replace <= Instant::now() {
if self.nodes.is_full() {
if self.status(Position(0)) == NodeStatus::Connected {
return None;
}
debug_assert!(self.first_connected_pos.is_none_or(|p| p > 0)); let inserted = pending.node.clone();
if pending.status == NodeStatus::Connected {
let evicted = Some(self.nodes.remove(0));
self.first_connected_pos = self
.first_connected_pos
.map_or_else(|| Some(self.nodes.len()), |p| p.checked_sub(1));
self.nodes.push(pending.node);
return Some(AppliedPending { inserted, evicted });
}
else if let Some(p) = self.first_connected_pos {
let insert_pos = p.checked_sub(1).expect("by (*)");
let evicted = Some(self.nodes.remove(0));
self.nodes.insert(insert_pos, pending.node);
return Some(AppliedPending { inserted, evicted });
} else {
let evicted = Some(self.nodes.remove(0));
self.nodes.push(pending.node);
return Some(AppliedPending { inserted, evicted });
}
} else {
let inserted = pending.node.clone();
match self.insert(pending.node, pending.status) {
InsertResult::Inserted => {
return Some(AppliedPending {
inserted,
evicted: None,
});
}
_ => unreachable!("Bucket is not full."),
}
}
} else {
self.pending = Some(pending);
}
}
None
}
pub fn update_pending(&mut self, status: NodeStatus) {
if let Some(pending) = &mut self.pending {
pending.status = status
}
}
pub fn remove_pending(&mut self) -> Option<PendingNode<TVal>> {
self.pending.take()
}
pub fn update(&mut self, key: &IdBytes, status: NodeStatus) {
if let Some((node, _status, pos)) = self.remove(key) {
if pos == Position(0) && status == NodeStatus::Connected {
self.pending = None
}
match self.insert(node, status) {
InsertResult::Inserted => {}
_ => unreachable!("The node is removed before being (re)inserted."),
}
}
}
pub fn insert(&mut self, node: Node<IdBytes, TVal>, status: NodeStatus) -> InsertResult {
match status {
NodeStatus::Connected => {
if self.nodes.is_full() {
if self.first_connected_pos == Some(0) || self.pending.is_some() {
return InsertResult::Full;
} else {
self.pending = Some(PendingNode {
node,
status: NodeStatus::Connected,
replace: Instant::now() + self.pending_timeout,
});
return InsertResult::Pending {
disconnected: self.nodes[0].key,
};
}
}
let pos = self.nodes.len();
self.first_connected_pos = self.first_connected_pos.or(Some(pos));
self.nodes.push(node);
InsertResult::Inserted
}
NodeStatus::Disconnected => {
if self.nodes.is_full() {
return InsertResult::Full;
}
if let Some(ref mut p) = self.first_connected_pos {
self.nodes.insert(*p, node);
*p += 1;
} else {
self.nodes.push(node);
}
InsertResult::Inserted
}
}
}
pub fn remove(&mut self, key: &IdBytes) -> Option<(Node<IdBytes, TVal>, NodeStatus, Position)> {
if let Some(pos) = self.position(key) {
let status = self.status(pos);
let node = self.nodes.remove(pos.0);
match status {
NodeStatus::Connected => {
if (self.first_connected_pos == Some(pos.0)) && pos.0 == self.nodes.len() {
self.first_connected_pos = None
}
}
NodeStatus::Disconnected => {
if let Some(ref mut p) = self.first_connected_pos {
*p -= 1;
}
}
}
Some((node, status, pos))
} else {
None
}
}
pub fn status(&self, pos: Position) -> NodeStatus {
if self.first_connected_pos.is_some_and(|i| pos.0 >= i) {
NodeStatus::Connected
} else {
NodeStatus::Disconnected
}
}
#[expect(unused)] pub fn is_connected(&self, pos: Position) -> bool {
self.status(pos) == NodeStatus::Connected
}
pub fn num_entries(&self) -> usize {
self.nodes.len()
}
pub fn num_connected(&self) -> usize {
self.first_connected_pos.map_or(0, |i| self.nodes.len() - i)
}
#[expect(unused)] pub fn num_disconnected(&self) -> usize {
self.nodes.len() - self.num_connected()
}
pub fn position(&self, key: &IdBytes) -> Option<Position> {
self.nodes
.iter()
.position(|p| p.key.as_ref() == key.as_ref())
.map(Position)
}
pub fn get_mut(&mut self, key: &IdBytes) -> Option<&mut Node<IdBytes, TVal>> {
self.nodes
.iter_mut()
.find(move |p| p.key.as_ref() == key.as_ref())
}
}