Skip to main content

kcode_k1_launch_nodes/
lib.rs

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}