unb-core 2.0.3

Core unb protocol types: envelope, session, routing, taxonomy
Documentation
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, TargetPath};

#[derive(Clone)]
pub struct NodeCore {
    routes: RouteTable,
    catalog: BTreeMap<String, Value>,
    revision: u64,
    route_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,
            route_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.catalog.keys().map(String::as_str))
    }

    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);
                let target_path = TargetPath::application(self.node(), subject)
                    .expect("installed catalog subjects are path-safe")
                    .to_string();
                if detail_full {
                    let mut entry = described.clone();
                    entry["subject"] = Value::String(subject.clone());
                    entry["target_path"] = Value::String(target_path);
                    entry
                } else {
                    json!({
                        "subject": subject,
                        "target_path": target_path,
                        "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, _)| TargetPath::application(self.node(), subject).is_ok())
            .collect();
        if self.catalog == capabilities {
            return false;
        }
        self.catalog = capabilities;
        self.revision = self
            .revision
            .checked_add(1)
            .expect("catalog revision overflow");
        true
    }

    pub fn resolve(&self, destination: &str) -> Resolution {
        self.routes.resolve(destination)
    }

    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 {
            destination: node.clone(),
            owner: node.clone(),
            owner_instance: self.instance_id.clone(),
            owner_epoch: self.epoch,
            owner_revision: self.route_revision,
            distance: 0,
            path: vec![node.clone()],
        }];
        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 {
                destination: advertisement.destination.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(destination: &str, owner: &str, path: &[&str]) -> RouteAdvertisement {
        RouteAdvertisement {
            destination: destination.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("leaf-a", "leaf-a", &["leaf-a"])]),
        )
        .unwrap();
        core.apply_snapshot(
            "sess-2",
            "leaf-b",
            &RouteSnapshot::canonical(1, vec![advertisement("leaf-b", "leaf-b", &["leaf-b"])]),
        )
        .unwrap();
        core.leave("sess-1");
        assert_eq!(core.resolve("leaf-a"), Resolution::Unknown);
        assert_eq!(core.resolve("leaf-b"), 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, 0);
        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("owner-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("owner-far"), Resolution::Route(_)));
        let export = core.export_for("elsewhere");
        assert_eq!(export.len(), 1, "the overlong transit route is omitted");
        assert_eq!(export[0].destination, "hub");
    }
}