use std::{
collections::HashMap,
hash::Hash,
path::Path,
sync::{LockResult, Mutex, RwLock, RwLockReadGuard},
time::{Duration, Instant},
};
pub use kcode_k1_access_format::{AccessAction, OwnerWitness};
pub use kcode_k1_access_types::{
AccessCheck, AccessId, AccessRevision, Authorizations, GroupId, ModelId, OwnerSubject,
RequestPrincipal, SubsystemId, Target, TxId, UserId, ViewerSubject,
};
pub use kcode_k1_txn_ordering::K1TxnOrdering;
use kcode_k1_access_store::{Store, StoreMutation, StoredAccess};
use kcode_k1_transaction::Transaction;
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum ApplyOutcome {
Applied(AccessRevision),
Unchanged(AccessRevision),
Rejected(String),
}
pub struct Projection {
store: Store,
state: RwLock<State>,
apply: Mutex<()>,
}
struct State {
available: bool,
cursor: Option<TxId>,
objects: HashMap<AccessId, StoredAccess>,
targets: HashMap<Target, AccessId>,
}
struct Prepared {
outcome: ApplyOutcome,
mutation: StoreMutation,
delta: Delta,
}
enum Delta {
None,
Create(AccessId, Target, StoredAccess),
Replace(AccessId, StoredAccess),
}
fn locked<T>(result: LockResult<T>, message: &str) -> Result<T, String> {
result.map_err(|_| message.to_string())
}
fn reserve<K: Eq + Hash, V>(
map: &mut HashMap<K, V>,
additional: usize,
kind: &str,
) -> Result<(), String> {
map.try_reserve(additional)
.map_err(|error| format!("unable to reserve {kind} index: {error}"))
}
impl State {
fn finish(&mut self, callback_txid: TxId, delta: Delta) -> Result<(), String> {
let error = match delta {
Delta::None => None,
Delta::Create(access_id, target, stored) => {
if self.objects.contains_key(&access_id) || self.targets.contains_key(&target) {
Some("projection create contradiction after commit")
} else {
self.targets.insert(target, access_id);
self.objects.insert(access_id, stored);
None
}
}
Delta::Replace(access_id, stored) => {
let consistent = self.targets.get(stored.target()) == Some(&access_id)
&& self
.objects
.get(&access_id)
.is_some_and(|current| current.target() == stored.target());
if !consistent {
Some("projection replace contradiction after commit")
} else if let Some(slot) = self.objects.get_mut(&access_id) {
*slot = stored;
None
} else {
Some("projection replace disappeared after commit")
}
}
};
if let Some(error) = error {
self.available = false;
return Err(error.to_string());
}
self.cursor = Some(callback_txid);
Ok(())
}
}
impl Projection {
pub fn open(root: &Path, ordering: &K1TxnOrdering) -> Result<(Self, Option<TxId>), String> {
let started = Instant::now();
let result = Self::open_inner(root, ordering);
let elapsed = started.elapsed();
if elapsed > Duration::from_millis(100) {
eprintln!(
"level=warn module=kcode-k1-access-projection operation=open elapsed_us={} outcome={}",
elapsed.as_micros(),
if result.is_ok() { "ready" } else { "error" }
);
}
result
}
fn open_inner(root: &Path, ordering: &K1TxnOrdering) -> Result<(Self, Option<TxId>), String> {
let store = Store::open(root)?;
let (cursor, rows) = store.snapshot()?.into_parts();
let subsystem = SubsystemId::from_str("k1-access-subsystem")?;
if !Self::valid_snapshot(ordering, subsystem, cursor, &rows)? {
return Self::reset(store);
}
let mut objects = HashMap::new();
let mut targets = HashMap::new();
reserve(&mut objects, rows.len(), "access")?;
reserve(&mut targets, rows.len(), "target")?;
for row in rows {
let access_id = row.access_id();
let target = row.target().clone();
if objects.insert(access_id, row).is_some()
|| targets.insert(target, access_id).is_some()
{
return Self::reset(store);
}
}
let projection = Self::new(store, cursor, objects, targets);
Ok((projection, cursor))
}
fn new(
store: Store,
cursor: Option<TxId>,
objects: HashMap<AccessId, StoredAccess>,
targets: HashMap<Target, AccessId>,
) -> Self {
Self {
store,
state: RwLock::new(State {
available: true,
cursor,
objects,
targets,
}),
apply: Mutex::new(()),
}
}
fn empty(store: Store) -> Self {
Self::new(store, None, HashMap::new(), HashMap::new())
}
fn reset(store: Store) -> Result<(Self, Option<TxId>), String> {
store.clear()?;
Ok((Self::empty(store), None))
}
fn valid_snapshot(
ordering: &K1TxnOrdering,
subsystem: SubsystemId,
cursor: Option<TxId>,
rows: &[StoredAccess],
) -> Result<bool, String> {
if !rows.is_empty() && cursor.is_none() {
return Ok(false);
}
if let Some(cursor) = cursor
&& Self::load_action(ordering, subsystem, cursor)?.is_none()
{
return Ok(false);
}
for row in rows {
let access_id = row.access_id();
let Some(create) = Self::load_action(ordering, subsystem, access_id.txid())? else {
return Ok(false);
};
let created = match create {
AccessAction::Create {
target,
authorizations,
} if &target == row.target() => authorizations,
_ => return Ok(false),
};
if row.revision() == access_id.txid() {
if &created != row.authorizations() {
return Ok(false);
}
} else if !matches!(
Self::load_action(ordering, subsystem, row.revision())?,
Some(AccessAction::Replace {
access_id: replaced,
authorizations,
..
}) if replaced == access_id && &authorizations == row.authorizations()
) {
return Ok(false);
}
}
Ok(true)
}
fn load_action(
ordering: &K1TxnOrdering,
subsystem: SubsystemId,
txid: TxId,
) -> Result<Option<AccessAction>, String> {
let Some(bytes) = ordering
.get_txn(txid)
.map_err(|error| format!("KTO transaction lookup failed: {error}"))?
else {
return Ok(None);
};
let Ok(transaction) = Transaction::parse(&bytes) else {
return Ok(None);
};
if transaction.subsystem() != subsystem {
return Ok(None);
}
Ok(kcode_k1_access_format::decode(transaction.payload())
.ok()
.map(|(_, action)| action))
}
pub fn apply(&self, callback_txid: TxId, action: AccessAction) -> Result<ApplyOutcome, String> {
let _lane = locked(self.apply.lock(), "projection apply lock poisoned")?;
let Prepared {
outcome,
mutation,
delta,
} = {
let mut state = locked(self.state.write(), "projection state lock poisoned")?;
if !state.available {
return Err("projection unavailable".to_string());
}
Self::prepare(&mut state, callback_txid, action)?
};
if let Err(error) = self.store.commit(callback_txid, &mutation) {
if let Ok(mut state) = self.state.write() {
state.available = false;
}
return Err(error);
}
let mut state = locked(
self.state.write(),
"projection state lock poisoned after commit",
)?;
if !state.available {
return Err("projection became unavailable after commit".to_string());
}
state.finish(callback_txid, delta)?;
Ok(outcome)
}
fn prepare(state: &mut State, txid: TxId, action: AccessAction) -> Result<Prepared, String> {
match action {
AccessAction::Create {
target,
authorizations,
} => {
let access_id = AccessId::new(txid);
if state.objects.contains_key(&access_id) {
return Ok(Self::rejected("access ID already exists"));
}
if state.targets.contains_key(&target) {
return Ok(Self::rejected("target already exists"));
}
reserve(&mut state.objects, 1, "access")?;
reserve(&mut state.targets, 1, "target")?;
let store =
StoredAccess::new(access_id, target.clone(), txid, authorizations.clone());
let state = StoredAccess::new(access_id, target.clone(), txid, authorizations);
Ok(Prepared {
outcome: ApplyOutcome::Applied(AccessRevision::new(access_id, txid)),
mutation: StoreMutation::Create(store),
delta: Delta::Create(access_id, target, state),
})
}
AccessAction::Replace {
access_id,
actor,
groups_revision: _,
witness,
authorizations,
} => {
let Some(current) = state.objects.get(&access_id) else {
return Ok(Self::rejected("unknown access ID"));
};
if !Self::valid_witness(current.authorizations(), actor, &witness) {
return Ok(Self::rejected("owner witness is invalid"));
}
if current.authorizations() == &authorizations {
return Ok(Prepared {
outcome: ApplyOutcome::Unchanged(AccessRevision::new(
access_id,
current.revision(),
)),
mutation: StoreMutation::CursorOnly,
delta: Delta::None,
});
}
let stored = StoredAccess::new(
access_id,
current.target().clone(),
txid,
authorizations.clone(),
);
Ok(Prepared {
outcome: ApplyOutcome::Applied(AccessRevision::new(access_id, txid)),
mutation: StoreMutation::Replace {
access_id,
revision: txid,
authorizations,
},
delta: Delta::Replace(access_id, stored),
})
}
}
}
fn rejected(reason: &str) -> Prepared {
Prepared {
outcome: ApplyOutcome::Rejected(reason.to_string()),
mutation: StoreMutation::CursorOnly,
delta: Delta::None,
}
}
fn valid_witness(auth: &Authorizations, actor: UserId, witness: &OwnerWitness) -> bool {
auth.owners().iter().any(|owner| match (owner, witness) {
(OwnerSubject::User(owner), OwnerWitness::User) => *owner == actor,
(OwnerSubject::Group(owner), OwnerWitness::Group(group)) => owner == group,
_ => false,
})
}
fn readable(&self) -> Result<RwLockReadGuard<'_, State>, String> {
let state = locked(self.state.read(), "projection state lock poisoned")?;
if state.available {
Ok(state)
} else {
Err("projection unavailable".to_string())
}
}
pub fn owner_witness(
&self,
access_id: AccessId,
user: UserId,
user_groups: &[GroupId],
) -> Result<Option<OwnerWitness>, String> {
let state = self.readable()?;
let Some(stored) = state.objects.get(&access_id) else {
return Ok(None);
};
let owners = stored.authorizations().owners();
if owners
.iter()
.any(|owner| matches!(owner, OwnerSubject::User(owner) if *owner == user))
{
return Ok(Some(OwnerWitness::User));
}
Ok(owners
.iter()
.filter_map(|owner| match owner {
OwnerSubject::Group(group) if user_groups.contains(group) => Some(*group),
_ => None,
})
.min()
.map(OwnerWitness::Group))
}
pub fn check(
&self,
principal: RequestPrincipal,
access_id: AccessId,
expected_subsystem: SubsystemId,
user_groups: &[GroupId],
model_groups: &[GroupId],
groups_revision: Option<TxId>,
) -> Result<AccessCheck, String> {
let state = self.readable()?;
let Some(stored) = state.objects.get(&access_id) else {
return Self::hidden_check();
};
if stored.target().subsystem() != expected_subsystem {
return Self::hidden_check();
}
let auth = stored.authorizations();
let user_owner = Self::user_owner(auth, principal.user(), user_groups);
let user_view = user_owner
|| auth.viewers().iter().any(|viewer| match viewer {
ViewerSubject::User(user) => *user == principal.user(),
ViewerSubject::Group(group) => user_groups.contains(group),
ViewerSubject::Model(_) => false,
});
let model_view = auth.owners().iter().any(|owner| match owner {
OwnerSubject::Group(group) => model_groups.contains(group),
OwnerSubject::User(_) => false,
}) || auth.viewers().iter().any(|viewer| match viewer {
ViewerSubject::Model(model) => *model == principal.model(),
ViewerSubject::Group(group) => model_groups.contains(group),
ViewerSubject::User(_) => false,
});
let can_manage = user_owner;
let can_view = user_view && model_view;
let evidence = can_view || can_manage;
AccessCheck::new(
can_view,
can_manage,
can_view.then(|| stored.target().clone()),
evidence.then_some(stored.revision()),
evidence.then_some(groups_revision).flatten(),
)
}
fn user_owner(auth: &Authorizations, user: UserId, groups: &[GroupId]) -> bool {
auth.owners().iter().any(|owner| match owner {
OwnerSubject::User(owner) => *owner == user,
OwnerSubject::Group(group) => groups.contains(group),
})
}
fn hidden_check() -> Result<AccessCheck, String> {
AccessCheck::new(false, false, None, None, None)
}
pub fn clear(&self) -> Result<(), String> {
let _lane = locked(self.apply.lock(), "projection apply lock poisoned")?;
{
let mut state = locked(self.state.write(), "projection state lock poisoned")?;
state.available = false;
state.cursor = None;
state.objects.clear();
state.targets.clear();
}
self.store.clear()
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::fs;
#[test]
fn owner_semantics_canary() -> Result<(), String> {
let id = |byte| TxId::from_bytes([byte; 12]);
let base = std::env::temp_dir().join(format!("k1-access-canary-{}", std::process::id()));
let _ = fs::remove_dir_all(&base);
let ordering = K1TxnOrdering::open(&base.join("ordering"))?;
let (projection, _) = Projection::open(&base.join("projection"), &ordering)?;
let direct = UserId::from_tx_id(id(1));
let outsider = UserId::from_tx_id(id(2));
let low = GroupId::new(id(3));
let high = GroupId::new(id(4));
let model = ModelId::from_bytes([5; 32]);
let auth = Authorizations::new(
vec![
OwnerSubject::User(direct),
OwnerSubject::Group(high),
OwnerSubject::Group(low),
],
vec![ViewerSubject::User(outsider)],
)?;
let access_id = AccessId::new(id(6));
let subsystem = SubsystemId::from_str("canary")?;
projection.apply(
id(6),
AccessAction::Create {
target: Target::new(subsystem, vec![7]),
authorizations: auth.clone(),
},
)?;
assert_eq!(
projection.owner_witness(access_id, direct, &[high, low])?,
Some(OwnerWitness::User)
);
assert_eq!(
projection.owner_witness(access_id, outsider, &[high, low])?,
Some(OwnerWitness::Group(low))
);
assert!(Projection::valid_witness(
&auth,
direct,
&OwnerWitness::User
));
assert!(!Projection::valid_witness(
&auth,
outsider,
&OwnerWitness::User
));
assert!(!Projection::valid_witness(
&auth,
outsider,
&OwnerWitness::Group(GroupId::new(id(8)))
));
let view = projection.check(
RequestPrincipal::new(outsider, model),
access_id,
subsystem,
&[],
&[low],
None,
)?;
assert!(view.can_view() && !view.can_manage());
let manage = projection.check(
RequestPrincipal::new(direct, model),
access_id,
subsystem,
&[],
&[],
None,
)?;
assert!(!manage.can_view() && manage.can_manage());
drop(projection);
drop(ordering);
fs::remove_dir_all(base).map_err(|error| error.to_string())
}
}