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