kcode-k1-launch-nodes 0.4.0

KTO-authoritative launch-node bindings with a durable startup projection
Documentation
use std::{
    collections::HashMap,
    fs::{File, OpenOptions},
    io::{Read, Seek, SeekFrom, Write},
    path::Path,
    sync::{Arc, Mutex},
};

pub use kcode_k1_launch_node_codec::{
    Authority, GroupId, NodeId, TargetId, TargetName, TxId, UserId,
};
use kcode_k1_launch_node_codec::{
    PROJECTION_HEADER, ProjectionDecode, ProjectionRecord, SetAction, decode_projection_record,
};
use kcode_k1_peering::K1Peering;
use kcode_k1_txn_ordering::{K1TxnOrdering, Subsystem as OrderingSubsystem, SubsystemId};

struct Pending {
    payload: Vec<u8>,
    evidence: Option<TxId>,
}

struct Inner {
    file: File,
    bindings: HashMap<TargetId, NodeId>,
    checkpoint: Option<TxId>,
    pending: HashMap<[u8; 16], Pending>,
    failed: bool,
}

struct Shared {
    inner: Mutex<Inner>,
}

struct Driver {
    shared: Arc<Shared>,
}

pub struct LaunchNodes {
    shared: Arc<Shared>,
    peering: Arc<K1Peering>,
}

impl LaunchNodes {
    pub fn open(
        path: &Path,
        ordering: Arc<K1TxnOrdering>,
        peering: Arc<K1Peering>,
    ) -> Result<Self, String> {
        let mut projection = open_projection(path)?;
        if let Some(checkpoint) = projection.checkpoint.to_owned() {
            match ordering.get_txn(checkpoint) {
                Ok(Some(_)) => {}
                Ok(None) => {
                    reset_projection(&mut projection.file, path)?;
                    projection.bindings.clear();
                    projection.checkpoint = None;
                }
                Err(error) => return Err(format!("validate launch nodes checkpoint: {error}")),
            }
        }
        let after = projection.checkpoint.to_owned();
        let shared = Arc::new(Shared {
            inner: Mutex::new(Inner {
                file: projection.file,
                bindings: projection.bindings,
                checkpoint: projection.checkpoint,
                pending: HashMap::new(),
                failed: false,
            }),
        });
        ordering.register_subsystem(
            subsystem()?,
            after,
            Arc::new(Driver {
                shared: shared.clone(),
            }),
        )?;
        Ok(Self { shared, peering })
    }

    pub fn get(&self, target: &TargetId) -> Result<Option<NodeId>, String> {
        let inner = self
            .shared
            .inner
            .lock()
            .map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
        if inner.failed {
            return Err("k1-launch-nodes facade is unavailable".into());
        }
        Ok(inner.bindings.get(target).cloned())
    }

    pub fn set(&self, target: TargetId, node: NodeId) -> Result<(), String> {
        let (operation_id, payload) = loop {
            let mut operation_id = [0_u8; 16];
            getrandom::fill(&mut operation_id)
                .map_err(|error| format!("generate launch nodes operation ID: {error}"))?;
            let payload = SetAction::new(operation_id, target.clone(), node).encode();
            let mut inner = self
                .shared
                .inner
                .lock()
                .map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
            if inner.failed {
                return Err("k1-launch-nodes facade is unavailable".into());
            }
            if inner.pending.contains_key(&operation_id) {
                continue;
            }
            inner.pending.insert(
                operation_id,
                Pending {
                    payload: payload.clone(),
                    evidence: None,
                },
            );
            break (operation_id, payload);
        };
        let submission = self.peering.submit_txn(subsystem()?, &payload);
        let mut inner = self
            .shared
            .inner
            .lock()
            .map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
        if inner.failed {
            return Err("k1-launch-nodes facade is unavailable".into());
        }
        let Some(pending) = inner.pending.remove(&operation_id) else {
            return Err(fault(
                &mut inner,
                "k1-launch-nodes callback evidence is missing",
            ));
        };
        match (submission, pending.evidence) {
            (Ok(submitted), Some(callback)) if submitted == callback => Ok(()),
            (Err(_), Some(_)) => Ok(()),
            (Err(error), None) => Err(format!("submit k1-launch-nodes transaction: {error}")),
            (Ok(_), None) => Err(fault(
                &mut inner,
                "k1-launch-nodes callback evidence is missing",
            )),
            (Ok(_), Some(_)) => Err(fault(
                &mut inner,
                "k1-launch-nodes callback transaction does not match submission",
            )),
        }
    }
}

impl OrderingSubsystem for Driver {
    fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String> {
        let action = match SetAction::decode(payload) {
            Ok(action) => action,
            Err(error) => {
                return self
                    .shared
                    .fail(format!("decode k1-launch-nodes Set action: {error}"));
            }
        };
        let operation_id = action.operation_id().to_owned();
        let canonical_payload = action.encode();
        let record =
            ProjectionRecord::new(id, action.target().to_owned(), action.node().to_owned())
                .encode();
        let mut inner = self
            .shared
            .inner
            .lock()
            .map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
        if inner.failed {
            return Err("k1-launch-nodes facade is unavailable".into());
        }
        if inner.checkpoint == Some(id) {
            return Err(fault(
                &mut inner,
                "k1-launch-nodes callback transaction is duplicated",
            ));
        }
        if let Some(pending) = inner.pending.get(&operation_id) {
            if pending.payload != canonical_payload {
                return Err(fault(
                    &mut inner,
                    "k1-launch-nodes callback action does not match reservation",
                ));
            }
            if pending.evidence.is_some() {
                return Err(fault(
                    &mut inner,
                    "k1-launch-nodes callback evidence is duplicated",
                ));
            }
        }
        let append = inner
            .file
            .seek(SeekFrom::End(0))
            .and_then(|_| inner.file.write_all(&record))
            .and_then(|()| inner.file.sync_data());
        if let Err(error) = append {
            return Err(fault(
                &mut inner,
                format!("append k1-launch-nodes projection record: {error}"),
            ));
        }
        inner
            .bindings
            .insert(action.target().to_owned(), action.node().to_owned());
        inner.checkpoint = Some(id);
        if let Some(pending) = inner.pending.get_mut(&operation_id) {
            pending.evidence = Some(id);
        }
        Ok(())
    }

    fn reorg(&self) -> Result<(), String> {
        let mut inner = self
            .shared
            .inner
            .lock()
            .map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
        inner.failed = true;
        inner.pending.clear();
        Ok(())
    }
}

impl Shared {
    fn fail(&self, message: String) -> Result<(), String> {
        let mut inner = self
            .inner
            .lock()
            .map_err(|_| "k1-launch-nodes lock poisoned".to_owned())?;
        Err(fault(&mut inner, message))
    }
}

struct Projection {
    file: File,
    bindings: HashMap<TargetId, NodeId>,
    checkpoint: Option<TxId>,
}

fn open_projection(path: &Path) -> Result<Projection, String> {
    let mut file = OpenOptions::new()
        .create(true)
        .truncate(false)
        .read(true)
        .write(true)
        .open(path)
        .map_err(|error| format!("open launch nodes projection: {error}"))?;
    if file
        .metadata()
        .map_err(|error| format!("inspect launch nodes projection: {error}"))?
        .len()
        == 0
    {
        reset_projection(&mut file, path)?;
    }
    file.seek(SeekFrom::Start(0))
        .map_err(|error| format!("seek launch nodes projection: {error}"))?;
    let mut bytes = Vec::new();
    file.read_to_end(&mut bytes)
        .map_err(|error| format!("read launch nodes projection: {error}"))?;
    if bytes.len() < PROJECTION_HEADER.len()
        || &bytes[..PROJECTION_HEADER.len()] != PROJECTION_HEADER
    {
        return Err("unsupported launch nodes projection format; no migration is available".into());
    }
    let mut bindings = HashMap::new();
    let mut checkpoint = None;
    let mut offset = PROJECTION_HEADER.len();
    let mut corrupt = false;
    while offset < bytes.len() {
        match decode_projection_record(&bytes[offset..]) {
            Ok(ProjectionDecode::Complete { record, consumed }) => {
                bindings.insert(record.target().to_owned(), record.node().to_owned());
                let callback: [u8; 12] = bytes[offset..offset + 12]
                    .try_into()
                    .map_err(|_| "launch nodes projection callback is unavailable".to_owned())?;
                checkpoint = Some(TxId::from_bytes(callback));
                offset += consumed;
            }
            Ok(ProjectionDecode::Incomplete) => {
                repair_tail(&mut file, offset)?;
                break;
            }
            Err(_) => {
                corrupt = true;
                break;
            }
        }
    }
    if corrupt {
        reset_projection(&mut file, path)?;
        bindings.clear();
        checkpoint = None;
    }
    file.seek(SeekFrom::End(0))
        .map_err(|error| format!("seek launch nodes append position: {error}"))?;
    Ok(Projection {
        file,
        bindings,
        checkpoint,
    })
}

fn reset_projection(file: &mut File, path: &Path) -> Result<(), String> {
    file.set_len(0)
        .and_then(|()| file.seek(SeekFrom::Start(0)).map(|_| ()))
        .and_then(|()| file.write_all(PROJECTION_HEADER))
        .and_then(|()| file.sync_data())
        .map_err(|error| format!("reset launch nodes projection: {error}"))?;
    sync_parent(path)
}

fn repair_tail(file: &mut File, offset: usize) -> Result<(), String> {
    file.set_len(offset as u64)
        .and_then(|()| file.sync_data())
        .map_err(|error| format!("repair launch nodes projection tail: {error}"))
}

fn sync_parent(path: &Path) -> Result<(), String> {
    let parent = path
        .parent()
        .filter(|parent| !parent.as_os_str().is_empty())
        .unwrap_or_else(|| Path::new("."));
    File::open(parent)
        .and_then(|directory| directory.sync_all())
        .map_err(|error| format!("synchronize launch nodes projection directory: {error}"))
}

fn subsystem() -> Result<SubsystemId, String> {
    SubsystemId::from_str("k1-launch-nodes")
}

fn fault(inner: &mut Inner, message: impl Into<String>) -> String {
    inner.failed = true;
    inner.pending.clear();
    message.into()
}

#[cfg(test)]
mod tests;