kcode-k1-web-code-workspace 0.1.0

KTO-backed unpublished K1 Web code workspaces
Documentation
#![forbid(unsafe_code)]

use getrandom::fill;
use kcode_k1_peering::K1Peering;
pub use kcode_k1_transaction_id::TxId;
use kcode_k1_txn_ordering::{K1TxnOrdering, Subsystem, SubsystemId, TxId as KtoTxId};
pub use kcode_k1_web_code_document::CodeDocument;
pub use kcode_k1_web_code_workspace_format::Workspace;
use kcode_k1_web_code_workspace_format::{Action, OperationId, Projection};
pub use kcode_k1_web_package::WebFamily;
use std::{
    collections::{HashMap, hash_map::Entry},
    sync::{Arc, Mutex, MutexGuard},
    time::{Duration, Instant},
};

type Evidence = (KtoTxId, Result<Workspace, String>);

struct Pending {
    action: Action,
    evidence: Option<Evidence>,
}

struct State {
    projection: Projection,
    pending: HashMap<OperationId, Pending>,
    fault: Option<String>,
}

impl State {
    fn ready(&self) -> Result<(), String> {
        self.fault.clone().map_or(Ok(()), Err)
    }

    fn poison(&mut self, message: impl Into<String>) -> String {
        self.fault.get_or_insert_with(|| message.into()).clone()
    }

    fn record(&mut self, action: &Action, evidence: Evidence) -> Result<(), String> {
        let operation = action.operation_id();
        let problem = match self.pending.get_mut(&operation) {
            None => None,
            Some(pending) if pending.action != *action => {
                Some("workspace callback action mismatch")
            }
            Some(pending) if pending.evidence.is_some() => {
                Some("duplicate workspace callback evidence")
            }
            Some(pending) => {
                pending.evidence = Some(evidence);
                None
            }
        };
        problem.map_or(Ok(()), |message| Err(self.poison(message)))
    }
}

struct Core(Mutex<State>);

impl Core {
    fn lock(&self) -> Result<MutexGuard<'_, State>, String> {
        self.0.lock().map_err(|_| "workspace lock poisoned".into())
    }

    fn apply(&self, transaction: KtoTxId, payload: &[u8]) -> Result<(), String> {
        let action = Action::decode(payload)?;
        let workspace_transaction = TxId::from_bytes(*transaction.as_bytes());
        let mut state = self.lock()?;
        state.ready()?;
        let outcome = match state
            .projection
            .apply(workspace_transaction, action.clone())
        {
            Ok(value) => value,
            Err(error) => {
                let error = state.poison(format!("workspace projection failed: {error}"));
                return Err(error);
            }
        };
        let evidence = (
            transaction,
            outcome.result().cloned().map_err(str::to_owned),
        );
        state.record(&action, evidence)
    }
}

impl Subsystem for Core {
    fn submit_txn(&self, id: KtoTxId, payload: &[u8]) -> Result<(), String> {
        self.apply(id, payload)
    }

    fn reorg(&self) -> Result<(), String> {
        let mut state = self.lock()?;
        state.pending.clear();
        Err(state.poison("Web code workspaces unavailable after reorganization"))
    }
}

#[derive(Clone)]
pub struct K1WebCodeWorkspace {
    core: Arc<Core>,
    peering: Arc<K1Peering>,
    subsystem: SubsystemId,
}

impl K1WebCodeWorkspace {
    pub fn open(ordering: Arc<K1TxnOrdering>, peering: Arc<K1Peering>) -> Result<Self, String> {
        let began = Instant::now();
        let core = Arc::new(Core(Mutex::new(State {
            projection: Projection::new(),
            pending: HashMap::new(),
            fault: None,
        })));
        let subsystem = SubsystemId::from_str("k1-web-ws")?;
        ordering.register_subsystem(subsystem, None, core.clone())?;
        core.lock()?.ready()?;
        if began.elapsed() > Duration::from_millis(100) {
            eprintln!("{{\"level\":\"warning\",\"event\":\"k1_web_code_workspace_open_slow\"}}");
        }
        Ok(Self {
            core,
            peering,
            subsystem,
        })
    }

    pub fn create(&self, document: CodeDocument) -> Result<Workspace, String> {
        let action = self.reserve(|operation| Action::create(operation, document.clone()))?;
        self.submit(action)
    }

    pub fn branch(&self, document: CodeDocument) -> Result<Workspace, String> {
        let action = self.reserve(|operation| Action::branch(operation, document.clone()))?;
        self.submit(action)
    }

    pub fn overwrite(
        &self,
        workspace: TxId,
        expected: TxId,
        document: CodeDocument,
    ) -> Result<Workspace, String> {
        let action = self.reserve(|operation| {
            Action::overwrite(operation, workspace, expected, document.clone())
        })?;
        self.submit(action)
    }

    pub fn get(&self, workspace: TxId) -> Result<Option<Workspace>, String> {
        let state = self.core.lock()?;
        state.ready()?;
        Ok(state.projection.get(workspace))
    }

    pub fn latest(&self, family: &WebFamily) -> Result<Option<Workspace>, String> {
        let state = self.core.lock()?;
        state.ready()?;
        Ok(state.projection.latest(family))
    }

    fn reserve(&self, mut build: impl FnMut(OperationId) -> Action) -> Result<Action, String> {
        loop {
            let mut operation = [0; 16];
            fill(&mut operation).map_err(|error| error.to_string())?;
            let action = build(operation);
            let mut state = self.core.lock()?;
            state.ready()?;
            if let Entry::Vacant(slot) = state.pending.entry(operation) {
                slot.insert(Pending {
                    action: action.clone(),
                    evidence: None,
                });
                return Ok(action);
            }
        }
    }

    fn submit(&self, action: Action) -> Result<Workspace, String> {
        let operation = action.operation_id();
        let submitted = self.peering.submit_txn(self.subsystem, &action.encode());
        let mut state = self.core.lock()?;
        let Some(pending) = state.pending.remove(&operation) else {
            return Err(state.poison("missing workspace pending operation"));
        };
        let result = correlate(submitted, pending.evidence).map_err(|error| state.poison(error))?;
        state.ready()?;
        result
    }
}

fn correlate(
    submitted: Result<KtoTxId, String>,
    evidence: Option<Evidence>,
) -> Result<Result<Workspace, String>, &'static str> {
    match (submitted, evidence) {
        (Ok(id), Some((seen, outcome))) if id == seen => Ok(outcome),
        (Ok(_), Some(_)) => Err("workspace callback transaction mismatch"),
        (Ok(_), None) => Err("missing workspace callback evidence"),
        (Err(_), Some((_, outcome))) => Ok(outcome),
        (Err(error), None) => Ok(Err(error)),
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use kcode_k1_web_code_workspace_testkit::{Candidate, run_tests};

    struct Subject;

    impl Candidate for Subject {
        type Driver = K1WebCodeWorkspace;
        fn open(
            ordering: Arc<K1TxnOrdering>,
            peering: Arc<K1Peering>,
        ) -> Result<Self::Driver, String> {
            K1WebCodeWorkspace::open(ordering, peering)
        }
        fn create(driver: &Self::Driver, document: CodeDocument) -> Result<Workspace, String> {
            driver.create(document)
        }
        fn branch(driver: &Self::Driver, document: CodeDocument) -> Result<Workspace, String> {
            driver.branch(document)
        }
        fn overwrite(
            driver: &Self::Driver,
            workspace: TxId,
            expected: TxId,
            document: CodeDocument,
        ) -> Result<Workspace, String> {
            driver.overwrite(workspace, expected, document)
        }
        fn get(driver: &Self::Driver, workspace: TxId) -> Result<Option<Workspace>, String> {
            driver.get(workspace)
        }
        fn latest(driver: &Self::Driver, family: &WebFamily) -> Result<Option<Workspace>, String> {
            driver.latest(family)
        }
    }

    #[test]
    fn conformance() {
        run_tests::<Subject>();
    }
}