use std::collections::BTreeMap;
use serde_json::{json, Value};
use crate::discover::NodeCatalogSnapshot;
use crate::route_control::{RouteAdvertisement, RouteSnapshot};
use crate::route_table::RouteTable;
use crate::{Envelope, Resolution, RouteError};
#[derive(Clone)]
pub struct NodeCore {
routes: RouteTable,
catalog: BTreeMap<String, Value>,
revision: u64,
instance_id: String,
epoch: u64,
proof: Value,
}
impl NodeCore {
pub fn new(node: &str) -> NodeCore {
NodeCore {
routes: RouteTable::new(node),
catalog: BTreeMap::new(),
revision: 0,
instance_id: node.to_string(),
epoch: 0,
proof: Value::Null,
}
}
pub fn set_identity(&mut self, instance_id: &str, epoch: u64) {
self.instance_id = instance_id.to_string();
self.epoch = epoch;
}
pub fn set_node_identity(&mut self, identity: crate::NodeIdentity) {
self.instance_id = identity.instance_id;
self.epoch = identity.epoch;
self.proof = identity.proof;
}
pub fn node(&self) -> &str {
self.routes.node()
}
pub fn identity(&self) -> crate::NodeIdentity {
crate::NodeIdentity {
node_id: self.node().to_string(),
instance_id: self.instance_id.clone(),
epoch: self.epoch,
proof: self.proof.clone(),
}
}
pub fn catalog_revision(&self) -> u64 {
self.revision
}
pub fn fingerprint(&self) -> String {
fingerprint_of(self.routes.local_subjects())
}
pub fn catalog_subjects(&self, detail_full: bool) -> Vec<Value> {
self.catalog
.iter()
.map(|(subject, described)| {
let one_line = described.get("one_line").cloned().unwrap_or(Value::Null);
if detail_full {
let mut entry = described.clone();
entry["subject"] = Value::String(subject.clone());
entry
} else {
json!({ "subject": subject, "one_line": one_line })
}
})
.collect()
}
pub fn catalog(&self, detail_full: bool) -> Value {
json!({ "node": self.node(), "subjects": self.catalog_subjects(detail_full) })
}
pub fn catalog_snapshot(&self, detail_full: bool) -> NodeCatalogSnapshot {
NodeCatalogSnapshot {
node: self.node().to_string(),
instance_id: self.instance_id.clone(),
revision: self.revision,
fingerprint: self.fingerprint(),
subjects: self.catalog_subjects(detail_full),
}
}
pub fn install_local_capabilities(&mut self, capabilities: BTreeMap<String, Value>) -> bool {
let capabilities: BTreeMap<_, _> = capabilities
.into_iter()
.filter(|(subject, _)| {
!subject.is_empty() && subject.len() <= crate::route_control::MAX_SUBJECT_LEN
})
.collect();
let current: Vec<_> = self.routes.local_subjects().map(str::to_owned).collect();
let subjects: Vec<_> = capabilities.keys().cloned().collect();
if self.catalog == capabilities && current == subjects {
return false;
}
for subject in current {
if !capabilities.contains_key(&subject) {
self.routes.unregister_local(&subject);
}
}
for subject in capabilities.keys() {
self.routes.register_local(subject);
}
self.catalog = capabilities;
self.revision = self
.revision
.checked_add(1)
.expect("catalog revision overflow");
true
}
pub fn resolve(&self, subject: &str) -> Resolution {
self.routes.resolve(subject)
}
pub fn apply_snapshot(
&mut self,
session: &str,
advertiser: &str,
snapshot: &RouteSnapshot,
) -> Result<Vec<String>, RouteError> {
self.routes.apply_snapshot(session, advertiser, snapshot)
}
pub fn apply_delta(
&mut self,
session: &str,
advertiser: &str,
delta: &crate::route_control::RouteDelta,
) -> Result<Vec<String>, RouteError> {
self.routes.apply_delta(session, advertiser, delta)
}
pub fn leave(&mut self, session: &str) -> Vec<String> {
self.routes.leave(session)
}
pub fn applied_generation(&self, session: &str) -> Option<u64> {
self.routes.applied_generation(session)
}
pub fn export_for(&self, peer_node: &str) -> Vec<RouteAdvertisement> {
let node = self.routes.node().to_string();
let mut routes: Vec<RouteAdvertisement> = self
.routes
.local_subjects()
.map(|subject| RouteAdvertisement {
subject: subject.to_string(),
owner: node.clone(),
owner_instance: self.instance_id.clone(),
owner_epoch: self.epoch,
owner_revision: self.revision,
distance: 0,
path: vec![node.clone()],
})
.collect();
routes.extend(self.routes.selected_transit().filter_map(|candidate| {
let advertisement = &candidate.advertisement;
if advertisement.path.iter().any(|hop| hop == peer_node) {
return None;
}
if advertisement.path.len() >= crate::route_control::MAX_ROUTE_PATH {
return None;
}
let mut path = advertisement.path.clone();
path.push(node.clone());
Some(RouteAdvertisement {
subject: advertisement.subject.clone(),
owner: advertisement.owner.clone(),
owner_instance: advertisement.owner_instance.clone(),
owner_epoch: advertisement.owner_epoch,
owner_revision: advertisement.owner_revision,
distance: advertisement.distance.checked_add(1)?,
path,
})
}));
routes.sort();
routes
}
pub fn reachable_names(&self) -> Vec<String> {
self.routes.reachable_names()
}
pub fn forward(&self, envelope: Envelope) -> Result<(String, Envelope), RouteError> {
self.routes.forward(envelope)
}
pub fn annotate_error(&self, envelope: Envelope) -> Envelope {
self.routes.annotate_error(envelope)
}
}
fn fingerprint_of<'a>(names: impl Iterator<Item = &'a str>) -> String {
let mut hash: u64 = 0xcbf29ce484222325;
for name in names {
for byte in name.bytes() {
hash ^= u64::from(byte);
hash = hash.wrapping_mul(0x100000001b3);
}
hash ^= 0xff;
hash = hash.wrapping_mul(0x100000001b3);
}
format!("{hash:016x}")
}
#[cfg(test)]
mod tests {
use super::*;
use crate::route_control::MAX_ROUTE_PATH;
use serde_json::json;
fn advertisement(subject: &str, owner: &str, path: &[&str]) -> RouteAdvertisement {
RouteAdvertisement {
subject: subject.into(),
owner: owner.into(),
owner_instance: format!("{owner}-inst"),
owner_epoch: 1,
owner_revision: 0,
distance: (path.len() - 1) as u32,
path: path.iter().map(|s| s.to_string()).collect(),
}
}
#[test]
fn session_close_removes_exactly_its_routes() {
let mut core = NodeCore::new("hub");
core.apply_snapshot(
"sess-1",
"leaf-a",
&RouteSnapshot::canonical(
1,
vec![
advertisement("chess", "leaf-a", &["leaf-a"]),
advertisement("checkers", "leaf-a", &["leaf-a"]),
],
),
)
.unwrap();
core.apply_snapshot(
"sess-2",
"leaf-b",
&RouteSnapshot::canonical(1, vec![advertisement("go", "leaf-b", &["leaf-b"])]),
)
.unwrap();
core.leave("sess-1");
assert_eq!(core.resolve("chess"), Resolution::Unknown);
assert_eq!(core.resolve("checkers"), Resolution::Unknown);
assert_eq!(core.resolve("go"), Resolution::Route("leaf-b".into()));
}
#[test]
fn exports_carry_the_local_identity_and_revision() {
let mut core = NodeCore::new("hub");
core.set_identity("hub-7", 3);
core.install_local_capabilities(BTreeMap::from([("chess".into(), json!({}))]));
let export = core.export_for("anyone");
assert_eq!(export.len(), 1);
assert_eq!(export[0].owner, "hub");
assert_eq!(export[0].owner_instance, "hub-7");
assert_eq!(export[0].owner_epoch, 3);
assert_eq!(export[0].owner_revision, core.catalog_revision());
assert_eq!(export[0].path, vec!["hub"]);
}
#[test]
fn a_transit_route_at_the_path_limit_is_not_exported() {
let mut core = NodeCore::new("hub");
let mut path: Vec<String> = (0..MAX_ROUTE_PATH - 1).map(|i| format!("n{i}")).collect();
path.insert(0, "owner-far".to_string());
let path_refs: Vec<&str> = path.iter().map(String::as_str).collect();
let mut long = advertisement("far", "owner-far", &path_refs);
long.path[MAX_ROUTE_PATH - 1] = "adv".into();
core.apply_snapshot("sess-1", "adv", &RouteSnapshot::canonical(1, vec![long]))
.unwrap();
assert!(matches!(core.resolve("far"), Resolution::Route(_)));
assert!(
core.export_for("elsewhere").is_empty(),
"appending this node would exceed the path limit"
);
}
}