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}