Skip to main content

kcode_k1_kmap/
lib.rs

1use std::{
2    path::Path,
3    sync::{
4        Arc,
5        atomic::{AtomicBool, Ordering},
6    },
7};
8
9use kcode_k1_kmap_projection::Projection;
10use kcode_k1_peering::K1Peering;
11use kcode_k1_transaction::SubsystemId;
12use kcode_k1_txn_ordering::{K1TxnOrdering, Subsystem};
13
14pub use kcode_k1_kmap_format::{
15    Connection, ConnectionMeasurement, ConnectionSpec, ConnectionTier, KmapAction, KmapError,
16    MeasurementImportance, MeasurementOutcome, Node, NodeId, Weight,
17};
18pub use kcode_k1_kmap_loader::{
19    DEPTH_DECAY, LoadedNode, NARRATIVE_COST, OpenMode, OpenResult, PREVIEW_COST,
20};
21pub use kcode_k1_txn_ordering::TxId;
22
23const SUBSYSTEM_NAME: &str = "k1-kmap-subsystem";
24
25struct Availability {
26    available: AtomicBool,
27}
28
29impl Availability {
30    fn new() -> Self {
31        Self {
32            available: AtomicBool::new(true),
33        }
34    }
35
36    fn ensure(&self) -> Result<(), String> {
37        if self.available.load(Ordering::Acquire) {
38            Ok(())
39        } else {
40            Err("k1 Kmap instance unavailable".to_owned())
41        }
42    }
43
44    fn fail(&self) {
45        self.available.store(false, Ordering::Release);
46    }
47}
48
49struct KmapSubsystem {
50    projection: Arc<Projection>,
51    availability: Arc<Availability>,
52}
53
54impl Subsystem for KmapSubsystem {
55    fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String> {
56        self.availability.ensure()?;
57        let action = match KmapAction::decode(payload) {
58            Ok(action) => action,
59            Err(error) => {
60                self.availability.fail();
61                return Err(error.to_string());
62            }
63        };
64        match self.projection.apply(id, action) {
65            Ok(_) => Ok(()),
66            Err(error) => {
67                self.availability.fail();
68                Err(error)
69            }
70        }
71    }
72
73    fn reorg(&self) -> Result<(), String> {
74        self.availability.fail();
75        self.projection.clear()
76    }
77}
78
79pub struct K1Kmap {
80    projection: Arc<Projection>,
81    _ordering: Arc<K1TxnOrdering>,
82    peering: Arc<K1Peering>,
83    availability: Arc<Availability>,
84    subsystem: SubsystemId,
85}
86
87impl K1Kmap {
88    pub fn open(
89        root: &Path,
90        ordering: Arc<K1TxnOrdering>,
91        peering: Arc<K1Peering>,
92    ) -> Result<Self, String> {
93        let (projection, checkpoint) = Projection::open(root)?;
94        let projection = Arc::new(projection);
95        let availability = Arc::new(Availability::new());
96        let subsystem = SubsystemId::from_str(SUBSYSTEM_NAME)?;
97        let callback = Arc::new(KmapSubsystem {
98            projection: projection.clone(),
99            availability: availability.clone(),
100        });
101        if let Err(error) = ordering.register_subsystem(subsystem, checkpoint, callback) {
102            availability.fail();
103            return Err(error);
104        }
105        Ok(Self {
106            projection,
107            _ordering: ordering,
108            peering,
109            availability,
110            subsystem,
111        })
112    }
113
114    pub fn create_node(
115        &self,
116        title: impl Into<String>,
117        navigation_hint: impl Into<String>,
118        narrative: impl Into<String>,
119        connections: Vec<ConnectionSpec>,
120    ) -> Result<TxId, String> {
121        self.submit_action(KmapAction::CreateNode {
122            title: title.into(),
123            navigation_hint: navigation_hint.into(),
124            narrative: narrative.into(),
125            connections,
126        })
127    }
128
129    pub fn update_node(
130        &self,
131        node_id: NodeId,
132        title: Option<String>,
133        navigation_hint: Option<String>,
134        narrative: Option<String>,
135        connection_updates: Vec<ConnectionSpec>,
136    ) -> Result<TxId, String> {
137        self.submit_action(KmapAction::UpdateNode {
138            node_id,
139            title,
140            navigation_hint,
141            narrative,
142            connection_updates,
143        })
144    }
145
146    pub fn apply_measurements(
147        &self,
148        measurements: Vec<ConnectionMeasurement>,
149    ) -> Result<TxId, String> {
150        self.submit_action(KmapAction::ApplyMeasurements { measurements })
151    }
152
153    pub fn get_node(&self, node_id: NodeId) -> Result<Option<Node>, String> {
154        self.availability.ensure()?;
155        self.read_projection(node_id)
156    }
157
158    pub fn open_node(
159        &self,
160        node_id: NodeId,
161        budget: f64,
162        temperature: f64,
163        mode: OpenMode,
164        access_filter: impl FnMut(NodeId) -> bool,
165    ) -> Result<OpenResult, String> {
166        self.availability.ensure()?;
167        kcode_k1_kmap_loader::open_node(
168            node_id,
169            budget,
170            temperature,
171            mode,
172            |id| self.read_projection(id),
173            access_filter,
174        )
175    }
176
177    fn submit_action(&self, action: KmapAction) -> Result<TxId, String> {
178        self.availability.ensure()?;
179        let payload = action.encode().map_err(|error| error.to_string())?;
180        self.peering.submit_txn(self.subsystem, &payload)
181    }
182
183    fn read_projection(&self, node_id: NodeId) -> Result<Option<Node>, String> {
184        match self.projection.get(node_id) {
185            Ok(node) => Ok(node),
186            Err(error) => {
187                self.availability.fail();
188                Err(error)
189            }
190        }
191    }
192}