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::{
DEPTH_DECAY, 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,
access_filter: impl FnMut(NodeId) -> bool,
) -> Result<OpenResult, String> {
self.availability.ensure()?;
kcode_k1_kmap_loader::open_node(
node_id,
budget,
temperature,
mode,
|id| self.read_projection(id),
access_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)
}
}
}
}