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}