1use std::{
2 collections::HashMap,
3 fs::{File, OpenOptions},
4 io::{Read, Seek, SeekFrom, Write},
5 path::Path,
6 sync::{Arc, Mutex},
7};
8
9pub use kcode_k1_launch_node_codec::{
10 Authority, GroupId, NodeId, TargetId, TargetName, TxId, UserId,
11};
12use kcode_k1_launch_node_codec::{
13 PROJECTION_HEADER, ProjectionDecode, ProjectionRecord, SetAction, decode_projection_record,
14};
15use kcode_k1_peering::K1Peering;
16use kcode_k1_txn_ordering::{K1TxnOrdering, Subsystem as OrderingSubsystem, SubsystemId};
17
18struct Pending {
19 payload: Vec<u8>,
20 evidence: Option<TxId>,
21}
22
23struct Inner {
24 file: File,
25 bindings: HashMap<TargetId, NodeId>,
26 checkpoint: Option<TxId>,
27 pending: HashMap<[u8; 16], Pending>,
28 failed: bool,
29}
30
31struct Shared {
32 inner: Mutex<Inner>,
33}
34
35struct Driver {
36 shared: Arc<Shared>,
37}
38
39pub struct LaunchNodes {
40 shared: Arc<Shared>,
41 peering: Arc<K1Peering>,
42}
43
44impl LaunchNodes {
45 pub fn open(
46 path: &Path,
47 ordering: Arc<K1TxnOrdering>,
48 peering: Arc<K1Peering>,
49 ) -> Result<Self, String> {
50 let mut projection = open_projection(path)?;
51 if let Some(checkpoint) = projection.checkpoint.to_owned() {
52 match ordering.get_txn(checkpoint) {
53 Ok(Some(_)) => {}
54 Ok(None) => {
55 reset_projection(&mut projection.file, path)?;
56 projection.bindings.clear();
57 projection.checkpoint = None;
58 }
59 Err(error) => return Err(format!("validate launch nodes checkpoint: {error}")),
60 }
61 }
62 let after = projection.checkpoint.to_owned();
63 let shared = Arc::new(Shared {
64 inner: Mutex::new(Inner {
65 file: projection.file,
66 bindings: projection.bindings,
67 checkpoint: projection.checkpoint,
68 pending: HashMap::new(),
69 failed: false,
70 }),
71 });
72 ordering.register_subsystem(
73 subsystem()?,
74 after,
75 Arc::new(Driver {
76 shared: shared.clone(),
77 }),
78 )?;
79 Ok(Self { shared, peering })
80 }
81
82 pub fn get(&self, target: &TargetId) -> Result<Option<NodeId>, String> {
83 let inner = self
84 .shared
85 .inner
86 .lock()
87 .map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
88 if inner.failed {
89 return Err("k1-launch-nodes facade is unavailable".into());
90 }
91 Ok(inner.bindings.get(target).cloned())
92 }
93
94 pub fn set(&self, target: TargetId, node: NodeId) -> Result<(), String> {
95 let (operation_id, payload) = loop {
96 let mut operation_id = [0_u8; 16];
97 getrandom::fill(&mut operation_id)
98 .map_err(|error| format!("generate launch nodes operation ID: {error}"))?;
99 let payload = SetAction::new(operation_id, target.clone(), node).encode();
100 let mut inner = self
101 .shared
102 .inner
103 .lock()
104 .map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
105 if inner.failed {
106 return Err("k1-launch-nodes facade is unavailable".into());
107 }
108 if inner.pending.contains_key(&operation_id) {
109 continue;
110 }
111 inner.pending.insert(
112 operation_id,
113 Pending {
114 payload: payload.clone(),
115 evidence: None,
116 },
117 );
118 break (operation_id, payload);
119 };
120 let submission = self.peering.submit_txn(subsystem()?, &payload);
121 let mut inner = self
122 .shared
123 .inner
124 .lock()
125 .map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
126 if inner.failed {
127 return Err("k1-launch-nodes facade is unavailable".into());
128 }
129 let Some(pending) = inner.pending.remove(&operation_id) else {
130 return Err(fault(
131 &mut inner,
132 "k1-launch-nodes callback evidence is missing",
133 ));
134 };
135 match (submission, pending.evidence) {
136 (Ok(submitted), Some(callback)) if submitted == callback => Ok(()),
137 (Err(_), Some(_)) => Ok(()),
138 (Err(error), None) => Err(format!("submit k1-launch-nodes transaction: {error}")),
139 (Ok(_), None) => Err(fault(
140 &mut inner,
141 "k1-launch-nodes callback evidence is missing",
142 )),
143 (Ok(_), Some(_)) => Err(fault(
144 &mut inner,
145 "k1-launch-nodes callback transaction does not match submission",
146 )),
147 }
148 }
149}
150
151impl OrderingSubsystem for Driver {
152 fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String> {
153 let action = match SetAction::decode(payload) {
154 Ok(action) => action,
155 Err(error) => {
156 return self
157 .shared
158 .fail(format!("decode k1-launch-nodes Set action: {error}"));
159 }
160 };
161 let operation_id = action.operation_id().to_owned();
162 let canonical_payload = action.encode();
163 let record =
164 ProjectionRecord::new(id, action.target().to_owned(), action.node().to_owned())
165 .encode();
166 let mut inner = self
167 .shared
168 .inner
169 .lock()
170 .map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
171 if inner.failed {
172 return Err("k1-launch-nodes facade is unavailable".into());
173 }
174 if inner.checkpoint == Some(id) {
175 return Err(fault(
176 &mut inner,
177 "k1-launch-nodes callback transaction is duplicated",
178 ));
179 }
180 if let Some(pending) = inner.pending.get(&operation_id) {
181 if pending.payload != canonical_payload {
182 return Err(fault(
183 &mut inner,
184 "k1-launch-nodes callback action does not match reservation",
185 ));
186 }
187 if pending.evidence.is_some() {
188 return Err(fault(
189 &mut inner,
190 "k1-launch-nodes callback evidence is duplicated",
191 ));
192 }
193 }
194 let append = inner
195 .file
196 .seek(SeekFrom::End(0))
197 .and_then(|_| inner.file.write_all(&record))
198 .and_then(|()| inner.file.sync_data());
199 if let Err(error) = append {
200 return Err(fault(
201 &mut inner,
202 format!("append k1-launch-nodes projection record: {error}"),
203 ));
204 }
205 inner
206 .bindings
207 .insert(action.target().to_owned(), action.node().to_owned());
208 inner.checkpoint = Some(id);
209 if let Some(pending) = inner.pending.get_mut(&operation_id) {
210 pending.evidence = Some(id);
211 }
212 Ok(())
213 }
214
215 fn reorg(&self) -> Result<(), String> {
216 let mut inner = self
217 .shared
218 .inner
219 .lock()
220 .map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
221 inner.failed = true;
222 inner.pending.clear();
223 Ok(())
224 }
225}
226
227impl Shared {
228 fn fail(&self, message: String) -> Result<(), String> {
229 let mut inner = self
230 .inner
231 .lock()
232 .map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
233 Err(fault(&mut inner, message))
234 }
235}
236
237struct Projection {
238 file: File,
239 bindings: HashMap<TargetId, NodeId>,
240 checkpoint: Option<TxId>,
241}
242
243fn open_projection(path: &Path) -> Result<Projection, String> {
244 let mut file = OpenOptions::new()
245 .create(true)
246 .truncate(false)
247 .read(true)
248 .write(true)
249 .open(path)
250 .map_err(|error| format!("open launch nodes projection: {error}"))?;
251 if file
252 .metadata()
253 .map_err(|error| format!("inspect launch nodes projection: {error}"))?
254 .len()
255 == 0
256 {
257 reset_projection(&mut file, path)?;
258 }
259 file.seek(SeekFrom::Start(0))
260 .map_err(|error| format!("seek launch nodes projection: {error}"))?;
261 let mut bytes = Vec::new();
262 file.read_to_end(&mut bytes)
263 .map_err(|error| format!("read launch nodes projection: {error}"))?;
264 if bytes.len() < PROJECTION_HEADER.len()
265 || &bytes[..PROJECTION_HEADER.len()] != PROJECTION_HEADER
266 {
267 return Err("unsupported launch nodes projection format; no migration is available".into());
268 }
269 let mut bindings = HashMap::new();
270 let mut checkpoint = None;
271 let mut offset = PROJECTION_HEADER.len();
272 let mut corrupt = false;
273 while offset < bytes.len() {
274 match decode_projection_record(&bytes[offset..]) {
275 Ok(ProjectionDecode::Complete { record, consumed }) => {
276 bindings.insert(record.target().to_owned(), record.node().to_owned());
277 let callback: [u8; 12] = bytes[offset..offset + 12]
278 .try_into()
279 .map_err(|_| "launch nodes projection callback is unavailable".to_owned())?;
280 checkpoint = Some(TxId::from_bytes(callback));
281 offset += consumed;
282 }
283 Ok(ProjectionDecode::Incomplete) => {
284 repair_tail(&mut file, offset)?;
285 break;
286 }
287 Err(_) => {
288 corrupt = true;
289 break;
290 }
291 }
292 }
293 if corrupt {
294 reset_projection(&mut file, path)?;
295 bindings.clear();
296 checkpoint = None;
297 }
298 file.seek(SeekFrom::End(0))
299 .map_err(|error| format!("seek launch nodes append position: {error}"))?;
300 Ok(Projection {
301 file,
302 bindings,
303 checkpoint,
304 })
305}
306
307fn reset_projection(file: &mut File, path: &Path) -> Result<(), String> {
308 file.set_len(0)
309 .and_then(|()| file.seek(SeekFrom::Start(0)).map(|_| ()))
310 .and_then(|()| file.write_all(PROJECTION_HEADER))
311 .and_then(|()| file.sync_data())
312 .map_err(|error| format!("reset launch nodes projection: {error}"))?;
313 sync_parent(path)
314}
315
316fn repair_tail(file: &mut File, offset: usize) -> Result<(), String> {
317 file.set_len(offset as u64)
318 .and_then(|()| file.sync_data())
319 .map_err(|error| format!("repair launch nodes projection tail: {error}"))
320}
321
322fn sync_parent(path: &Path) -> Result<(), String> {
323 let parent = path
324 .parent()
325 .filter(|parent| !parent.as_os_str().is_empty())
326 .unwrap_or_else(|| Path::new("."));
327 File::open(parent)
328 .and_then(|directory| directory.sync_all())
329 .map_err(|error| format!("synchronize launch nodes projection directory: {error}"))
330}
331
332fn subsystem() -> Result<SubsystemId, String> {
333 SubsystemId::from_str("k1-launch-nodes")
334}
335
336fn fault(inner: &mut Inner, message: impl Into<String>) -> String {
337 inner.failed = true;
338 inner.pending.clear();
339 message.into()
340}
341
342#[cfg(test)]
343mod tests {
344 use super::*;
345 use std::{
346 fs,
347 path::{Path, PathBuf},
348 sync::atomic::{AtomicU64, Ordering},
349 };
350
351 static NEXT: AtomicU64 = AtomicU64::new(0);
352
353 struct Temp {
354 root: PathBuf,
355 }
356
357 impl Temp {
358 fn new() -> Self {
359 let root = std::env::temp_dir().join(format!(
360 "k1-launch-nodes-{}-{}",
361 std::process::id(),
362 NEXT.fetch_add(1, Ordering::Relaxed)
363 ));
364 fs::create_dir(&root).unwrap();
365 Self { root }
366 }
367 }
368
369 impl Drop for Temp {
370 fn drop(&mut self) {
371 fs::remove_dir_all(&self.root).unwrap();
372 }
373 }
374
375 fn services(root: &Path) -> (Arc<K1TxnOrdering>, Arc<K1Peering>) {
376 fs::create_dir_all(root).unwrap();
377 let ordering = Arc::new(K1TxnOrdering::open(&root.join("ordering")).unwrap());
378 let peering = Arc::new(K1Peering::open(&root.join("peering"), ordering.clone()).unwrap());
379 (ordering, peering)
380 }
381
382 fn target(authority: Authority, name: &str) -> TargetId {
383 TargetId::new(authority, TargetName::new(name.to_owned()).unwrap())
384 }
385
386 fn user(byte: u8) -> Authority {
387 Authority::User(UserId::from_tx_id(TxId::from_bytes([byte; 12])))
388 }
389
390 fn group(byte: u8) -> Authority {
391 Authority::Group(GroupId::new(TxId::from_bytes([byte; 12])))
392 }
393
394 fn record_count(path: &Path) -> usize {
395 let bytes = fs::read(path).unwrap();
396 let mut offset = PROJECTION_HEADER.len();
397 let mut count = 0;
398 while offset < bytes.len() {
399 match decode_projection_record(&bytes[offset..]).unwrap() {
400 ProjectionDecode::Complete { consumed, .. } => {
401 count += 1;
402 offset += consumed;
403 }
404 ProjectionDecode::Incomplete => panic!("unexpected incomplete record"),
405 }
406 }
407 count
408 }
409
410 #[test]
411 fn callbacks_preserve_authorities_and_last_canonical_write() {
412 let temp = Temp::new();
413 let (ordering, peering) = services(&temp.root.join("service"));
414 let store = LaunchNodes::open(&temp.root.join("projection"), ordering, peering).unwrap();
415 let user_target = target(user(1), "default/chat");
416 let group_target = target(group(1), "default/chat");
417 store.set(user_target.clone(), NodeId([1; 12])).unwrap();
418 store.set(group_target.clone(), NodeId([2; 12])).unwrap();
419 store.set(user_target.clone(), NodeId([3; 12])).unwrap();
420 assert_eq!(store.get(&user_target).unwrap(), Some(NodeId([3; 12])));
421 assert_eq!(store.get(&group_target).unwrap(), Some(NodeId([2; 12])));
422 }
423
424 #[test]
425 fn peering_error_without_callback_does_not_bind() {
426 let temp = Temp::new();
427 let ordering = Arc::new(K1TxnOrdering::open(&temp.root.join("ordering-a")).unwrap());
428 let other = Arc::new(K1TxnOrdering::open(&temp.root.join("ordering-b")).unwrap());
429 let peering = Arc::new(K1Peering::open(&temp.root.join("peering-b"), other).unwrap());
430 let store = LaunchNodes::open(&temp.root.join("projection"), ordering, peering).unwrap();
431 let key = target(user(2), "model/x");
432 assert!(store.set(key.clone(), NodeId([4; 12])).is_err());
433 assert_eq!(store.get(&key).unwrap(), None);
434 assert_eq!(record_count(&temp.root.join("projection")), 0);
435 }
436
437 #[test]
438 fn projection_repairs_rebuilds_and_replays_same_value_sets() {
439 let temp = Temp::new();
440 let service = temp.root.join("service");
441 let projection = temp.root.join("projection");
442 let key = target(group(3), "same/value");
443 let node = NodeId([5; 12]);
444 let (ordering, peering) = services(&service);
445 let store = LaunchNodes::open(&projection, ordering.clone(), peering.clone()).unwrap();
446 store.set(key.clone(), node).unwrap();
447 store.set(key.clone(), node).unwrap();
448 drop((store, peering, ordering));
449
450 let clean = fs::read(&projection).unwrap();
451 let mut incomplete = clean.clone();
452 incomplete.extend_from_slice(&[1, 2, 3]);
453 fs::write(&projection, incomplete).unwrap();
454 let (ordering, peering) = services(&service);
455 let store = LaunchNodes::open(&projection, ordering.clone(), peering.clone()).unwrap();
456 assert_eq!(store.get(&key).unwrap(), Some(node));
457 assert_eq!(fs::read(&projection).unwrap(), clean);
458 drop((store, peering, ordering));
459
460 let mut corrupt = fs::read(&projection).unwrap();
461 let last = corrupt.last_mut().unwrap();
462 *last ^= 0xff;
463 fs::write(&projection, corrupt).unwrap();
464 let (ordering, peering) = services(&service);
465 let store = LaunchNodes::open(&projection, ordering, peering).unwrap();
466 assert_eq!(store.get(&key).unwrap(), Some(node));
467 assert_eq!(record_count(&projection), 2);
468 }
469
470 #[test]
471 fn non_v3_projection_is_rejected_without_migration() {
472 let temp = Temp::new();
473 let projection = temp.root.join("projection");
474 fs::write(&projection, b"K1LNV2\0\0").unwrap();
475 let (ordering, peering) = services(&temp.root.join("service"));
476 let error = LaunchNodes::open(&projection, ordering, peering)
477 .err()
478 .unwrap();
479 assert_eq!(
480 error,
481 "unsupported launch nodes projection format; no migration is available"
482 );
483 assert_eq!(fs::read(projection).unwrap(), b"K1LNV2\0\0");
484 }
485
486 #[test]
487 fn malformed_callback_faults_facade() {
488 let temp = Temp::new();
489 let (ordering, peering) = services(&temp.root.join("malformed"));
490 let store =
491 LaunchNodes::open(&temp.root.join("malformed-projection"), ordering, peering).unwrap();
492 let driver = Driver {
493 shared: store.shared.clone(),
494 };
495 assert!(
496 driver
497 .submit_txn(TxId::from_bytes([8; 12]), b"malformed")
498 .is_err()
499 );
500 assert!(store.get(&target(user(8), "faulted")).is_err());
501 }
502
503 #[test]
504 fn complete_package_stays_below_the_managed_limit() {
505 let files = [
506 include_str!("../Cargo.toml"),
507 include_str!("../Documentation.md"),
508 include_str!("lib.rs"),
509 ];
510 let count = files
511 .iter()
512 .flat_map(|file| file.lines())
513 .filter(|line| !line.trim().is_empty())
514 .count();
515 assert!(count < 500, "complete package has {count} nonblank lines");
516 }
517}