Skip to main content

kcode_k1_persons/
lib.rs

1use std::{
2    collections::HashMap,
3    path::Path,
4    sync::{Arc, Mutex, MutexGuard},
5};
6
7use kcode_k1_peering::K1Peering;
8pub use kcode_k1_persons_projection::PersonView;
9use kcode_k1_persons_projection::{ApplyOutcome, Projection};
10pub use kcode_k1_persons_wire::PersonId;
11use kcode_k1_persons_wire::{PersonsWire, decode};
12pub use kcode_k1_txn_ordering::TxId;
13use kcode_k1_txn_ordering::{K1TxnOrdering, Subsystem, SubsystemId};
14
15const SUBSYSTEM_NAME: &str = "k1-persons-subsystem";
16const UNAVAILABLE: &str = "persons facade is unavailable";
17
18pub struct K1Persons {
19    _ordering: Arc<K1TxnOrdering>,
20    peering: Arc<K1Peering>,
21    projection: Arc<Projection>,
22    state: Arc<Mutex<State>>,
23    subsystem: SubsystemId,
24}
25
26struct PersonSubsystem {
27    projection: Arc<Projection>,
28    state: Arc<Mutex<State>>,
29}
30
31struct State {
32    available: bool,
33    pending: HashMap<Vec<u8>, usize>,
34    outcomes: HashMap<TxId, CapturedOutcome>,
35}
36
37struct CapturedOutcome {
38    payload: Vec<u8>,
39    outcome: ApplyOutcome,
40}
41
42impl State {
43    fn new() -> Self {
44        Self {
45            available: true,
46            pending: HashMap::new(),
47            outcomes: HashMap::new(),
48        }
49    }
50
51    fn fault(&mut self) {
52        self.available = false;
53        self.pending.clear();
54        self.outcomes.clear();
55    }
56}
57
58impl K1Persons {
59    pub fn open(
60        root: &Path,
61        ordering: Arc<K1TxnOrdering>,
62        peering: Arc<K1Peering>,
63    ) -> Result<Self, String> {
64        let (projection, checkpoint) = Projection::open(root, &ordering)?;
65        let projection = Arc::new(projection);
66        let state = Arc::new(Mutex::new(State::new()));
67        let subsystem = SubsystemId::from_str(SUBSYSTEM_NAME)?;
68        let handler = Arc::new(PersonSubsystem {
69            projection: projection.clone(),
70            state: state.clone(),
71        });
72        if let Err(error) = ordering.register_subsystem(subsystem, checkpoint, handler) {
73            mark_fault(&state);
74            return Err(error);
75        }
76        Ok(Self {
77            _ordering: ordering,
78            peering,
79            projection,
80            state,
81            subsystem,
82        })
83    }
84
85    pub fn create(&self, name: String) -> Result<PersonId, String> {
86        let wire = PersonsWire::create(name)?;
87        let (transaction, outcome) = self.submit(wire.as_bytes())?;
88        match outcome {
89            ApplyOutcome::Applied(person) if person.as_tx_id() == transaction => Ok(person),
90            ApplyOutcome::Rejected(reason) => Err(reason),
91            ApplyOutcome::Applied(_) | ApplyOutcome::Unchanged(_) => {
92                mark_fault(&self.state);
93                Err(UNAVAILABLE.to_owned())
94            }
95        }
96    }
97
98    pub fn update(&self, person: PersonId, name: String) -> Result<(), String> {
99        let wire = PersonsWire::update(person, name)?;
100        let (_, outcome) = self.submit(wire.as_bytes())?;
101        match outcome {
102            ApplyOutcome::Applied(_) | ApplyOutcome::Unchanged(_) => Ok(()),
103            ApplyOutcome::Rejected(reason) => Err(reason),
104        }
105    }
106
107    pub fn resolve(&self, canonical: PersonId, alias: PersonId) -> Result<PersonId, String> {
108        let wire = PersonsWire::resolve(canonical, alias);
109        let (_, outcome) = self.submit(wire.as_bytes())?;
110        match outcome {
111            ApplyOutcome::Applied(root) | ApplyOutcome::Unchanged(root) => Ok(root),
112            ApplyOutcome::Rejected(reason) => Err(reason),
113        }
114    }
115
116    pub fn read(&self, person: PersonId) -> Result<Option<PersonView>, String> {
117        ensure_available(&self.state)?;
118        let result = self.projection.read(person);
119        ensure_available(&self.state)?;
120        result
121    }
122
123    pub fn get(&self, person: PersonId) -> Result<Option<String>, String> {
124        ensure_available(&self.state)?;
125        let result = self.projection.get(person);
126        ensure_available(&self.state)?;
127        result
128    }
129
130    fn submit(&self, payload: &[u8]) -> Result<(TxId, ApplyOutcome), String> {
131        self.reserve(payload)?;
132        match self.peering.submit_txn(self.subsystem, payload) {
133            Ok(transaction) => self
134                .claim(payload, transaction)
135                .map(|outcome| (transaction, outcome)),
136            Err(error) => Err(self.finish_submission_error(payload, error)),
137        }
138    }
139
140    fn reserve(&self, payload: &[u8]) -> Result<(), String> {
141        let mut state = lock_state(&self.state)?;
142        if !state.available {
143            return Err(UNAVAILABLE.to_owned());
144        }
145        if let Some(count) = state.pending.get_mut(payload) {
146            if let Some(next) = count.checked_add(1) {
147                *count = next;
148                return Ok(());
149            }
150            state.fault();
151            return Err(UNAVAILABLE.to_owned());
152        }
153        if state.pending.try_reserve(1).is_err() {
154            state.fault();
155            return Err(UNAVAILABLE.to_owned());
156        }
157        let key = match copy_payload(payload) {
158            Ok(key) => key,
159            Err(()) => {
160                state.fault();
161                return Err(UNAVAILABLE.to_owned());
162            }
163        };
164        state.pending.insert(key, 1);
165        Ok(())
166    }
167
168    fn claim(&self, payload: &[u8], transaction: TxId) -> Result<ApplyOutcome, String> {
169        let mut state = lock_state(&self.state)?;
170        if !state.available {
171            return Err(UNAVAILABLE.to_owned());
172        }
173        let captured = match state.outcomes.remove(&transaction) {
174            Some(captured) => captured,
175            None => {
176                state.fault();
177                return Err(UNAVAILABLE.to_owned());
178            }
179        };
180        if captured.payload != payload || !release_pending(&mut state, payload) {
181            state.fault();
182            return Err(UNAVAILABLE.to_owned());
183        }
184        Ok(captured.outcome)
185    }
186
187    fn finish_submission_error(&self, payload: &[u8], error: String) -> String {
188        let mut state = match lock_state(&self.state) {
189            Ok(state) => state,
190            Err(error) => return error,
191        };
192        if !state.available {
193            return UNAVAILABLE.to_owned();
194        }
195        if !release_pending(&mut state, payload) {
196            state.fault();
197            return UNAVAILABLE.to_owned();
198        }
199        error
200    }
201}
202
203impl Subsystem for PersonSubsystem {
204    fn submit_txn(&self, transaction: TxId, payload: &[u8]) -> Result<(), String> {
205        ensure_available(&self.state)?;
206        let action = match decode(payload) {
207            Ok(action) => action,
208            Err(error) => {
209                mark_fault(&self.state);
210                return Err(error);
211            }
212        };
213        let outcome = match self.projection.apply(transaction, action) {
214            Ok(outcome) => outcome,
215            Err(error) => {
216                mark_fault(&self.state);
217                return Err(error);
218            }
219        };
220        self.capture(transaction, payload, outcome)
221    }
222
223    fn reorg(&self) -> Result<(), String> {
224        let lock_error = match self.state.lock() {
225            Ok(mut state) => {
226                state.fault();
227                None
228            }
229            Err(error) => {
230                error.into_inner().fault();
231                Some(UNAVAILABLE.to_owned())
232            }
233        };
234        self.projection.clear()?;
235        match lock_error {
236            Some(error) => Err(error),
237            None => Ok(()),
238        }
239    }
240}
241
242impl PersonSubsystem {
243    fn capture(
244        &self,
245        transaction: TxId,
246        payload: &[u8],
247        outcome: ApplyOutcome,
248    ) -> Result<(), String> {
249        let mut state = lock_state(&self.state)?;
250        if !state.available {
251            return Err(UNAVAILABLE.to_owned());
252        }
253        match state.pending.get(payload).copied() {
254            None => return Ok(()),
255            Some(0) => {
256                state.fault();
257                return Err(UNAVAILABLE.to_owned());
258            }
259            Some(_) => {}
260        }
261        if state.outcomes.contains_key(&transaction) || state.outcomes.try_reserve(1).is_err() {
262            state.fault();
263            return Err(UNAVAILABLE.to_owned());
264        }
265        let exact_payload = match copy_payload(payload) {
266            Ok(payload) => payload,
267            Err(()) => {
268                state.fault();
269                return Err(UNAVAILABLE.to_owned());
270            }
271        };
272        if state
273            .outcomes
274            .insert(
275                transaction,
276                CapturedOutcome {
277                    payload: exact_payload,
278                    outcome,
279                },
280            )
281            .is_some()
282        {
283            state.fault();
284            return Err(UNAVAILABLE.to_owned());
285        }
286        Ok(())
287    }
288}
289
290fn copy_payload(payload: &[u8]) -> Result<Vec<u8>, ()> {
291    let mut copy = Vec::new();
292    copy.try_reserve_exact(payload.len()).map_err(|_| ())?;
293    copy.extend_from_slice(payload);
294    Ok(copy)
295}
296
297fn release_pending(state: &mut State, payload: &[u8]) -> bool {
298    let remove = match state.pending.get_mut(payload) {
299        Some(count) if *count > 1 => {
300            *count -= 1;
301            false
302        }
303        Some(count) if *count == 1 => true,
304        _ => return false,
305    };
306    if remove {
307        state.pending.remove(payload);
308        state
309            .outcomes
310            .retain(|_, captured| captured.payload.as_slice() != payload);
311    }
312    true
313}
314
315fn ensure_available(state: &Mutex<State>) -> Result<(), String> {
316    let state = lock_state(state)?;
317    if state.available {
318        Ok(())
319    } else {
320        Err(UNAVAILABLE.to_owned())
321    }
322}
323
324fn lock_state(state: &Mutex<State>) -> Result<MutexGuard<'_, State>, String> {
325    match state.lock() {
326        Ok(state) => Ok(state),
327        Err(error) => {
328            error.into_inner().fault();
329            Err(UNAVAILABLE.to_owned())
330        }
331    }
332}
333
334fn mark_fault(state: &Mutex<State>) {
335    match state.lock() {
336        Ok(mut state) => state.fault(),
337        Err(error) => error.into_inner().fault(),
338    }
339}