kcode_k1_persons_projection/
lib.rs1use 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}