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}