use std::collections::{BTreeMap, BTreeSet};
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex, RwLock};
use rand_core::RngCore;
use super::engine::{
AccessDecision, AccessPrincipal, AccessResource, evaluate_access, groups_for_subject,
principal_may_perform,
};
use super::model::{
ACTION_AGENT_VIEW, AccessConfigError, AccessControlConfig, AccessGroup, AccessRule,
validate_access_config,
};
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct AgentResourceAttributes {
pub identity: String,
pub agent_id: Option<String>,
pub role: Option<String>,
pub labels: BTreeMap<String, String>,
}
pub(crate) struct AccessEditPrecondition {
pub owner_instance: String,
pub expected_revision: u64,
}
pub(crate) enum AccessMutation {
Replace(AccessControlConfig),
UpsertRule(AccessRule),
DeleteRule(String),
SetGroup(String, AccessGroup),
DeleteGroup(String),
SetEnabled(bool),
}
impl AccessMutation {
fn apply(self, config: &mut AccessControlConfig) -> Result<(), AccessConfigError> {
match self {
Self::Replace(replacement) => *config = replacement,
Self::UpsertRule(rule) => {
match config
.rules
.iter_mut()
.find(|existing| existing.id == rule.id)
{
Some(existing) => *existing = rule,
None => config.rules.push(rule),
}
}
Self::DeleteRule(id) => {
let before = config.rules.len();
config.rules.retain(|rule| rule.id != id);
if config.rules.len() == before {
return Err(AccessConfigError::UnknownRule(id));
}
}
Self::SetGroup(name, group) => {
config.groups.insert(name, group);
}
Self::DeleteGroup(name) => {
config.groups.remove(&name);
}
Self::SetEnabled(enabled) => config.enabled = enabled,
}
Ok(())
}
}
#[derive(Debug)]
pub(crate) enum AccessEditError {
Denied,
Unavailable,
OwnerChanged,
RevisionConflict { expected: u64, actual: u64 },
Config(AccessConfigError),
}
#[cfg(test)]
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum CheckedMutationStage {
Attempt,
Acquired,
}
#[cfg(test)]
type CheckedMutationProbe = dyn Fn(CheckedMutationStage) + Send + Sync;
fn new_edit_identity() -> Option<String> {
let mut bytes = [0_u8; 16];
rand_core::OsRng.try_fill_bytes(&mut bytes).ok()?;
Some(format!("{:032x}", u128::from_be_bytes(bytes)))
}
struct AccessState {
config: Arc<AccessControlConfig>,
revision: u64,
}
struct AccessControllerInner {
owner_instance: Option<String>,
#[cfg(test)]
checked_probe: Mutex<Option<Arc<CheckedMutationProbe>>>,
state: RwLock<AccessState>,
persist_path: RwLock<Option<PathBuf>>,
attributes: RwLock<BTreeMap<String, Arc<AgentResourceAttributes>>>,
mutation: Mutex<()>,
}
#[derive(Clone)]
pub struct AccessController {
inner: Arc<AccessControllerInner>,
}
impl std::fmt::Debug for AccessController {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let (config, revision) = self.snapshot();
f.debug_struct("AccessController")
.field("enabled", &config.enabled)
.field("revision", &revision)
.field("rules", &config.rules.len())
.finish()
}
}
impl AccessController {
pub fn new(config: AccessControlConfig) -> Result<Self, AccessConfigError> {
Self::with_edit_identity(config, new_edit_identity())
}
fn with_edit_identity(
mut config: AccessControlConfig,
owner_instance: Option<String>,
) -> Result<Self, AccessConfigError> {
super::model::normalize_access_config_for_memory_actions(&mut config);
validate_access_config(&config)?;
Ok(Self {
inner: Arc::new(AccessControllerInner {
owner_instance,
#[cfg(test)]
checked_probe: Mutex::new(None),
state: RwLock::new(AccessState {
config: Arc::new(config),
revision: 0,
}),
persist_path: RwLock::new(None),
attributes: RwLock::new(BTreeMap::new()),
mutation: Mutex::new(()),
}),
})
}
pub fn disabled() -> Self {
Self::new(AccessControlConfig::default()).unwrap_or_else(|_| unreachable!())
}
pub fn load_or_default(path: impl Into<PathBuf>) -> Result<Self, AccessConfigError> {
let path = path.into();
let config = if path.is_file() {
let raw = std::fs::read_to_string(&path)
.map_err(|err| AccessConfigError::Io(err.to_string()))?;
toml::from_str::<AccessControlConfig>(&raw)
.map_err(|err| AccessConfigError::Parse(err.to_string()))?
} else {
AccessControlConfig::default()
};
let controller = Self::new(config)?;
*controller
.inner
.persist_path
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(path);
Ok(controller)
}
pub fn with_persist_path(self, path: impl Into<PathBuf>) -> Self {
*self
.inner
.persist_path
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(path.into());
self
}
pub fn snapshot(&self) -> (Arc<AccessControlConfig>, u64) {
let state = self
.inner
.state
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner);
(Arc::clone(&state.config), state.revision)
}
pub fn enabled(&self) -> bool {
self.snapshot().0.enabled
}
pub fn replace_config(&self, config: AccessControlConfig) -> Result<u64, AccessConfigError> {
self.mutate(AccessMutation::Replace(config))
}
pub fn upsert_rule(&self, rule: AccessRule) -> Result<u64, AccessConfigError> {
self.mutate(AccessMutation::UpsertRule(rule))
}
pub fn delete_rule(&self, rule_id: &str) -> Result<u64, AccessConfigError> {
self.mutate(AccessMutation::DeleteRule(rule_id.to_string()))
}
pub fn set_group(&self, name: &str, group: AccessGroup) -> Result<u64, AccessConfigError> {
self.mutate(AccessMutation::SetGroup(name.to_string(), group))
}
pub fn delete_group(&self, name: &str) -> Result<u64, AccessConfigError> {
self.mutate(AccessMutation::DeleteGroup(name.to_string()))
}
pub fn set_enabled(&self, enabled: bool) -> Result<u64, AccessConfigError> {
self.mutate(AccessMutation::SetEnabled(enabled))
}
pub(crate) fn edit_identity(&self) -> Option<&str> {
self.inner.owner_instance.as_deref()
}
pub(crate) fn admin_config(
&self,
subject: Option<&str>,
) -> Result<(Arc<AccessControlConfig>, u64), AccessEditError> {
let (view, revision) = self.view_and_revision(subject);
if !view.can_administer() {
return Err(AccessEditError::Denied);
}
Ok((view.config, revision))
}
pub(crate) fn admin_preview(
&self,
caller: Option<&str>,
target: Option<&str>,
) -> Result<AccessView, AccessEditError> {
let (config, _) = self.snapshot();
let caller_view = self.view_from_config(Arc::clone(&config), caller);
if !caller_view.can_administer() {
return Err(AccessEditError::Denied);
}
Ok(self.view_from_config(config, target))
}
pub(crate) fn mutate_admin(
&self,
subject: Option<&str>,
expected: Option<AccessEditPrecondition>,
mutation: AccessMutation,
) -> Result<u64, AccessEditError> {
#[cfg(test)]
self.observe_checked_mutation(CheckedMutationStage::Attempt);
let _mutation = self
.inner
.mutation
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
#[cfg(test)]
self.observe_checked_mutation(CheckedMutationStage::Acquired);
let (config, revision) = self.admin_config(subject)?;
if let Some(expected) = expected {
let instance = self.edit_identity().ok_or(AccessEditError::Unavailable)?;
if instance != expected.owner_instance {
return Err(AccessEditError::OwnerChanged);
}
if revision != expected.expected_revision {
return Err(AccessEditError::RevisionConflict {
expected: expected.expected_revision,
actual: revision,
});
}
}
self.edit_and_commit((*config).clone(), mutation)
.map_err(AccessEditError::Config)
}
#[cfg(test)]
fn observe_checked_mutation(&self, stage: CheckedMutationStage) {
let probe = self
.inner
.checked_probe
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
if let Some(probe) = probe {
probe(stage);
}
}
fn mutate(&self, mutation: AccessMutation) -> Result<u64, AccessConfigError> {
let _mutation = self
.inner
.mutation
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
self.edit_and_commit((*self.snapshot().0).clone(), mutation)
}
fn edit_and_commit(
&self,
mut config: AccessControlConfig,
mutation: AccessMutation,
) -> Result<u64, AccessConfigError> {
mutation.apply(&mut config)?;
super::model::normalize_access_config_for_memory_actions(&mut config);
validate_access_config(&config)?;
self.commit(config)
}
fn commit(&self, config: AccessControlConfig) -> Result<u64, AccessConfigError> {
let next_revision = self
.snapshot()
.1
.checked_add(1)
.ok_or(AccessConfigError::RevisionExhausted)?;
let persist_path = self
.inner
.persist_path
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
if let Some(path) = persist_path {
persist_config(&path, &config)?;
}
let mut state = self
.inner
.state
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner);
state.config = Arc::new(config);
state.revision = next_revision;
Ok(state.revision)
}
pub fn view_for_subject(&self, subject: Option<&str>) -> AccessView {
self.view_and_revision(subject).0
}
pub(crate) fn view_and_revision(&self, subject: Option<&str>) -> (AccessView, u64) {
let (config, revision) = self.snapshot();
(self.view_from_config(config, subject), revision)
}
fn view_from_config(
&self,
config: Arc<AccessControlConfig>,
subject: Option<&str>,
) -> AccessView {
let principal = match subject {
Some(subject) => AccessPrincipal {
subject: Some(subject.to_string()),
groups: groups_for_subject(&config, subject),
},
None => AccessPrincipal::anonymous(),
};
let is_admin = principal
.subject
.as_deref()
.is_some_and(|subject| config.admins.iter().any(|admin| admin == subject));
AccessView {
inner: Arc::clone(&self.inner),
config,
principal,
is_admin,
}
}
pub fn record_agent_attributes(&self, attributes: AgentResourceAttributes) {
if attributes.identity.is_empty() {
return;
}
let mut cache = self
.inner
.attributes
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner);
cache.insert(attributes.identity.clone(), Arc::new(attributes));
}
pub fn replace_agent_attributes(
&self,
attributes: impl IntoIterator<Item = AgentResourceAttributes>,
) {
let next: BTreeMap<String, Arc<AgentResourceAttributes>> = attributes
.into_iter()
.filter(|entry| !entry.identity.is_empty())
.map(|entry| (entry.identity.clone(), Arc::new(entry)))
.collect();
if next.is_empty() {
return;
}
*self
.inner
.attributes
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = next;
}
}
fn persist_config(path: &Path, config: &AccessControlConfig) -> Result<(), AccessConfigError> {
let rendered =
toml::to_string_pretty(config).map_err(|err| AccessConfigError::Parse(err.to_string()))?;
if let Some(parent) = path.parent()
&& !parent.as_os_str().is_empty()
{
std::fs::create_dir_all(parent).map_err(|err| AccessConfigError::Io(err.to_string()))?;
}
let header = "# MobKit access control. Managed by the console Access panel;\n# hand edits are preserved until the next console save.\n\n";
let mut tmp = path.to_path_buf();
let mut tmp_name = path
.file_name()
.map(std::ffi::OsString::from)
.ok_or_else(|| {
AccessConfigError::Io(format!(
"access config path has no file name: {}",
path.display()
))
})?;
tmp_name.push(".tmp");
tmp.set_file_name(tmp_name);
std::fs::write(&tmp, format!("{header}{rendered}"))
.map_err(|err| AccessConfigError::Io(err.to_string()))?;
std::fs::rename(&tmp, path).map_err(|err| AccessConfigError::Io(err.to_string()))
}
struct LineageLink {
identity: String,
attributes: Option<Arc<AgentResourceAttributes>>,
}
#[derive(Clone)]
pub struct AccessView {
inner: Arc<AccessControllerInner>,
config: Arc<AccessControlConfig>,
principal: AccessPrincipal,
is_admin: bool,
}
impl AccessView {
pub fn enforced(&self) -> bool {
self.config.enabled
}
pub fn subject(&self) -> Option<&str> {
self.principal.subject.as_deref()
}
pub fn groups(&self) -> &BTreeSet<String> {
&self.principal.groups
}
pub fn is_admin(&self) -> bool {
self.is_admin
}
pub fn decide(&self, action: &str, resource: &AccessResource<'_>) -> AccessDecision {
evaluate_access(&self.config, &self.principal, action, resource)
}
pub fn allows(&self, action: &str) -> bool {
self.decide(action, &AccessResource::none()).is_allow()
}
pub fn may_perform_anywhere(&self, action: &str) -> bool {
self.is_admin || principal_may_perform(&self.config, &self.principal, action)
}
pub fn allows_agent(&self, action: &str, identity: &str) -> bool {
self.decide_agent(action, identity).is_allow()
}
pub fn decide_agent(&self, action: &str, identity: &str) -> AccessDecision {
if !self.config.enabled {
return AccessDecision::Allow;
}
let lineage = self.lineage_for(identity);
if lineage.is_empty() {
return self.decide(action, &AccessResource::for_identity(identity));
}
self.decide_agent_lineage(action, identity, &lineage)
}
pub(crate) fn decide_agent_with_attributes(
&self,
action: &str,
attributes: &AgentResourceAttributes,
) -> AccessDecision {
if !self.config.enabled {
return AccessDecision::Allow;
}
let lineage = self.lineage_for_attributes(attributes);
self.decide_agent_lineage(action, attributes.identity.as_str(), &lineage)
}
fn decide_agent_lineage(
&self,
action: &str,
fallback_identity: &str,
lineage: &[LineageLink],
) -> AccessDecision {
let resources = lineage
.iter()
.enumerate()
.map(|(index, link)| match link.attributes.as_deref() {
Some(attributes) => AccessResource {
identity: Some(attributes.identity.as_str()),
agent_id: attributes
.agent_id
.as_deref()
.or((index == 0).then_some(fallback_identity)),
role: attributes.role.as_deref(),
labels: Some(&attributes.labels),
},
None => AccessResource::for_identity(link.identity.as_str()),
})
.collect::<Vec<_>>();
super::engine::evaluate_access_lineage(&self.config, &self.principal, action, &resources)
}
fn lineage_for(&self, identity: &str) -> Vec<LineageLink> {
let cache = self
.inner
.attributes
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let resolve = |key: &str| {
cache.get(key).cloned().or_else(|| {
cache
.values()
.find(|attributes| attributes.agent_id.as_deref() == Some(key))
.cloned()
})
};
let Some(own) = resolve(identity) else {
return Vec::new();
};
Self::lineage_from_attributes(own, &resolve)
}
fn lineage_for_attributes(&self, attributes: &AgentResourceAttributes) -> Vec<LineageLink> {
let cache = self
.inner
.attributes
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let resolve = |key: &str| {
cache.get(key).cloned().or_else(|| {
cache
.values()
.find(|candidate| candidate.agent_id.as_deref() == Some(key))
.cloned()
})
};
Self::lineage_from_attributes(Arc::new(attributes.clone()), &resolve)
}
fn lineage_from_attributes(
own: Arc<AgentResourceAttributes>,
resolve: &impl Fn(&str) -> Option<Arc<AgentResourceAttributes>>,
) -> Vec<LineageLink> {
const MAX_LINEAGE_DEPTH: usize = 8;
let mut visited = BTreeSet::from([own.identity.clone()]);
let mut lineage = vec![LineageLink {
identity: own.identity.clone(),
attributes: Some(own),
}];
while lineage.len() < MAX_LINEAGE_DEPTH {
let Some(parent) = lineage
.last()
.and_then(|link| link.attributes.as_deref())
.and_then(|attributes| attributes.labels.get("spawned_by"))
.map(|parent| parent.trim().to_string())
.filter(|parent| !parent.is_empty())
else {
break;
};
let attributes = resolve(&parent);
let parent_identity = attributes
.as_deref()
.map(|attributes| attributes.identity.clone())
.unwrap_or(parent);
if !visited.insert(parent_identity.clone()) {
break;
}
lineage.push(LineageLink {
identity: parent_identity,
attributes,
});
}
lineage
}
pub fn can_view_agent(&self, identity: &str) -> bool {
self.allows_agent(ACTION_AGENT_VIEW, identity)
}
pub fn knows_agent(&self, identity: &str) -> bool {
let cache = self
.inner
.attributes
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner);
cache.contains_key(identity)
|| cache
.values()
.any(|attributes| attributes.agent_id.as_deref() == Some(identity))
}
pub fn can_administer(&self) -> bool {
if self.is_admin {
return true;
}
if !self.config.enabled {
return self.config.admins.is_empty();
}
self.decide(super::model::ACTION_ACCESS_ADMIN, &AccessResource::none())
.is_allow()
}
}
impl std::fmt::Debug for AccessView {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("AccessView")
.field("subject", &self.principal.subject)
.field("groups", &self.principal.groups)
.field("is_admin", &self.is_admin)
.field("enforced", &self.config.enabled)
.finish()
}
}
#[cfg(test)]
#[allow(clippy::expect_used, clippy::unwrap_used)]
mod tests {
use super::*;
use crate::access::model::AccessEffect;
fn enabled_config() -> AccessControlConfig {
AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
groups: BTreeMap::from([(
"ops".to_string(),
AccessGroup {
description: None,
members: vec!["alice@example.test".to_string()],
},
)]),
rules: vec![AccessRule {
id: "ops-view-all".to_string(),
groups: vec!["ops".to_string()],
actions: vec!["agent.view".to_string()],
..AccessRule::default()
}],
}
}
#[test]
fn view_resolves_groups_and_admin_flag() {
let controller = AccessController::new(enabled_config()).expect("controller");
let alice = controller.view_for_subject(Some("alice@example.test"));
assert!(alice.groups().contains("ops"));
assert!(!alice.is_admin());
assert!(alice.can_view_agent("identity:scout-1"));
assert!(!alice.allows_agent("agent.send", "identity:scout-1"));
let root = controller.view_for_subject(Some("root@example.test"));
assert!(root.is_admin());
assert!(root.allows("access.admin"));
}
#[test]
fn live_mutations_bump_revision_and_apply() {
let controller = AccessController::new(enabled_config()).expect("controller");
let bob = controller.view_for_subject(Some("bob@example.test"));
assert!(!bob.can_view_agent("identity:scout-1"));
let revision = controller
.set_group(
"ops",
AccessGroup {
description: None,
members: vec![
"alice@example.test".to_string(),
"bob@example.test".to_string(),
],
},
)
.expect("set group");
assert_eq!(revision, 1);
let bob_after = controller.view_for_subject(Some("bob@example.test"));
assert!(bob_after.can_view_agent("identity:scout-1"));
assert!(!bob.can_view_agent("identity:scout-1"));
}
#[test]
fn delete_rule_unknown_id_errors() {
let controller = AccessController::new(enabled_config()).expect("controller");
assert_eq!(
controller.delete_rule("missing"),
Err(AccessConfigError::UnknownRule("missing".to_string()))
);
controller.delete_rule("ops-view-all").expect("delete");
let (config, revision) = controller.snapshot();
assert!(config.rules.is_empty());
assert_eq!(revision, 1);
}
#[test]
fn attribute_cache_feeds_label_selectors() {
let mut config = enabled_config();
config.rules.push(AccessRule {
id: "bob-payments".to_string(),
subjects: vec!["bob@example.test".to_string()],
actions: vec!["agent.view".to_string()],
match_labels: BTreeMap::from([("org".to_string(), "payments".to_string())]),
..AccessRule::default()
});
let controller = AccessController::new(config).expect("controller");
let bob = controller.view_for_subject(Some("bob@example.test"));
assert!(!bob.can_view_agent("identity:pay-1"));
controller.record_agent_attributes(AgentResourceAttributes {
identity: "identity:pay-1".to_string(),
agent_id: Some("pay-1".to_string()),
role: Some("analyst".to_string()),
labels: BTreeMap::from([("org".to_string(), "payments".to_string())]),
});
assert!(bob.can_view_agent("identity:pay-1"));
assert!(!bob.can_view_agent("identity:other"));
}
#[test]
fn exact_event_attributes_override_newer_alias_cache_entry() {
let controller = AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: vec![
AccessRule {
id: "view-all".to_string(),
actions: vec!["agent.view".to_string()],
agents: vec!["*".to_string()],
..AccessRule::default()
},
AccessRule {
id: "deny-secret".to_string(),
effect: AccessEffect::Deny,
actions: vec!["agent.view".to_string()],
match_labels: BTreeMap::from([("org".to_string(), "secret".to_string())]),
..AccessRule::default()
},
],
..AccessControlConfig::default()
})
.expect("controller");
controller.record_agent_attributes(AgentResourceAttributes {
identity: "reused-alias".to_string(),
agent_id: Some("reused-alias".to_string()),
role: Some("lead".to_string()),
labels: BTreeMap::from([("org".to_string(), "public".to_string())]),
});
let historical_secret = AgentResourceAttributes {
identity: "reused-alias".to_string(),
agent_id: Some("reused-alias".to_string()),
role: Some("lead".to_string()),
labels: BTreeMap::from([("org".to_string(), "secret".to_string())]),
};
let view = controller.view_for_subject(None);
assert!(view.can_view_agent("reused-alias"));
assert!(
!view
.decide_agent_with_attributes(ACTION_AGENT_VIEW, &historical_secret)
.is_allow(),
"the newer public cache entry must not authorize the historical secret event"
);
}
#[test]
fn knows_agent_detects_cold_cache_so_label_deny_can_be_made_fail_closed() {
let mut config = enabled_config();
config.rules.push(AccessRule {
id: "anon-view-all".to_string(),
actions: vec!["agent.view".to_string()],
agents: vec!["*".to_string()],
..AccessRule::default()
});
config.rules.push(AccessRule {
id: "deny-secret".to_string(),
effect: AccessEffect::Deny,
actions: vec!["agent.view".to_string()],
match_labels: BTreeMap::from([("org".to_string(), "secret".to_string())]),
..AccessRule::default()
});
let controller = AccessController::new(config).expect("controller");
let view = controller.view_for_subject(None);
assert!(!view.knows_agent("identity:secret-1"));
assert!(
view.can_view_agent("identity:secret-1"),
"cold cache currently fails OPEN — this is what knows_agent() detects"
);
controller.record_agent_attributes(AgentResourceAttributes {
identity: "identity:secret-1".to_string(),
agent_id: Some("secret-1".to_string()),
role: Some("worker".to_string()),
labels: BTreeMap::from([("org".to_string(), "secret".to_string())]),
});
assert!(view.knows_agent("identity:secret-1"));
assert!(
!view.can_view_agent("identity:secret-1"),
"after re-prime the label-scoped deny must hide the member"
);
}
#[test]
fn persistence_round_trips() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("config").join("access.toml");
let controller = AccessController::load_or_default(&path).expect("load default");
assert!(!controller.enabled());
let mut config = enabled_config();
config.rules.push(AccessRule {
id: "deny-secret".to_string(),
effect: AccessEffect::Deny,
actions: vec!["agent.*".to_string()],
agents: vec!["identity:secret".to_string()],
..AccessRule::default()
});
controller.replace_config(config.clone()).expect("replace");
super::super::model::normalize_access_config_for_memory_actions(&mut config);
let reloaded = AccessController::load_or_default(&path).expect("reload");
let (reloaded_config, _) = reloaded.snapshot();
assert_eq!(*reloaded_config, config);
}
#[test]
fn lockout_protected_on_live_surface() {
let controller = AccessController::new(enabled_config()).expect("controller");
let mut config = (*controller.snapshot().0).clone();
config.admins.clear();
assert_eq!(
controller.replace_config(config),
Err(AccessConfigError::EnabledWithoutAdmins)
);
}
#[test]
fn concurrent_rule_upserts_do_not_lose_updates() {
let controller = AccessController::new(enabled_config()).expect("controller");
let base_rules = controller.snapshot().0.rules.len();
let threads: usize = 16;
let handles: Vec<_> = (0..threads)
.map(|i| {
let controller = controller.clone();
std::thread::spawn(move || {
controller
.upsert_rule(AccessRule {
id: format!("rule-{i}"),
actions: vec!["agent.view".to_string()],
agents: vec![format!("identity:agent-{i}")],
..AccessRule::default()
})
.expect("upsert");
})
})
.collect();
for handle in handles {
handle.join().expect("thread");
}
let (config, revision) = controller.snapshot();
assert_eq!(
config.rules.len(),
base_rules + threads,
"every rule survived: {config:#?}"
);
assert_eq!(revision, threads as u64, "revision counts every commit");
for i in 0..threads {
assert!(
config
.rules
.iter()
.any(|rule| rule.id == format!("rule-{i}")),
"rule-{i} missing"
);
}
}
#[test]
fn spawn_lineage_inherits_parent_permissions() {
let mut config = enabled_config();
config.rules.push(AccessRule {
id: "bob-ops-lead".to_string(),
subjects: vec!["bob@example.test".to_string()],
actions: vec!["agent.view".to_string(), "agent.send".to_string()],
agents: vec!["ops-lead".to_string()],
..AccessRule::default()
});
let controller = AccessController::new(config).expect("controller");
controller.record_agent_attributes(AgentResourceAttributes {
identity: "ops-lead".to_string(),
agent_id: Some("ops-lead".to_string()),
role: Some("orchestrator".to_string()),
labels: BTreeMap::new(),
});
controller.record_agent_attributes(AgentResourceAttributes {
identity: "worker-3".to_string(),
agent_id: Some("worker-3".to_string()),
role: Some("person-worker".to_string()),
labels: BTreeMap::from([("spawned_by".to_string(), "ops-lead".to_string())]),
});
controller.record_agent_attributes(AgentResourceAttributes {
identity: "worker-3-sub".to_string(),
agent_id: Some("worker-3-sub".to_string()),
role: Some("helper".to_string()),
labels: BTreeMap::from([("spawned_by".to_string(), "worker-3".to_string())]),
});
controller.record_agent_attributes(AgentResourceAttributes {
identity: "scout-1".to_string(),
agent_id: Some("scout-1".to_string()),
role: Some("scout".to_string()),
labels: BTreeMap::new(),
});
let bob = controller.view_for_subject(Some("bob@example.test"));
assert!(bob.can_view_agent("ops-lead"));
assert!(
bob.can_view_agent("worker-3"),
"a member spawned by ops-lead inherits ops-lead's visibility"
);
assert!(
bob.allows_agent("agent.send", "worker-3"),
"permission inheritance covers every agent action, not just view"
);
assert!(
bob.can_view_agent("worker-3-sub"),
"spawn lineage inheritance is transitive"
);
assert!(
!bob.can_view_agent("scout-1"),
"agents outside the spawn lineage stay denied"
);
}
#[test]
fn spawn_lineage_deny_on_parent_overrides_descendants() {
let mut config = enabled_config();
config.rules.push(AccessRule {
id: "bob-view-all".to_string(),
subjects: vec!["bob@example.test".to_string()],
actions: vec!["agent.view".to_string()],
agents: vec!["*".to_string()],
..AccessRule::default()
});
config.rules.push(AccessRule {
id: "hide-secret-lead".to_string(),
effect: AccessEffect::Deny,
actions: vec!["agent.*".to_string()],
agents: vec!["secret-lead".to_string()],
..AccessRule::default()
});
let controller = AccessController::new(config).expect("controller");
controller.record_agent_attributes(AgentResourceAttributes {
identity: "secret-lead".to_string(),
agent_id: Some("secret-lead".to_string()),
role: None,
labels: BTreeMap::new(),
});
controller.record_agent_attributes(AgentResourceAttributes {
identity: "covert-worker".to_string(),
agent_id: Some("covert-worker".to_string()),
role: None,
labels: BTreeMap::from([("spawned_by".to_string(), "secret-lead".to_string())]),
});
let bob = controller.view_for_subject(Some("bob@example.test"));
assert!(!bob.can_view_agent("secret-lead"));
assert!(
!bob.can_view_agent("covert-worker"),
"a deny on the spawning parent must propagate to its descendants"
);
}
#[test]
fn spawn_lineage_cycles_terminate_and_fail_closed() {
let controller = AccessController::new(enabled_config()).expect("controller");
controller.record_agent_attributes(AgentResourceAttributes {
identity: "loop-a".to_string(),
agent_id: Some("loop-a".to_string()),
role: None,
labels: BTreeMap::from([("spawned_by".to_string(), "loop-b".to_string())]),
});
controller.record_agent_attributes(AgentResourceAttributes {
identity: "loop-b".to_string(),
agent_id: Some("loop-b".to_string()),
role: None,
labels: BTreeMap::from([("spawned_by".to_string(), "loop-a".to_string())]),
});
let bob = controller.view_for_subject(Some("bob@example.test"));
assert!(
!bob.can_view_agent("loop-a"),
"lineage cycles must terminate and deny by default"
);
}
#[test]
fn persist_is_atomic_via_temp_rename() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("access.toml");
let controller = AccessController::load_or_default(&path).expect("load");
controller
.replace_config(enabled_config())
.expect("replace");
assert!(path.is_file(), "target written");
assert!(
!dir.path().join("access.toml.tmp").exists(),
"temp file cleaned up by rename"
);
let reloaded = AccessController::load_or_default(&path).expect("reload");
assert!(reloaded.enabled());
}
mod checked_http {
use super::*;
use axum::body::{Body, to_bytes};
use axum::http::{Request, StatusCode, header};
use serde_json::{Value, json};
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Condvar, mpsc};
use std::thread::JoinHandle;
use std::time::Duration;
use tower::ServiceExt;
const A: &str = "root@example.test";
const B: &str = "alice@example.test";
const PRIVATE: &str = "PRIVATE_QUEUED_CONFIG_CANARY";
const ISSUER: &str = "https://trusted.mobkit.localhost";
fn app(controller: &AccessController) -> axum::Router {
let decisions = crate::build_runtime_decision_state(crate::RuntimeDecisionInputs {
bigquery: crate::BigQueryNaming { dataset: "access_dataset".into(), table: "access_table".into() },
trusted_mobkit_toml: "[[modules]]\nid = \"router\"\ncommand = \"router-bin\"\nargs = []\nrestart_policy = \"always\"\n".into(),
auth: crate::AuthPolicy { default_provider: crate::AuthProvider::GoogleOAuth,
email_allowlist: vec![A.into(), B.into(), "carol@example.test".into()] },
trusted_oidc: crate::TrustedOidcRuntimeConfig {
discovery_json: json!({"issuer": ISSUER, "jwks_uri": format!("{ISSUER}/.well-known/jwks.json")}).to_string(),
jwks_json: r#"{"keys":[{"kid":"kid-current","kty":"oct","alg":"HS256","k":"cGhhc2U3LXRydXN0ZWQtY3VycmVudC1zZWNyZXQ"}]}"#.into(),
audience: "meerkat-console".into(), require_verified_email: false,
},
console: crate::ConsolePolicy { require_app_auth: true, ..Default::default() },
ops: crate::RuntimeOpsPolicy::default(),
release_metadata_json: include_str!("../../assets/release-targets.json").into(),
}).expect("real decisions");
crate::console_json_router_with_aggregator_and_access(
decisions,
crate::MobKitConsoleAggregator::new(Arc::new(
crate::InMemoryConsoleLogStore::default(),
)),
Some(controller.clone()),
)
}
fn rpc(app: axum::Router, subject: &str, method: &str, params: Value) -> Value {
let mut h = jsonwebtoken::Header::new(jsonwebtoken::Algorithm::HS256);
h.kid = Some("kid-current".into());
let jwt = jsonwebtoken::encode(
&h,
&json!({"iss": ISSUER, "aud": "meerkat-console",
"sub": subject, "email": subject, "provider": "google_oauth",
"exp": chrono::Utc::now().timestamp() + 300}),
&jsonwebtoken::EncodingKey::from_secret(b"phase7-trusted-current-secret"),
)
.expect("JWT");
let request = Request::builder()
.method("POST")
.uri("/console/rpc")
.header(header::CONTENT_TYPE, "application/json")
.header(header::AUTHORIZATION, format!("Bearer {jwt}"))
.body(Body::from(
json!({"jsonrpc":"2.0","id":"queued-admin","method":method,"params":params})
.to_string(),
))
.expect("request");
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("runtime")
.block_on(async {
let response = app.oneshot(request).await.expect("real HTTP route");
assert_eq!(response.status(), StatusCode::OK);
let bytes = to_bytes(response.into_body(), 1024 * 1024)
.await
.expect("body");
let value: Value = serde_json::from_slice(&bytes).expect("JSON");
assert_eq!(value["id"], json!("queued-admin"));
value
})
}
fn config() -> AccessControlConfig {
AccessControlConfig {
enabled: true,
admins: vec![A.into(), B.into()],
rules: vec![AccessRule {
id: "private-rule".into(),
description: Some(PRIVATE.into()),
actions: vec!["agent.send".into()],
subjects: vec![B.into()],
..Default::default()
}],
..Default::default()
}
}
fn payload(owner: &str, revision: u64, config: &AccessControlConfig) -> Value {
json!({"checked_v1":{"owner_instance":owner,"expected_revision":revision,"config":config}})
}
#[derive(Default)]
struct Gate {
released: Mutex<bool>,
wake: Condvar,
used: AtomicBool,
timed_out: AtomicBool,
}
impl Gate {
fn release(&self) {
*self
.released
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = true;
self.wake.notify_all();
}
fn park_once(&self) {
if self.used.swap(true, Ordering::SeqCst) {
return;
}
let locked = self
.released
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let (released, _) = self
.wake
.wait_timeout_while(locked, Duration::from_secs(10), |ready| !*ready)
.unwrap_or_else(std::sync::PoisonError::into_inner);
if !*released {
self.timed_out.store(true, Ordering::SeqCst);
}
}
}
struct Workers {
gates: [Arc<Gate>; 2],
joins: Vec<JoinHandle<Value>>,
}
impl Drop for Workers {
fn drop(&mut self) {
for gate in &self.gates {
gate.release();
}
for task in self.joins.drain(..) {
let _ = task.join();
}
}
}
fn checkpoint(
controller: &AccessController,
path: &Path,
) -> (AccessControlConfig, u64, Vec<u8>) {
let (config, revision) = controller.snapshot();
let bytes = std::fs::read(path).expect("persisted bytes");
let disk: AccessControlConfig =
toml::from_str(std::str::from_utf8(&bytes).expect("UTF8")).expect("TOML");
assert_eq!(disk, *config);
((*config).clone(), revision, bytes)
}
#[test]
fn queued_revoked_administrator_is_denied_before_matching_checked_write() {
run_queued_revocation(true);
}
#[test]
fn queued_revoked_administrator_is_denied_before_legacy_write() {
run_queued_revocation(false);
}
fn run_queued_revocation(checked: bool) {
let dir = tempfile::tempdir().expect("directory");
let path = dir.path().join("access.toml");
let initial = config();
std::fs::write(&path, toml::to_string_pretty(&initial).expect("seed"))
.expect("persist seed");
let owner = AccessController::load_or_default(&path).expect("stored owner");
let route = app(&owner);
let read = rpc(route.clone(), A, "mobkit/access/get", json!({}));
let instance = if checked {
read["result"]["owner_instance"]
.as_str()
.expect("instance")
.to_string()
} else {
String::new()
};
let revision = read["result"]["revision"].as_u64().expect("revision");
assert_eq!(revision, 0);
let write = |revision, config: &AccessControlConfig| {
if checked {
payload(&instance, revision, config)
} else {
json!({"config": config})
}
};
let a_gate = Arc::new(Gate::default());
let b_gate = Arc::new(Gate::default());
let (tx, rx) = mpsc::channel();
let mut workers = Workers {
gates: [Arc::clone(&a_gate), Arc::clone(&b_gate)],
joins: Vec::new(),
};
let ag = Arc::clone(&a_gate);
let bg = Arc::clone(&b_gate);
let acquired = AtomicUsize::new(0);
*owner.inner.checked_probe.lock().expect("probe slot") = Some(Arc::new(move |stage| {
let _ = tx.send(stage);
if stage == CheckedMutationStage::Acquired {
match acquired.fetch_add(1, Ordering::SeqCst) {
0 => bg.park_once(),
1 => ag.park_once(),
_ => {}
}
}
}));
let mut revoked = initial.clone();
revoked.admins.retain(|admin| admin != A);
let b_route = route.clone();
let b_body = write(revision, &revoked);
workers.joins.push(std::thread::spawn(move || {
rpc(b_route, B, "mobkit/access/set", b_body)
}));
assert_eq!(
rx.recv_timeout(Duration::from_secs(5)).expect("B attempt"),
CheckedMutationStage::Attempt
);
assert_eq!(
rx.recv_timeout(Duration::from_secs(5))
.expect("B holds real mutation lock"),
CheckedMutationStage::Acquired
);
assert_eq!(checkpoint(&owner, &path).0, initial);
let a_route = route.clone();
let a_body = write(revision + 1, &initial);
workers.joins.push(std::thread::spawn(move || {
rpc(a_route, A, "mobkit/access/set", a_body)
}));
assert_eq!(
rx.recv_timeout(Duration::from_secs(5))
.expect("A reached actual mutex acquisition"),
CheckedMutationStage::Attempt
);
assert!(matches!(rx.try_recv(), Err(mpsc::TryRecvError::Empty)));
assert!(!workers.joins[1].is_finished());
b_gate.release();
assert_eq!(
rx.recv_timeout(Duration::from_secs(5))
.expect("A acquired after B commit"),
CheckedMutationStage::Acquired
);
let b_response = workers.joins.remove(0).join().expect("B worker");
assert_eq!(b_response["error"], Value::Null, "{b_response}");
assert_eq!(b_response["result"]["revision"], json!(revision + 1));
let after_b = checkpoint(&owner, &path);
assert_eq!(after_b.0, revoked);
assert_eq!(after_b.1, revision + 1);
assert!(!owner.view_for_subject(Some(A)).can_administer());
a_gate.release();
let denied = workers.joins.remove(0).join().expect("A worker");
assert_eq!(denied["error"]["code"], json!(-32030), "{denied}");
assert_eq!(denied["error"]["data"], json!({"kind":"access_denied"}));
assert_eq!(denied["result"], Value::Null);
for private in [
A,
B,
PRIVATE,
instance.as_str(),
path.to_str().expect("path"),
]
.into_iter()
.filter(|value| !value.is_empty())
{
assert!(!denied["error"].to_string().contains(private));
}
assert_eq!(checkpoint(&owner, &path), after_b);
assert!(!a_gate.timed_out.load(Ordering::SeqCst));
assert!(!b_gate.timed_out.load(Ordering::SeqCst));
let healthy = rpc(route, B, "mobkit/access/set", write(revision + 1, &initial));
assert_eq!(healthy["error"], Value::Null, "{healthy}");
assert_eq!(healthy["result"]["revision"], json!(revision + 2));
let final_state = checkpoint(&owner, &path);
assert_eq!(final_state.0, initial);
assert_eq!(final_state.1, revision + 2);
}
#[test]
fn unavailable_edit_identity_keeps_reads_and_legacy_edits_without_false_capability() {
let dir = tempfile::tempdir().expect("directory");
let path = dir.path().join("access.toml");
let initial = config();
std::fs::write(&path, toml::to_string_pretty(&initial).expect("seed"))
.expect("seed bytes");
let owner = AccessController::with_edit_identity(initial.clone(), None)
.expect("constructor failure branch")
.with_persist_path(&path);
let cloned = owner.clone();
let route = app(&cloned);
for method in ["mobkit/access/status", "mobkit/access/get"] {
let read = rpc(route.clone(), A, method, json!({}));
assert_eq!(read["error"], Value::Null, "{read}");
assert!(read["result"].get("conditional_mutations").is_none());
assert!(read["result"].get("owner_instance").is_none());
if method.ends_with("/get") {
assert_eq!(read["result"]["config"], json!(initial));
}
}
let before = checkpoint(&owner, &path);
let checked = payload("not-an-owner-token", 0, &initial);
let denied = rpc(
route.clone(),
"carol@example.test",
"mobkit/access/set",
checked.clone(),
);
assert_eq!(denied["error"]["code"], json!(-32030));
assert_eq!(denied["error"]["data"], json!({"kind":"access_denied"}));
let unavailable = rpc(route.clone(), A, "mobkit/access/set", checked);
assert_eq!(unavailable["error"]["code"], json!(-32004));
assert_eq!(
unavailable["error"]["data"],
json!({"kind":"access_mutation_unavailable"})
);
for private in [
A,
PRIVATE,
"not-an-owner-token",
path.to_str().expect("path"),
] {
assert!(!unavailable["error"].to_string().contains(private));
}
assert_eq!(checkpoint(&owner, &path), before);
cloned
.set_group(
"legacy-direct",
AccessGroup {
members: vec![B.into()],
..Default::default()
},
)
.expect("existing direct path");
assert_eq!(checkpoint(&owner, &path).1, 1);
assert!(owner.snapshot().0.groups.contains_key("legacy-direct"));
let legacy = rpc(
route,
A,
"mobkit/access/groups/set",
json!({
"name": "legacy-http", "group": {"members": [B]}
}),
);
assert_eq!(legacy["error"], Value::Null, "{legacy}");
assert_eq!(legacy["result"]["revision"], json!(2));
assert_eq!(checkpoint(&owner, &path).1, 2);
assert!(owner.snapshot().0.groups.contains_key("legacy-http"));
}
#[test]
fn exhausted_revision_refuses_direct_and_http_without_mutation() {
let dir = tempfile::tempdir().expect("directory");
let path = dir.path().join("access.toml");
let initial = config();
std::fs::write(&path, toml::to_string_pretty(&initial).expect("seed"))
.expect("seed bytes");
let owner = AccessController::load_or_default(&path).expect("stored owner");
owner.inner.state.write().expect("state").revision = u64::MAX;
let before = checkpoint(&owner, &path);
let route = app(&owner);
let read = rpc(route.clone(), A, "mobkit/access/get", json!({}));
let instance = read["result"]["owner_instance"].as_str().expect("instance");
assert_eq!(read["result"]["revision"], json!(u64::MAX));
let mut changed = initial.clone();
changed
.groups
.insert("must-not-publish".into(), AccessGroup::default());
assert_eq!(
owner.replace_config(changed.clone()),
Err(AccessConfigError::RevisionExhausted)
);
assert_eq!(checkpoint(&owner, &path), before);
for body in [
payload(instance, u64::MAX, &changed),
json!({"config": changed}),
] {
let response = rpc(route.clone(), A, "mobkit/access/set", body);
assert_eq!(response["error"]["code"], json!(-32004), "{response}");
assert_eq!(
response["error"]["data"],
json!({"kind":"access_mutation_unavailable"})
);
assert_eq!(checkpoint(&owner, &path), before);
assert!(!path.with_file_name("access.toml.tmp").exists());
}
let denied = rpc(
route,
"carol@example.test",
"mobkit/access/set",
payload("wrong-instance", 0, &initial),
);
assert_eq!(denied["error"]["code"], json!(-32030));
assert_eq!(denied["error"]["data"], json!({"kind":"access_denied"}));
assert_eq!(checkpoint(&owner, &path), before);
}
#[test]
fn checked_http_persistence_failure_is_finite_and_does_not_publish() {
let dir = tempfile::tempdir().expect("directory");
let path = dir.path().join("PRIVATE_PERSIST_PATH.toml");
let initial = config();
std::fs::write(&path, toml::to_string_pretty(&initial).expect("seed"))
.expect("seed bytes");
let owner = AccessController::load_or_default(&path).expect("stored owner");
let route = app(&owner);
let read = rpc(route.clone(), A, "mobkit/access/get", json!({}));
let instance = read["result"]["owner_instance"].as_str().expect("instance");
let before = checkpoint(&owner, &path);
let temp_path = path.with_file_name("PRIVATE_PERSIST_PATH.toml.tmp");
std::fs::create_dir(&temp_path).expect("block actual temp write");
let mut changed = initial;
changed
.groups
.insert("new-group".into(), AccessGroup::default());
let failed = rpc(
route.clone(),
A,
"mobkit/access/set",
payload(instance, 0, &changed),
);
assert_eq!(failed["error"]["code"], json!(-32000), "{failed}");
assert_eq!(
failed["error"]["data"],
json!({"kind":"access_persistence_failed"})
);
assert_eq!(
failed["error"]["message"],
json!("Access configuration could not be saved.")
);
for private in [
A,
PRIVATE,
instance,
path.to_str().expect("path"),
"PRIVATE_PERSIST_PATH",
] {
assert!(!failed["error"].to_string().contains(private));
}
assert_eq!(checkpoint(&owner, &path), before);
std::fs::remove_dir(&temp_path).expect("remove fault");
let healthy = rpc(
route,
A,
"mobkit/access/set",
payload(instance, 0, &changed),
);
assert_eq!(healthy["error"], Value::Null, "{healthy}");
assert_eq!(healthy["result"]["revision"], json!(1));
let after = checkpoint(&owner, &path);
assert_eq!(after.0, changed);
assert_eq!(after.1, 1);
assert!(!temp_path.exists());
}
}
}