Skip to main content

kcode_k1_persons_projection/
lib.rs

1use std::{
2    collections::HashMap,
3    path::Path,
4    sync::{
5        Mutex, MutexGuard, RwLock, RwLockReadGuard, RwLockWriteGuard,
6        atomic::{AtomicBool, Ordering},
7    },
8    time::{Duration, Instant},
9};
10
11pub use kcode_k1_person_types::PersonId;
12use kcode_k1_persons_store::{Store, StoreChange, StoredSnapshot};
13use kcode_k1_txn_ordering::K1TxnOrdering;
14pub use kcode_k1_txn_ordering::TxId;
15
16#[derive(Clone, Debug, Eq, PartialEq)]
17pub struct PersonView {
18    pub person_id: PersonId,
19    pub name: String,
20}
21
22#[derive(Clone, Debug, Eq, PartialEq)]
23pub struct PersonAction(Action);
24
25#[derive(Clone, Debug, Eq, PartialEq)]
26enum Action {
27    Create(String),
28    Update(PersonId, String),
29    Resolve(PersonId, PersonId),
30}
31
32impl PersonAction {
33    pub fn create(name: String) -> Result<Self, String> {
34        validate_name(&name)?;
35        Ok(Self(Action::Create(name)))
36    }
37
38    pub fn update(person: PersonId, name: String) -> Result<Self, String> {
39        validate_name(&name)?;
40        Ok(Self(Action::Update(person, name)))
41    }
42
43    pub fn resolve(canonical: PersonId, alias: PersonId) -> Self {
44        Self(Action::Resolve(canonical, alias))
45    }
46}
47
48fn validate_name(name: &str) -> Result<(), String> {
49    if !(1..=128).contains(&name.len()) {
50        return Err("person name must be 1 through 128 UTF-8 bytes".to_owned());
51    }
52    if name.chars().any(char::is_control) {
53        return Err("person name must not contain control characters".to_owned());
54    }
55    if !name.chars().any(|character| !character.is_whitespace()) {
56        return Err("person name must contain a non-whitespace character".to_owned());
57    }
58    Ok(())
59}
60
61#[derive(Clone, Debug, Eq, PartialEq)]
62pub enum ApplyOutcome {
63    Applied(PersonId),
64    Unchanged(PersonId),
65    Rejected(String),
66}
67
68struct Entry {
69    root: PersonId,
70    name: Option<String>,
71}
72
73#[derive(Default)]
74struct State {
75    entries: HashMap<PersonId, Entry>,
76    classes: HashMap<PersonId, Vec<PersonId>>,
77    checkpoint: Option<TxId>,
78}
79
80enum Change {
81    Create(PersonId, String),
82    Update(PersonId, String),
83    Resolve(PersonId, PersonId, usize),
84}
85
86impl State {
87    fn from_snapshot(snapshot: StoredSnapshot) -> Result<Self, ()> {
88        let mut state = Self {
89            checkpoint: snapshot.checkpoint(),
90            ..Self::default()
91        };
92        for person in snapshot.persons() {
93            if person
94                .name()
95                .is_some_and(|name| validate_name(name).is_err())
96                || state
97                    .entries
98                    .insert(
99                        person.id(),
100                        Entry {
101                            root: person.root(),
102                            name: person.name().map(str::to_owned),
103                        },
104                    )
105                    .is_some()
106            {
107                return Err(());
108            }
109        }
110        for (&id, entry) in &state.entries {
111            let root = state.entries.get(&entry.root).ok_or(())?;
112            if root.root != entry.root
113                || root.name.is_none()
114                || (id == entry.root) != entry.name.is_some()
115            {
116                return Err(());
117            }
118            state.classes.entry(entry.root).or_default().push(id);
119        }
120        Ok(state)
121    }
122
123    fn plan(&self, callback: TxId, action: PersonAction) -> (ApplyOutcome, Option<Change>) {
124        match action.0 {
125            Action::Create(name) => {
126                let id = PersonId::from_tx_id(callback);
127                if self.entries.contains_key(&id) {
128                    rejected("person already exists")
129                } else {
130                    (ApplyOutcome::Applied(id), Some(Change::Create(id, name)))
131                }
132            }
133            Action::Update(id, name) => {
134                let Some(entry) = self.entries.get(&id) else {
135                    return rejected("person is unknown");
136                };
137                let root = entry.root;
138                if self.entries[&root].name.as_deref() == Some(&name) {
139                    (ApplyOutcome::Unchanged(root), None)
140                } else {
141                    (
142                        ApplyOutcome::Applied(root),
143                        Some(Change::Update(root, name)),
144                    )
145                }
146            }
147            Action::Resolve(canonical, alias) => {
148                let Some(canonical) = self.entries.get(&canonical).map(|entry| entry.root) else {
149                    return rejected("canonical person is unknown");
150                };
151                let Some(alias) = self.entries.get(&alias).map(|entry| entry.root) else {
152                    return rejected("alias person is unknown");
153                };
154                if canonical == alias {
155                    (ApplyOutcome::Unchanged(canonical), None)
156                } else {
157                    (
158                        ApplyOutcome::Applied(canonical),
159                        Some(Change::Resolve(
160                            alias,
161                            canonical,
162                            self.classes[&alias].len(),
163                        )),
164                    )
165                }
166            }
167        }
168    }
169
170    fn read(&self, person: PersonId) -> Option<PersonView> {
171        let root = self.entries.get(&person)?.root;
172        let name = self.entries.get(&root)?.name.clone()?;
173        Some(PersonView {
174            person_id: root,
175            name,
176        })
177    }
178
179    fn publish(
180        &mut self,
181        callback: TxId,
182        outcome: ApplyOutcome,
183        change: Option<Change>,
184    ) -> ApplyOutcome {
185        match change {
186            Some(Change::Create(id, name)) => {
187                self.entries.insert(
188                    id,
189                    Entry {
190                        root: id,
191                        name: Some(name),
192                    },
193                );
194                self.classes.insert(id, vec![id]);
195            }
196            Some(Change::Update(root, name)) => {
197                self.entries.get_mut(&root).unwrap().name = Some(name)
198            }
199            Some(Change::Resolve(from, to, _)) => {
200                let members = self.classes.remove(&from).unwrap();
201                for id in &members {
202                    self.entries.get_mut(id).unwrap().root = to;
203                }
204                self.entries.get_mut(&from).unwrap().name = None;
205                self.classes.get_mut(&to).unwrap().extend(members);
206            }
207            None => {}
208        }
209        self.checkpoint = Some(callback);
210        outcome
211    }
212}
213
214fn rejected(reason: &str) -> (ApplyOutcome, Option<Change>) {
215    (ApplyOutcome::Rejected(reason.to_owned()), None)
216}
217
218fn store_change(change: &Change) -> StoreChange {
219    match change {
220        Change::Create(id, name) => StoreChange::create(*id, name.clone()),
221        Change::Update(id, name) => StoreChange::update(*id, name.clone()),
222        Change::Resolve(from, to, count) => StoreChange::resolve(*from, *to, *count),
223    }
224}
225
226pub struct Projection {
227    state: RwLock<State>,
228    apply_lane: Mutex<()>,
229    store: Mutex<Store>,
230    unavailable: AtomicBool,
231}
232
233impl Projection {
234    pub fn open(root: &Path, ordering: &K1TxnOrdering) -> Result<(Self, Option<TxId>), String> {
235        let started = Instant::now();
236        let result = (|| {
237            let (mut store, snapshot) = Store::open(root, ordering)?;
238            let state = match State::from_snapshot(snapshot) {
239                Ok(state) => state,
240                Err(()) => {
241                    store.clear()?;
242                    State::default()
243                }
244            };
245            let checkpoint = state.checkpoint;
246            Ok((
247                Self {
248                    state: RwLock::new(state),
249                    apply_lane: Mutex::new(()),
250                    store: Mutex::new(store),
251                    unavailable: AtomicBool::new(false),
252                },
253                checkpoint,
254            ))
255        })();
256        if started.elapsed() > Duration::from_millis(100) {
257            let outcome = if result.is_ok() { "ready" } else { "error" };
258            eprintln!(
259                "level=warn module=kcode-k1-persons-projection operation=open elapsed_us={} outcome={outcome}",
260                started.elapsed().as_micros()
261            );
262        }
263        result
264    }
265
266    pub fn apply(&self, callback: TxId, action: PersonAction) -> Result<ApplyOutcome, String> {
267        self.ensure_available()?;
268        let _lane = self.apply_lock()?;
269        self.ensure_available()?;
270        let (outcome, change) = self.state_read()?.plan(callback, action);
271        let store_change = change.as_ref().map(store_change);
272        if let Err(error) = self.store_lock()?.commit(callback, store_change.as_ref()) {
273            self.unavailable.store(true, Ordering::SeqCst);
274            return Err(error);
275        }
276        Ok(self.state_write()?.publish(callback, outcome, change))
277    }
278
279    pub fn read(&self, person: PersonId) -> Result<Option<PersonView>, String> {
280        self.ensure_available()?;
281        let state = self.state_read()?;
282        self.ensure_available()?;
283        Ok(state.read(person))
284    }
285
286    pub fn get(&self, person: PersonId) -> Result<Option<String>, String> {
287        Ok(self.read(person)?.map(|view| view.name))
288    }
289
290    pub fn clear(&self) -> Result<(), String> {
291        self.ensure_available()?;
292        let _lane = self.apply_lock()?;
293        self.ensure_available()?;
294        if let Err(error) = self.store_lock()?.clear() {
295            self.unavailable.store(true, Ordering::SeqCst);
296            return Err(error);
297        }
298        *self.state_write()? = State::default();
299        Ok(())
300    }
301
302    fn ensure_available(&self) -> Result<(), String> {
303        (!self.unavailable.load(Ordering::SeqCst))
304            .then_some(())
305            .ok_or_else(|| "projection is unavailable until reopen".to_owned())
306    }
307
308    fn apply_lock(&self) -> Result<MutexGuard<'_, ()>, String> {
309        self.apply_lane
310            .lock()
311            .map_err(|_| self.failed("projection apply lane is unavailable until reopen"))
312    }
313
314    fn store_lock(&self) -> Result<MutexGuard<'_, Store>, String> {
315        self.store
316            .lock()
317            .map_err(|_| self.failed("projection store lane is unavailable until reopen"))
318    }
319
320    fn state_read(&self) -> Result<RwLockReadGuard<'_, State>, String> {
321        self.state
322            .read()
323            .map_err(|_| self.failed("projection state is unavailable until reopen"))
324    }
325
326    fn state_write(&self) -> Result<RwLockWriteGuard<'_, State>, String> {
327        self.state
328            .write()
329            .map_err(|_| self.failed("projection state is unavailable until reopen"))
330    }
331
332    fn failed(&self, message: &str) -> String {
333        self.unavailable.store(true, Ordering::SeqCst);
334        message.to_owned()
335    }
336}