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}