Skip to main content

kcode_k1_peering/
lib.rs

1use ed25519_dalek::{Signer, SigningKey};
2use kcode_k1_txn_ordering::{K1TxnOrdering, SubsystemId};
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<(), 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(|_| ())
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, TxId};
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        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
317        let bytes = ordering.get_txn(entries[0].0).unwrap().unwrap();
318        assert_eq!(&bytes[..12], GENESIS_PARENT.as_bytes());
319        assert_eq!(&bytes[20..52], &peering.public_key);
320        assert_eq!(&bytes[52..72], target.as_bytes());
321        assert_eq!(&bytes[72..bytes.len() - 64], b"object update");
322
323        let signature_bytes: &[u8; 64] = bytes[bytes.len() - 64..].try_into().unwrap();
324        let signature = Signature::from_bytes(signature_bytes);
325        peering
326            .signing_key
327            .verifying_key()
328            .verify_strict(&bytes[..bytes.len() - 64], &signature)
329            .unwrap();
330    }
331
332    #[test]
333    fn concurrent_submissions_form_one_linear_callback_sequence() {
334        let roots = TempRoots::new("concurrent");
335        let ordering = open_ordering(&roots);
336        let target = subsystem(b'b');
337        let handler = Arc::new(RecordingSubsystem::new());
338        ordering
339            .register_subsystem(target, None, handler.clone())
340            .unwrap();
341        let peering = Arc::new(K1Peering::open(&roots.peering(), ordering.clone()).unwrap());
342
343        let threads: Vec<_> = (0_u8..12)
344            .map(|value| {
345                let peering = peering.clone();
346                std::thread::spawn(move || peering.submit_txn(target, &[value]).unwrap())
347            })
348            .collect();
349
350        for thread in threads {
351            thread.join().unwrap();
352        }
353
354        let entries = handler.entries();
355        assert_eq!(entries.len(), 12);
356
357        let mut expected_parent = GENESIS_PARENT;
358        for (id, _) in &entries {
359            let bytes = ordering.get_txn(*id).unwrap().unwrap();
360            assert_eq!(&bytes[..12], expected_parent.as_bytes());
361            expected_parent = *id;
362        }
363
364        assert_eq!(ordering.tip(), Some(expected_parent));
365    }
366
367    #[test]
368    fn rejects_unregistered_target_before_mutation() {
369        let roots = TempRoots::new("unregistered");
370        let ordering = open_ordering(&roots);
371        let peering = K1Peering::open(&roots.peering(), ordering.clone()).unwrap();
372
373        assert!(
374            peering
375                .submit_txn(subsystem(b'c'), b"not committed")
376                .is_err()
377        );
378        assert_eq!(ordering.tip(), None);
379    }
380
381    #[test]
382    fn callback_failure_is_reported_as_committed_without_retry() {
383        let roots = TempRoots::new("callback-failure");
384        let ordering = open_ordering(&roots);
385        let target = subsystem(b'd');
386        let handler = Arc::new(RecordingSubsystem::new());
387        handler.fail.store(true, Ordering::Relaxed);
388        ordering
389            .register_subsystem(target, None, handler.clone())
390            .unwrap();
391        let peering = K1Peering::open(&roots.peering(), ordering.clone()).unwrap();
392
393        let error = peering.submit_txn(target, b"committed once").unwrap_err();
394
395        assert!(error.contains("committed"));
396        assert!(ordering.tip().is_some());
397        assert_eq!(handler.entries().len(), 1);
398    }
399
400    #[test]
401    fn restart_and_checkpoint_replay_deliver_only_newer_updates() {
402        let roots = TempRoots::new("replay");
403        let ordering = open_ordering(&roots);
404        let target = subsystem(b'e');
405        let first_handler = Arc::new(RecordingSubsystem::new());
406        ordering
407            .register_subsystem(target, None, first_handler.clone())
408            .unwrap();
409        let peering = K1Peering::open(&roots.peering(), ordering.clone()).unwrap();
410
411        peering.submit_txn(target, b"first").unwrap();
412        peering.submit_txn(target, b"second").unwrap();
413        let original = first_handler.entries();
414        assert_eq!(original.len(), 2);
415
416        drop(peering);
417        drop(first_handler);
418        drop(ordering);
419
420        let reopened = open_ordering(&roots);
421        let reopened_peering = K1Peering::open(&roots.peering(), reopened.clone()).unwrap();
422        let replay_handler = Arc::new(RecordingSubsystem::new());
423        reopened
424            .register_subsystem(target, Some(original[0].0), replay_handler.clone())
425            .unwrap();
426
427        assert_eq!(
428            replay_handler.entries(),
429            vec![(original[1].0, b"second".to_vec())]
430        );
431        reopened_peering.submit_txn(target, b"third").unwrap();
432        assert_eq!(replay_handler.entries().len(), 2);
433    }
434
435    #[test]
436    fn rejects_truncated_and_mismatched_identity_material() {
437        let truncated = TempRoots::new("truncated");
438        let ordering = open_ordering(&truncated);
439        write_identity(&truncated.peering().join("identity.key"), &[1_u8; 63]);
440        assert!(K1Peering::open(&truncated.peering(), ordering).is_err());
441
442        let mismatched = TempRoots::new("mismatched");
443        let ordering = open_ordering(&mismatched);
444        write_identity(
445            &mismatched.peering().join("identity.key"),
446            &[0_u8; IDENTITY_BYTES],
447        );
448        assert!(K1Peering::open(&mismatched.peering(), ordering).is_err());
449    }
450
451    #[cfg(unix)]
452    #[test]
453    fn rejects_over_permissive_and_symbolic_link_identity_files() {
454        use std::os::unix::fs::{PermissionsExt, symlink};
455
456        let permissive = TempRoots::new("permissive");
457        let ordering = open_ordering(&permissive);
458        let identity = permissive.peering().join("identity.key");
459        let signing_key = SigningKey::from_bytes(&[7_u8; PRIVATE_KEY_BYTES]);
460        let mut bytes = [0_u8; IDENTITY_BYTES];
461        bytes[..PRIVATE_KEY_BYTES].fill(7);
462        bytes[PRIVATE_KEY_BYTES..].copy_from_slice(&signing_key.verifying_key().to_bytes());
463        write_identity(&identity, &bytes);
464        fs::set_permissions(&identity, fs::Permissions::from_mode(0o644)).unwrap();
465        assert!(K1Peering::open(&permissive.peering(), ordering).is_err());
466
467        let linked = TempRoots::new("linked-identity");
468        let ordering = open_ordering(&linked);
469        fs::create_dir_all(linked.peering()).unwrap();
470        let target = linked.base.join("identity-target");
471        write_identity(&target, &bytes);
472        symlink(&target, linked.peering().join("identity.key")).unwrap();
473        assert!(K1Peering::open(&linked.peering(), ordering).is_err());
474    }
475
476    #[cfg(unix)]
477    #[test]
478    fn rejects_symbolic_link_peering_root() {
479        use std::os::unix::fs::symlink;
480
481        let roots = TempRoots::new("linked-root");
482        let ordering = open_ordering(&roots);
483        let actual = roots.base.join("actual-peering");
484        fs::create_dir_all(&actual).unwrap();
485        let linked = roots.base.join("linked-peering");
486        symlink(&actual, &linked).unwrap();
487
488        assert!(K1Peering::open(&linked, ordering).is_err());
489    }
490
491    #[test]
492    fn missing_identity_with_canonical_history_fails_closed() {
493        let roots = TempRoots::new("missing-with-history");
494        let ordering = open_ordering(&roots);
495        let target = subsystem(b'f');
496        ordering
497            .register_subsystem(target, None, Arc::new(RecordingSubsystem::new()))
498            .unwrap();
499        let peering = K1Peering::open(&roots.peering(), ordering.clone()).unwrap();
500        peering.submit_txn(target, b"durable").unwrap();
501
502        drop(peering);
503        drop(ordering);
504        fs::remove_file(roots.peering().join("identity.key")).unwrap();
505
506        let reopened = open_ordering(&roots);
507        assert!(K1Peering::open(&roots.peering(), reopened).is_err());
508        assert!(!roots.peering().join("identity.key").exists());
509    }
510}