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}