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