use std::{
collections::HashMap,
path::Path,
sync::{LockResult, Mutex, RwLock, RwLockReadGuard},
};
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_discovery_store::DiscoveryStore;
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,
discovery: DiscoveryStore,
state: RwLock<State>,
apply: Mutex<()>,
}
struct State {
available: bool,
objects: HashMap<AccessId, StoredAccess>,
targets: HashMap<Target, AccessId>,
}
struct Prepared {
outcome: ApplyOutcome,
mutation: StoreMutation,
delta: Delta,
discovery: Option<(AccessId, Authorizations)>,
}
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_owned())
}
fn subjects(authorizations: &Authorizations) -> (Vec<UserId>, Vec<GroupId>) {
let mut users = Vec::new();
let mut groups = Vec::new();
for owner in authorizations.owners() {
match owner {
OwnerSubject::User(user) => users.push(*user),
OwnerSubject::Group(group) => groups.push(*group),
}
}
for viewer in authorizations.viewers() {
match viewer {
ViewerSubject::User(user) => users.push(*user),
ViewerSubject::Group(group) => groups.push(*group),
ViewerSubject::Model(_) => {}
}
}
(users, groups)
}
impl State {
fn finish(&mut self, delta: Delta) -> Result<(), String> {
let failure = 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(|row| row.target() == stored.target());
if !consistent {
Some("projection replace contradiction after commit")
} else if let Some(row) = self.objects.get_mut(&access_id) {
*row = stored;
None
} else {
Some("projection replace disappeared after commit")
}
}
};
if let Some(failure) = failure {
self.available = false;
return Err(failure.to_owned());
}
Ok(())
}
}
impl Projection {
pub fn open(root: &Path, ordering: &K1TxnOrdering) -> Result<(Self, Option<TxId>), String> {
let store = Store::open(root)?;
let discovery = DiscoveryStore::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)?
|| discovery.checkpoint()? != cursor
{
return Self::reset(store, discovery, root);
}
for row in &rows {
let (users, groups) = subjects(row.authorizations());
if discovery.contains_missing(&users, &groups, row.access_id())? {
return Self::reset(store, discovery, root);
}
}
let mut objects = HashMap::new();
let mut targets = HashMap::new();
objects
.try_reserve(rows.len())
.map_err(|error| format!("unable to reserve access index: {error}"))?;
targets
.try_reserve(rows.len())
.map_err(|error| format!("unable to reserve target index: {error}"))?;
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, discovery, root);
}
}
Ok((Self::new(store, discovery, objects, targets), cursor))
}
fn new(
store: Store,
discovery: DiscoveryStore,
objects: HashMap<AccessId, StoredAccess>,
targets: HashMap<Target, AccessId>,
) -> Self {
Self {
store,
discovery,
state: RwLock::new(State {
available: true,
objects,
targets,
}),
apply: Mutex::new(()),
}
}
fn reset(
store: Store,
discovery: DiscoveryStore,
root: &Path,
) -> Result<(Self, Option<TxId>), String> {
store.clear()?;
discovery.discard()?;
let discovery = DiscoveryStore::open(root)?;
Ok((
Self::new(store, discovery, HashMap::new(), HashMap::new()),
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 valid = match create {
AccessAction::Create {
target,
authorizations,
} if &target == row.target() => {
if row.revision() == access_id.txid() {
&authorizations == row.authorizations()
} else {
matches!(Self::load_action(ordering, subsystem, row.revision())?, Some(AccessAction::Replace { access_id: replaced, authorizations, .. }) if replaced == access_id && &authorizations == row.authorizations())
}
}
_ => false,
};
if !valid {
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 = {
let mut state = locked(self.state.write(), "projection state lock poisoned")?;
if !state.available {
return Err("projection unavailable".to_owned());
}
Self::prepare(&mut state, callback_txid, action)?
};
if let Err(error) = self.store.commit(callback_txid, &prepared.mutation) {
self.unavailable();
return Err(error);
}
let committed = if let Some((access_id, authorizations)) = &prepared.discovery {
let (users, groups) = subjects(authorizations);
self.discovery
.commit(callback_txid, &users, &groups, *access_id)
} else {
self.discovery
.commit(callback_txid, &[], &[], AccessId::new(callback_txid))
};
if let Err(error) = committed {
self.unavailable();
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_owned());
}
state.finish(prepared.delta)?;
Ok(prepared.outcome)
}
fn unavailable(&self) {
if let Ok(mut state) = self.state.write() {
state.available = false;
}
}
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"));
}
state
.objects
.try_reserve(1)
.map_err(|error| format!("unable to reserve access index: {error}"))?;
state
.targets
.try_reserve(1)
.map_err(|error| format!("unable to reserve target index: {error}"))?;
let stored =
StoredAccess::new(access_id, target.clone(), txid, authorizations.clone());
Ok(Prepared {
outcome: ApplyOutcome::Applied(AccessRevision::new(access_id, txid)),
mutation: StoreMutation::Create(stored.clone()),
delta: Delta::Create(access_id, target, stored),
discovery: Some((access_id, authorizations)),
})
}
AccessAction::Replace {
access_id,
actor,
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,
discovery: Some((access_id, authorizations)),
});
}
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: authorizations.clone(),
},
delta: Delta::Replace(access_id, stored),
discovery: Some((access_id, authorizations)),
})
}
AccessAction::EnsureDiscovery { access_id } => {
let Some(stored) = state.objects.get(&access_id) else {
return Ok(Self::rejected("unknown access ID"));
};
Ok(Prepared {
outcome: ApplyOutcome::Unchanged(AccessRevision::new(
access_id,
stored.revision(),
)),
mutation: StoreMutation::CursorOnly,
delta: Delta::None,
discovery: Some((access_id, stored.authorizations().clone())),
})
}
}
}
fn rejected(reason: &str) -> Prepared {
Prepared {
outcome: ApplyOutcome::Rejected(reason.to_owned()),
mutation: StoreMutation::CursorOnly,
delta: Delta::None,
discovery: 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_owned())
}
}
pub fn discovery_missing(
&self,
access_id: AccessId,
expected_subsystem: SubsystemId,
) -> Result<bool, String> {
let state = self.readable()?;
let Some(stored) = state.objects.get(&access_id) else {
return Ok(false);
};
if stored.target().subsystem() != expected_subsystem {
return Ok(false);
}
let (users, groups) = subjects(stored.authorizations());
self.discovery.contains_missing(&users, &groups, access_id)
}
pub fn discovered_for_user(
&self,
user: UserId,
subsystem: SubsystemId,
) -> Result<Vec<AccessId>, String> {
self.discovered(self.discovery.list_user(user)?, subsystem)
}
pub fn discovered_for_group(
&self,
group: GroupId,
subsystem: SubsystemId,
) -> Result<Vec<AccessId>, String> {
self.discovered(self.discovery.list_group(group)?, subsystem)
}
fn discovered(
&self,
ids: Vec<AccessId>,
subsystem: SubsystemId,
) -> Result<Vec<AccessId>, String> {
let state = self.readable()?;
Ok(ids
.into_iter()
.filter(|id| {
state
.objects
.get(id)
.is_some_and(|stored| stored.target().subsystem() == subsystem)
})
.collect())
}
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")?;
self.unavailable();
self.store.clear().and(self.discovery.discard())
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::{
fs,
sync::atomic::{AtomicU64, Ordering},
};
static NEXT: AtomicU64 = AtomicU64::new(0);
fn root() -> std::path::PathBuf {
let root = std::env::temp_dir().join(format!(
"access-projection-{}-{}",
std::process::id(),
NEXT.fetch_add(1, Ordering::Relaxed)
));
fs::create_dir(&root).unwrap();
root
}
fn tx(value: u8) -> TxId {
TxId::from_bytes([value; 12])
}
fn user(value: u8) -> UserId {
UserId::from_tx_id(tx(value))
}
fn group(value: u8) -> GroupId {
GroupId::new(tx(value))
}
fn target(name: &str) -> Target {
Target::new(SubsystemId::from_str(name).unwrap(), vec![1])
}
fn auth() -> Authorizations {
Authorizations::new(
vec![OwnerSubject::User(user(1)), OwnerSubject::Group(group(2))],
vec![
ViewerSubject::User(user(3)),
ViewerSubject::Group(group(4)),
ViewerSubject::Model(ModelId::from_bytes([5; 32])),
],
)
.unwrap()
}
#[test]
fn create_replace_ensure_queries_and_authorization_fan_out_without_duplicates()
-> Result<(), String> {
let root = root();
let ordering = K1TxnOrdering::open(&root.join("ordering"))?;
let (projection, _) = Projection::open(&root.join("projection"), &ordering)?;
let access_id = AccessId::new(tx(6));
let subsystem = SubsystemId::from_str("one")?;
assert_eq!(
projection.apply(
tx(6),
AccessAction::Create {
target: target("one"),
authorizations: auth()
}
)?,
ApplyOutcome::Applied(AccessRevision::new(access_id, tx(6)))
);
assert_eq!(
projection.discovered_for_user(user(1), subsystem)?,
vec![access_id]
);
assert_eq!(
projection.discovered_for_group(group(4), subsystem)?,
vec![access_id]
);
assert!(!projection.discovery_missing(access_id, subsystem)?);
assert_eq!(
projection.apply(tx(7), AccessAction::EnsureDiscovery { access_id })?,
ApplyOutcome::Unchanged(AccessRevision::new(access_id, tx(6)))
);
assert_eq!(
projection.discovered_for_user(user(1), subsystem)?,
vec![access_id]
);
let replacement = Authorizations::new(
vec![OwnerSubject::User(user(3))],
vec![ViewerSubject::Group(group(2))],
)?;
assert_eq!(
projection.apply(
tx(8),
AccessAction::Replace {
access_id,
actor: user(1),
groups_revision: None,
witness: OwnerWitness::User,
authorizations: replacement
}
)?,
ApplyOutcome::Applied(AccessRevision::new(access_id, tx(8)))
);
assert_eq!(
projection.discovered_for_user(user(3), subsystem)?,
vec![access_id]
);
assert_eq!(
projection.discovered_for_group(group(2), subsystem)?,
vec![access_id]
);
assert!(
projection
.discovered_for_user(user(1), SubsystemId::from_str("two")?)?
.is_empty()
);
assert_eq!(
projection.apply(
tx(9),
AccessAction::EnsureDiscovery {
access_id: AccessId::new(tx(9))
}
)?,
ApplyOutcome::Rejected("unknown access ID".to_owned())
);
let view = projection.check(
RequestPrincipal::new(user(1), ModelId::from_bytes([5; 32])),
access_id,
subsystem,
&[group(2)],
&[group(2)],
None,
)?;
assert!(view.can_view() && !view.can_manage());
assert_eq!(
projection.owner_witness(access_id, user(3), &[group(2)])?,
Some(OwnerWitness::User)
);
fs::remove_dir_all(root).map_err(|error| error.to_string())
}
#[test]
fn mismatched_discovery_checkpoint_resets_both_derived_stores() -> Result<(), String> {
let root = root();
let projection_root = root.join("projection");
let ordering = K1TxnOrdering::open(&root.join("ordering"))?;
kcode_k1_access_discovery_store::DiscoveryStore::open(&projection_root)?.commit(
tx(1),
&[user(2)],
&[],
AccessId::new(tx(3)),
)?;
let (projection, cursor) = Projection::open(&projection_root, &ordering)?;
assert_eq!(cursor, None);
assert!(
projection
.discovered_for_user(user(2), SubsystemId::from_str("one")?)?
.is_empty()
);
projection.clear()?;
fs::remove_dir_all(root).map_err(|error| error.to_string())
}
}