Skip to main content

kcode_k1_access_projection/
lib.rs

1use std::{
2    collections::HashMap,
3    hash::Hash,
4    path::Path,
5    sync::{LockResult, Mutex, RwLock, RwLockReadGuard},
6    time::{Duration, Instant},
7};
8
9pub use kcode_k1_access_format::{AccessAction, OwnerWitness};
10pub use kcode_k1_access_types::{
11    AccessCheck, AccessId, AccessRevision, Authorizations, GroupId, ModelId, OwnerSubject,
12    RequestPrincipal, SubsystemId, Target, TxId, UserId, ViewerSubject,
13};
14pub use kcode_k1_txn_ordering::K1TxnOrdering;
15
16use kcode_k1_access_store::{Store, StoreMutation, StoredAccess};
17use kcode_k1_transaction::Transaction;
18
19#[derive(Clone, Debug, Eq, PartialEq)]
20pub enum ApplyOutcome {
21    Applied(AccessRevision),
22    Unchanged(AccessRevision),
23    Rejected(String),
24}
25
26pub struct Projection {
27    store: Store,
28    state: RwLock<State>,
29    apply: Mutex<()>,
30}
31
32struct State {
33    available: bool,
34    cursor: Option<TxId>,
35    objects: HashMap<AccessId, StoredAccess>,
36    targets: HashMap<Target, AccessId>,
37}
38
39struct Prepared {
40    outcome: ApplyOutcome,
41    mutation: StoreMutation,
42    delta: Delta,
43}
44
45enum Delta {
46    None,
47    Create(AccessId, Target, StoredAccess),
48    Replace(AccessId, StoredAccess),
49}
50
51fn locked<T>(result: LockResult<T>, message: &str) -> Result<T, String> {
52    result.map_err(|_| message.to_string())
53}
54
55fn reserve<K: Eq + Hash, V>(
56    map: &mut HashMap<K, V>,
57    additional: usize,
58    kind: &str,
59) -> Result<(), String> {
60    map.try_reserve(additional)
61        .map_err(|error| format!("unable to reserve {kind} index: {error}"))
62}
63
64impl State {
65    fn finish(&mut self, callback_txid: TxId, delta: Delta) -> Result<(), String> {
66        let error = match delta {
67            Delta::None => None,
68            Delta::Create(access_id, target, stored) => {
69                if self.objects.contains_key(&access_id) || self.targets.contains_key(&target) {
70                    Some("projection create contradiction after commit")
71                } else {
72                    self.targets.insert(target, access_id);
73                    self.objects.insert(access_id, stored);
74                    None
75                }
76            }
77            Delta::Replace(access_id, stored) => {
78                let consistent = self.targets.get(stored.target()) == Some(&access_id)
79                    && self
80                        .objects
81                        .get(&access_id)
82                        .is_some_and(|current| current.target() == stored.target());
83                if !consistent {
84                    Some("projection replace contradiction after commit")
85                } else if let Some(slot) = self.objects.get_mut(&access_id) {
86                    *slot = stored;
87                    None
88                } else {
89                    Some("projection replace disappeared after commit")
90                }
91            }
92        };
93        if let Some(error) = error {
94            self.available = false;
95            return Err(error.to_string());
96        }
97        self.cursor = Some(callback_txid);
98        Ok(())
99    }
100}
101
102impl Projection {
103    pub fn open(root: &Path, ordering: &K1TxnOrdering) -> Result<(Self, Option<TxId>), String> {
104        let started = Instant::now();
105        let result = Self::open_inner(root, ordering);
106        let elapsed = started.elapsed();
107        if elapsed > Duration::from_millis(100) {
108            eprintln!(
109                "level=warn module=kcode-k1-access-projection operation=open elapsed_us={} outcome={}",
110                elapsed.as_micros(),
111                if result.is_ok() { "ready" } else { "error" }
112            );
113        }
114        result
115    }
116
117    fn open_inner(root: &Path, ordering: &K1TxnOrdering) -> Result<(Self, Option<TxId>), String> {
118        let store = Store::open(root)?;
119        let (cursor, rows) = store.snapshot()?.into_parts();
120        let subsystem = SubsystemId::from_str("k1-access-subsystem")?;
121        if !Self::valid_snapshot(ordering, subsystem, cursor, &rows)? {
122            return Self::reset(store);
123        }
124        let mut objects = HashMap::new();
125        let mut targets = HashMap::new();
126        reserve(&mut objects, rows.len(), "access")?;
127        reserve(&mut targets, rows.len(), "target")?;
128        for row in rows {
129            let access_id = row.access_id();
130            let target = row.target().clone();
131            if objects.insert(access_id, row).is_some()
132                || targets.insert(target, access_id).is_some()
133            {
134                return Self::reset(store);
135            }
136        }
137        let projection = Self::new(store, cursor, objects, targets);
138        Ok((projection, cursor))
139    }
140
141    fn new(
142        store: Store,
143        cursor: Option<TxId>,
144        objects: HashMap<AccessId, StoredAccess>,
145        targets: HashMap<Target, AccessId>,
146    ) -> Self {
147        Self {
148            store,
149            state: RwLock::new(State {
150                available: true,
151                cursor,
152                objects,
153                targets,
154            }),
155            apply: Mutex::new(()),
156        }
157    }
158
159    fn empty(store: Store) -> Self {
160        Self::new(store, None, HashMap::new(), HashMap::new())
161    }
162
163    fn reset(store: Store) -> Result<(Self, Option<TxId>), String> {
164        store.clear()?;
165        Ok((Self::empty(store), None))
166    }
167
168    fn valid_snapshot(
169        ordering: &K1TxnOrdering,
170        subsystem: SubsystemId,
171        cursor: Option<TxId>,
172        rows: &[StoredAccess],
173    ) -> Result<bool, String> {
174        if !rows.is_empty() && cursor.is_none() {
175            return Ok(false);
176        }
177        if let Some(cursor) = cursor
178            && Self::load_action(ordering, subsystem, cursor)?.is_none()
179        {
180            return Ok(false);
181        }
182        for row in rows {
183            let access_id = row.access_id();
184            let Some(create) = Self::load_action(ordering, subsystem, access_id.txid())? else {
185                return Ok(false);
186            };
187            let created = match create {
188                AccessAction::Create {
189                    target,
190                    authorizations,
191                } if &target == row.target() => authorizations,
192                _ => return Ok(false),
193            };
194            if row.revision() == access_id.txid() {
195                if &created != row.authorizations() {
196                    return Ok(false);
197                }
198            } else if !matches!(
199                Self::load_action(ordering, subsystem, row.revision())?,
200                Some(AccessAction::Replace {
201                    access_id: replaced,
202                    authorizations,
203                    ..
204                }) if replaced == access_id && &authorizations == row.authorizations()
205            ) {
206                return Ok(false);
207            }
208        }
209        Ok(true)
210    }
211
212    fn load_action(
213        ordering: &K1TxnOrdering,
214        subsystem: SubsystemId,
215        txid: TxId,
216    ) -> Result<Option<AccessAction>, String> {
217        let Some(bytes) = ordering
218            .get_txn(txid)
219            .map_err(|error| format!("KTO transaction lookup failed: {error}"))?
220        else {
221            return Ok(None);
222        };
223        let Ok(transaction) = Transaction::parse(&bytes) else {
224            return Ok(None);
225        };
226        if transaction.subsystem() != subsystem {
227            return Ok(None);
228        }
229        Ok(kcode_k1_access_format::decode(transaction.payload())
230            .ok()
231            .map(|(_, action)| action))
232    }
233
234    pub fn apply(&self, callback_txid: TxId, action: AccessAction) -> Result<ApplyOutcome, String> {
235        let _lane = locked(self.apply.lock(), "projection apply lock poisoned")?;
236        let Prepared {
237            outcome,
238            mutation,
239            delta,
240        } = {
241            let mut state = locked(self.state.write(), "projection state lock poisoned")?;
242            if !state.available {
243                return Err("projection unavailable".to_string());
244            }
245            Self::prepare(&mut state, callback_txid, action)?
246        };
247        if let Err(error) = self.store.commit(callback_txid, &mutation) {
248            if let Ok(mut state) = self.state.write() {
249                state.available = false;
250            }
251            return Err(error);
252        }
253        let mut state = locked(
254            self.state.write(),
255            "projection state lock poisoned after commit",
256        )?;
257        if !state.available {
258            return Err("projection became unavailable after commit".to_string());
259        }
260        state.finish(callback_txid, delta)?;
261        Ok(outcome)
262    }
263
264    fn prepare(state: &mut State, txid: TxId, action: AccessAction) -> Result<Prepared, String> {
265        match action {
266            AccessAction::Create {
267                target,
268                authorizations,
269            } => {
270                let access_id = AccessId::new(txid);
271                if state.objects.contains_key(&access_id) {
272                    return Ok(Self::rejected("access ID already exists"));
273                }
274                if state.targets.contains_key(&target) {
275                    return Ok(Self::rejected("target already exists"));
276                }
277                reserve(&mut state.objects, 1, "access")?;
278                reserve(&mut state.targets, 1, "target")?;
279                let store =
280                    StoredAccess::new(access_id, target.clone(), txid, authorizations.clone());
281                let state = StoredAccess::new(access_id, target.clone(), txid, authorizations);
282                Ok(Prepared {
283                    outcome: ApplyOutcome::Applied(AccessRevision::new(access_id, txid)),
284                    mutation: StoreMutation::Create(store),
285                    delta: Delta::Create(access_id, target, state),
286                })
287            }
288            AccessAction::Replace {
289                access_id,
290                actor,
291                groups_revision: _,
292                witness,
293                authorizations,
294            } => {
295                let Some(current) = state.objects.get(&access_id) else {
296                    return Ok(Self::rejected("unknown access ID"));
297                };
298                if !Self::valid_witness(current.authorizations(), actor, &witness) {
299                    return Ok(Self::rejected("owner witness is invalid"));
300                }
301                if current.authorizations() == &authorizations {
302                    return Ok(Prepared {
303                        outcome: ApplyOutcome::Unchanged(AccessRevision::new(
304                            access_id,
305                            current.revision(),
306                        )),
307                        mutation: StoreMutation::CursorOnly,
308                        delta: Delta::None,
309                    });
310                }
311                let stored = StoredAccess::new(
312                    access_id,
313                    current.target().clone(),
314                    txid,
315                    authorizations.clone(),
316                );
317                Ok(Prepared {
318                    outcome: ApplyOutcome::Applied(AccessRevision::new(access_id, txid)),
319                    mutation: StoreMutation::Replace {
320                        access_id,
321                        revision: txid,
322                        authorizations,
323                    },
324                    delta: Delta::Replace(access_id, stored),
325                })
326            }
327        }
328    }
329
330    fn rejected(reason: &str) -> Prepared {
331        Prepared {
332            outcome: ApplyOutcome::Rejected(reason.to_string()),
333            mutation: StoreMutation::CursorOnly,
334            delta: Delta::None,
335        }
336    }
337
338    fn valid_witness(auth: &Authorizations, actor: UserId, witness: &OwnerWitness) -> bool {
339        auth.owners().iter().any(|owner| match (owner, witness) {
340            (OwnerSubject::User(owner), OwnerWitness::User) => *owner == actor,
341            (OwnerSubject::Group(owner), OwnerWitness::Group(group)) => owner == group,
342            _ => false,
343        })
344    }
345
346    fn readable(&self) -> Result<RwLockReadGuard<'_, State>, String> {
347        let state = locked(self.state.read(), "projection state lock poisoned")?;
348        if state.available {
349            Ok(state)
350        } else {
351            Err("projection unavailable".to_string())
352        }
353    }
354
355    pub fn owner_witness(
356        &self,
357        access_id: AccessId,
358        user: UserId,
359        user_groups: &[GroupId],
360    ) -> Result<Option<OwnerWitness>, String> {
361        let state = self.readable()?;
362        let Some(stored) = state.objects.get(&access_id) else {
363            return Ok(None);
364        };
365        let owners = stored.authorizations().owners();
366        if owners
367            .iter()
368            .any(|owner| matches!(owner, OwnerSubject::User(owner) if *owner == user))
369        {
370            return Ok(Some(OwnerWitness::User));
371        }
372        Ok(owners
373            .iter()
374            .filter_map(|owner| match owner {
375                OwnerSubject::Group(group) if user_groups.contains(group) => Some(*group),
376                _ => None,
377            })
378            .min()
379            .map(OwnerWitness::Group))
380    }
381
382    pub fn check(
383        &self,
384        principal: RequestPrincipal,
385        access_id: AccessId,
386        expected_subsystem: SubsystemId,
387        user_groups: &[GroupId],
388        model_groups: &[GroupId],
389        groups_revision: Option<TxId>,
390    ) -> Result<AccessCheck, String> {
391        let state = self.readable()?;
392        let Some(stored) = state.objects.get(&access_id) else {
393            return Self::hidden_check();
394        };
395        if stored.target().subsystem() != expected_subsystem {
396            return Self::hidden_check();
397        }
398        let auth = stored.authorizations();
399        let user_owner = Self::user_owner(auth, principal.user(), user_groups);
400        let user_view = user_owner
401            || auth.viewers().iter().any(|viewer| match viewer {
402                ViewerSubject::User(user) => *user == principal.user(),
403                ViewerSubject::Group(group) => user_groups.contains(group),
404                ViewerSubject::Model(_) => false,
405            });
406        let model_view = auth.owners().iter().any(|owner| match owner {
407            OwnerSubject::Group(group) => model_groups.contains(group),
408            OwnerSubject::User(_) => false,
409        }) || auth.viewers().iter().any(|viewer| match viewer {
410            ViewerSubject::Model(model) => *model == principal.model(),
411            ViewerSubject::Group(group) => model_groups.contains(group),
412            ViewerSubject::User(_) => false,
413        });
414        let can_manage = user_owner;
415        let can_view = user_view && model_view;
416        let evidence = can_view || can_manage;
417        AccessCheck::new(
418            can_view,
419            can_manage,
420            can_view.then(|| stored.target().clone()),
421            evidence.then_some(stored.revision()),
422            evidence.then_some(groups_revision).flatten(),
423        )
424    }
425
426    fn user_owner(auth: &Authorizations, user: UserId, groups: &[GroupId]) -> bool {
427        auth.owners().iter().any(|owner| match owner {
428            OwnerSubject::User(owner) => *owner == user,
429            OwnerSubject::Group(group) => groups.contains(group),
430        })
431    }
432
433    fn hidden_check() -> Result<AccessCheck, String> {
434        AccessCheck::new(false, false, None, None, None)
435    }
436
437    pub fn clear(&self) -> Result<(), String> {
438        let _lane = locked(self.apply.lock(), "projection apply lock poisoned")?;
439        {
440            let mut state = locked(self.state.write(), "projection state lock poisoned")?;
441            state.available = false;
442            state.cursor = None;
443            state.objects.clear();
444            state.targets.clear();
445        }
446        self.store.clear()
447    }
448}
449
450#[cfg(test)]
451mod tests {
452    use super::*;
453    use std::fs;
454    #[test]
455    fn owner_semantics_canary() -> Result<(), String> {
456        let id = |byte| TxId::from_bytes([byte; 12]);
457        let base = std::env::temp_dir().join(format!("k1-access-canary-{}", std::process::id()));
458        let _ = fs::remove_dir_all(&base);
459        let ordering = K1TxnOrdering::open(&base.join("ordering"))?;
460        let (projection, _) = Projection::open(&base.join("projection"), &ordering)?;
461        let direct = UserId::from_tx_id(id(1));
462        let outsider = UserId::from_tx_id(id(2));
463        let low = GroupId::new(id(3));
464        let high = GroupId::new(id(4));
465        let model = ModelId::from_bytes([5; 32]);
466        let auth = Authorizations::new(
467            vec![
468                OwnerSubject::User(direct),
469                OwnerSubject::Group(high),
470                OwnerSubject::Group(low),
471            ],
472            vec![ViewerSubject::User(outsider)],
473        )?;
474        let access_id = AccessId::new(id(6));
475        let subsystem = SubsystemId::from_str("canary")?;
476        projection.apply(
477            id(6),
478            AccessAction::Create {
479                target: Target::new(subsystem, vec![7]),
480                authorizations: auth.clone(),
481            },
482        )?;
483        assert_eq!(
484            projection.owner_witness(access_id, direct, &[high, low])?,
485            Some(OwnerWitness::User)
486        );
487        assert_eq!(
488            projection.owner_witness(access_id, outsider, &[high, low])?,
489            Some(OwnerWitness::Group(low))
490        );
491        assert!(Projection::valid_witness(
492            &auth,
493            direct,
494            &OwnerWitness::User
495        ));
496        assert!(!Projection::valid_witness(
497            &auth,
498            outsider,
499            &OwnerWitness::User
500        ));
501        assert!(!Projection::valid_witness(
502            &auth,
503            outsider,
504            &OwnerWitness::Group(GroupId::new(id(8)))
505        ));
506        let view = projection.check(
507            RequestPrincipal::new(outsider, model),
508            access_id,
509            subsystem,
510            &[],
511            &[low],
512            None,
513        )?;
514        assert!(view.can_view() && !view.can_manage());
515        let manage = projection.check(
516            RequestPrincipal::new(direct, model),
517            access_id,
518            subsystem,
519            &[],
520            &[],
521            None,
522        )?;
523        assert!(!manage.can_view() && manage.can_manage());
524        drop(projection);
525        drop(ordering);
526        fs::remove_dir_all(base).map_err(|error| error.to_string())
527    }
528}