#![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>();
}
}