use ed25519_dalek::{Signer, SigningKey};
use kcode_k1_txn_ordering::{K1TxnOrdering, SubsystemId};
use std::fs::{self, File, OpenOptions};
use std::io::{Read, Write};
use std::path::Path;
use std::sync::Arc;
use std::time::{SystemTime, UNIX_EPOCH};
const IDENTITY_BYTES: usize = 64;
const PRIVATE_KEY_BYTES: usize = 32;
const PUBLIC_KEY_BYTES: usize = 32;
pub struct K1Peering {
ordering: Arc<K1TxnOrdering>,
signing_key: SigningKey,
public_key: [u8; PUBLIC_KEY_BYTES],
}
impl K1Peering {
pub fn open(root: &Path, ordering: Arc<K1TxnOrdering>) -> Result<Self, String> {
prepare_root(root)?;
let identity_path = root.join("identity.key");
let (signing_key, public_key) = match fs::symlink_metadata(&identity_path) {
Ok(metadata) => load_identity(&identity_path, metadata)?,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
if ordering.tip().is_some() {
return Err(
"identity.key is missing while canonical transaction history exists"
.to_owned(),
);
}
create_identity(&identity_path)?
}
Err(error) => return Err(format!("cannot inspect identity.key: {error}")),
};
Ok(Self {
ordering,
signing_key,
public_key,
})
}
pub fn submit_txn(&self, subsystem: SubsystemId, payload: &[u8]) -> Result<(), String> {
let timestamp = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_err(|_| "system clock is before the Unix epoch".to_owned())?
.as_secs();
self.ordering
.submit_local_txn(
timestamp,
self.public_key,
subsystem,
payload,
|bytes| Ok(self.signing_key.sign(bytes).to_bytes()),
|_| Ok(()),
)
.map(|_| ())
}
}
fn prepare_root(root: &Path) -> Result<(), String> {
match fs::symlink_metadata(root) {
Ok(metadata) if metadata.file_type().is_dir() => return Ok(()),
Ok(_) => return Err("peering root is not a real directory".to_owned()),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
Err(error) => return Err(format!("cannot inspect peering root: {error}")),
}
fs::create_dir_all(root).map_err(|error| format!("cannot create peering root: {error}"))?;
let metadata = fs::symlink_metadata(root)
.map_err(|error| format!("cannot inspect created peering root: {error}"))?;
if !metadata.file_type().is_dir() {
return Err("created peering root is not a real directory".to_owned());
}
Ok(())
}
fn load_identity(
path: &Path,
metadata: fs::Metadata,
) -> Result<(SigningKey, [u8; PUBLIC_KEY_BYTES]), String> {
if !metadata.file_type().is_file() {
return Err("identity.key is not a real regular file".to_owned());
}
if metadata.len() != IDENTITY_BYTES as u64 {
return Err("identity.key must contain exactly 64 bytes".to_owned());
}
validate_identity_permissions(&metadata)?;
let mut file =
File::open(path).map_err(|error| format!("cannot open identity.key: {error}"))?;
let mut bytes = [0_u8; IDENTITY_BYTES];
file.read_exact(&mut bytes)
.map_err(|error| format!("cannot read identity.key: {error}"))?;
let private_seed: [u8; PRIVATE_KEY_BYTES] = bytes[..PRIVATE_KEY_BYTES]
.try_into()
.expect("fixed private-key range");
let stored_public_key: [u8; PUBLIC_KEY_BYTES] = bytes[PRIVATE_KEY_BYTES..]
.try_into()
.expect("fixed public-key range");
let signing_key = SigningKey::from_bytes(&private_seed);
let derived_public_key = signing_key.verifying_key().to_bytes();
if stored_public_key != derived_public_key {
return Err("identity.key public key does not match its private seed".to_owned());
}
Ok((signing_key, derived_public_key))
}
fn create_identity(path: &Path) -> Result<(SigningKey, [u8; PUBLIC_KEY_BYTES]), String> {
let mut private_seed = [0_u8; PRIVATE_KEY_BYTES];
getrandom::fill(&mut private_seed)
.map_err(|error| format!("cannot generate peering identity: {error}"))?;
let signing_key = SigningKey::from_bytes(&private_seed);
let public_key = signing_key.verifying_key().to_bytes();
let mut bytes = [0_u8; IDENTITY_BYTES];
bytes[..PRIVATE_KEY_BYTES].copy_from_slice(&private_seed);
bytes[PRIVATE_KEY_BYTES..].copy_from_slice(&public_key);
let mut options = OpenOptions::new();
options.write(true).create_new(true);
set_owner_only_creation_mode(&mut options);
let mut file = options
.open(path)
.map_err(|error| format!("cannot create identity.key: {error}"))?;
file.write_all(&bytes)
.map_err(|error| format!("cannot write identity.key: {error}"))?;
file.sync_all()
.map_err(|error| format!("cannot synchronize identity.key: {error}"))?;
Ok((signing_key, public_key))
}
#[cfg(unix)]
fn set_owner_only_creation_mode(options: &mut OpenOptions) {
use std::os::unix::fs::OpenOptionsExt;
options.mode(0o600);
}
#[cfg(not(unix))]
fn set_owner_only_creation_mode(_: &mut OpenOptions) {}
#[cfg(unix)]
fn validate_identity_permissions(metadata: &fs::Metadata) -> Result<(), String> {
use std::os::unix::fs::PermissionsExt;
if metadata.permissions().mode() & 0o077 != 0 {
return Err("identity.key grants group or other permissions".to_owned());
}
Ok(())
}
#[cfg(not(unix))]
fn validate_identity_permissions(_: &fs::Metadata) -> Result<(), String> {
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use ed25519_dalek::Signature;
use kcode_k1_txn_ordering::{GENESIS_PARENT, Subsystem, TxId};
use std::path::PathBuf;
use std::sync::Mutex;
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
static NEXT_ROOT: AtomicU64 = AtomicU64::new(0);
struct TempRoots {
base: PathBuf,
}
impl TempRoots {
fn new(label: &str) -> Self {
let sequence = NEXT_ROOT.fetch_add(1, Ordering::Relaxed);
let base = std::env::temp_dir().join(format!(
"kcode-k1-peering-{}-{}-{}",
std::process::id(),
sequence,
label
));
let _ = fs::remove_dir_all(&base);
Self { base }
}
fn peering(&self) -> PathBuf {
self.base.join("peering")
}
fn ordering(&self) -> PathBuf {
self.base.join("ordering")
}
}
impl Drop for TempRoots {
fn drop(&mut self) {
let _ = fs::remove_dir_all(&self.base);
}
}
struct RecordingSubsystem {
submissions: Mutex<Vec<(TxId, Vec<u8>)>>,
fail: AtomicBool,
reorgs: AtomicUsize,
}
impl RecordingSubsystem {
fn new() -> Self {
Self {
submissions: Mutex::new(Vec::new()),
fail: AtomicBool::new(false),
reorgs: AtomicUsize::new(0),
}
}
fn entries(&self) -> Vec<(TxId, Vec<u8>)> {
self.submissions.lock().unwrap().clone()
}
}
impl Subsystem for RecordingSubsystem {
fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String> {
self.submissions
.lock()
.unwrap()
.push((id, payload.to_vec()));
if self.fail.load(Ordering::Relaxed) {
Err("integration failed".to_owned())
} else {
Ok(())
}
}
fn reorg(&self) -> Result<(), String> {
self.reorgs.fetch_add(1, Ordering::Relaxed);
Ok(())
}
}
fn subsystem(value: u8) -> SubsystemId {
SubsystemId::from_bytes([value; 20]).unwrap()
}
fn open_ordering(roots: &TempRoots) -> Arc<K1TxnOrdering> {
Arc::new(K1TxnOrdering::open(&roots.ordering()).unwrap())
}
fn write_identity(path: &Path, bytes: &[u8]) {
fs::create_dir_all(path.parent().unwrap()).unwrap();
let mut options = OpenOptions::new();
options.write(true).create_new(true);
set_owner_only_creation_mode(&mut options);
let mut file = options.open(path).unwrap();
file.write_all(bytes).unwrap();
file.sync_all().unwrap();
}
#[test]
fn creates_owner_only_identity_and_reopens_stably() {
let roots = TempRoots::new("identity");
let ordering = open_ordering(&roots);
let first = K1Peering::open(&roots.peering(), ordering.clone()).unwrap();
let first_bytes = fs::read(roots.peering().join("identity.key")).unwrap();
assert_eq!(first_bytes.len(), IDENTITY_BYTES);
assert_eq!(&first_bytes[PRIVATE_KEY_BYTES..], &first.public_key);
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let mode = fs::metadata(roots.peering().join("identity.key"))
.unwrap()
.permissions()
.mode();
assert_eq!(mode & 0o077, 0);
}
drop(first);
let second = K1Peering::open(&roots.peering(), ordering).unwrap();
let second_bytes = fs::read(roots.peering().join("identity.key")).unwrap();
assert_eq!(second_bytes, first_bytes);
assert_eq!(second.public_key, first_bytes[PRIVATE_KEY_BYTES..]);
}
#[test]
fn submits_exact_signed_transaction_and_delivers_once() {
let roots = TempRoots::new("signed");
let ordering = open_ordering(&roots);
let target = subsystem(b'a');
let handler = Arc::new(RecordingSubsystem::new());
ordering
.register_subsystem(target, None, handler.clone())
.unwrap();
let peering = K1Peering::open(&roots.peering(), ordering.clone()).unwrap();
peering.submit_txn(target, b"object update").unwrap();
let entries = handler.entries();
assert_eq!(entries.len(), 1);
assert_eq!(entries[0].1, b"object update");
let bytes = ordering.get_txn(entries[0].0).unwrap().unwrap();
assert_eq!(&bytes[..12], GENESIS_PARENT.as_bytes());
assert_eq!(&bytes[20..52], &peering.public_key);
assert_eq!(&bytes[52..72], target.as_bytes());
assert_eq!(&bytes[72..bytes.len() - 64], b"object update");
let signature_bytes: &[u8; 64] = bytes[bytes.len() - 64..].try_into().unwrap();
let signature = Signature::from_bytes(signature_bytes);
peering
.signing_key
.verifying_key()
.verify_strict(&bytes[..bytes.len() - 64], &signature)
.unwrap();
}
#[test]
fn concurrent_submissions_form_one_linear_callback_sequence() {
let roots = TempRoots::new("concurrent");
let ordering = open_ordering(&roots);
let target = subsystem(b'b');
let handler = Arc::new(RecordingSubsystem::new());
ordering
.register_subsystem(target, None, handler.clone())
.unwrap();
let peering = Arc::new(K1Peering::open(&roots.peering(), ordering.clone()).unwrap());
let threads: Vec<_> = (0_u8..12)
.map(|value| {
let peering = peering.clone();
std::thread::spawn(move || peering.submit_txn(target, &[value]).unwrap())
})
.collect();
for thread in threads {
thread.join().unwrap();
}
let entries = handler.entries();
assert_eq!(entries.len(), 12);
let mut expected_parent = GENESIS_PARENT;
for (id, _) in &entries {
let bytes = ordering.get_txn(*id).unwrap().unwrap();
assert_eq!(&bytes[..12], expected_parent.as_bytes());
expected_parent = *id;
}
assert_eq!(ordering.tip(), Some(expected_parent));
}
#[test]
fn rejects_unregistered_target_before_mutation() {
let roots = TempRoots::new("unregistered");
let ordering = open_ordering(&roots);
let peering = K1Peering::open(&roots.peering(), ordering.clone()).unwrap();
assert!(
peering
.submit_txn(subsystem(b'c'), b"not committed")
.is_err()
);
assert_eq!(ordering.tip(), None);
}
#[test]
fn callback_failure_is_reported_as_committed_without_retry() {
let roots = TempRoots::new("callback-failure");
let ordering = open_ordering(&roots);
let target = subsystem(b'd');
let handler = Arc::new(RecordingSubsystem::new());
handler.fail.store(true, Ordering::Relaxed);
ordering
.register_subsystem(target, None, handler.clone())
.unwrap();
let peering = K1Peering::open(&roots.peering(), ordering.clone()).unwrap();
let error = peering.submit_txn(target, b"committed once").unwrap_err();
assert!(error.contains("committed"));
assert!(ordering.tip().is_some());
assert_eq!(handler.entries().len(), 1);
}
#[test]
fn restart_and_checkpoint_replay_deliver_only_newer_updates() {
let roots = TempRoots::new("replay");
let ordering = open_ordering(&roots);
let target = subsystem(b'e');
let first_handler = Arc::new(RecordingSubsystem::new());
ordering
.register_subsystem(target, None, first_handler.clone())
.unwrap();
let peering = K1Peering::open(&roots.peering(), ordering.clone()).unwrap();
peering.submit_txn(target, b"first").unwrap();
peering.submit_txn(target, b"second").unwrap();
let original = first_handler.entries();
assert_eq!(original.len(), 2);
drop(peering);
drop(first_handler);
drop(ordering);
let reopened = open_ordering(&roots);
let reopened_peering = K1Peering::open(&roots.peering(), reopened.clone()).unwrap();
let replay_handler = Arc::new(RecordingSubsystem::new());
reopened
.register_subsystem(target, Some(original[0].0), replay_handler.clone())
.unwrap();
assert_eq!(
replay_handler.entries(),
vec![(original[1].0, b"second".to_vec())]
);
reopened_peering.submit_txn(target, b"third").unwrap();
assert_eq!(replay_handler.entries().len(), 2);
}
#[test]
fn rejects_truncated_and_mismatched_identity_material() {
let truncated = TempRoots::new("truncated");
let ordering = open_ordering(&truncated);
write_identity(&truncated.peering().join("identity.key"), &[1_u8; 63]);
assert!(K1Peering::open(&truncated.peering(), ordering).is_err());
let mismatched = TempRoots::new("mismatched");
let ordering = open_ordering(&mismatched);
write_identity(
&mismatched.peering().join("identity.key"),
&[0_u8; IDENTITY_BYTES],
);
assert!(K1Peering::open(&mismatched.peering(), ordering).is_err());
}
#[cfg(unix)]
#[test]
fn rejects_over_permissive_and_symbolic_link_identity_files() {
use std::os::unix::fs::{PermissionsExt, symlink};
let permissive = TempRoots::new("permissive");
let ordering = open_ordering(&permissive);
let identity = permissive.peering().join("identity.key");
let signing_key = SigningKey::from_bytes(&[7_u8; PRIVATE_KEY_BYTES]);
let mut bytes = [0_u8; IDENTITY_BYTES];
bytes[..PRIVATE_KEY_BYTES].fill(7);
bytes[PRIVATE_KEY_BYTES..].copy_from_slice(&signing_key.verifying_key().to_bytes());
write_identity(&identity, &bytes);
fs::set_permissions(&identity, fs::Permissions::from_mode(0o644)).unwrap();
assert!(K1Peering::open(&permissive.peering(), ordering).is_err());
let linked = TempRoots::new("linked-identity");
let ordering = open_ordering(&linked);
fs::create_dir_all(linked.peering()).unwrap();
let target = linked.base.join("identity-target");
write_identity(&target, &bytes);
symlink(&target, linked.peering().join("identity.key")).unwrap();
assert!(K1Peering::open(&linked.peering(), ordering).is_err());
}
#[cfg(unix)]
#[test]
fn rejects_symbolic_link_peering_root() {
use std::os::unix::fs::symlink;
let roots = TempRoots::new("linked-root");
let ordering = open_ordering(&roots);
let actual = roots.base.join("actual-peering");
fs::create_dir_all(&actual).unwrap();
let linked = roots.base.join("linked-peering");
symlink(&actual, &linked).unwrap();
assert!(K1Peering::open(&linked, ordering).is_err());
}
#[test]
fn missing_identity_with_canonical_history_fails_closed() {
let roots = TempRoots::new("missing-with-history");
let ordering = open_ordering(&roots);
let target = subsystem(b'f');
ordering
.register_subsystem(target, None, Arc::new(RecordingSubsystem::new()))
.unwrap();
let peering = K1Peering::open(&roots.peering(), ordering.clone()).unwrap();
peering.submit_txn(target, b"durable").unwrap();
drop(peering);
drop(ordering);
fs::remove_file(roots.peering().join("identity.key")).unwrap();
let reopened = open_ordering(&roots);
assert!(K1Peering::open(&roots.peering(), reopened).is_err());
assert!(!roots.peering().join("identity.key").exists());
}
}