kcode-k1-kmap 0.5.0

K1 Kmap transaction facade and attention-budgeted loader
Documentation
use std::{
    path::Path,
    sync::{
        Arc,
        atomic::{AtomicBool, Ordering},
    },
};

use kcode_k1_kmap_projection::Projection;
use kcode_k1_peering::K1Peering;
use kcode_k1_transaction::SubsystemId;
use kcode_k1_txn_ordering::{K1TxnOrdering, Subsystem};

pub use kcode_k1_kmap_format::{
    Connection, ConnectionMeasurement, ConnectionSpec, ConnectionTier, KmapAction, KmapError,
    MeasurementImportance, MeasurementOutcome, Node, NodeId, Weight,
};
pub use kcode_k1_kmap_loader::{LoadedNode, NARRATIVE_COST, OpenMode, OpenResult, PREVIEW_COST};
pub use kcode_k1_txn_ordering::TxId;

const SUBSYSTEM_NAME: &str = "k1-kmap-subsystem";

struct Availability {
    available: AtomicBool,
}

impl Availability {
    fn new() -> Self {
        Self {
            available: AtomicBool::new(true),
        }
    }

    fn ensure(&self) -> Result<(), String> {
        if self.available.load(Ordering::Acquire) {
            Ok(())
        } else {
            Err("k1 Kmap instance unavailable".to_owned())
        }
    }

    fn fail(&self) {
        self.available.store(false, Ordering::Release);
    }
}

struct KmapSubsystem {
    projection: Arc<Projection>,
    availability: Arc<Availability>,
}

impl Subsystem for KmapSubsystem {
    fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String> {
        self.availability.ensure()?;
        let action = match KmapAction::decode(payload) {
            Ok(action) => action,
            Err(error) => {
                self.availability.fail();
                return Err(error.to_string());
            }
        };
        match self.projection.apply(id, action) {
            Ok(_) => Ok(()),
            Err(error) => {
                self.availability.fail();
                Err(error)
            }
        }
    }

    fn reorg(&self) -> Result<(), String> {
        self.availability.fail();
        self.projection.clear()
    }
}

pub struct K1Kmap {
    projection: Arc<Projection>,
    _ordering: Arc<K1TxnOrdering>,
    peering: Arc<K1Peering>,
    availability: Arc<Availability>,
    subsystem: SubsystemId,
}

impl K1Kmap {
    pub fn open(
        root: &Path,
        ordering: Arc<K1TxnOrdering>,
        peering: Arc<K1Peering>,
    ) -> Result<Self, String> {
        let (projection, checkpoint) = Projection::open(root)?;
        let projection = Arc::new(projection);
        let availability = Arc::new(Availability::new());
        let subsystem = SubsystemId::from_str(SUBSYSTEM_NAME)?;
        let callback = Arc::new(KmapSubsystem {
            projection: projection.clone(),
            availability: availability.clone(),
        });
        if let Err(error) = ordering.register_subsystem(subsystem, checkpoint, callback) {
            availability.fail();
            return Err(error);
        }
        Ok(Self {
            projection,
            _ordering: ordering,
            peering,
            availability,
            subsystem,
        })
    }

    pub fn create_node(
        &self,
        title: impl Into<String>,
        navigation_hint: impl Into<String>,
        narrative: impl Into<String>,
        connections: Vec<ConnectionSpec>,
    ) -> Result<TxId, String> {
        self.submit_action(KmapAction::CreateNode {
            title: title.into(),
            navigation_hint: navigation_hint.into(),
            narrative: narrative.into(),
            connections,
        })
    }

    pub fn update_node(
        &self,
        node_id: NodeId,
        title: Option<String>,
        navigation_hint: Option<String>,
        narrative: Option<String>,
        connection_updates: Vec<ConnectionSpec>,
    ) -> Result<TxId, String> {
        self.submit_action(KmapAction::UpdateNode {
            node_id,
            title,
            navigation_hint,
            narrative,
            connection_updates,
        })
    }

    pub fn apply_measurements(
        &self,
        measurements: Vec<ConnectionMeasurement>,
    ) -> Result<TxId, String> {
        self.submit_action(KmapAction::ApplyMeasurements { measurements })
    }

    pub fn get_node(&self, node_id: NodeId) -> Result<Option<Node>, String> {
        self.availability.ensure()?;
        self.read_projection(node_id)
    }

    pub fn open_node(
        &self,
        node_id: NodeId,
        budget: f64,
        temperature: f64,
        mode: OpenMode,
        candidate_filter: impl FnMut(&[NodeId]) -> Result<Vec<NodeId>, String>,
    ) -> Result<OpenResult, String> {
        self.availability.ensure()?;
        kcode_k1_kmap_loader::open_node(
            node_id,
            budget,
            temperature,
            mode,
            |id| self.read_projection(id),
            candidate_filter,
        )
    }

    fn submit_action(&self, action: KmapAction) -> Result<TxId, String> {
        self.availability.ensure()?;
        let payload = action.encode().map_err(|error| error.to_string())?;
        self.peering.submit_txn(self.subsystem, &payload)
    }

    fn read_projection(&self, node_id: NodeId) -> Result<Option<Node>, String> {
        match self.projection.get(node_id) {
            Ok(node) => Ok(node),
            Err(error) => {
                self.availability.fail();
                Err(error)
            }
        }
    }
}