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 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 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 get(&self, person: PersonId) -> Result<Option<String>, String> {
self.ensure_available()?;
let state = self.state_read()?;
self.ensure_available()?;
Ok(state
.entries
.get(&person)
.and_then(|entry| state.entries[&entry.root].name.clone()))
}
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()
}
}