Skip to main content

kcode_k1_peering/
lib.rs

1use ed25519_dalek::{Signer, SigningKey};
2use kcode_k1_transaction_id::TxId;
3use kcode_k1_txn_ordering::{K1TxnOrdering, SubsystemId};
4use std::fs::{self, File, OpenOptions};
5use std::io::{Read, Write};
6use std::path::Path;
7use std::sync::Arc;
8use std::time::{SystemTime, UNIX_EPOCH};
9
10const IDENTITY_BYTES: usize = 64;
11const PRIVATE_KEY_BYTES: usize = 32;
12const PUBLIC_KEY_BYTES: usize = 32;
13
14pub struct K1Peering {
15    ordering: Arc<K1TxnOrdering>,
16    signing_key: SigningKey,
17    public_key: [u8; PUBLIC_KEY_BYTES],
18}
19
20impl K1Peering {
21    pub fn open(root: &Path, ordering: Arc<K1TxnOrdering>) -> Result<Self, String> {
22        prepare_root(root)?;
23        let identity_path = root.join("identity.key");
24
25        let (signing_key, public_key) = match fs::symlink_metadata(&identity_path) {
26            Ok(metadata) => load_identity(&identity_path, metadata)?,
27            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
28                if ordering.tip().is_some() {
29                    return Err(
30                        "identity.key is missing while canonical transaction history exists"
31                            .to_owned(),
32                    );
33                }
34                create_identity(&identity_path)?
35            }
36            Err(error) => return Err(format!("cannot inspect identity.key: {error}")),
37        };
38
39        Ok(Self {
40            ordering,
41            signing_key,
42            public_key,
43        })
44    }
45
46    pub fn submit_txn(&self, subsystem: SubsystemId, payload: &[u8]) -> Result<TxId, String> {
47        let timestamp = SystemTime::now()
48            .duration_since(UNIX_EPOCH)
49            .map_err(|_| "system clock is before the Unix epoch".to_owned())?
50            .as_secs();
51
52        self.ordering
53            .submit_local_txn(
54                timestamp,
55                self.public_key,
56                subsystem,
57                payload,
58                |bytes| Ok(self.signing_key.sign(bytes).to_bytes()),
59                |_| Ok(()),
60            )
61            .map(|(id, _)| id)
62    }
63}
64
65fn prepare_root(root: &Path) -> Result<(), String> {
66    match fs::symlink_metadata(root) {
67        Ok(metadata) if metadata.file_type().is_dir() => return Ok(()),
68        Ok(_) => return Err("peering root is not a real directory".to_owned()),
69        Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
70        Err(error) => return Err(format!("cannot inspect peering root: {error}")),
71    }
72
73    fs::create_dir_all(root).map_err(|error| format!("cannot create peering root: {error}"))?;
74
75    let metadata = fs::symlink_metadata(root)
76        .map_err(|error| format!("cannot inspect created peering root: {error}"))?;
77
78    if !metadata.file_type().is_dir() {
79        return Err("created peering root is not a real directory".to_owned());
80    }
81
82    Ok(())
83}
84
85fn load_identity(
86    path: &Path,
87    metadata: fs::Metadata,
88) -> Result<(SigningKey, [u8; PUBLIC_KEY_BYTES]), String> {
89    if !metadata.file_type().is_file() {
90        return Err("identity.key is not a real regular file".to_owned());
91    }
92
93    if metadata.len() != IDENTITY_BYTES as u64 {
94        return Err("identity.key must contain exactly 64 bytes".to_owned());
95    }
96
97    validate_identity_permissions(&metadata)?;
98
99    let mut file =
100        File::open(path).map_err(|error| format!("cannot open identity.key: {error}"))?;
101    let mut bytes = [0_u8; IDENTITY_BYTES];
102    file.read_exact(&mut bytes)
103        .map_err(|error| format!("cannot read identity.key: {error}"))?;
104
105    let private_seed: [u8; PRIVATE_KEY_BYTES] = bytes[..PRIVATE_KEY_BYTES]
106        .try_into()
107        .expect("fixed private-key range");
108    let stored_public_key: [u8; PUBLIC_KEY_BYTES] = bytes[PRIVATE_KEY_BYTES..]
109        .try_into()
110        .expect("fixed public-key range");
111    let signing_key = SigningKey::from_bytes(&private_seed);
112    let derived_public_key = signing_key.verifying_key().to_bytes();
113
114    if stored_public_key != derived_public_key {
115        return Err("identity.key public key does not match its private seed".to_owned());
116    }
117
118    Ok((signing_key, derived_public_key))
119}
120
121fn create_identity(path: &Path) -> Result<(SigningKey, [u8; PUBLIC_KEY_BYTES]), String> {
122    let mut private_seed = [0_u8; PRIVATE_KEY_BYTES];
123    getrandom::fill(&mut private_seed)
124        .map_err(|error| format!("cannot generate peering identity: {error}"))?;
125
126    let signing_key = SigningKey::from_bytes(&private_seed);
127    let public_key = signing_key.verifying_key().to_bytes();
128    let mut bytes = [0_u8; IDENTITY_BYTES];
129    bytes[..PRIVATE_KEY_BYTES].copy_from_slice(&private_seed);
130    bytes[PRIVATE_KEY_BYTES..].copy_from_slice(&public_key);
131
132    let mut options = OpenOptions::new();
133    options.write(true).create_new(true);
134    set_owner_only_creation_mode(&mut options);
135
136    let mut file = options
137        .open(path)
138        .map_err(|error| format!("cannot create identity.key: {error}"))?;
139    file.write_all(&bytes)
140        .map_err(|error| format!("cannot write identity.key: {error}"))?;
141    file.sync_all()
142        .map_err(|error| format!("cannot synchronize identity.key: {error}"))?;
143
144    Ok((signing_key, public_key))
145}
146
147#[cfg(unix)]
148fn set_owner_only_creation_mode(options: &mut OpenOptions) {
149    use std::os::unix::fs::OpenOptionsExt;
150    options.mode(0o600);
151}
152
153#[cfg(not(unix))]
154fn set_owner_only_creation_mode(_: &mut OpenOptions) {}
155
156#[cfg(unix)]
157fn validate_identity_permissions(metadata: &fs::Metadata) -> Result<(), String> {
158    use std::os::unix::fs::PermissionsExt;
159
160    if metadata.permissions().mode() & 0o077 != 0 {
161        return Err("identity.key grants group or other permissions".to_owned());
162    }
163
164    Ok(())
165}
166
167#[cfg(not(unix))]
168fn validate_identity_permissions(_: &fs::Metadata) -> Result<(), String> {
169    Ok(())
170}
171
172#[cfg(test)]
173mod tests {
174    use super::*;
175    use ed25519_dalek::Signature;
176    use kcode_k1_txn_ordering::{GENESIS_PARENT, Subsystem};
177    use std::path::PathBuf;
178    use std::sync::Mutex;
179    use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
180
181    static NEXT_ROOT: AtomicU64 = AtomicU64::new(0);
182
183    struct TempRoots {
184        base: PathBuf,
185    }
186
187    impl TempRoots {
188        fn new(label: &str) -> Self {
189            let sequence = NEXT_ROOT.fetch_add(1, Ordering::Relaxed);
190            let base = std::env::temp_dir().join(format!(
191                "kcode-k1-peering-{}-{}-{}",
192                std::process::id(),
193                sequence,
194                label
195            ));
196            let _ = fs::remove_dir_all(&base);
197            Self { base }
198        }
199
200        fn peering(&self) -> PathBuf {
201            self.base.join("peering")
202        }
203
204        fn ordering(&self) -> PathBuf {
205            self.base.join("ordering")
206        }
207    }
208
209    impl Drop for TempRoots {
210        fn drop(&mut self) {
211            let _ = fs::remove_dir_all(&self.base);
212        }
213    }
214
215    struct RecordingSubsystem {
216        submissions: Mutex<Vec<(TxId, Vec<u8>)>>,
217        fail: AtomicBool,
218        reorgs: AtomicUsize,
219    }
220
221    impl RecordingSubsystem {
222        fn new() -> Self {
223            Self {
224                submissions: Mutex::new(Vec::new()),
225                fail: AtomicBool::new(false),
226                reorgs: AtomicUsize::new(0),
227            }
228        }
229
230        fn entries(&self) -> Vec<(TxId, Vec<u8>)> {
231            self.submissions.lock().unwrap().clone()
232        }
233    }
234
235    impl Subsystem for RecordingSubsystem {
236        fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String> {
237            self.submissions
238                .lock()
239                .unwrap()
240                .push((id, payload.to_vec()));
241
242            if self.fail.load(Ordering::Relaxed) {
243                Err("integration failed".to_owned())
244            } else {
245                Ok(())
246            }
247        }
248
249        fn reorg(&self) -> Result<(), String> {
250            self.reorgs.fetch_add(1, Ordering::Relaxed);
251            Ok(())
252        }
253    }
254
255    fn subsystem(value: u8) -> SubsystemId {
256        SubsystemId::from_bytes([value; 20]).unwrap()
257    }
258
259    fn open_ordering(roots: &TempRoots) -> Arc<K1TxnOrdering> {
260        Arc::new(K1TxnOrdering::open(&roots.ordering()).unwrap())
261    }
262
263    fn write_identity(path: &Path, bytes: &[u8]) {
264        fs::create_dir_all(path.parent().unwrap()).unwrap();
265        let mut options = OpenOptions::new();
266        options.write(true).create_new(true);
267        set_owner_only_creation_mode(&mut options);
268        let mut file = options.open(path).unwrap();
269        file.write_all(bytes).unwrap();
270        file.sync_all().unwrap();
271    }
272
273    #[test]
274    fn creates_owner_only_identity_and_reopens_stably() {
275        let roots = TempRoots::new("identity");
276        let ordering = open_ordering(&roots);
277        let first = K1Peering::open(&roots.peering(), ordering.clone()).unwrap();
278        let first_bytes = fs::read(roots.peering().join("identity.key")).unwrap();
279
280        assert_eq!(first_bytes.len(), IDENTITY_BYTES);
281        assert_eq!(&first_bytes[PRIVATE_KEY_BYTES..], &first.public_key);
282
283        #[cfg(unix)]
284        {
285            use std::os::unix::fs::PermissionsExt;
286            let mode = fs::metadata(roots.peering().join("identity.key"))
287                .unwrap()
288                .permissions()
289                .mode();
290            assert_eq!(mode & 0o077, 0);
291        }
292
293        drop(first);
294        let second = K1Peering::open(&roots.peering(), ordering).unwrap();
295        let second_bytes = fs::read(roots.peering().join("identity.key")).unwrap();
296
297        assert_eq!(second_bytes, first_bytes);
298        assert_eq!(second.public_key, first_bytes[PRIVATE_KEY_BYTES..]);
299    }
300
301    #[test]
302    fn submits_exact_signed_transaction_and_delivers_once() {
303        let roots = TempRoots::new("signed");
304        let ordering = open_ordering(&roots);
305        let target = subsystem(b'a');
306        let handler = Arc::new(RecordingSubsystem::new());
307        ordering
308            .register_subsystem(target, None, handler.clone())
309            .unwrap();
310        let peering = K1Peering::open(&roots.peering(), ordering.clone()).unwrap();
311
312        let returned_id = peering.submit_txn(target, b"object update").unwrap();
313
314        let entries = handler.entries();
315        assert_eq!(entries.len(), 1);
316        assert_eq!(entries[0].1, b"object update");
317        assert_eq!(returned_id, entries[0].0);
318        assert_eq!(ordering.tip(), Some(returned_id));
319
320        let bytes = ordering.get_txn(returned_id).unwrap().unwrap();
321        assert_eq!(TxId::for_transaction(&bytes), returned_id);
322        assert_eq!(&bytes[..12], GENESIS_PARENT.as_bytes());
323        assert_eq!(&bytes[20..52], &peering.public_key);
324        assert_eq!(&bytes[52..72], target.as_bytes());
325        assert_eq!(&bytes[72..bytes.len() - 64], b"object update");
326
327        let signature_bytes: &[u8; 64] = bytes[bytes.len() - 64..].try_into().unwrap();
328        let signature = Signature::from_bytes(signature_bytes);
329        peering
330            .signing_key
331            .verifying_key()
332            .verify_strict(&bytes[..bytes.len() - 64], &signature)
333            .unwrap();
334    }
335
336    #[test]
337    fn concurrent_submissions_form_one_linear_callback_sequence() {
338        let roots = TempRoots::new("concurrent");
339        let ordering = open_ordering(&roots);
340        let target = subsystem(b'b');
341        let handler = Arc::new(RecordingSubsystem::new());
342        ordering
343            .register_subsystem(target, None, handler.clone())
344            .unwrap();
345        let peering = Arc::new(K1Peering::open(&roots.peering(), ordering.clone()).unwrap());
346
347        let threads: Vec<_> = (0_u8..12)
348            .map(|value| {
349                let peering = peering.clone();
350                std::thread::spawn(move || peering.submit_txn(target, &[value]).unwrap())
351            })
352            .collect();
353
354        for thread in threads {
355            thread.join().unwrap();
356        }
357
358        let entries = handler.entries();
359        assert_eq!(entries.len(), 12);
360
361        let mut expected_parent = GENESIS_PARENT;
362        for (id, _) in &entries {
363            let bytes = ordering.get_txn(*id).unwrap().unwrap();
364            assert_eq!(&bytes[..12], expected_parent.as_bytes());
365            expected_parent = *id;
366        }
367
368        assert_eq!(ordering.tip(), Some(expected_parent));
369    }
370
371    #[test]
372    fn rejects_unregistered_target_before_mutation() {
373        let roots = TempRoots::new("unregistered");
374        let ordering = open_ordering(&roots);
375        let peering = K1Peering::open(&roots.peering(), ordering.clone()).unwrap();
376
377        assert!(
378            peering
379                .submit_txn(subsystem(b'c'), b"not committed")
380                .is_err()
381        );
382        assert_eq!(ordering.tip(), None);
383    }
384
385    #[test]
386    fn callback_failure_is_reported_as_committed_without_retry() {
387        let roots = TempRoots::new("callback-failure");
388        let ordering = open_ordering(&roots);
389        let target = subsystem(b'd');
390        let handler = Arc::new(RecordingSubsystem::new());
391        handler.fail.store(true, Ordering::Relaxed);
392        ordering
393            .register_subsystem(target, None, handler.clone())
394            .unwrap();
395        let peering = K1Peering::open(&roots.peering(), ordering.clone()).unwrap();
396
397        let error = peering.submit_txn(target, b"committed once").unwrap_err();
398
399        assert!(error.contains("committed"));
400        assert!(ordering.tip().is_some());
401        assert_eq!(handler.entries().len(), 1);
402    }
403
404    #[test]
405    fn restart_and_checkpoint_replay_deliver_only_newer_updates() {
406        let roots = TempRoots::new("replay");
407        let ordering = open_ordering(&roots);
408        let target = subsystem(b'e');
409        let first_handler = Arc::new(RecordingSubsystem::new());
410        ordering
411            .register_subsystem(target, None, first_handler.clone())
412            .unwrap();
413        let peering = K1Peering::open(&roots.peering(), ordering.clone()).unwrap();
414
415        peering.submit_txn(target, b"first").unwrap();
416        peering.submit_txn(target, b"second").unwrap();
417        let original = first_handler.entries();
418        assert_eq!(original.len(), 2);
419
420        drop(peering);
421        drop(first_handler);
422        drop(ordering);
423
424        let reopened = open_ordering(&roots);
425        let reopened_peering = K1Peering::open(&roots.peering(), reopened.clone()).unwrap();
426        let replay_handler = Arc::new(RecordingSubsystem::new());
427        reopened
428            .register_subsystem(target, Some(original[0].0), replay_handler.clone())
429            .unwrap();
430
431        assert_eq!(
432            replay_handler.entries(),
433            vec![(original[1].0, b"second".to_vec())]
434        );
435        reopened_peering.submit_txn(target, b"third").unwrap();
436        assert_eq!(replay_handler.entries().len(), 2);
437    }
438
439    #[test]
440    fn rejects_truncated_and_mismatched_identity_material() {
441        let truncated = TempRoots::new("truncated");
442        let ordering = open_ordering(&truncated);
443        write_identity(&truncated.peering().join("identity.key"), &[1_u8; 63]);
444        assert!(K1Peering::open(&truncated.peering(), ordering).is_err());
445
446        let mismatched = TempRoots::new("mismatched");
447        let ordering = open_ordering(&mismatched);
448        write_identity(
449            &mismatched.peering().join("identity.key"),
450            &[0_u8; IDENTITY_BYTES],
451        );
452        assert!(K1Peering::open(&mismatched.peering(), ordering).is_err());
453    }
454
455    #[cfg(unix)]
456    #[test]
457    fn rejects_over_permissive_and_symbolic_link_identity_files() {
458        use std::os::unix::fs::{PermissionsExt, symlink};
459
460        let permissive = TempRoots::new("permissive");
461        let ordering = open_ordering(&permissive);
462        let identity = permissive.peering().join("identity.key");
463        let signing_key = SigningKey::from_bytes(&[7_u8; PRIVATE_KEY_BYTES]);
464        let mut bytes = [0_u8; IDENTITY_BYTES];
465        bytes[..PRIVATE_KEY_BYTES].fill(7);
466        bytes[PRIVATE_KEY_BYTES..].copy_from_slice(&signing_key.verifying_key().to_bytes());
467        write_identity(&identity, &bytes);
468        fs::set_permissions(&identity, fs::Permissions::from_mode(0o644)).unwrap();
469        assert!(K1Peering::open(&permissive.peering(), ordering).is_err());
470
471        let linked = TempRoots::new("linked-identity");
472        let ordering = open_ordering(&linked);
473        fs::create_dir_all(linked.peering()).unwrap();
474        let target = linked.base.join("identity-target");
475        write_identity(&target, &bytes);
476        symlink(&target, linked.peering().join("identity.key")).unwrap();
477        assert!(K1Peering::open(&linked.peering(), ordering).is_err());
478    }
479
480    #[cfg(unix)]
481    #[test]
482    fn rejects_symbolic_link_peering_root() {
483        use std::os::unix::fs::symlink;
484
485        let roots = TempRoots::new("linked-root");
486        let ordering = open_ordering(&roots);
487        let actual = roots.base.join("actual-peering");
488        fs::create_dir_all(&actual).unwrap();
489        let linked = roots.base.join("linked-peering");
490        symlink(&actual, &linked).unwrap();
491
492        assert!(K1Peering::open(&linked, ordering).is_err());
493    }
494
495    #[test]
496    fn missing_identity_with_canonical_history_fails_closed() {
497        let roots = TempRoots::new("missing-with-history");
498        let ordering = open_ordering(&roots);
499        let target = subsystem(b'f');
500        ordering
501            .register_subsystem(target, None, Arc::new(RecordingSubsystem::new()))
502            .unwrap();
503        let peering = K1Peering::open(&roots.peering(), ordering.clone()).unwrap();
504        peering.submit_txn(target, b"durable").unwrap();
505
506        drop(peering);
507        drop(ordering);
508        fs::remove_file(roots.peering().join("identity.key")).unwrap();
509
510        let reopened = open_ordering(&roots);
511        assert!(K1Peering::open(&roots.peering(), reopened).is_err());
512        assert!(!roots.peering().join("identity.key").exists());
513    }
514}