#![allow(
clippy::expect_used,
clippy::unwrap_used,
clippy::panic,
clippy::uninlined_format_args,
clippy::too_many_lines
)]
use std::sync::Arc;
use std::time::Duration;
use meerkat_client::TestClient;
use meerkat_mob::ids::AgentIdentity as MeerkatId;
use meerkat_mob::{MobDefinition, ProfileName, SpawnMemberSpec};
use meerkat_mobkit::memory::events::MemoryTimelineEvent;
use meerkat_mobkit::memory::records::RecordStatus;
use meerkat_mobkit::memory::records::{
InjectionLogEntry, InjectionSurface, MemoryAuthor, MemoryKind,
};
use meerkat_mobkit::memory::staged::StagedBatchKind;
use meerkat_mobkit::{
AccessControlConfig, AccessController, AccessEffect, AccessRule, AgentMemoryProvider,
AuthPolicy, BigQueryNaming, ConsolePolicy, DiscoverySpec, MemoryPanelStore, MemoryScope,
MobKitConfig, NewMemoryRecord, PreSpawnData, RuntimeDecisionInputs, RuntimeOpsPolicy,
SqliteAgentMemoryStore, StagedMemoryStore, StagedMutationBatch, StagedOp, StewardStore,
TaintableStore, TrustTier, TrustedOidcRuntimeConfig, UnifiedRuntime,
build_runtime_decision_state,
};
use serde_json::json;
const REALM: &str = "default";
const MOB: &str = "memory-e2e-mob";
const MOB_TOML: &str = r#"
[mob]
id = "memory-e2e-mob"
[profiles.lead]
model = "gpt-5.5"
external_addressable = true
[profiles.lead.tools]
comms = true
"#;
struct QuarantineAgentWrites;
impl meerkat_mobkit::memory::taint::LlmWriteGate for QuarantineAgentWrites {
fn quarantine_reason(
&self,
author: &MemoryAuthor,
_kind: StagedBatchKind,
_evidence: &[meerkat_mobkit::memory::records::EvidenceRef],
) -> Option<String> {
matches!(author, MemoryAuthor::Agent { .. })
.then(|| "matches the 'credential-assignment' secret pattern class (§10.4)".to_string())
}
}
struct SeededIds {
chain_tip: String,
quarantined: String,
delivery_fact: String,
dream_run: String,
}
fn now_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let state_dir = match std::env::var("MOBKIT_MEMORY_E2E_STATE") {
Ok(path) => std::path::PathBuf::from(path),
Err(_) => {
let dir =
std::env::temp_dir().join(format!("mobkit-memory-e2e-{}", std::process::id()));
std::fs::create_dir_all(&dir)?;
dir
}
};
let store_path = state_dir.join("agent-memory");
std::fs::create_dir_all(&store_path)?;
let store = SqliteAgentMemoryStore::open(&store_path)?;
let ids = seed_store(&store).await;
verify_seeded_store(&store, &ids).await;
if std::env::var("MOBKIT_MEMORY_E2E_SEED_ONLY").ok().as_deref() == Some("1") {
println!("memory e2e fixture: seed-only mode, exiting before serve");
return Ok(());
}
let definition = MobDefinition::from_toml(MOB_TOML)
.map_err(|e| std::io::Error::other(format!("bad mob definition: {e}")))?;
let mut runtime = Box::pin(
UnifiedRuntime::builder()
.definition(definition)
.default_llm_client(Arc::new(TestClient::default()))
.module_config(MobKitConfig {
modules: vec![],
discovery: DiscoverySpec {
namespace: "memory-e2e".to_string(),
modules: vec![],
},
pre_spawn: Vec::<PreSpawnData>::new(),
})
.timeout(Duration::from_secs(5))
.build(),
)
.await?;
runtime
.reconcile(vec![
SpawnMemberSpec::new(ProfileName::from("lead"), MeerkatId::from("router")),
SpawnMemberSpec::new(ProfileName::from("lead"), MeerkatId::from("delivery")),
])
.await?;
let access_mode =
std::env::var("MOBKIT_MEMORY_E2E_ACCESS").unwrap_or_else(|_| "open".to_string());
if let Some(controller) = access_controller(&access_mode) {
runtime.set_access_controller(controller);
}
println!("memory e2e fixture: access mode '{access_mode}'");
runtime.set_memory_panel_store(Arc::new(store.clone()));
emit_timeline_events(&runtime, &ids);
let decisions = build_runtime_decision_state(RuntimeDecisionInputs {
bigquery: BigQueryNaming {
dataset: "memory_e2e_dataset".to_string(),
table: "memory_e2e_table".to_string(),
},
trusted_mobkit_toml: r#"
[[modules]]
id = "router"
command = "router-bin"
args = []
restart_policy = "always"
"#
.to_string(),
auth: AuthPolicy::default(),
trusted_oidc: trusted_oidc(),
console: ConsolePolicy {
require_app_auth: false,
..ConsolePolicy::default()
},
ops: RuntimeOpsPolicy::default(),
release_metadata_json: include_str!("../assets/release-targets.json").to_string(),
})
.map_err(|err| std::io::Error::other(format!("failed to build console decisions: {err:?}")))?;
let listen_addr =
std::env::var("MOBKIT_MEMORY_E2E_ADDR").unwrap_or_else(|_| "127.0.0.1:3230".to_string());
let listener = tokio::net::TcpListener::bind(&listen_addr).await?;
println!("memory e2e fixture listening on http://{listen_addr}");
let run_report = runtime
.run(listener, decisions, async {
let _ = tokio::signal::ctrl_c().await;
})
.await;
run_report
.shutdown
.mob_stop
.map_err(|err| std::io::Error::other(format!("failed to stop mob runtime: {err}")))?;
run_report.serve_result?;
Ok(())
}
async fn seed_store(store: &SqliteAgentMemoryStore) -> SeededIds {
let record = |title: &str, body: &str| NewMemoryRecord {
kind: MemoryKind::Fact,
title: title.to_string(),
description: format!("{title} — seeded for the memory console e2e"),
body: body.to_string(),
tags: Vec::new(),
evidence: Vec::new(),
verification: None,
};
let identity_scope = |identity: &str| MemoryScope::Identity {
realm: REALM.to_string(),
identity: identity.to_string(),
};
let mob_scope = MemoryScope::Mob {
realm: REALM.to_string(),
mob: MOB.to_string(),
};
let filler_ops: Vec<StagedOp> = (1..=60u64)
.map(|index| StagedOp::Create {
id: None,
scope: identity_scope("router"),
record: record(
&format!("Router ops note {index:02}"),
&format!("Routine operational note number {index} for paging."),
),
trust: TrustTier::AgentObserved,
derived_from: Vec::new(),
rationale: Some("seed filler".to_string()),
created_at_ms: Some(1_000 + index),
updated_at_ms: Some(1_000 + index),
})
.collect();
let filler_token = store
.stage(StagedMutationBatch {
kind: StagedBatchKind::FreshWrite,
realm: REALM.to_string(),
author: MemoryAuthor::Application,
ops: filler_ops,
})
.await
.expect("stage filler batch");
store
.commit(filler_token)
.await
.expect("commit filler batch");
let chain_root = store
.remember_authored(
&identity_scope("router"),
record(
"Router deploy cadence",
"Deploys happen ad hoc whenever a fix lands.",
),
MemoryAuthor::Operator,
)
.await
.expect("chain root");
let chain_mid = store
.supersede_authored(
&identity_scope("router"),
&chain_root.memory_id,
record(
"Router deploy cadence",
"Deploys moved to a weekly Tuesday window.",
),
MemoryAuthor::Operator,
)
.await
.expect("chain mid");
let mut tip_record = record(
"Router deploy cadence",
"Deploys are gated on the release train: Tuesdays, after CI green.",
);
tip_record.evidence = vec![meerkat_mobkit::memory::records::EvidenceRef {
session_id: "sess-router-archived-1".to_string(),
generation: 2,
revision: None,
range: Some((3, 9)),
}];
let chain_tip = store
.supersede_authored(
&identity_scope("router"),
&chain_mid.memory_id,
tip_record,
MemoryAuthor::Operator,
)
.await
.expect("chain tip");
for (title, body) in [
(
"Router escalation contact",
"Escalations for routing incidents go to the on-call channel first.",
),
(
"Router retry budget",
"Retries are capped at 2 with 25ms backoff before dead-lettering.",
),
] {
store
.remember_authored(
&identity_scope("router"),
record(title, body),
MemoryAuthor::Operator,
)
.await
.expect("router fact");
}
let delivery_fact = store
.remember_authored(
&identity_scope("delivery"),
record(
"Delivery sink preference",
"The delivery sink prefers batched flushes over per-message writes.",
),
MemoryAuthor::Operator,
)
.await
.expect("delivery fact");
store
.remember_authored(
&identity_scope("delivery"),
record(
"Delivery quiet hours",
"No non-urgent delivery notifications between 22:00 and 07:00.",
),
MemoryAuthor::Operator,
)
.await
.expect("delivery fact 2");
for (title, body) in [
(
"Mob review convention",
"Cross-agent changes need a second reviewer from another domain.",
),
(
"Mob incident channel",
"Incidents coordinate in the shared mob channel, not DMs.",
),
] {
store
.remember_authored(&mob_scope, record(title, body), MemoryAuthor::Operator)
.await
.expect("mob fact");
}
store
.remember_authored(
&MemoryScope::Realm {
realm: REALM.to_string(),
},
record(
"Realm data residency",
"All realm data stays in the EU region unless a ticket says otherwise.",
),
MemoryAuthor::Operator,
)
.await
.expect("realm fact");
store
.remember_authored(
&MemoryScope::Operator {
realm: REALM.to_string(),
operator: "op-luka".to_string(),
},
record(
"Operator briefing preference",
"Prefers terse morning summaries with links over long prose.",
),
MemoryAuthor::Operator,
)
.await
.expect("operator fact");
store.set_llm_write_gate(Arc::new(QuarantineAgentWrites));
let quarantined = store
.remember_authored(
&identity_scope("router"),
record(
"Router upstream credential",
"The upstream accepts a static token configured at deploy time.",
),
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(
REALM,
&[
InjectionLogEntry {
record_id: chain_tip.memory_id.clone(),
identity: "router".to_string(),
session_key: Some("sess-router-1".to_string()),
surface: InjectionSurface::Build,
at_ms: 1_000,
},
InjectionLogEntry {
record_id: delivery_fact.memory_id.clone(),
identity: "delivery".to_string(),
session_key: Some("sess-delivery-1".to_string()),
surface: InjectionSurface::Turn,
at_ms: 2_000,
},
],
)
.await
.expect("injection rows");
let dream_run = "run-dream-e2e-1".to_string();
let token = store
.stage(StagedMutationBatch {
kind: StagedBatchKind::FreshWrite,
realm: REALM.to_string(),
author: MemoryAuthor::Steward {
run_id: dream_run.clone(),
},
ops: vec![
StagedOp::Create {
id: None,
scope: mob_scope.clone(),
record: record(
"Consolidated deploy learnings",
"Three deploy retros distilled: gate on CI, announce in channel.",
),
trust: TrustTier::AgentObserved,
derived_from: Vec::new(),
rationale: Some("consolidated during dream".to_string()),
created_at_ms: None,
updated_at_ms: None,
},
StagedOp::Create {
id: None,
scope: mob_scope.clone(),
record: record(
"Consolidated escalation map",
"Escalation paths across router/delivery merged into one map.",
),
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, "Pending gated promotion payload.");
async move {
store
.stage(StagedMutationBatch {
kind: StagedBatchKind::FreshWrite,
realm: REALM.to_string(),
author: MemoryAuthor::Steward {
run_id: "run-gate-e2e-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(mob_scope.clone(), "Promoted mob claim").await;
store
.record_pending_promotion(
REALM,
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: MOB.to_string(),
rationale: Some("steward: mob-wide convention".to_string()),
status: "pending".to_string(),
created_at_ms: now_ms(),
},
)
.await
.expect("mob promotion row");
let delivery_stage =
stage_promotion(identity_scope("delivery"), "Promoted delivery claim").await;
store
.record_pending_promotion(
REALM,
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: now_ms(),
},
)
.await
.expect("delivery promotion row");
SeededIds {
chain_tip: chain_tip.memory_id,
quarantined: quarantined.memory_id,
delivery_fact: delivery_fact.memory_id,
dream_run,
}
}
async fn verify_seeded_store(store: &SqliteAgentMemoryStore, ids: &SeededIds) {
let page = store
.records_page(REALM, None, None, None, 500, None)
.await
.expect("records page");
let active = page
.records
.iter()
.filter(|record| matches!(record.status, RecordStatus::Active))
.count();
let superseded = page
.records
.iter()
.filter(|record| matches!(record.status, RecordStatus::Superseded { .. }))
.count();
let quarantined_count = page
.records
.iter()
.filter(|record| matches!(record.status, RecordStatus::Quarantined { .. }))
.count();
assert_eq!(page.records.len(), 74, "total seeded records");
assert_eq!(active, 71, "active records");
assert_eq!(superseded, 2, "superseded chain records");
assert_eq!(quarantined_count, 1, "quarantined record");
let quarantined = store
.quarantined_records(REALM, 10)
.await
.expect("quarantined records");
assert_eq!(quarantined.len(), 1);
assert_eq!(quarantined[0].id, ids.quarantined);
let promotions = store
.pending_promotions(REALM)
.await
.expect("pending promotions");
assert_eq!(promotions.len(), 2, "pending promotions");
let dreams = store.dream_history(REALM, 10).await.expect("dream history");
assert!(
dreams.iter().any(|run| run.run_id == ids.dream_run),
"dream audit rows missing run {}: {dreams:?}",
ids.dream_run
);
assert_eq!(dreams.len(), 1, "exactly one dream run: {dreams:?}");
let tip_injections = store
.injection_log_for_record(REALM, &ids.chain_tip, 10)
.await
.expect("tip injections");
assert_eq!(tip_injections.len(), 1, "chain-tip injection row");
let delivery_injections = store
.injection_log_for_record(REALM, &ids.delivery_fact, 10)
.await
.expect("delivery injections");
assert_eq!(delivery_injections.len(), 1, "delivery injection row");
println!(
"memory e2e fixture seeded: {} records ({active} active, {superseded} superseded, \
{quarantined_count} quarantined), {} pending promotions, dream run '{}', 2 injection rows",
page.records.len(),
promotions.len(),
ids.dream_run,
);
}
fn emit_timeline_events(runtime: &UnifiedRuntime, ids: &SeededIds) {
let sink = runtime.memory_event_sink();
for event in [
MemoryTimelineEvent::DreamStarted {
realm: REALM.to_string(),
run_id: ids.dream_run.clone(),
},
MemoryTimelineEvent::DreamCompleted {
realm: REALM.to_string(),
run_id: ids.dream_run.clone(),
ops_committed: 2,
detail: json!({ "phase": "consolidation", "verdicts": { "release": 0 } }),
},
MemoryTimelineEvent::QuarantinedWrite {
realm: REALM.to_string(),
author: "agent router".to_string(),
reason: "matches the 'credential-assignment' secret pattern class (§10.4)".to_string(),
},
MemoryTimelineEvent::QuarantineVerdict {
realm: REALM.to_string(),
record_id: ids.quarantined.clone(),
verdict: "unverifiable".to_string(),
rationale: Some("no evidence span reaches the claimed credential".to_string()),
},
MemoryTimelineEvent::QuarantineReleaseBlocked {
realm: REALM.to_string(),
record_id: ids.quarantined.clone(),
verdict: "release".to_string(),
class: "credential-assignment".to_string(),
},
MemoryTimelineEvent::PromotionPendingGate {
realm: REALM.to_string(),
pending_id: "gate-mob-promotion".to_string(),
record_id: ids.quarantined.clone(),
scope_kind: "mob".to_string(),
scope_key: MOB.to_string(),
},
MemoryTimelineEvent::ConflictSignal {
realm: REALM.to_string(),
entity: "router".to_string(),
topic: "deploy cadence".to_string(),
reason: "weekly window contradicts ad-hoc deploy claim".to_string(),
},
MemoryTimelineEvent::TaintTransition {
identity: Some("router".to_string()),
session_key: "sess-router-1".to_string(),
kind: "tainted".to_string(),
source: "web_search".to_string(),
},
] {
sink.emit(event);
}
}
fn access_controller(mode: &str) -> Option<AccessController> {
let everyone = |id: &str, actions: &[&str]| AccessRule {
id: id.to_string(),
actions: actions
.iter()
.map(std::string::ToString::to_string)
.collect(),
..AccessRule::default()
};
let mut rules = vec![
everyone("everyone-views-agents", &["agent.view"]),
everyone("everyone-observes-mob", &["mob.observe"]),
];
match mode {
"open" => return None,
"reader" => rules.push(everyone(
"everyone-reads-memory",
&[
"agent.memory.read",
"mob.memory.read",
"operator.memory.read",
],
)),
"partial" => rules.push(everyone(
"everyone-reads-agent-and-mob-memory",
&["agent.memory.read", "mob.memory.read"],
)),
"scoped" => rules.push(AccessRule {
id: "router-only-memory-read".to_string(),
actions: vec!["agent.memory.read".to_string()],
agents: vec!["router".to_string()],
..AccessRule::default()
}),
"none" => rules.push(AccessRule {
id: "no-memory-access".to_string(),
effect: AccessEffect::Deny,
actions: vec![
"agent.memory.read".to_string(),
"mob.memory.read".to_string(),
"operator.memory.read".to_string(),
"memory.quarantine.review".to_string(),
],
..AccessRule::default()
}),
other => panic!("unknown MOBKIT_MEMORY_E2E_ACCESS mode: {other}"),
}
Some(
AccessController::new(AccessControlConfig {
enabled: true,
admins: vec!["root@memory-e2e.test".to_string()],
rules,
..AccessControlConfig::default()
})
.expect("access controller"),
)
}
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(),
}
}