#![allow(clippy::expect_used, clippy::unwrap_used, clippy::panic)]
use std::collections::BTreeMap;
use std::sync::Arc;
use std::time::Duration;
use axum::body::{Body, to_bytes};
use axum::http::{Request, StatusCode, header};
use futures::StreamExt;
use meerkat::{AgentFactory, Config, build_ephemeral_service};
use meerkat_client::TestClient;
use meerkat_core::types::HandlingMode;
use meerkat_mob::event::MobEventKind;
use meerkat_mob::ids::AgentIdentity as MeerkatId;
use meerkat_mob::{MobDefinition, MobStorage, SpawnMemberSpec};
use meerkat_mobkit::runtime::ConsoleMember;
use meerkat_mobkit::{
AccessControlConfig, AccessController, AccessGroup, AccessRule, AgentResourceAttributes,
AuthPolicy, BigQueryNaming, ConsoleAccessRequest, ConsoleLiveSnapshot,
ConsoleModelCapabilities, ConsolePolicy, ConsoleRestJsonRequest, ConsoleVisibilityPolicy,
DiscoverySpec, MobBootstrapOptions, MobBootstrapSpec, MobKitConfig, RuntimeDecisionInputs,
RuntimeOpsPolicy, TopologyControlMode, TopologyControlPolicy, TrustedOidcRuntimeConfig,
UnifiedRuntime, build_runtime_decision_state,
handle_console_rest_json_route_with_snapshot_and_access,
};
use meerkat_mobkit::{StewardStore, TaintableStore};
use serde_json::{Value, json};
use sha2::{Digest, Sha256};
use tower::ServiceExt;
fn trusted_toml() -> String {
r#"
[[modules]]
id = "router"
command = "router-bin"
args = []
restart_policy = "always"
"#
.to_string()
}
fn trusted_oidc() -> TrustedOidcRuntimeConfig {
TrustedOidcRuntimeConfig {
discovery_json:
r#"{"issuer":"https://trusted.mobkit.local","jwks_uri":"https://trusted.mobkit.local/.well-known/jwks.json"}"#
.to_string(),
jwks_json: r#"{"keys":[{"kid":"kid-current","kty":"oct","alg":"HS256","k":"cGhhc2U3LXRydXN0ZWQtY3VycmVudC1zZWNyZXQ"}]}"#
.to_string(),
audience: "meerkat-console".to_string(),
}
}
fn decision_state(require_app_auth: bool) -> meerkat_mobkit::RuntimeDecisionState {
build_runtime_decision_state(RuntimeDecisionInputs {
bigquery: BigQueryNaming {
dataset: "access_dataset".to_string(),
table: "access_table".to_string(),
},
trusted_mobkit_toml: trusted_toml(),
auth: AuthPolicy {
default_provider: meerkat_mobkit::AuthProvider::GoogleOAuth,
email_allowlist: vec![
"root@example.test".to_string(),
"alice@example.test".to_string(),
"carol@example.test".to_string(),
],
},
trusted_oidc: trusted_oidc(),
console: ConsolePolicy {
require_app_auth,
..ConsolePolicy::default()
},
ops: RuntimeOpsPolicy::default(),
release_metadata_json: include_str!("../assets/release-targets.json").to_string(),
})
.expect("decision state builds")
}
fn member(identity: &str, role: &str, labels: &[(&str, &str)]) -> ConsoleMember {
ConsoleMember {
agent_identity: identity.to_string(),
role: role.to_string(),
state: "active".to_string(),
model_capabilities: ConsoleModelCapabilities::default(),
runtime_mode: None,
session_id: None,
wired_to: Vec::new(),
labels: labels
.iter()
.map(|(key, value)| (key.to_string(), value.to_string()))
.collect(),
progress: None,
}
}
fn snapshot_with_members(members: Vec<ConsoleMember>) -> ConsoleLiveSnapshot {
ConsoleLiveSnapshot::new(
Some("access-test-runtime".to_string()),
true,
Vec::new(),
Vec::new(),
members,
true,
)
}
fn ops_access_config() -> AccessControlConfig {
AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
groups: BTreeMap::from([(
"ops".to_string(),
AccessGroup {
description: Some("Operations".to_string()),
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()
},
AccessRule {
id: "ops-send-lead".to_string(),
groups: vec!["ops".to_string()],
actions: vec!["agent.send".to_string()],
agents: vec!["ops-lead".to_string()],
..AccessRule::default()
},
],
}
}
fn experience_for(
controller: &AccessController,
subject: &str,
snapshot: &ConsoleLiveSnapshot,
) -> Value {
let decisions = decision_state(true);
let response = handle_console_rest_json_route_with_snapshot_and_access(
&decisions,
&ConsoleRestJsonRequest {
method: "GET".to_string(),
path: "/console/experience".to_string(),
auth: Some(ConsoleAccessRequest {
provider: meerkat_mobkit::AuthProvider::GoogleOAuth,
email: subject.to_string(),
}),
},
Some(snapshot),
Some(controller),
);
assert_eq!(response.status, 200, "experience: {:?}", response.body);
response.body
}
fn sidebar_identities(experience: &Value) -> Vec<String> {
experience["agent_sidebar"]["live_snapshot"]["agents"]
.as_array()
.expect("sidebar agents")
.iter()
.map(|agent| {
agent["identity"]
.as_str()
.filter(|value| !value.is_empty())
.or_else(|| agent["agent_id"].as_str())
.unwrap_or_default()
.to_string()
})
.collect()
}
#[test]
fn experience_is_filtered_per_principal() {
let controller = AccessController::new(ops_access_config()).expect("controller");
let snapshot = snapshot_with_members(vec![
member("ops-lead", "lead", &[]),
member("scout-1", "scout", &[]),
]);
let admin_experience = experience_for(&controller, "root@example.test", &snapshot);
assert_eq!(
sidebar_identities(&admin_experience),
["ops-lead", "scout-1"]
);
assert_eq!(admin_experience["access"]["can_administer"], json!(true));
assert_eq!(admin_experience["access"]["enabled"], json!(true));
let ops_experience = experience_for(&controller, "alice@example.test", &snapshot);
assert_eq!(sidebar_identities(&ops_experience), ["ops-lead", "scout-1"]);
let agents = ops_experience["agent_sidebar"]["live_snapshot"]["agents"]
.as_array()
.expect("agents");
let affordance = |identity: &str, key: &str| -> bool {
agents
.iter()
.find(|agent| agent["identity"] == identity)
.and_then(|agent| agent["affordances"][key].as_bool())
.unwrap_or(false)
};
assert!(affordance("ops-lead", "can_send_message"));
assert!(!affordance("scout-1", "can_send_message"));
assert!(!affordance("scout-1", "can_retire"));
assert_eq!(ops_experience["access"]["can_administer"], json!(false));
assert_eq!(
ops_experience["access"]["groups"],
json!(["ops"]),
"groups surface in the access section"
);
assert_eq!(
ops_experience["runtime_capabilities"]["can_send_messages"],
json!(true)
);
assert_eq!(
ops_experience["runtime_capabilities"]["can_spawn_members"],
json!(false)
);
let outsider_experience = experience_for(&controller, "carol@example.test", &snapshot);
assert_eq!(
sidebar_identities(&outsider_experience),
Vec::<String>::new(),
"an ungranted authenticated subject sees no agents"
);
assert_eq!(
outsider_experience["access"]["can_administer"],
json!(false)
);
let decisions = decision_state(true);
let response = handle_console_rest_json_route_with_snapshot_and_access(
&decisions,
&ConsoleRestJsonRequest {
method: "GET".to_string(),
path: "/console/experience".to_string(),
auth: Some(ConsoleAccessRequest {
provider: meerkat_mobkit::AuthProvider::GoogleOAuth,
email: "alice@example.test".to_string(),
}),
},
Some(&snapshot_with_members(vec![member(
"hidden-only",
"lead",
&[],
)])),
Some(
&AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
..AccessControlConfig::default()
})
.expect("deny-all controller"),
),
);
assert_eq!(
sidebar_identities(&response.body),
Vec::<String>::new(),
"deny-by-default hides every agent"
);
}
#[test]
fn module_fallback_rows_are_filtered_for_denied_callers() {
let controller = AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
..AccessControlConfig::default()
})
.expect("deny-all controller");
let snapshot = ConsoleLiveSnapshot::new(
Some("access-test-runtime".to_string()),
true,
vec!["router".to_string()],
Vec::new(),
Vec::new(),
false,
);
let denied = experience_for(&controller, "alice@example.test", &snapshot);
assert_eq!(
sidebar_identities(&denied),
Vec::<String>::new(),
"module fallback rows must not leak to denied callers"
);
let admin = experience_for(&controller, "root@example.test", &snapshot);
assert_eq!(sidebar_identities(&admin).len(), 1, "admin keeps modules");
let partial = AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: vec![AccessRule {
id: "carol-views-router".to_string(),
subjects: vec!["carol@example.test".to_string()],
actions: vec!["agent.view".to_string()],
agents: vec!["router".to_string()],
..AccessRule::default()
}],
..AccessControlConfig::default()
})
.expect("partial controller");
let two_modules = ConsoleLiveSnapshot::new(
Some("access-test-runtime".to_string()),
true,
vec!["router".to_string(), "delivery".to_string()],
Vec::new(),
Vec::new(),
false,
);
let scoped = experience_for(&partial, "carol@example.test", &two_modules);
assert_eq!(
sidebar_identities(&scoped),
["router"],
"only the granted module row is visible"
);
}
#[test]
fn label_selector_rules_filter_experience() {
let config = AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: vec![AccessRule {
id: "payments-only".to_string(),
subjects: vec!["alice@example.test".to_string()],
actions: vec!["agent.view".to_string()],
match_labels: BTreeMap::from([("org".to_string(), "payments".to_string())]),
..AccessRule::default()
}],
..AccessControlConfig::default()
};
let controller = AccessController::new(config).expect("controller");
let snapshot = snapshot_with_members(vec![
member("pay-analyst", "analyst", &[("org", "payments")]),
member("hr-analyst", "analyst", &[("org", "people")]),
]);
let experience = experience_for(&controller, "alice@example.test", &snapshot);
assert_eq!(sidebar_identities(&experience), ["pay-analyst"]);
}
#[test]
fn identity_first_console_identity_drives_filtering() {
let config = AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: vec![
AccessRule {
id: "alice-views-lead-by-identity".to_string(),
subjects: vec!["alice@example.test".to_string()],
actions: vec!["agent.view".to_string()],
agents: vec!["identity:ops-lead".to_string()],
..AccessRule::default()
},
AccessRule {
id: "alice-sends-lead-by-runtime-id".to_string(),
subjects: vec!["alice@example.test".to_string()],
actions: vec!["agent.send".to_string()],
agents: vec!["member-runtime-7".to_string()],
..AccessRule::default()
},
],
..AccessControlConfig::default()
};
let controller = AccessController::new(config).expect("controller");
let snapshot = snapshot_with_members(vec![
member(
"member-runtime-7",
"lead",
&[("agent_identity", "identity:ops-lead")],
),
member(
"member-runtime-8",
"scout",
&[("agent_identity", "identity:scout-9")],
),
]);
let experience = experience_for(&controller, "alice@example.test", &snapshot);
assert_eq!(sidebar_identities(&experience), ["identity:ops-lead"]);
let agents = experience["agent_sidebar"]["live_snapshot"]["agents"]
.as_array()
.expect("agents");
let lead = agents
.iter()
.find(|agent| agent["identity"] == "identity:ops-lead")
.expect("lead row");
assert_eq!(lead["affordances"]["can_send_message"], json!(true));
}
#[test]
fn disabled_controller_changes_nothing() {
let controller = AccessController::disabled();
let snapshot = snapshot_with_members(vec![
member("ops-lead", "lead", &[]),
member("scout-1", "scout", &[]),
]);
let experience = experience_for(&controller, "alice@example.test", &snapshot);
assert_eq!(sidebar_identities(&experience), ["ops-lead", "scout-1"]);
assert_eq!(experience["access"]["enabled"], json!(false));
assert_eq!(experience["access"]["can_administer"], json!(true));
}
async fn build_access_runtime_fixture() -> (tempfile::TempDir, UnifiedRuntime) {
let temp_dir = tempfile::tempdir().expect("temp dir");
let session_path = temp_dir.path().join("sessions");
std::fs::create_dir_all(&session_path).expect("session path");
let factory = AgentFactory::new(&session_path).comms(true);
let session_service = Arc::new(build_ephemeral_service(factory, Config::default(), 16));
static NEXT_ACCESS_MOB: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
let definition = MobDefinition::from_toml(&format!(
r#"
[mob]
id = "access-control-mob-{}"
[profiles.lead]
model = "gpt-5.5"
external_addressable = true
[profiles.lead.tools]
comms = true
"#,
NEXT_ACCESS_MOB.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
))
.expect("definition");
let mob_spec = MobBootstrapSpec::new(definition, MobStorage::in_memory(), session_service)
.with_options(MobBootstrapOptions {
allow_ephemeral_sessions: true,
notify_orchestrator_on_resume: true,
default_llm_client: Some(Arc::new(TestClient::default())),
});
let module_config = MobKitConfig {
modules: vec![],
discovery: DiscoverySpec {
namespace: "access-control".to_string(),
modules: vec![],
},
pre_spawn: vec![],
};
let runtime = UnifiedRuntime::bootstrap(mob_spec, module_config, Duration::from_secs(2))
.await
.expect("bootstrap runtime");
for member_id in ["router", "delivery"] {
runtime
.spawn(SpawnMemberSpec::from_wire(
"lead".to_string(),
MeerkatId::from(member_id).to_string(),
Some(format!("You are {member_id}.").into()),
None,
None,
))
.await
.expect("spawn member");
}
(temp_dir, runtime)
}
async fn rpc(app: &axum::Router, method: &str, params: Value) -> Value {
let payload = json!({
"jsonrpc": "2.0",
"id": "test",
"method": method,
"params": params,
});
let response = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/console/rpc")
.header(header::CONTENT_TYPE, "application/json")
.body(Body::from(payload.to_string()))
.expect("rpc request"),
)
.await
.expect("rpc response");
assert_eq!(response.status(), StatusCode::OK);
let body = to_bytes(response.into_body(), 1024 * 1024)
.await
.expect("rpc body");
serde_json::from_slice(&body).expect("rpc json")
}
async fn get_status(app: &axum::Router, uri: &str) -> StatusCode {
app.clone()
.oneshot(
Request::builder()
.method("GET")
.uri(uri)
.body(Body::empty())
.expect("request"),
)
.await
.expect("response")
.status()
}
async fn retire_historical_secret_member(runtime: &UnifiedRuntime, identity: &str) -> u64 {
let before_spawn = runtime
.mob_handle()
.events()
.latest_cursor()
.await
.expect("latest cursor before historical member");
let mut spec = SpawnMemberSpec::from_wire(
"lead".to_string(),
identity.to_string(),
Some("historical secret member".into()),
None,
None,
);
spec.labels = Some(BTreeMap::from([("org".to_string(), "secret".to_string())]));
runtime.spawn(spec).await.expect("spawn historical member");
runtime
.mob_handle()
.retire(MeerkatId::from(identity))
.await
.expect("retire historical member");
tokio::time::timeout(Duration::from_secs(2), async {
loop {
let still_present = runtime
.mob_handle()
.list_members_including_retiring()
.await
.iter()
.any(|member| member.agent_identity.as_str() == identity);
if !still_present {
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("retired member should leave the operational roster");
before_spawn
}
async fn assert_sse_replay_emits_no_data(response: axum::response::Response) {
assert_eq!(response.status(), StatusCode::OK);
let mut stream = response.into_body().into_data_stream();
match tokio::time::timeout(Duration::from_millis(200), stream.next()).await {
Err(_) | Ok(None) => {}
Ok(Some(Err(error))) => panic!("unexpected SSE body error: {error}"),
Ok(Some(Ok(bytes))) => panic!(
"historical hidden member leaked through SSE replay: {}",
String::from_utf8_lossy(&bytes)
),
}
}
async fn assert_sse_emits_no_identity(response: axum::response::Response, identity: &str) {
assert_eq!(response.status(), StatusCode::OK);
let mut stream = response.into_body().into_data_stream();
let deadline = tokio::time::Instant::now() + Duration::from_millis(500);
let mut observed = String::new();
loop {
let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
if remaining.is_zero() {
break;
}
match tokio::time::timeout(remaining, stream.next()).await {
Err(_) | Ok(None) => break,
Ok(Some(Err(error))) => panic!("unexpected SSE body error: {error}"),
Ok(Some(Ok(bytes))) => {
observed.push_str(&String::from_utf8_lossy(&bytes));
assert!(
!observed.contains(identity),
"hidden identity {identity} leaked through SSE: {observed}"
);
}
}
}
}
async fn assert_structural_sse_excludes_and_includes_cursor(
response: axum::response::Response,
excluded_cursor: u64,
included_cursor: u64,
) {
assert_eq!(response.status(), StatusCode::OK);
let mut stream = response.into_body().into_data_stream();
let excluded = format!("id: mob-evt-{excluded_cursor}");
let included = format!("id: mob-evt-{included_cursor}");
let deadline = tokio::time::Instant::now() + Duration::from_secs(2);
let mut observed = String::new();
while !observed.contains(&included) {
let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
assert!(
!remaining.is_zero(),
"expected current public cursor {included_cursor} in structural SSE: {observed}"
);
match tokio::time::timeout(remaining, stream.next()).await {
Err(_) | Ok(None) => panic!(
"expected current public cursor {included_cursor} in structural SSE: {observed}"
),
Ok(Some(Err(error))) => panic!("unexpected SSE body error: {error}"),
Ok(Some(Ok(bytes))) => {
observed.push_str(&String::from_utf8_lossy(&bytes));
assert!(
!observed.contains(&excluded),
"historical secret cursor {excluded_cursor} leaked through structural SSE: {observed}"
);
}
}
}
}
fn anonymous_router_only_config() -> AccessControlConfig {
AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
groups: BTreeMap::new(),
rules: vec![
AccessRule {
id: "everyone-views-router".to_string(),
actions: vec!["agent.view".to_string()],
agents: vec!["router".to_string()],
..AccessRule::default()
},
AccessRule {
id: "everyone-sends-router".to_string(),
actions: vec!["agent.send".to_string()],
agents: vec!["router".to_string()],
..AccessRule::default()
},
],
}
}
#[tokio::test]
async fn http_router_enforces_access_end_to_end() {
let (_temp_dir, mut runtime) = build_access_runtime_fixture().await;
let controller = AccessController::new(anonymous_router_only_config()).expect("controller");
runtime.set_access_controller(controller.clone());
let app = runtime.build_reference_app_router(decision_state(false));
let response = app
.clone()
.oneshot(
Request::builder()
.method("GET")
.uri("/console/experience")
.body(Body::empty())
.expect("experience request"),
)
.await
.expect("experience response");
assert_eq!(response.status(), StatusCode::OK);
let body = to_bytes(response.into_body(), 1024 * 1024)
.await
.expect("experience body");
let experience: Value = serde_json::from_slice(&body).expect("experience json");
assert_eq!(sidebar_identities(&experience), ["router"]);
assert_eq!(experience["access"]["enabled"], json!(true));
let members = rpc(&app, "mobkit/list_members", json!({})).await;
let member_rows = members["result"].as_array().expect("members array");
assert_eq!(member_rows.len(), 1, "members: {member_rows:#?}");
assert_eq!(
member_rows[0]["state"],
json!("active"),
"member rows must carry the `state` wire key: {member_rows:#?}"
);
let denied = rpc(
&app,
"mobkit/console/send",
json!({
"identity": "delivery",
"content": "hello",
"origin": "test",
"idempotency_key": "denied-send",
}),
)
.await;
assert_eq!(
denied["error"]["code"],
json!(-32030),
"denied: {denied:#?}"
);
assert_eq!(denied["error"]["data"]["kind"], json!("access_denied"));
let allowed = rpc(
&app,
"mobkit/console/send",
json!({
"identity": "router",
"content": "hello",
"origin": "test",
"idempotency_key": "allowed-send",
}),
)
.await;
assert_ne!(
allowed["error"]["code"],
json!(-32030),
"allowed send must not be access-denied: {allowed:#?}"
);
let retire = rpc(&app, "mobkit/retire", json!({ "identity": "router" })).await;
assert_eq!(retire["error"]["code"], json!(-32030));
let labels = rpc(
&app,
"mobkit/mob_labels/set",
json!({ "labels": { "a": "b" } }),
)
.await;
assert_eq!(labels["error"]["code"], json!(-32030));
for method in [
"mobkit/routing/routes/list",
"mobkit/delivery/history",
"mobkit/cross_mob/directory",
"mobkit/mob_labels/get",
"mobkit/list_runs",
"mobkit/list_flows",
] {
let denied = rpc(&app, method, json!({})).await;
assert_eq!(
denied["error"]["code"],
json!(-32030),
"{method} must be access-gated: {denied:#?}"
);
}
let status = rpc(&app, "mobkit/access/status", json!({})).await;
assert_eq!(status["result"]["available"], json!(true));
assert_eq!(status["result"]["enabled"], json!(true));
assert_eq!(status["result"]["can_administer"], json!(false));
let get_config = rpc(&app, "mobkit/access/get", json!({})).await;
assert_eq!(get_config["error"]["code"], json!(-32030));
assert_eq!(
get_status(&app, "/agents/router/events").await,
StatusCode::OK
);
assert_eq!(
get_status(&app, "/agents/delivery/events").await,
StatusCode::FORBIDDEN
);
assert_eq!(get_status(&app, "/mob/events").await, StatusCode::FORBIDDEN);
assert_eq!(
get_status(&app, "/mobkit/mob_events/stream").await,
StatusCode::FORBIDDEN
);
controller
.upsert_rule(AccessRule {
id: "everyone-views-delivery".to_string(),
actions: vec!["agent.view".to_string()],
agents: vec!["delivery".to_string()],
..AccessRule::default()
})
.expect("live rule update");
let response = app
.clone()
.oneshot(
Request::builder()
.method("GET")
.uri("/console/experience")
.body(Body::empty())
.expect("experience request"),
)
.await
.expect("experience response");
let body = to_bytes(response.into_body(), 1024 * 1024)
.await
.expect("experience body");
let experience: Value = serde_json::from_slice(&body).expect("experience json");
assert_eq!(sidebar_identities(&experience), ["delivery", "router"]);
assert_eq!(
get_status(&app, "/agents/delivery/events").await,
StatusCode::OK
);
let _ = runtime.mob_handle().stop().await;
}
#[tokio::test]
async fn structural_sse_replay_fails_closed_when_historical_agent_attributes_are_unknown() {
let (_temp_dir, mut runtime) = build_access_runtime_fixture().await;
let identity = "historical-secret";
let after_seq = retire_historical_secret_member(&runtime, identity).await;
let controller = AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: vec![
AccessRule {
id: "observe-mob".to_string(),
actions: vec!["mob.observe".to_string()],
..AccessRule::default()
},
AccessRule {
id: "view-all".to_string(),
actions: vec!["agent.view".to_string()],
agents: vec!["*".to_string()],
..AccessRule::default()
},
AccessRule {
id: "deny-secret-label".to_string(),
effect: meerkat_mobkit::AccessEffect::Deny,
actions: vec!["agent.view".to_string()],
match_labels: BTreeMap::from([("org".to_string(), "secret".to_string())]),
..AccessRule::default()
},
],
..AccessControlConfig::default()
})
.expect("access controller");
assert!(
!controller.view_for_subject(None).knows_agent(identity),
"the replay regression requires a genuinely cold historical identity"
);
runtime.set_access_controller(controller);
let app = runtime.build_reference_app_router(decision_state(false));
for method in ["mobkit/mob_events/query", "mobkit/mob_events/subscribe"] {
let page = rpc(
&app,
method,
json!({ "after_seq": after_seq, "identity": identity }),
)
.await;
assert_eq!(page["error"], Value::Null, "{method}: {page:#?}");
assert!(
page["result"]["events"]
.as_array()
.is_some_and(Vec::is_empty),
"cold historical attributes must fail closed in {method}: {page:#?}"
);
}
let response = app
.clone()
.oneshot(
Request::builder()
.method("GET")
.uri(format!(
"/mobkit/mob_events/stream?after_seq={after_seq}&identity={identity}"
))
.body(Body::empty())
.expect("historical structural SSE request"),
)
.await
.expect("historical structural SSE response");
assert_sse_replay_emits_no_data(response).await;
runtime.shutdown().await;
}
#[tokio::test]
async fn structural_sse_replay_fails_closed_without_live_visibility_projection() {
#[derive(Debug)]
struct HideHistoricalSecret;
impl ConsoleVisibilityPolicy for HideHistoricalSecret {
fn member_visible(&self, member: &ConsoleMember) -> bool {
member.agent_identity != "historical-secret"
}
}
let (_temp_dir, runtime) = build_access_runtime_fixture().await;
let identity = "historical-secret";
let after_seq = retire_historical_secret_member(&runtime, identity).await;
let app = runtime.build_reference_app_router_with_console_visibility_policy(
decision_state(false),
Arc::new(HideHistoricalSecret),
);
for method in ["mobkit/mob_events/query", "mobkit/mob_events/subscribe"] {
let page = rpc(
&app,
method,
json!({ "after_seq": after_seq, "identity": identity }),
)
.await;
assert_eq!(page["error"], Value::Null, "{method}: {page:#?}");
assert!(
page["result"]["events"]
.as_array()
.is_some_and(Vec::is_empty),
"missing live visibility projection must fail closed in {method}: {page:#?}"
);
}
let response = app
.clone()
.oneshot(
Request::builder()
.method("GET")
.uri(format!(
"/mobkit/mob_events/stream?after_seq={after_seq}&identity={identity}"
))
.body(Body::empty())
.expect("historical visibility SSE request"),
)
.await
.expect("historical visibility SSE response");
assert_sse_replay_emits_no_data(response).await;
runtime.shutdown().await;
}
#[tokio::test]
async fn long_lived_sse_reauthorizes_stale_alias_when_new_generation_appears() {
let (_temp_dir, mut runtime) = build_access_runtime_fixture().await;
let identity = "label-transition";
let controller = AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: vec![
AccessRule {
id: "observe-mob".to_string(),
actions: vec!["mob.observe".to_string()],
..AccessRule::default()
},
AccessRule {
id: "view-all".to_string(),
actions: vec!["agent.view".to_string()],
agents: vec!["*".to_string()],
..AccessRule::default()
},
AccessRule {
id: "deny-secret-label".to_string(),
effect: meerkat_mobkit::AccessEffect::Deny,
actions: vec!["agent.view".to_string()],
match_labels: BTreeMap::from([("org".to_string(), "secret".to_string())]),
..AccessRule::default()
},
],
..AccessControlConfig::default()
})
.expect("access controller");
controller.record_agent_attributes(AgentResourceAttributes {
identity: identity.to_string(),
agent_id: Some(identity.to_string()),
role: Some("lead".to_string()),
labels: BTreeMap::from([("org".to_string(), "public".to_string())]),
});
runtime.set_access_controller(controller.clone());
let app = runtime.build_reference_app_router(decision_state(false));
assert!(controller.view_for_subject(None).knows_agent(identity));
let after_seq = runtime
.mob_handle()
.events()
.latest_cursor()
.await
.expect("cursor before new generation");
let mob_response = app
.clone()
.oneshot(
Request::builder()
.method("GET")
.uri("/mob/events")
.body(Body::empty())
.expect("mob SSE request"),
)
.await
.expect("mob SSE response");
let structural_response = app
.clone()
.oneshot(
Request::builder()
.method("GET")
.uri(format!(
"/mobkit/mob_events/stream?after_seq={after_seq}&identity={identity}"
))
.body(Body::empty())
.expect("structural SSE request"),
)
.await
.expect("structural SSE response");
let mut secret_spec = SpawnMemberSpec::from_wire(
"lead".to_string(),
identity.to_string(),
Some("secret member".into()),
None,
None,
);
secret_spec.labels = Some(BTreeMap::from([("org".to_string(), "secret".to_string())]));
runtime
.spawn(secret_spec)
.await
.expect("spawn secret member");
let member_id = MeerkatId::from(identity);
let mut proof_stream = runtime
.mob_handle()
.subscribe_agent_events(&member_id)
.await
.expect("proof event stream");
runtime
.mob_handle()
.member(&member_id)
.await
.expect("secret member handle")
.send("emit a test event", HandlingMode::Queue)
.await
.expect("secret member turn");
tokio::time::timeout(Duration::from_secs(2), proof_stream.next())
.await
.expect("secret member should emit an event")
.expect("proof stream should remain open");
assert_sse_emits_no_identity(mob_response, identity).await;
assert_sse_emits_no_identity(structural_response, identity).await;
runtime.shutdown().await;
}
#[tokio::test]
async fn event_surfaces_do_not_reauthorize_secret_generation_as_public_same_alias() {
let (_temp_dir, mut runtime) = build_access_runtime_fixture().await;
let identity = "secret-then-public";
let controller = AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: vec![
AccessRule {
id: "observe-mob".to_string(),
actions: vec!["mob.observe".to_string()],
..AccessRule::default()
},
AccessRule {
id: "view-all".to_string(),
actions: vec!["agent.view".to_string()],
agents: vec!["*".to_string()],
..AccessRule::default()
},
AccessRule {
id: "deny-secret-label".to_string(),
effect: meerkat_mobkit::AccessEffect::Deny,
actions: vec!["agent.view".to_string()],
match_labels: BTreeMap::from([("org".to_string(), "secret".to_string())]),
..AccessRule::default()
},
],
..AccessControlConfig::default()
})
.expect("access controller");
runtime.set_access_controller(controller.clone());
let app = runtime.build_reference_app_router(decision_state(false));
let events_view = runtime.mob_handle().events();
let before_secret = events_view
.latest_cursor()
.await
.expect("cursor before secret generation");
let mut secret_spec = SpawnMemberSpec::from_wire(
"lead".to_string(),
identity.to_string(),
Some("secret generation".into()),
None,
None,
);
secret_spec.labels = Some(BTreeMap::from([("org".to_string(), "secret".to_string())]));
runtime
.spawn(secret_spec)
.await
.expect("spawn secret generation");
let secret_spawn_cursor = events_view
.poll_strict(before_secret, 128)
.await
.expect("secret spawn events")
.into_iter()
.find_map(|event| match event.kind {
MobEventKind::MemberSpawned(spawned) if spawned.agent_identity.as_str() == identity => {
Some(event.cursor)
}
_ => None,
})
.expect("secret member_spawned cursor");
let mob_response = app
.clone()
.oneshot(
Request::builder()
.method("GET")
.uri("/mob/events")
.body(Body::empty())
.expect("mob SSE request"),
)
.await
.expect("mob SSE response");
let member_id = MeerkatId::from(identity);
let mut proof_stream = runtime
.mob_handle()
.subscribe_agent_events(&member_id)
.await
.expect("secret proof event stream");
runtime
.mob_handle()
.member(&member_id)
.await
.expect("secret member handle")
.send("emit a secret event", HandlingMode::Queue)
.await
.expect("secret member turn");
tokio::time::timeout(Duration::from_secs(2), proof_stream.next())
.await
.expect("secret member should emit an event")
.expect("proof stream should remain open");
let before_public = events_view
.latest_cursor()
.await
.expect("cursor before public alias projection");
let public_runtime_identity = "zz-public-incarnation";
let mut public_spec = SpawnMemberSpec::from_wire(
"lead".to_string(),
public_runtime_identity.to_string(),
Some("public alias incarnation".into()),
None,
None,
);
public_spec.labels = Some(BTreeMap::from([
("agent_identity".to_string(), identity.to_string()),
("org".to_string(), "public".to_string()),
]));
runtime
.mob_handle()
.spawn_spec(public_spec)
.await
.expect("spawn exact public binding for the reused alias");
let public_spawn_cursor = events_view
.poll_strict(before_public, 128)
.await
.expect("public alias events")
.into_iter()
.find_map(|event| match event.kind {
MobEventKind::MemberSpawned(spawned)
if spawned.agent_identity.as_str() == public_runtime_identity =>
{
Some(event.cursor)
}
_ => None,
})
.expect("public alias member_spawned cursor");
for method in ["mobkit/mob_events/query", "mobkit/mob_events/subscribe"] {
let page = rpc(
&app,
method,
json!({
"after_seq": before_secret,
"event_types": ["member_spawned"],
"limit": 1,
}),
)
.await;
assert_eq!(page["error"], Value::Null, "{method}: {page:#?}");
let cursors = page["result"]["events"]
.as_array()
.expect("event array")
.iter()
.filter_map(|event| event["cursor"].as_u64())
.collect::<Vec<_>>();
assert!(
!cursors.contains(&secret_spawn_cursor),
"historical secret generation leaked through {method}: {page:#?}"
);
assert!(
cursors.contains(&public_spawn_cursor),
"current public binding should remain visible in {method}: {page:#?}"
);
assert_eq!(
cursors.len(),
1,
"authorization-denied rows must not consume the visible page limit: {page:#?}"
);
assert_eq!(page["result"]["next_after_seq"], public_spawn_cursor);
}
let structural_response = app
.clone()
.oneshot(
Request::builder()
.method("GET")
.uri(format!(
"/mobkit/mob_events/stream?after_seq={before_secret}"
))
.body(Body::empty())
.expect("structural SSE request"),
)
.await
.expect("structural SSE response");
assert_structural_sse_excludes_and_includes_cursor(
structural_response,
secret_spawn_cursor,
public_spawn_cursor,
)
.await;
let before_latest_secret = events_view
.latest_cursor()
.await
.expect("cursor before latest secret alias projection");
let latest_secret_runtime_identity = "zzz-secret-incarnation";
let mut latest_secret_spec = SpawnMemberSpec::from_wire(
"lead".to_string(),
latest_secret_runtime_identity.to_string(),
Some("latest secret alias incarnation".into()),
None,
None,
);
latest_secret_spec.labels = Some(BTreeMap::from([
("agent_identity".to_string(), identity.to_string()),
("org".to_string(), "secret".to_string()),
]));
runtime
.mob_handle()
.spawn_spec(latest_secret_spec)
.await
.expect("spawn latest exact secret binding");
let latest_secret_cursor = events_view
.poll_strict(before_latest_secret, 128)
.await
.expect("latest secret alias events")
.into_iter()
.find_map(|event| match event.kind {
MobEventKind::MemberSpawned(spawned)
if spawned.agent_identity.as_str() == latest_secret_runtime_identity =>
{
Some(event.cursor)
}
_ => None,
})
.expect("latest secret alias member_spawned cursor");
let latest_secret_frontier = events_view
.latest_cursor()
.await
.expect("raw frontier after latest secret alias projection");
let backward = rpc(
&app,
"mobkit/mob_events/query",
json!({ "event_types": ["member_spawned"], "limit": 1 }),
)
.await;
assert_eq!(backward["error"], Value::Null, "{backward:#?}");
assert_eq!(
backward["result"]["events"][0]["cursor"],
public_spawn_cursor
);
assert_eq!(backward["result"]["events"].as_array().unwrap().len(), 1);
let backward_next = backward["result"]["next_after_seq"]
.as_u64()
.expect("backward next_after_seq");
let raw_frontier_after_backward = events_view
.latest_cursor()
.await
.expect("raw frontier after backward query");
assert!(
backward_next >= latest_secret_frontier && backward_next <= raw_frontier_after_backward,
"backwards snapshot must expose the raw ledger frontier (captured \
{latest_secret_frontier}, reported {backward_next}, raw now \
{raw_frontier_after_backward})"
);
let empty_subscribe = rpc(
&app,
"mobkit/mob_events/subscribe",
json!({
"after_seq": public_spawn_cursor,
"event_types": ["member_spawned"],
"limit": 1,
}),
)
.await;
assert_eq!(
empty_subscribe["error"],
Value::Null,
"{empty_subscribe:#?}"
);
assert!(
empty_subscribe["result"]["events"]
.as_array()
.is_some_and(Vec::is_empty),
"hidden-only snapshot should remain empty: {empty_subscribe:#?}"
);
let empty_next = empty_subscribe["result"]["next_after_seq"]
.as_u64()
.expect("empty subscribe next_after_seq");
let raw_frontier_after_subscribe = events_view
.latest_cursor()
.await
.expect("raw frontier after empty subscribe");
assert!(
empty_next >= latest_secret_frontier && empty_next <= raw_frontier_after_subscribe,
"empty page must advance over the hidden raw row to a real ledger \
frontier (captured {latest_secret_frontier}, reported {empty_next}, \
raw now {raw_frontier_after_subscribe})"
);
let before_next_public = events_view
.latest_cursor()
.await
.expect("cursor before next public alias projection");
let next_public_runtime_identity = "zzzz-public-incarnation";
let mut next_public_spec = SpawnMemberSpec::from_wire(
"lead".to_string(),
next_public_runtime_identity.to_string(),
Some("next public alias incarnation".into()),
None,
None,
);
next_public_spec.labels = Some(BTreeMap::from([
("agent_identity".to_string(), identity.to_string()),
("org".to_string(), "public".to_string()),
]));
runtime
.mob_handle()
.spawn_spec(next_public_spec)
.await
.expect("spawn next exact public binding");
let next_public_cursor = events_view
.poll_strict(before_next_public, 128)
.await
.expect("next public alias events")
.into_iter()
.find_map(|event| match event.kind {
MobEventKind::MemberSpawned(spawned)
if spawned.agent_identity.as_str() == next_public_runtime_identity =>
{
Some(event.cursor)
}
_ => None,
})
.expect("next public alias member_spawned cursor");
let subscribe_url = empty_subscribe["result"]["subscribe_url"]
.as_str()
.expect("subscribe URL");
let resumed_response = app
.clone()
.oneshot(
Request::builder()
.method("GET")
.uri(subscribe_url)
.body(Body::empty())
.expect("resumed structural SSE request"),
)
.await
.expect("resumed structural SSE response");
assert_structural_sse_excludes_and_includes_cursor(
resumed_response,
latest_secret_cursor,
next_public_cursor,
)
.await;
assert_sse_emits_no_identity(mob_response, identity).await;
runtime.shutdown().await;
}
#[tokio::test]
async fn public_member_ingress_rejects_encoded_roster_ids_before_abac_and_resolution() {
let (_temp_dir, mut runtime) = build_access_runtime_fixture().await;
let controller = AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: vec![
AccessRule {
id: "everyone-views".to_string(),
actions: vec!["agent.view".to_string()],
..AccessRule::default()
},
AccessRule {
id: "deny-delivery-view".to_string(),
effect: meerkat_mobkit::AccessEffect::Deny,
actions: vec!["agent.view".to_string()],
agents: vec!["delivery".to_string()],
..AccessRule::default()
},
AccessRule {
id: "everyone-spawns".to_string(),
actions: vec!["agent.spawn".to_string()],
..AccessRule::default()
},
],
..AccessControlConfig::default()
})
.expect("controller");
runtime.set_access_controller(controller);
let app = runtime.build_reference_app_router(decision_state(false));
let encoded_rpc = rpc(
&app,
"mobkit/get_member",
json!({ "member_id": "mk--delivery" }),
)
.await;
assert_eq!(encoded_rpc["error"]["code"], json!(-32602));
assert!(
encoded_rpc["error"]["message"]
.as_str()
.is_some_and(|message| message.contains("encoded roster-id")),
"{encoded_rpc:#?}"
);
assert_eq!(
get_status(&app, "/agents/mk--delivery/events").await,
StatusCode::BAD_REQUEST,
"per-agent SSE must reject encoded roster ids before its agent.view decision"
);
assert_eq!(
get_status(&app, "/agents/%20delivery%20/events").await,
StatusCode::FORBIDDEN,
"per-agent SSE must trim the public alias before its targeted agent.view decision"
);
let forged_label = rpc(
&app,
"mobkit/ensure_member",
json!({
"role": "lead",
"agent_identity": "raw-imposter",
"labels": { "agent_identity": "delivery" },
}),
)
.await;
assert_eq!(forged_label["error"]["code"], json!(-32602));
assert!(
forged_label["error"]["message"]
.as_str()
.is_some_and(|message| message.contains("runtime-authoritative")),
"{forged_label:#?}"
);
#[derive(Debug)]
struct HideDelivery;
impl ConsoleVisibilityPolicy for HideDelivery {
fn member_visible(&self, member: &ConsoleMember) -> bool {
member.agent_identity != "delivery"
}
}
let hidden_app = runtime.build_reference_app_router_with_console_visibility_policy(
decision_state(false),
Arc::new(HideDelivery),
);
assert_eq!(
get_status(&hidden_app, "/agents/delivery/events").await,
StatusCode::NOT_FOUND,
"a visibility-hidden agent must be rejected before SSE subscription"
);
let _ = runtime.mob_handle().stop().await;
}
#[tokio::test]
async fn topology_plan_and_audit_capabilities_follow_endpoint_scoped_grants() {
let (_temp_dir, mut runtime) = build_access_runtime_fixture().await;
runtime
.set_topology_control_policy(TopologyControlPolicy {
mode: TopologyControlMode::ReadOnly,
..TopologyControlPolicy::default()
})
.expect("read-only topology policy");
let controller = AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: vec![AccessRule {
id: "anonymous-topology-view".to_string(),
actions: vec!["agent.view".to_string(), "topology.view".to_string()],
agents: vec!["*".to_string()],
..AccessRule::default()
}],
..AccessControlConfig::default()
})
.expect("topology controller");
runtime.set_access_controller(controller.clone());
let app = runtime.build_reference_app_router(decision_state(false));
let capabilities = rpc(&app, "mobkit/capabilities", json!({})).await;
let methods = capabilities["result"]["methods"]
.as_array()
.expect("capability methods");
assert!(
methods
.iter()
.any(|method| method == "mobkit/topology/query")
);
assert!(
!methods
.iter()
.any(|method| method == "mobkit/topology/plan"),
"view-only callers must not be offered a planning oracle: {capabilities:#?}"
);
assert!(
!methods
.iter()
.any(|method| method == "mobkit/topology/audit/query"),
"audit is an independent sensitive permission: {capabilities:#?}"
);
let query = rpc(&app, "mobkit/topology/query", json!({})).await;
let revision = query["result"]["revision"]
.as_u64()
.expect("topology revision");
let connect = json!({
"expected_revision": revision,
"operations": [{
"action": "connect",
"edge": {
"a": {"identity": "router"},
"b": {"identity": "delivery"}
}
}]
});
let denied_without_action = rpc(&app, "mobkit/topology/plan", connect.clone()).await;
assert_eq!(denied_without_action["error"]["code"], json!(-32030));
assert_eq!(
denied_without_action["error"]["data"]["action"],
json!("topology.connect")
);
let denied_audit = rpc(&app, "mobkit/topology/audit/query", json!({})).await;
assert_eq!(denied_audit["error"]["code"], json!(-32030));
assert_eq!(
denied_audit["error"]["data"]["action"],
json!("topology.audit")
);
controller
.upsert_rule(AccessRule {
id: "connect-router-only".to_string(),
actions: vec!["topology.connect".to_string()],
agents: vec!["router".to_string()],
..AccessRule::default()
})
.expect("live connect grant");
let capabilities = rpc(&app, "mobkit/capabilities", json!({})).await;
assert!(
capabilities["result"]["methods"]
.as_array()
.expect("capability methods")
.iter()
.any(|method| method == "mobkit/topology/plan")
);
let denied_other_endpoint = rpc(&app, "mobkit/topology/plan", connect).await;
assert_eq!(denied_other_endpoint["error"]["code"], json!(-32030));
assert_eq!(
denied_other_endpoint["error"]["data"]["resource"],
json!("delivery")
);
let _ = runtime.mob_handle().stop().await;
}
#[tokio::test]
async fn multipart_send_denied_does_not_write_blob_before_access_gate() {
let (_temp_dir, mut runtime) = build_access_runtime_fixture().await;
let controller = AccessController::new(anonymous_router_only_config()).expect("controller");
runtime.set_access_controller(controller);
let app = runtime.build_reference_app_router(decision_state(false));
let image_bytes = b"denied-multipart-png-bytes";
let boundary = "mobkit-access-boundary";
let payload = json!({
"jsonrpc": "2.0",
"id": "denied-multipart-send",
"method": "mobkit/console/send",
"params": {
"identity": "delivery",
"origin": "test",
"idempotency_key": "denied-multipart-send",
"content": [
{ "type": "image_upload", "upload_id": "u1", "media_type": "image/png" }
]
}
});
let mut body = Vec::new();
body.extend_from_slice(format!("--{boundary}\r\n").as_bytes());
body.extend_from_slice(b"Content-Disposition: form-data; name=\"payload\"\r\n");
body.extend_from_slice(b"Content-Type: application/json\r\n\r\n");
body.extend_from_slice(payload.to_string().as_bytes());
body.extend_from_slice(b"\r\n");
body.extend_from_slice(format!("--{boundary}\r\n").as_bytes());
body.extend_from_slice(
b"Content-Disposition: form-data; name=\"file:u1\"; filename=\"x.png\"\r\n",
);
body.extend_from_slice(b"Content-Type: image/png\r\n\r\n");
body.extend_from_slice(image_bytes);
body.extend_from_slice(b"\r\n");
body.extend_from_slice(format!("--{boundary}--\r\n").as_bytes());
let response = app
.clone()
.oneshot(
Request::builder()
.method("POST")
.uri("/console/rpc/multipart")
.header(
header::CONTENT_TYPE,
format!("multipart/form-data; boundary={boundary}"),
)
.body(Body::from(body))
.expect("multipart request"),
)
.await
.expect("multipart response");
assert_eq!(response.status(), StatusCode::OK);
let resp_body = to_bytes(response.into_body(), 1024 * 1024)
.await
.expect("multipart body");
let resp_json: Value = serde_json::from_slice(&resp_body).expect("multipart json");
assert_eq!(
resp_json["error"]["code"],
json!(-32030),
"denied multipart send must be access-denied: {resp_json:#?}"
);
assert_eq!(resp_json["error"]["data"]["kind"], json!("access_denied"));
let mut hasher = Sha256::new();
hasher.update(b"image/png");
hasher.update([0]);
hasher.update(image_bytes);
let blob_id = format!("sha256:{:x}", hasher.finalize());
let blob_status = get_status(&app, &format!("/blobs/{blob_id}")).await;
assert_eq!(
blob_status,
StatusCode::NOT_FOUND,
"denied upload bytes must not be persisted before the access gate"
);
let _ = runtime.mob_handle().stop().await;
}
#[tokio::test]
async fn bootstrap_path_configures_access_from_console_then_locks_down() {
let (_temp_dir, mut runtime) = build_access_runtime_fixture().await;
let controller = AccessController::disabled();
runtime.set_access_controller(controller.clone());
let app = runtime.build_reference_app_router(decision_state(false));
let status = rpc(&app, "mobkit/access/status", json!({})).await;
assert_eq!(status["result"]["enabled"], json!(false));
assert_eq!(status["result"]["can_administer"], json!(true));
let set = rpc(
&app,
"mobkit/access/set",
json!({ "config": {
"enabled": true,
"admins": ["root@example.test"],
"rules": [
{ "id": "everyone-views-router", "actions": ["agent.view"], "agents": ["router"] }
],
}}),
)
.await;
assert_eq!(set["error"], Value::Null, "set: {set:#?}");
assert_eq!(set["result"]["revision"], json!(1));
let status = rpc(&app, "mobkit/access/status", json!({})).await;
assert_eq!(status["result"]["enabled"], json!(true));
assert_eq!(status["result"]["can_administer"], json!(false));
let get_config = rpc(&app, "mobkit/access/get", json!({})).await;
assert_eq!(get_config["error"]["code"], json!(-32030));
let members = rpc(&app, "mobkit/list_members", json!({})).await;
assert_eq!(members["result"].as_array().expect("members").len(), 1);
assert!(
controller
.replace_config(AccessControlConfig {
enabled: true,
..AccessControlConfig::default()
})
.is_err()
);
let _ = runtime.mob_handle().stop().await;
}
#[tokio::test]
async fn timeline_role_rules_resolve_without_prior_experience() {
let (_temp_dir, mut runtime) = build_access_runtime_fixture().await;
let controller = AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: vec![AccessRule {
id: "lead-role-only".to_string(),
actions: vec!["agent.view".to_string(), "agent.send".to_string()],
roles: vec!["lead".to_string()],
..AccessRule::default()
}],
..AccessControlConfig::default()
})
.expect("controller");
runtime.set_access_controller(controller);
let app = runtime.build_reference_app_router(decision_state(false));
let send = rpc(
&app,
"mobkit/console/send",
json!({
"identity": "router",
"content": "ping",
"origin": "test",
"idempotency_key": "role-cold-send",
}),
)
.await;
assert_eq!(
send["error"],
Value::Null,
"role-keyed send must resolve cold: {send:#?}"
);
let page = rpc(
&app,
"mobkit/console/query_timeline",
json!({ "mode": "recent", "limit": 200 }),
)
.await;
assert_eq!(
page["error"],
Value::Null,
"cold timeline query must succeed: {page:#?}"
);
if let Some(frames) = page["result"]["frames"].as_array() {
assert!(
frames
.iter()
.all(|frame| frame["identity"] == json!("router")),
"only role-matched agents are visible: {frames:#?}"
);
}
assert_eq!(
get_status(&app, "/console/timeline/stream").await,
StatusCode::OK
);
let _ = runtime.mob_handle().stop().await;
}
#[tokio::test]
async fn rest_console_send_resolves_label_deny_on_cold_cache() {
let (_temp_dir, mut runtime) = build_access_runtime_fixture().await;
runtime
.spawn(
SpawnMemberSpec::from_wire(
"lead".to_string(),
MeerkatId::from("red-shadow").to_string(),
Some("You are red-shadow.".into()),
None,
None,
)
.with_labels(BTreeMap::from([("team".to_string(), "red".to_string())])),
)
.await
.expect("spawn labeled member");
let controller = AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: vec![
AccessRule {
id: "everyone-sends".to_string(),
actions: vec!["agent.send".to_string()],
..AccessRule::default()
},
AccessRule {
id: "deny-red-team-send".to_string(),
effect: meerkat_mobkit::AccessEffect::Deny,
actions: vec!["agent.send".to_string()],
match_labels: BTreeMap::from([("team".to_string(), "red".to_string())]),
..AccessRule::default()
},
],
..AccessControlConfig::default()
})
.expect("controller");
runtime.set_access_controller(controller);
let app = runtime.build_reference_app_router(decision_state(false));
let rest_send = |identity: &'static str, idempotency_key: &'static str| {
let app = app.clone();
async move {
app.oneshot(
Request::builder()
.method("POST")
.uri("/console/send")
.header(header::CONTENT_TYPE, "application/json")
.body(Body::from(
json!({
"identity": identity,
"content": "hello",
"origin": "test",
"idempotency_key": idempotency_key,
})
.to_string(),
))
.expect("send request"),
)
.await
.expect("send response")
.status()
}
};
assert_eq!(
rest_send("red-shadow", "cold-deny-send").await,
StatusCode::FORBIDDEN,
"label-scoped deny must resolve on a cold cache"
);
assert_ne!(
rest_send("router", "cold-allow-send").await,
StatusCode::FORBIDDEN,
"non-matching members must keep passing the send gate"
);
let _ = runtime.mob_handle().stop().await;
}
#[tokio::test]
async fn timeline_rpc_is_filtered_per_caller() {
let (_temp_dir, mut runtime) = build_access_runtime_fixture().await;
let controller = AccessController::new(anonymous_router_only_config()).expect("controller");
runtime.set_access_controller(controller);
let app = runtime.build_reference_app_router(decision_state(false));
let _ = rpc(
&app,
"mobkit/console/send",
json!({
"identity": "router",
"content": "ping",
"origin": "test",
"idempotency_key": "timeline-send",
}),
)
.await;
let page = rpc(
&app,
"mobkit/console/query_timeline",
json!({ "mode": "recent", "limit": 200 }),
)
.await;
let frames = page["result"]["frames"].as_array().expect("frames");
assert!(
frames
.iter()
.all(|frame| frame["identity"] == json!("router")),
"timeline must only contain visible identities: {frames:#?}"
);
let _ = runtime.mob_handle().stop().await;
}
#[tokio::test]
async fn system_memory_frames_respect_timeline_access_filtering() {
let (_temp_dir, mut runtime) = build_access_runtime_fixture().await;
runtime.memory_event_sink().emit(
meerkat_mobkit::memory::events::MemoryTimelineEvent::DreamCompleted {
realm: "default".to_string(),
run_id: "run-acl-dream-1".to_string(),
ops_committed: 1,
detail: json!({ "phase": "test" }),
},
);
let system_frames = |page: &Value| -> Vec<Value> {
page["result"]["frames"]
.as_array()
.expect("frames")
.iter()
.filter(|frame| frame["identity"] == json!("_system"))
.cloned()
.collect()
};
let open_controller = AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: vec![AccessRule {
id: "everyone-views-everything".to_string(),
actions: vec!["agent.view".to_string()],
..AccessRule::default()
}],
..AccessControlConfig::default()
})
.expect("open controller");
runtime.set_access_controller(open_controller);
let open_app = runtime.build_reference_app_router(decision_state(false));
let mut seen = Vec::new();
for _ in 0..80 {
let page = rpc(
&open_app,
"mobkit/console/query_timeline",
json!({ "mode": "recent", "limit": 200 }),
)
.await;
seen = system_frames(&page);
if !seen.is_empty() {
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
assert!(
seen.iter()
.any(|frame| frame["kind"] == json!("memory.dream.completed")),
"unscoped agent.view must see the _system memory event: {seen:#?}"
);
let scoped_controller =
AccessController::new(anonymous_router_only_config()).expect("scoped controller");
runtime.set_access_controller(scoped_controller);
let scoped_app = runtime.build_reference_app_router(decision_state(false));
let page = rpc(
&scoped_app,
"mobkit/console/query_timeline",
json!({ "mode": "recent", "limit": 200 }),
)
.await;
assert!(
system_frames(&page).is_empty(),
"agent-scoped viewers must not see _system frames: {page:#?}"
);
let _ = runtime.mob_handle().stop().await;
}
#[tokio::test]
async fn mob_observe_does_not_bypass_per_agent_view() {
let (_temp_dir, mut runtime) = build_access_runtime_fixture().await;
let controller = AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: vec![
AccessRule {
id: "everyone-observes".to_string(),
actions: vec!["mob.observe".to_string()],
..AccessRule::default()
},
AccessRule {
id: "everyone-views-router".to_string(),
actions: vec!["agent.view".to_string()],
agents: vec!["router".to_string()],
..AccessRule::default()
},
],
..AccessControlConfig::default()
})
.expect("controller");
runtime.set_access_controller(controller);
let app = runtime.build_reference_app_router(decision_state(false));
let page = rpc(&app, "mobkit/mob_events/query", json!({})).await;
assert_eq!(
page["error"],
Value::Null,
"observe surface open: {page:#?}"
);
let events = page["result"]["events"].as_array().expect("events");
assert!(
events.iter().all(|event| {
event["agent_identity"].is_null() || event["agent_identity"] == json!("router")
}),
"mob.observe must not surface denied agents' events: {events:#?}"
);
assert!(
events
.iter()
.any(|event| event["agent_identity"] == json!("router")),
"router's own ledger entries should still be visible: {events:#?}"
);
let _ = runtime.mob_handle().stop().await;
}
async fn seeded_memory_store(
root: &std::path::Path,
) -> (
meerkat_mobkit::SqliteAgentMemoryStore,
String,
String,
String,
String,
String,
) {
use meerkat_mobkit::memory::records::{
InjectionLogEntry, InjectionSurface, MemoryAuthor, MemoryKind,
};
use meerkat_mobkit::{
AgentMemoryProvider, MemoryScope, NewMemoryRecord, SqliteAgentMemoryStore,
StagedMemoryStore, StagedMutationBatch, StagedOp, TrustTier,
};
struct AlwaysQuarantine;
impl meerkat_mobkit::memory::taint::LlmWriteGate for AlwaysQuarantine {
fn quarantine_reason(
&self,
author: &MemoryAuthor,
_kind: meerkat_mobkit::memory::staged::StagedBatchKind,
_evidence: &[meerkat_mobkit::memory::records::EvidenceRef],
) -> Option<String> {
author.is_llm().then(|| "test taint".to_string())
}
}
let store = SqliteAgentMemoryStore::open(root).expect("open store");
let record = |title: &str| NewMemoryRecord {
kind: MemoryKind::Fact,
title: title.to_string(),
description: format!("{title} description"),
body: format!("{title} body"),
tags: Vec::new(),
evidence: Vec::new(),
verification: None,
};
let identity_scope = |identity: &str| MemoryScope::Identity {
realm: "default".to_string(),
identity: identity.to_string(),
};
let router_root = store
.remember_authored(
&identity_scope("router"),
record("Router root"),
MemoryAuthor::Operator,
)
.await
.expect("router root");
let router_tip = store
.supersede_authored(
&identity_scope("router"),
&router_root.memory_id,
record("Router tip"),
MemoryAuthor::Operator,
)
.await
.expect("router tip");
let delivery = store
.remember_authored(
&identity_scope("delivery"),
record("Delivery fact"),
MemoryAuthor::Operator,
)
.await
.expect("delivery record");
let mob_record = store
.remember_authored(
&MemoryScope::Mob {
realm: "default".to_string(),
mob: "access-control-mob".to_string(),
},
record("Mob convention"),
MemoryAuthor::Operator,
)
.await
.expect("mob record");
store
.remember_authored(
&MemoryScope::Realm {
realm: "default".to_string(),
},
record("Realm fact"),
MemoryAuthor::Operator,
)
.await
.expect("realm record");
let operator_record = store
.remember_authored(
&MemoryScope::Operator {
realm: "default".to_string(),
operator: "op-luka".to_string(),
},
record("Operator preference"),
MemoryAuthor::Operator,
)
.await
.expect("operator record");
store.set_llm_write_gate(Arc::new(AlwaysQuarantine));
let quarantined = store
.remember_authored(
&identity_scope("router"),
record("Router quarantined claim"),
MemoryAuthor::Agent {
identity: "router".to_string(),
},
)
.await
.expect("quarantined record");
assert!(
matches!(
quarantined.status,
meerkat_mobkit::memory::records::RecordStatus::Quarantined { .. }
),
"seed record should land quarantined: {quarantined:?}"
);
store
.log_injections(
"default",
&[InjectionLogEntry {
record_id: router_tip.memory_id.clone(),
identity: "router".to_string(),
session_key: Some("sess-1".to_string()),
surface: InjectionSurface::Build,
at_ms: 1,
}],
)
.await
.expect("injection row");
let token = store
.stage(StagedMutationBatch {
kind: meerkat_mobkit::memory::staged::StagedBatchKind::FreshWrite,
realm: "default".to_string(),
author: MemoryAuthor::Steward {
run_id: "run-dream-1".to_string(),
},
ops: vec![StagedOp::Create {
id: None,
scope: MemoryScope::Mob {
realm: "default".to_string(),
mob: "access-control-mob".to_string(),
},
record: record("Dream consolidated"),
trust: TrustTier::AgentObserved,
derived_from: Vec::new(),
rationale: Some("consolidated during dream".to_string()),
created_at_ms: None,
updated_at_ms: None,
}],
})
.await
.expect("stage dream batch");
store.commit(token).await.expect("commit dream batch");
let stage_promotion = |scope: MemoryScope, title: &str| {
let store = store.clone();
let record = record(title);
async move {
store
.stage(StagedMutationBatch {
kind: meerkat_mobkit::memory::staged::StagedBatchKind::FreshWrite,
realm: "default".to_string(),
author: MemoryAuthor::Steward {
run_id: "run-gate-1".to_string(),
},
ops: vec![StagedOp::Create {
id: None,
scope,
record,
trust: TrustTier::AgentObserved,
derived_from: Vec::new(),
rationale: Some("gated promotion".to_string()),
created_at_ms: None,
updated_at_ms: None,
}],
})
.await
.expect("stage promotion batch")
}
};
let mob_stage = stage_promotion(
MemoryScope::Mob {
realm: "default".to_string(),
mob: "access-control-mob".to_string(),
},
"Promoted mob claim",
)
.await;
store
.record_pending_promotion(
"default",
meerkat_mobkit::memory::PendingPromotion {
pending_id: "gate-mob-promotion".to_string(),
stage_token: mob_stage.token,
record_id: quarantined.memory_id.clone(),
scope_kind: "mob".to_string(),
scope_key: "access-control-mob".to_string(),
rationale: Some("steward: mob-wide convention".to_string()),
status: "pending".to_string(),
created_at_ms: 2,
},
)
.await
.expect("mob promotion row");
let delivery_stage =
stage_promotion(identity_scope("delivery"), "Promoted delivery claim").await;
store
.record_pending_promotion(
"default",
meerkat_mobkit::memory::PendingPromotion {
pending_id: "gate-delivery-promotion".to_string(),
stage_token: delivery_stage.token,
record_id: quarantined.memory_id.clone(),
scope_kind: "identity".to_string(),
scope_key: "delivery".to_string(),
rationale: Some("steward: delivery personal fact".to_string()),
status: "pending".to_string(),
created_at_ms: 3,
},
)
.await
.expect("delivery promotion row");
(
store,
router_tip.memory_id,
quarantined.memory_id,
delivery.memory_id,
mob_record.memory_id,
operator_record.memory_id,
)
}
#[tokio::test]
async fn memory_panel_reads_seeded_store_without_access_control() {
let (_temp_dir, runtime) = build_access_runtime_fixture().await;
let memory_dir = tempfile::tempdir().expect("memory dir");
let (store, tip_id, quarantined_id, _delivery_id, _mob_id, _operator_id) =
seeded_memory_store(memory_dir.path()).await;
runtime.set_memory_panel_store(Arc::new(store.clone()));
let app = runtime.build_reference_app_router(decision_state(false));
let capabilities = rpc(&app, "mobkit/capabilities", json!({})).await;
let methods = capabilities["result"]["methods"]
.as_array()
.expect("methods");
for method in [
"mobkit/memory/panel/records",
"mobkit/memory/panel/record",
"mobkit/memory/panel/quarantine",
"mobkit/memory/panel/dreams",
] {
assert!(
methods.iter().any(|value| value == method),
"{method} must be advertised: {methods:#?}"
);
}
let records = rpc(&app, "mobkit/memory/panel/records", json!({})).await;
assert_eq!(records["error"], Value::Null, "{records:#?}");
let rows = records["result"]["records"].as_array().expect("records");
assert!(rows.len() >= 5, "all seeded records: {rows:#?}");
assert!(
rows.iter().all(|row| row.get("body").is_none()),
"list rows must be body-free: {rows:#?}"
);
assert!(
rows.iter()
.any(|row| row["status"]["status"] == json!("quarantined")),
"quarantined row visible without enforcement: {rows:#?}"
);
let router_rows = rpc(
&app,
"mobkit/memory/panel/records",
json!({ "identity": "router" }),
)
.await;
let router_rows = router_rows["result"]["records"]
.as_array()
.expect("router rows")
.clone();
assert!(!router_rows.is_empty());
assert!(
router_rows
.iter()
.all(|row| row["scope"]["identity"] == json!("router")),
"identity filter leaked other scopes: {router_rows:#?}"
);
let detail = rpc(
&app,
"mobkit/memory/panel/record",
json!({ "memory_id": tip_id }),
)
.await;
assert_eq!(detail["error"], Value::Null, "{detail:#?}");
assert_eq!(detail["result"]["record"]["body"], json!("Router tip body"));
let chain = detail["result"]["chain"].as_array().expect("chain");
assert_eq!(chain.len(), 2, "root + tip: {chain:#?}");
assert_eq!(chain[0]["status"]["status"], json!("superseded"));
assert_eq!(chain[1]["id"], json!(tip_id));
let injections = detail["result"]["injections"]
.as_array()
.expect("injections");
assert_eq!(injections.len(), 1, "{injections:#?}");
assert_eq!(injections[0]["surface"], json!("build"));
let quarantine = rpc(&app, "mobkit/memory/panel/quarantine", json!({})).await;
assert_eq!(quarantine["error"], Value::Null, "{quarantine:#?}");
let queue = quarantine["result"]["records"].as_array().expect("queue");
assert!(
queue.iter().any(|row| row["id"] == json!(quarantined_id)),
"{queue:#?}"
);
let promotions = quarantine["result"]["pending_promotions"]
.as_array()
.expect("promotions");
assert_eq!(promotions.len(), 2, "{promotions:#?}");
assert!(
promotions
.iter()
.all(|row| row.get("stage_token").is_none()),
"stage_token is a commit capability and must never surface: {promotions:#?}"
);
let dreams = rpc(&app, "mobkit/memory/panel/dreams", json!({})).await;
assert_eq!(dreams["error"], Value::Null, "{dreams:#?}");
let runs = dreams["result"]["runs"].as_array().expect("runs");
assert_eq!(runs.len(), 1, "{runs:#?}");
assert_eq!(runs[0]["run_id"], json!("run-dream-1"));
assert_eq!(runs[0]["op_kinds"]["create"], json!(1));
let response = app
.clone()
.oneshot(
Request::builder()
.method("GET")
.uri("/console/experience")
.body(Body::empty())
.expect("experience request"),
)
.await
.expect("experience response");
let body = to_bytes(response.into_body(), 1024 * 1024)
.await
.expect("experience body");
let experience: Value = serde_json::from_slice(&body).expect("experience json");
assert_eq!(experience["memory"]["available"], json!(true));
assert_eq!(experience["memory"]["can_read"], json!(true));
assert_eq!(experience["memory"]["can_review_quarantine"], json!(true));
let _ = runtime.mob_handle().stop().await;
}
#[tokio::test]
async fn memory_panel_enforces_scope_actions_end_to_end() {
let (_temp_dir, mut runtime) = build_access_runtime_fixture().await;
let memory_dir = tempfile::tempdir().expect("memory dir");
let (store, tip_id, quarantined_id, delivery_id, mob_id, operator_id) =
seeded_memory_store(memory_dir.path()).await;
runtime.set_memory_panel_store(Arc::new(store.clone()));
let controller = AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: vec![
AccessRule {
id: "view-router".to_string(),
actions: vec!["agent.view".to_string()],
agents: vec!["router".to_string()],
..AccessRule::default()
},
AccessRule {
id: "read-router-memory".to_string(),
actions: vec!["agent.memory.read".to_string()],
agents: vec!["router".to_string()],
..AccessRule::default()
},
],
..AccessControlConfig::default()
})
.expect("controller");
runtime.set_access_controller(controller.clone());
let app = runtime.build_reference_app_router(decision_state(false));
let response = app
.clone()
.oneshot(
Request::builder()
.method("GET")
.uri("/console/experience")
.body(Body::empty())
.expect("experience request"),
)
.await
.expect("experience response");
let body = to_bytes(response.into_body(), 1024 * 1024)
.await
.expect("experience body");
let experience: Value = serde_json::from_slice(&body).expect("experience json");
assert_eq!(experience["memory"]["can_read"], json!(true));
assert_eq!(experience["memory"]["can_review_quarantine"], json!(false));
let records = rpc(&app, "mobkit/memory/panel/records", json!({})).await;
assert_eq!(records["error"], Value::Null, "{records:#?}");
let rows = records["result"]["records"].as_array().expect("records");
assert!(
rows.iter()
.all(|row| row["scope"]["identity"] == json!("router")),
"only router-scope rows may survive: {rows:#?}"
);
assert!(
rows.iter().all(|row| row["id"] != json!(quarantined_id)),
"quarantined rows need the review grant: {rows:#?}"
);
let denied = rpc(
&app,
"mobkit/memory/panel/records",
json!({ "identity": "delivery" }),
)
.await;
assert_eq!(denied["error"]["code"], json!(-32030), "{denied:#?}");
let allowed = rpc(
&app,
"mobkit/memory/panel/record",
json!({ "memory_id": tip_id }),
)
.await;
assert_eq!(allowed["error"], Value::Null, "{allowed:#?}");
for (memory_id, action) in [
(&delivery_id, "agent.memory.read"),
(&mob_id, "mob.memory.read"),
(&operator_id, "operator.memory.read"),
(&quarantined_id, "memory.quarantine.review"),
] {
let denied = rpc(
&app,
"mobkit/memory/panel/record",
json!({ "memory_id": memory_id }),
)
.await;
assert_eq!(denied["error"]["code"], json!(-32030), "{denied:#?}");
assert_eq!(
denied["error"]["data"]["action"],
json!(action),
"{denied:#?}"
);
}
let quarantine = rpc(&app, "mobkit/memory/panel/quarantine", json!({})).await;
assert_eq!(quarantine["error"]["code"], json!(-32030));
let dreams = rpc(&app, "mobkit/memory/panel/dreams", json!({})).await;
assert_eq!(dreams["error"]["code"], json!(-32030));
let recall_router = rpc(
&app,
"mobkit/agent_memory/recall",
json!({ "identity": "router", "selection": "always" }),
)
.await;
assert_ne!(
recall_router["error"]["code"],
json!(-32030),
"router recall passes the access gate: {recall_router:#?}"
);
let recall_delivery = rpc(
&app,
"mobkit/agent_memory/recall",
json!({ "identity": "delivery", "selection": "always" }),
)
.await;
assert_eq!(recall_delivery["error"]["code"], json!(-32030));
let capabilities = rpc(&app, "mobkit/capabilities", json!({})).await;
let methods = capabilities["result"]["methods"]
.as_array()
.expect("methods")
.clone();
assert!(
methods
.iter()
.any(|value| value == "mobkit/memory/panel/records")
);
for hidden in [
"mobkit/memory/panel/dreams",
"mobkit/memory/panel/quarantine",
] {
assert!(
methods.iter().all(|value| value != hidden),
"{hidden} requires a grant this caller lacks: {methods:#?}"
);
}
controller
.upsert_rule(AccessRule {
id: "reviewer".to_string(),
actions: vec![
"agent.memory.read".to_string(),
"mob.memory.read".to_string(),
"memory.quarantine.review".to_string(),
],
..AccessRule::default()
})
.expect("live reviewer grant");
let dreams = rpc(&app, "mobkit/memory/panel/dreams", json!({})).await;
assert_eq!(dreams["error"], Value::Null, "{dreams:#?}");
let quarantine = rpc(&app, "mobkit/memory/panel/quarantine", json!({})).await;
assert_eq!(quarantine["error"], Value::Null, "{quarantine:#?}");
let queue = quarantine["result"]["records"].as_array().expect("queue");
assert!(
queue.iter().any(|row| row["id"] == json!(quarantined_id)),
"{queue:#?}"
);
let promotions = quarantine["result"]["pending_promotions"]
.as_array()
.expect("promotions");
assert!(
promotions
.iter()
.any(|row| row["pending_id"] == json!("gate-mob-promotion")),
"{promotions:#?}"
);
assert!(
promotions
.iter()
.all(|row| row["pending_id"] != json!("gate-delivery-promotion")),
"identity promotions must not ride the unscoped read grant: {promotions:#?}"
);
let records = rpc(&app, "mobkit/memory/panel/records", json!({})).await;
let rows = records["result"]["records"].as_array().expect("records");
assert!(
rows.iter().all(|row| row["id"] != json!(operator_id)),
"operator rows must not ride the unscoped read grant: {rows:#?}"
);
let denied_operator = rpc(
&app,
"mobkit/memory/panel/record",
json!({ "memory_id": operator_id }),
)
.await;
assert_eq!(denied_operator["error"]["code"], json!(-32030));
let denied_scope = rpc(
&app,
"mobkit/memory/panel/records",
json!({ "scope": "operator" }),
)
.await;
assert_eq!(denied_scope["error"]["code"], json!(-32030));
controller
.upsert_rule(AccessRule {
id: "operator-reader".to_string(),
actions: vec!["operator.memory.read".to_string()],
..AccessRule::default()
})
.expect("live operator grant");
let operator_detail = rpc(
&app,
"mobkit/memory/panel/record",
json!({ "memory_id": operator_id }),
)
.await;
assert_eq!(
operator_detail["error"],
Value::Null,
"{operator_detail:#?}"
);
let _ = runtime.mob_handle().stop().await;
}
#[tokio::test]
async fn memory_panel_quarantine_queue_filters_rows_per_scope() {
let (_temp_dir, mut runtime) = build_access_runtime_fixture().await;
let memory_dir = tempfile::tempdir().expect("memory dir");
let (store, _tip_id, quarantined_id, _delivery_id, _mob_id, _operator_id) =
seeded_memory_store(memory_dir.path()).await;
runtime.set_memory_panel_store(Arc::new(store.clone()));
let controller = AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: Vec::new(),
..AccessControlConfig::default()
})
.expect("controller");
runtime.set_access_controller(controller.clone());
let app = runtime.build_reference_app_router(decision_state(false));
let denied = rpc(&app, "mobkit/memory/panel/quarantine", json!({})).await;
assert_eq!(denied["error"]["code"], json!(-32030), "{denied:#?}");
controller
.upsert_rule(AccessRule {
id: "reviewer".to_string(),
actions: vec!["memory.quarantine.review".to_string()],
..AccessRule::default()
})
.expect("reviewer grant");
let quarantine = rpc(&app, "mobkit/memory/panel/quarantine", json!({})).await;
assert_eq!(quarantine["error"], Value::Null, "{quarantine:#?}");
let queue = quarantine["result"]["records"].as_array().expect("queue");
assert!(
queue.is_empty(),
"reviewer-only must see no queue records: {queue:#?}"
);
let promotions = quarantine["result"]["pending_promotions"]
.as_array()
.expect("promotions");
assert!(
promotions.is_empty(),
"reviewer-only must see no promotions: {promotions:#?}"
);
controller
.upsert_rule(AccessRule {
id: "router-reader".to_string(),
actions: vec!["agent.memory.read".to_string(), "agent.view".to_string()],
agents: vec!["router".to_string()],
..AccessRule::default()
})
.expect("router grant");
let quarantine = rpc(&app, "mobkit/memory/panel/quarantine", json!({})).await;
let queue = quarantine["result"]["records"].as_array().expect("queue");
assert_eq!(queue.len(), 1, "{queue:#?}");
assert_eq!(queue[0]["id"], json!(quarantined_id));
let promotions = quarantine["result"]["pending_promotions"]
.as_array()
.expect("promotions");
assert!(
promotions.is_empty(),
"router grants must not expose mob/delivery promotions: {promotions:#?}"
);
controller
.upsert_rule(AccessRule {
id: "mob-reader".to_string(),
actions: vec!["mob.memory.read".to_string()],
..AccessRule::default()
})
.expect("mob grant");
let quarantine = rpc(&app, "mobkit/memory/panel/quarantine", json!({})).await;
let promotions = quarantine["result"]["pending_promotions"]
.as_array()
.expect("promotions");
assert_eq!(promotions.len(), 1, "{promotions:#?}");
assert_eq!(promotions[0]["pending_id"], json!("gate-mob-promotion"));
controller
.upsert_rule(AccessRule {
id: "delivery-reader".to_string(),
actions: vec!["agent.memory.read".to_string(), "agent.view".to_string()],
agents: vec!["delivery".to_string()],
..AccessRule::default()
})
.expect("delivery grant");
let quarantine = rpc(&app, "mobkit/memory/panel/quarantine", json!({})).await;
let promotions = quarantine["result"]["pending_promotions"]
.as_array()
.expect("promotions");
assert_eq!(promotions.len(), 2, "{promotions:#?}");
assert!(
promotions
.iter()
.any(|row| row["pending_id"] == json!("gate-delivery-promotion")),
"{promotions:#?}"
);
let _ = runtime.mob_handle().stop().await;
}
#[tokio::test]
async fn recall_read_action_migration_compat_rule_both_ways() {
let (_temp_dir, mut runtime) = build_access_runtime_fixture().await;
let naive = AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: vec![AccessRule {
id: "view-router".to_string(),
actions: vec!["agent.view".to_string()],
agents: vec!["router".to_string()],
..AccessRule::default()
}],
..AccessControlConfig::default()
})
.expect("naive controller");
let (config, _) = naive.snapshot();
assert!(
config.rules[0]
.actions
.contains(&"agent.memory.read".to_string()),
"compat rewrite materializes the read grant: {config:#?}"
);
runtime.set_access_controller(naive);
let app = runtime.build_reference_app_router(decision_state(false));
let recall_router = rpc(
&app,
"mobkit/agent_memory/recall",
json!({ "identity": "router", "selection": "always" }),
)
.await;
assert_ne!(
recall_router["error"]["code"],
json!(-32030),
"naive config keeps recall working on view grants: {recall_router:#?}"
);
let recall_delivery = rpc(
&app,
"mobkit/agent_memory/recall",
json!({ "identity": "delivery", "selection": "always" }),
)
.await;
assert_eq!(recall_delivery["error"]["code"], json!(-32030));
let explicit = AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@example.test".to_string()],
rules: vec![
AccessRule {
id: "view-router".to_string(),
actions: vec!["agent.view".to_string()],
agents: vec!["router".to_string()],
..AccessRule::default()
},
AccessRule {
id: "unrelated-memory-rule".to_string(),
actions: vec!["agent.memory.read".to_string()],
agents: vec!["someone-else".to_string()],
..AccessRule::default()
},
],
..AccessControlConfig::default()
})
.expect("explicit controller");
runtime.set_access_controller(explicit);
let app = runtime.build_reference_app_router(decision_state(false));
let denied = rpc(
&app,
"mobkit/agent_memory/recall",
json!({ "identity": "router", "selection": "always" }),
)
.await;
assert_eq!(denied["error"]["code"], json!(-32030), "{denied:#?}");
assert_eq!(
denied["error"]["data"]["action"],
json!("agent.memory.read"),
"explicit configs are taken literally: {denied:#?}"
);
let _ = runtime.mob_handle().stop().await;
}