Skip to main content

kcode_k1_peering/
lib.rs

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