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}