kcode-k1-persons-projection 0.2.0

Durable current-name projection for K1 persons
Documentation
use std::{
    collections::HashMap,
    path::Path,
    sync::{
        Mutex, MutexGuard, RwLock, RwLockReadGuard, RwLockWriteGuard,
        atomic::{AtomicBool, Ordering},
    },
    time::{Duration, Instant},
};

pub use kcode_k1_person_types::PersonId;
use kcode_k1_persons_store::{Store, StoreChange, StoredSnapshot};
use kcode_k1_txn_ordering::K1TxnOrdering;
pub use kcode_k1_txn_ordering::TxId;

#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PersonView {
    pub person_id: PersonId,
    pub name: String,
}

#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PersonAction(Action);

#[derive(Clone, Debug, Eq, PartialEq)]
enum Action {
    Create(String),
    Update(PersonId, String),
    Resolve(PersonId, PersonId),
}

impl PersonAction {
    pub fn create(name: String) -> Result<Self, String> {
        validate_name(&name)?;
        Ok(Self(Action::Create(name)))
    }

    pub fn update(person: PersonId, name: String) -> Result<Self, String> {
        validate_name(&name)?;
        Ok(Self(Action::Update(person, name)))
    }

    pub fn resolve(canonical: PersonId, alias: PersonId) -> Self {
        Self(Action::Resolve(canonical, alias))
    }
}

fn validate_name(name: &str) -> Result<(), String> {
    if !(1..=128).contains(&name.len()) {
        return Err("person name must be 1 through 128 UTF-8 bytes".to_owned());
    }
    if name.chars().any(char::is_control) {
        return Err("person name must not contain control characters".to_owned());
    }
    if !name.chars().any(|character| !character.is_whitespace()) {
        return Err("person name must contain a non-whitespace character".to_owned());
    }
    Ok(())
}

#[derive(Clone, Debug, Eq, PartialEq)]
pub enum ApplyOutcome {
    Applied(PersonId),
    Unchanged(PersonId),
    Rejected(String),
}

struct Entry {
    root: PersonId,
    name: Option<String>,
}

#[derive(Default)]
struct State {
    entries: HashMap<PersonId, Entry>,
    classes: HashMap<PersonId, Vec<PersonId>>,
    checkpoint: Option<TxId>,
}

enum Change {
    Create(PersonId, String),
    Update(PersonId, String),
    Resolve(PersonId, PersonId, usize),
}

impl State {
    fn from_snapshot(snapshot: StoredSnapshot) -> Result<Self, ()> {
        let mut state = Self {
            checkpoint: snapshot.checkpoint(),
            ..Self::default()
        };
        for person in snapshot.persons() {
            if person
                .name()
                .is_some_and(|name| validate_name(name).is_err())
                || state
                    .entries
                    .insert(
                        person.id(),
                        Entry {
                            root: person.root(),
                            name: person.name().map(str::to_owned),
                        },
                    )
                    .is_some()
            {
                return Err(());
            }
        }
        for (&id, entry) in &state.entries {
            let root = state.entries.get(&entry.root).ok_or(())?;
            if root.root != entry.root
                || root.name.is_none()
                || (id == entry.root) != entry.name.is_some()
            {
                return Err(());
            }
            state.classes.entry(entry.root).or_default().push(id);
        }
        Ok(state)
    }

    fn plan(&self, callback: TxId, action: PersonAction) -> (ApplyOutcome, Option<Change>) {
        match action.0 {
            Action::Create(name) => {
                let id = PersonId::from_tx_id(callback);
                if self.entries.contains_key(&id) {
                    rejected("person already exists")
                } else {
                    (ApplyOutcome::Applied(id), Some(Change::Create(id, name)))
                }
            }
            Action::Update(id, name) => {
                let Some(entry) = self.entries.get(&id) else {
                    return rejected("person is unknown");
                };
                let root = entry.root;
                if self.entries[&root].name.as_deref() == Some(&name) {
                    (ApplyOutcome::Unchanged(root), None)
                } else {
                    (
                        ApplyOutcome::Applied(root),
                        Some(Change::Update(root, name)),
                    )
                }
            }
            Action::Resolve(canonical, alias) => {
                let Some(canonical) = self.entries.get(&canonical).map(|entry| entry.root) else {
                    return rejected("canonical person is unknown");
                };
                let Some(alias) = self.entries.get(&alias).map(|entry| entry.root) else {
                    return rejected("alias person is unknown");
                };
                if canonical == alias {
                    (ApplyOutcome::Unchanged(canonical), None)
                } else {
                    (
                        ApplyOutcome::Applied(canonical),
                        Some(Change::Resolve(
                            alias,
                            canonical,
                            self.classes[&alias].len(),
                        )),
                    )
                }
            }
        }
    }

    fn read(&self, person: PersonId) -> Option<PersonView> {
        let root = self.entries.get(&person)?.root;
        let name = self.entries.get(&root)?.name.clone()?;
        Some(PersonView {
            person_id: root,
            name,
        })
    }

    fn publish(
        &mut self,
        callback: TxId,
        outcome: ApplyOutcome,
        change: Option<Change>,
    ) -> ApplyOutcome {
        match change {
            Some(Change::Create(id, name)) => {
                self.entries.insert(
                    id,
                    Entry {
                        root: id,
                        name: Some(name),
                    },
                );
                self.classes.insert(id, vec![id]);
            }
            Some(Change::Update(root, name)) => {
                self.entries.get_mut(&root).unwrap().name = Some(name)
            }
            Some(Change::Resolve(from, to, _)) => {
                let members = self.classes.remove(&from).unwrap();
                for id in &members {
                    self.entries.get_mut(id).unwrap().root = to;
                }
                self.entries.get_mut(&from).unwrap().name = None;
                self.classes.get_mut(&to).unwrap().extend(members);
            }
            None => {}
        }
        self.checkpoint = Some(callback);
        outcome
    }
}

fn rejected(reason: &str) -> (ApplyOutcome, Option<Change>) {
    (ApplyOutcome::Rejected(reason.to_owned()), None)
}

fn store_change(change: &Change) -> StoreChange {
    match change {
        Change::Create(id, name) => StoreChange::create(*id, name.clone()),
        Change::Update(id, name) => StoreChange::update(*id, name.clone()),
        Change::Resolve(from, to, count) => StoreChange::resolve(*from, *to, *count),
    }
}

pub struct Projection {
    state: RwLock<State>,
    apply_lane: Mutex<()>,
    store: Mutex<Store>,
    unavailable: AtomicBool,
}

impl Projection {
    pub fn open(root: &Path, ordering: &K1TxnOrdering) -> Result<(Self, Option<TxId>), String> {
        let started = Instant::now();
        let result = (|| {
            let (mut store, snapshot) = Store::open(root, ordering)?;
            let state = match State::from_snapshot(snapshot) {
                Ok(state) => state,
                Err(()) => {
                    store.clear()?;
                    State::default()
                }
            };
            let checkpoint = state.checkpoint;
            Ok((
                Self {
                    state: RwLock::new(state),
                    apply_lane: Mutex::new(()),
                    store: Mutex::new(store),
                    unavailable: AtomicBool::new(false),
                },
                checkpoint,
            ))
        })();
        if started.elapsed() > Duration::from_millis(100) {
            let outcome = if result.is_ok() { "ready" } else { "error" };
            eprintln!(
                "level=warn module=kcode-k1-persons-projection operation=open elapsed_us={} outcome={outcome}",
                started.elapsed().as_micros()
            );
        }
        result
    }

    pub fn apply(&self, callback: TxId, action: PersonAction) -> Result<ApplyOutcome, String> {
        self.ensure_available()?;
        let _lane = self.apply_lock()?;
        self.ensure_available()?;
        let (outcome, change) = self.state_read()?.plan(callback, action);
        let store_change = change.as_ref().map(store_change);
        if let Err(error) = self.store_lock()?.commit(callback, store_change.as_ref()) {
            self.unavailable.store(true, Ordering::SeqCst);
            return Err(error);
        }
        Ok(self.state_write()?.publish(callback, outcome, change))
    }

    pub fn read(&self, person: PersonId) -> Result<Option<PersonView>, String> {
        self.ensure_available()?;
        let state = self.state_read()?;
        self.ensure_available()?;
        Ok(state.read(person))
    }

    pub fn get(&self, person: PersonId) -> Result<Option<String>, String> {
        Ok(self.read(person)?.map(|view| view.name))
    }

    pub fn clear(&self) -> Result<(), String> {
        self.ensure_available()?;
        let _lane = self.apply_lock()?;
        self.ensure_available()?;
        if let Err(error) = self.store_lock()?.clear() {
            self.unavailable.store(true, Ordering::SeqCst);
            return Err(error);
        }
        *self.state_write()? = State::default();
        Ok(())
    }

    fn ensure_available(&self) -> Result<(), String> {
        (!self.unavailable.load(Ordering::SeqCst))
            .then_some(())
            .ok_or_else(|| "projection is unavailable until reopen".to_owned())
    }

    fn apply_lock(&self) -> Result<MutexGuard<'_, ()>, String> {
        self.apply_lane
            .lock()
            .map_err(|_| self.failed("projection apply lane is unavailable until reopen"))
    }

    fn store_lock(&self) -> Result<MutexGuard<'_, Store>, String> {
        self.store
            .lock()
            .map_err(|_| self.failed("projection store lane is unavailable until reopen"))
    }

    fn state_read(&self) -> Result<RwLockReadGuard<'_, State>, String> {
        self.state
            .read()
            .map_err(|_| self.failed("projection state is unavailable until reopen"))
    }

    fn state_write(&self) -> Result<RwLockWriteGuard<'_, State>, String> {
        self.state
            .write()
            .map_err(|_| self.failed("projection state is unavailable until reopen"))
    }

    fn failed(&self, message: &str) -> String {
        self.unavailable.store(true, Ordering::SeqCst);
        message.to_owned()
    }
}