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}