kcode-k1-launch-nodes 0.4.0

KTO-authoritative launch-node bindings with a durable startup projection
Documentation
use super::*;
use std::{
    fs,
    path::{Path, PathBuf},
    sync::atomic::{AtomicU64, Ordering},
};

static NEXT: AtomicU64 = AtomicU64::new(0);

struct Temp {
    root: PathBuf,
}

impl Temp {
    fn new() -> Self {
        let root = std::env::temp_dir().join(format!(
            "k1-launch-nodes-{}-{}",
            std::process::id(),
            NEXT.fetch_add(1, Ordering::Relaxed)
        ));
        fs::create_dir(&root).unwrap();
        Self { root }
    }
}

impl Drop for Temp {
    fn drop(&mut self) {
        fs::remove_dir_all(&self.root).unwrap();
    }
}

fn services(root: &Path) -> (Arc<K1TxnOrdering>, Arc<K1Peering>) {
    fs::create_dir_all(root).unwrap();
    let ordering = Arc::new(K1TxnOrdering::open(&root.join("ordering")).unwrap());
    let peering = Arc::new(K1Peering::open(&root.join("peering"), ordering.clone()).unwrap());
    (ordering, peering)
}

fn target(authority: Authority, name: &str) -> TargetId {
    TargetId::new(authority, TargetName::new(name.to_owned()).unwrap())
}

fn user(byte: u8) -> Authority {
    Authority::User(UserId::from_tx_id(TxId::from_bytes([byte; 12])))
}

fn group(byte: u8) -> Authority {
    Authority::Group(GroupId::new(TxId::from_bytes([byte; 12])))
}

fn record_count(path: &Path) -> usize {
    let bytes = fs::read(path).unwrap();
    let mut offset = PROJECTION_HEADER.len();
    let mut count = 0;
    while offset < bytes.len() {
        match decode_projection_record(&bytes[offset..]).unwrap() {
            ProjectionDecode::Complete { consumed, .. } => {
                count += 1;
                offset += consumed;
            }
            ProjectionDecode::Incomplete => panic!("unexpected incomplete record"),
        }
    }
    count
}

#[test]
fn callbacks_preserve_authorities_and_last_canonical_write() {
    let temp = Temp::new();
    let (ordering, peering) = services(&temp.root.join("service"));
    let store = LaunchNodes::open(&temp.root.join("projection"), ordering, peering).unwrap();
    let user_target = target(user(1), "default/chat");
    let group_target = target(group(1), "default/chat");
    store.set(user_target.clone(), NodeId([1; 12])).unwrap();
    store.set(group_target.clone(), NodeId([2; 12])).unwrap();
    store.set(user_target.clone(), NodeId([3; 12])).unwrap();
    assert_eq!(store.get(&user_target).unwrap(), Some(NodeId([3; 12])));
    assert_eq!(store.get(&group_target).unwrap(), Some(NodeId([2; 12])));
}

#[test]
fn peering_error_without_callback_does_not_bind() {
    let temp = Temp::new();
    let ordering = Arc::new(K1TxnOrdering::open(&temp.root.join("ordering-a")).unwrap());
    let other = Arc::new(K1TxnOrdering::open(&temp.root.join("ordering-b")).unwrap());
    let peering = Arc::new(K1Peering::open(&temp.root.join("peering-b"), other).unwrap());
    let store = LaunchNodes::open(&temp.root.join("projection"), ordering, peering).unwrap();
    let key = target(user(2), "model/x");
    assert!(store.set(key.clone(), NodeId([4; 12])).is_err());
    assert_eq!(store.get(&key).unwrap(), None);
    assert_eq!(record_count(&temp.root.join("projection")), 0);
}

#[test]
fn projection_repairs_rebuilds_and_replays_same_value_sets() {
    let temp = Temp::new();
    let service = temp.root.join("service");
    let projection = temp.root.join("projection");
    let key = target(group(3), "same/value");
    let node = NodeId([5; 12]);
    let (ordering, peering) = services(&service);
    let store = LaunchNodes::open(&projection, ordering.clone(), peering.clone()).unwrap();
    store.set(key.clone(), node).unwrap();
    store.set(key.clone(), node).unwrap();
    drop((store, peering, ordering));

    let clean = fs::read(&projection).unwrap();
    let mut incomplete = clean.clone();
    incomplete.extend_from_slice(&[1, 2, 3]);
    fs::write(&projection, incomplete).unwrap();
    let (ordering, peering) = services(&service);
    let store = LaunchNodes::open(&projection, ordering.clone(), peering.clone()).unwrap();
    assert_eq!(store.get(&key).unwrap(), Some(node));
    assert_eq!(fs::read(&projection).unwrap(), clean);
    drop((store, peering, ordering));

    let mut corrupt = fs::read(&projection).unwrap();
    let last = corrupt.last_mut().unwrap();
    *last ^= 0xff;
    fs::write(&projection, corrupt).unwrap();
    let (ordering, peering) = services(&service);
    let store = LaunchNodes::open(&projection, ordering, peering).unwrap();
    assert_eq!(store.get(&key).unwrap(), Some(node));
    assert_eq!(record_count(&projection), 2);
}

#[test]
fn non_v3_projection_is_rejected_without_migration() {
    let temp = Temp::new();
    let projection = temp.root.join("projection");
    fs::write(&projection, b"K1LNV2\0\0").unwrap();
    let (ordering, peering) = services(&temp.root.join("service"));
    let error = LaunchNodes::open(&projection, ordering, peering)
        .err()
        .unwrap();
    assert_eq!(
        error,
        "unsupported launch nodes projection format; no migration is available"
    );
    assert_eq!(fs::read(projection).unwrap(), b"K1LNV2\0\0");
}

#[test]
fn malformed_callback_faults_facade() {
    let temp = Temp::new();
    let (ordering, peering) = services(&temp.root.join("malformed"));
    let store =
        LaunchNodes::open(&temp.root.join("malformed-projection"), ordering, peering).unwrap();
    let driver = Driver {
        shared: store.shared.clone(),
    };
    assert!(
        driver
            .submit_txn(TxId::from_bytes([8; 12]), b"malformed")
            .is_err()
    );
    assert!(store.get(&target(user(8), "faulted")).is_err());
}

#[test]
fn complete_package_stays_below_the_managed_limit() {
    let files = [
        include_str!("../Cargo.toml"),
        include_str!("../Documentation.md"),
        include_str!("lib.rs"),
        include_str!("tests.rs"),
    ];
    let count = files
        .iter()
        .flat_map(|file| file.lines())
        .filter(|line| !line.trim().is_empty())
        .count();
    assert!(count < 500, "complete package has {count} nonblank lines");
}