kcode-k1-persons 0.1.0

Public K1 persons transaction facade
Documentation
use std::{
    collections::HashMap,
    path::Path,
    sync::{Arc, Mutex, MutexGuard},
};

use kcode_k1_peering::K1Peering;
use kcode_k1_persons_projection::{ApplyOutcome, Projection};
pub use kcode_k1_persons_wire::PersonId;
use kcode_k1_persons_wire::{PersonsWire, decode};
pub use kcode_k1_txn_ordering::TxId;
use kcode_k1_txn_ordering::{K1TxnOrdering, Subsystem, SubsystemId};

const SUBSYSTEM_NAME: &str = "k1-persons-subsystem";
const UNAVAILABLE: &str = "persons facade is unavailable";

pub struct K1Persons {
    _ordering: Arc<K1TxnOrdering>,
    peering: Arc<K1Peering>,
    projection: Arc<Projection>,
    state: Arc<Mutex<State>>,
    subsystem: SubsystemId,
}

struct PersonSubsystem {
    projection: Arc<Projection>,
    state: Arc<Mutex<State>>,
}

struct State {
    available: bool,
    pending: HashMap<Vec<u8>, usize>,
    outcomes: HashMap<TxId, CapturedOutcome>,
}

struct CapturedOutcome {
    payload: Vec<u8>,
    outcome: ApplyOutcome,
}

impl State {
    fn new() -> Self {
        Self {
            available: true,
            pending: HashMap::new(),
            outcomes: HashMap::new(),
        }
    }

    fn fault(&mut self) {
        self.available = false;
        self.pending.clear();
        self.outcomes.clear();
    }
}

impl K1Persons {
    pub fn open(
        root: &Path,
        ordering: Arc<K1TxnOrdering>,
        peering: Arc<K1Peering>,
    ) -> Result<Self, String> {
        let (projection, checkpoint) = Projection::open(root, &ordering)?;
        let projection = Arc::new(projection);
        let state = Arc::new(Mutex::new(State::new()));
        let subsystem = SubsystemId::from_str(SUBSYSTEM_NAME)?;
        let handler = Arc::new(PersonSubsystem {
            projection: projection.clone(),
            state: state.clone(),
        });
        if let Err(error) = ordering.register_subsystem(subsystem, checkpoint, handler) {
            mark_fault(&state);
            return Err(error);
        }
        Ok(Self {
            _ordering: ordering,
            peering,
            projection,
            state,
            subsystem,
        })
    }

    pub fn create(&self, name: String) -> Result<PersonId, String> {
        let wire = PersonsWire::create(name)?;
        let (transaction, outcome) = self.submit(wire.as_bytes())?;
        match outcome {
            ApplyOutcome::Applied(person) if person.as_tx_id() == transaction => Ok(person),
            ApplyOutcome::Rejected(reason) => Err(reason),
            ApplyOutcome::Applied(_) | ApplyOutcome::Unchanged(_) => {
                mark_fault(&self.state);
                Err(UNAVAILABLE.to_owned())
            }
        }
    }

    pub fn update(&self, person: PersonId, name: String) -> Result<(), String> {
        let wire = PersonsWire::update(person, name)?;
        let (_, outcome) = self.submit(wire.as_bytes())?;
        match outcome {
            ApplyOutcome::Applied(_) | ApplyOutcome::Unchanged(_) => Ok(()),
            ApplyOutcome::Rejected(reason) => Err(reason),
        }
    }

    pub fn resolve(&self, canonical: PersonId, alias: PersonId) -> Result<PersonId, String> {
        let wire = PersonsWire::resolve(canonical, alias);
        let (_, outcome) = self.submit(wire.as_bytes())?;
        match outcome {
            ApplyOutcome::Applied(root) | ApplyOutcome::Unchanged(root) => Ok(root),
            ApplyOutcome::Rejected(reason) => Err(reason),
        }
    }

    pub fn get(&self, person: PersonId) -> Result<Option<String>, String> {
        ensure_available(&self.state)?;
        let result = self.projection.get(person);
        ensure_available(&self.state)?;
        result
    }

    fn submit(&self, payload: &[u8]) -> Result<(TxId, ApplyOutcome), String> {
        self.reserve(payload)?;
        match self.peering.submit_txn(self.subsystem, payload) {
            Ok(transaction) => self
                .claim(payload, transaction)
                .map(|outcome| (transaction, outcome)),
            Err(error) => Err(self.finish_submission_error(payload, error)),
        }
    }

    fn reserve(&self, payload: &[u8]) -> Result<(), String> {
        let mut state = lock_state(&self.state)?;
        if !state.available {
            return Err(UNAVAILABLE.to_owned());
        }
        if let Some(count) = state.pending.get_mut(payload) {
            if let Some(next) = count.checked_add(1) {
                *count = next;
                return Ok(());
            }
            state.fault();
            return Err(UNAVAILABLE.to_owned());
        }
        if state.pending.try_reserve(1).is_err() {
            state.fault();
            return Err(UNAVAILABLE.to_owned());
        }
        let key = match copy_payload(payload) {
            Ok(key) => key,
            Err(()) => {
                state.fault();
                return Err(UNAVAILABLE.to_owned());
            }
        };
        state.pending.insert(key, 1);
        Ok(())
    }

    fn claim(&self, payload: &[u8], transaction: TxId) -> Result<ApplyOutcome, String> {
        let mut state = lock_state(&self.state)?;
        if !state.available {
            return Err(UNAVAILABLE.to_owned());
        }
        let captured = match state.outcomes.remove(&transaction) {
            Some(captured) => captured,
            None => {
                state.fault();
                return Err(UNAVAILABLE.to_owned());
            }
        };
        if captured.payload != payload || !release_pending(&mut state, payload) {
            state.fault();
            return Err(UNAVAILABLE.to_owned());
        }
        Ok(captured.outcome)
    }

    fn finish_submission_error(&self, payload: &[u8], error: String) -> String {
        let mut state = match lock_state(&self.state) {
            Ok(state) => state,
            Err(error) => return error,
        };
        if !state.available {
            return UNAVAILABLE.to_owned();
        }
        if !release_pending(&mut state, payload) {
            state.fault();
            return UNAVAILABLE.to_owned();
        }
        error
    }
}

impl Subsystem for PersonSubsystem {
    fn submit_txn(&self, transaction: TxId, payload: &[u8]) -> Result<(), String> {
        ensure_available(&self.state)?;
        let action = match decode(payload) {
            Ok(action) => action,
            Err(error) => {
                mark_fault(&self.state);
                return Err(error);
            }
        };
        let outcome = match self.projection.apply(transaction, action) {
            Ok(outcome) => outcome,
            Err(error) => {
                mark_fault(&self.state);
                return Err(error);
            }
        };
        self.capture(transaction, payload, outcome)
    }

    fn reorg(&self) -> Result<(), String> {
        let lock_error = match self.state.lock() {
            Ok(mut state) => {
                state.fault();
                None
            }
            Err(error) => {
                error.into_inner().fault();
                Some(UNAVAILABLE.to_owned())
            }
        };
        self.projection.clear()?;
        match lock_error {
            Some(error) => Err(error),
            None => Ok(()),
        }
    }
}

impl PersonSubsystem {
    fn capture(
        &self,
        transaction: TxId,
        payload: &[u8],
        outcome: ApplyOutcome,
    ) -> Result<(), String> {
        let mut state = lock_state(&self.state)?;
        if !state.available {
            return Err(UNAVAILABLE.to_owned());
        }
        match state.pending.get(payload).copied() {
            None => return Ok(()),
            Some(0) => {
                state.fault();
                return Err(UNAVAILABLE.to_owned());
            }
            Some(_) => {}
        }
        if state.outcomes.contains_key(&transaction) || state.outcomes.try_reserve(1).is_err() {
            state.fault();
            return Err(UNAVAILABLE.to_owned());
        }
        let exact_payload = match copy_payload(payload) {
            Ok(payload) => payload,
            Err(()) => {
                state.fault();
                return Err(UNAVAILABLE.to_owned());
            }
        };
        if state
            .outcomes
            .insert(
                transaction,
                CapturedOutcome {
                    payload: exact_payload,
                    outcome,
                },
            )
            .is_some()
        {
            state.fault();
            return Err(UNAVAILABLE.to_owned());
        }
        Ok(())
    }
}

fn copy_payload(payload: &[u8]) -> Result<Vec<u8>, ()> {
    let mut copy = Vec::new();
    copy.try_reserve_exact(payload.len()).map_err(|_| ())?;
    copy.extend_from_slice(payload);
    Ok(copy)
}

fn release_pending(state: &mut State, payload: &[u8]) -> bool {
    let remove = match state.pending.get_mut(payload) {
        Some(count) if *count > 1 => {
            *count -= 1;
            false
        }
        Some(count) if *count == 1 => true,
        _ => return false,
    };
    if remove {
        state.pending.remove(payload);
        state
            .outcomes
            .retain(|_, captured| captured.payload.as_slice() != payload);
    }
    true
}

fn ensure_available(state: &Mutex<State>) -> Result<(), String> {
    let state = lock_state(state)?;
    if state.available {
        Ok(())
    } else {
        Err(UNAVAILABLE.to_owned())
    }
}

fn lock_state(state: &Mutex<State>) -> Result<MutexGuard<'_, State>, String> {
    match state.lock() {
        Ok(state) => Ok(state),
        Err(error) => {
            error.into_inner().fault();
            Err(UNAVAILABLE.to_owned())
        }
    }
}

fn mark_fault(state: &Mutex<State>) {
    match state.lock() {
        Ok(mut state) => state.fault(),
        Err(error) => error.into_inner().fault(),
    }
}