use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use std::time::{SystemTime, UNIX_EPOCH};
use meerkat_core::event::AgentEvent;
use serde::{Deserialize, Serialize};
use crate::identity_first::agent_memory::AgentMemoryLlmWrites;
use crate::memory::records::{EvidenceRef, MemoryAuthor};
use crate::memory::staged::StagedBatchKind;
const ALWAYS_UNTRUSTED_TOOL_NAMES: &[&str] = &["web_search", "web_fetch", "fetch", "http_request"];
const MCP_QUALIFIED_PREFIX: &str = "mcp__";
const MAX_TRACKED_TAINTED_SESSIONS: usize = 4096;
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ContentTrustConfig {
#[serde(default)]
pub trusted_mcp_servers: Vec<String>,
#[serde(default)]
pub untrusted_tools: Vec<String>,
#[serde(default)]
pub trusted_tools: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ToolContentTrust {
Trusted,
Untrusted { source: String },
}
impl ContentTrustConfig {
pub fn from_json_value(value: &serde_json::Value) -> Result<Self, String> {
let object = value
.as_object()
.ok_or_else(|| "content_trust must be an object".to_string())?;
let supported = ["trusted_mcp_servers", "untrusted_tools", "trusted_tools"];
let unsupported = object
.keys()
.filter(|key| !supported.contains(&key.as_str()))
.map(String::as_str)
.collect::<Vec<_>>();
if !unsupported.is_empty() {
return Err(format!(
"unsupported content_trust fields: {}",
unsupported.join(", ")
));
}
let parse_names = |key: &str| -> Result<Vec<String>, String> {
match object.get(key) {
None => Ok(Vec::new()),
Some(value) => {
let entries = value
.as_array()
.ok_or_else(|| format!("content_trust.{key} must be an array"))?;
entries
.iter()
.map(|entry| {
entry
.as_str()
.map(str::trim)
.filter(|name| !name.is_empty())
.map(ToString::to_string)
.ok_or_else(|| {
format!("content_trust.{key} entries must be non-empty strings")
})
})
.collect()
}
}
};
Ok(Self {
trusted_mcp_servers: parse_names("trusted_mcp_servers")?,
untrusted_tools: parse_names("untrusted_tools")?,
trusted_tools: parse_names("trusted_tools")?,
})
}
pub fn classify_tool(&self, name: &str) -> ToolContentTrust {
if ALWAYS_UNTRUSTED_TOOL_NAMES.contains(&name) {
return ToolContentTrust::Untrusted {
source: format!("web tool '{name}'"),
};
}
if self.untrusted_tools.iter().any(|tool| tool == name) {
return ToolContentTrust::Untrusted {
source: format!("configured untrusted tool '{name}'"),
};
}
if self.trusted_tools.iter().any(|tool| tool == name) {
return ToolContentTrust::Trusted;
}
if let Some(rest) = name.strip_prefix(MCP_QUALIFIED_PREFIX) {
let server = rest.split("__").next().unwrap_or(rest);
if self
.trusted_mcp_servers
.iter()
.any(|trusted| trusted == server)
{
return ToolContentTrust::Trusted;
}
return ToolContentTrust::Untrusted {
source: format!("MCP server '{server}' (tool '{name}')"),
};
}
ToolContentTrust::Trusted
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TaintState {
pub tainted_at_ms: u64,
pub source: String,
}
#[derive(Default)]
struct TaintInner {
current_session: HashMap<String, String>,
tainted: HashMap<String, TaintState>,
pending_identity_taint: HashMap<String, TaintState>,
reset_boundaries: HashMap<String, u64>,
}
pub type OutboundTaintDeclarer =
Arc<dyn Fn(&str, Option<meerkat_core::comms::SenderContentTaint>) + Send + Sync>;
#[derive(Clone, Default)]
pub struct SessionTaintTracker {
config: Arc<ContentTrustConfig>,
inner: Arc<Mutex<TaintInner>>,
event_sink: Arc<Mutex<Option<Arc<dyn crate::memory::events::MemoryEventSink>>>>,
outbound_declarer: Arc<Mutex<Option<OutboundTaintDeclarer>>>,
}
impl SessionTaintTracker {
pub fn new(config: ContentTrustConfig) -> Self {
Self {
config: Arc::new(config),
inner: Arc::new(Mutex::new(TaintInner::default())),
event_sink: Arc::new(Mutex::new(None)),
outbound_declarer: Arc::new(Mutex::new(None)),
}
}
pub fn set_outbound_taint_declarer(&self, declarer: OutboundTaintDeclarer) {
*self
.outbound_declarer
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(declarer);
}
fn declare_outbound(
&self,
identity: &str,
taint: Option<meerkat_core::comms::SenderContentTaint>,
) {
if let Some(declarer) = self
.outbound_declarer
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.as_ref()
{
declarer(identity, taint);
}
}
pub fn set_event_sink(&self, sink: Arc<dyn crate::memory::events::MemoryEventSink>) {
*self
.event_sink
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(sink);
}
fn emit_event(&self, event: crate::memory::events::MemoryTimelineEvent) {
if let Some(sink) = self
.event_sink
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.as_ref()
{
sink.emit(event);
}
}
pub fn observe_agent_event(&self, identity: &str, event: &AgentEvent) {
match event {
AgentEvent::RunStarted { session_id, input } => {
let session_key = session_id.to_string();
self.note_current_session(identity, &session_key);
if let Some(text) = input.prompt_text() {
self.observe_inbound_peer_content(identity, &session_key, &text);
}
}
AgentEvent::PeerContentIngested {
peer, sender_taint, ..
} => {
if *sender_taint == Some(meerkat_core::comms::SenderContentTaint::Tainted) {
let sender = peer
.as_ref()
.and_then(|peer| peer.display_name.clone())
.unwrap_or_else(|| "peer".to_string());
self.mark_identity_tainted(
identity,
format!("peer content declared tainted by sender '{sender}'"),
);
}
}
AgentEvent::ToolResultReceived { name, .. }
| AgentEvent::ToolExecutionCompleted { name, .. } => {
if let ToolContentTrust::Untrusted { source } = self.config.classify_tool(name) {
self.mark_identity_tainted(identity, source);
}
}
AgentEvent::ServerToolContent { kind, .. } => {
self.mark_identity_tainted(
identity,
format!("provider server tool '{}'", kind.provider_name()),
);
}
_ => {}
}
}
pub fn note_current_session(&self, identity: &str, session_key: &str) {
let mut inner = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let pending = inner.pending_identity_taint.remove(identity);
let previous = inner
.current_session
.insert(identity.to_string(), session_key.to_string());
let rotated_away_from_tainted = previous.as_deref() != Some(session_key)
&& previous.is_some_and(|prior| inner.tainted.contains_key(&prior));
if rotated_away_from_tainted {
tracing::warn!(
identity,
session_key,
"agent memory taint: session rotated away from a tainted session; \
new session starts clean"
);
self.emit_event(
crate::memory::events::MemoryTimelineEvent::TaintTransition {
identity: Some(identity.to_string()),
session_key: session_key.to_string(),
kind: "rotated_clean".to_string(),
source: "session rotation away from tainted session".to_string(),
},
);
}
let pending_reapplied = pending.is_some();
if let Some(state) = pending {
self.insert_taint(&mut inner, session_key.to_string(), state);
}
if rotated_away_from_tainted && !pending_reapplied {
drop(inner);
self.declare_outbound(identity, None);
}
}
pub fn clear_identity(&self, identity: &str) {
let mut inner = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
inner.pending_identity_taint.remove(identity);
inner.current_session.remove(identity);
drop(inner);
self.declare_outbound(identity, None);
}
pub fn mark_reset_boundary(&self, session_key: &str) {
let mut inner = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if inner
.reset_boundaries
.insert(session_key.to_string(), now_ms())
.is_none()
{
tracing::warn!(
session_key,
"agent memory taint: reset boundary marked; distillates over this \
session will land quarantined pending steward review (§8.4)"
);
self.emit_event(
crate::memory::events::MemoryTimelineEvent::TaintTransition {
identity: None,
session_key: session_key.to_string(),
kind: "reset_boundary".to_string(),
source: "reset() boundary (§8.4)".to_string(),
},
);
}
if inner.reset_boundaries.len() > MAX_TRACKED_TAINTED_SESSIONS
&& let Some(oldest) = inner
.reset_boundaries
.iter()
.min_by_key(|(_, at_ms)| **at_ms)
.map(|(key, _)| key.clone())
{
inner.reset_boundaries.remove(&oldest);
}
}
pub fn evidence_quarantine_reason(&self, session_key: &str) -> Option<String> {
let inner = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(state) = inner.tainted.get(session_key) {
return Some(format!(
"evidence session tainted by {} (session-tainted ⇒ range-tainted)",
state.source
));
}
if inner.reset_boundaries.contains_key(session_key) {
return Some(
"evidence session closed at a reset boundary; distillates quarantine \
pending steward review (§8.4)"
.to_string(),
);
}
None
}
pub fn session_taint(&self, session_key: &str) -> Option<TaintState> {
let inner = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
inner.tainted.get(session_key).cloned()
}
pub fn identity_taint(&self, identity: &str) -> Option<TaintState> {
let inner = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(state) = inner.pending_identity_taint.get(identity) {
return Some(state.clone());
}
inner
.current_session
.get(identity)
.and_then(|session| inner.tainted.get(session))
.cloned()
}
fn observe_inbound_peer_content(&self, identity: &str, session_key: &str, text: &str) {
for line in text.lines() {
let Some(sender_identity) = peer_projection_sender_identity(line) else {
continue;
};
if sender_identity == identity {
continue;
}
let source = {
let inner = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
inner
.pending_identity_taint
.get(sender_identity)
.or_else(|| {
inner
.current_session
.get(sender_identity)
.and_then(|session| inner.tainted.get(session))
})
.map(|state| state.source.clone())
};
if let Some(source) = source {
let state = TaintState {
tainted_at_ms: now_ms(),
source: format!(
"peer message from tainted sender '{sender_identity}' \
(sender session tainted by {source})"
),
};
let mut inner = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
self.insert_taint(&mut inner, session_key.to_string(), state);
}
}
}
fn mark_identity_tainted(&self, identity: &str, source: String) {
let state = TaintState {
tainted_at_ms: now_ms(),
source,
};
let mut inner = self
.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
match inner.current_session.get(identity).cloned() {
Some(session) => self.insert_taint(&mut inner, session, state),
None => {
match inner.pending_identity_taint.entry(identity.to_string()) {
std::collections::hash_map::Entry::Vacant(slot) => {
tracing::warn!(
identity,
source = %state.source,
"agent memory taint: untrusted ingestion observed before session \
attribution; holding identity-sticky taint"
);
slot.insert(state);
}
std::collections::hash_map::Entry::Occupied(mut slot) => {
slot.insert(state);
}
}
}
}
drop(inner);
self.declare_outbound(
identity,
Some(meerkat_core::comms::SenderContentTaint::Tainted),
);
}
fn insert_taint(&self, inner: &mut TaintInner, session: String, state: TaintState) {
match inner.tainted.entry(session) {
std::collections::hash_map::Entry::Occupied(_) => return,
std::collections::hash_map::Entry::Vacant(slot) => {
tracing::warn!(
session_key = %slot.key(),
source = %state.source,
"agent memory taint: session ingested untrusted content; LLM-authored \
memory writes from this session will quarantine until a fresh-context \
boundary (reset/respawn/fresh spawn)"
);
self.emit_event(
crate::memory::events::MemoryTimelineEvent::TaintTransition {
identity: None,
session_key: slot.key().clone(),
kind: "tainted".to_string(),
source: state.source.clone(),
},
);
slot.insert(state);
}
}
if inner.tainted.len() > MAX_TRACKED_TAINTED_SESSIONS
&& let Some(oldest) = inner
.tainted
.iter()
.min_by_key(|(_, state)| state.tainted_at_ms)
.map(|(key, _)| key.clone())
{
inner.tainted.remove(&oldest);
}
}
}
pub trait LlmWriteGate: Send + Sync {
fn quarantine_reason(
&self,
author: &MemoryAuthor,
kind: StagedBatchKind,
evidence: &[EvidenceRef],
) -> Option<String>;
}
pub struct TaintLlmWriteGate {
tracker: Option<SessionTaintTracker>,
llm_writes: AgentMemoryLlmWrites,
}
impl TaintLlmWriteGate {
pub fn new(tracker: Option<SessionTaintTracker>, llm_writes: AgentMemoryLlmWrites) -> Self {
Self {
tracker,
llm_writes,
}
}
}
impl LlmWriteGate for TaintLlmWriteGate {
fn quarantine_reason(
&self,
author: &MemoryAuthor,
kind: StagedBatchKind,
evidence: &[EvidenceRef],
) -> Option<String> {
if !author.is_llm() {
return None;
}
if self.llm_writes == AgentMemoryLlmWrites::Quarantined
&& kind != StagedBatchKind::ReviewVerdict
{
return Some("llm_writes=quarantined policy".to_string());
}
let tracker = self.tracker.as_ref()?;
if let MemoryAuthor::Agent { identity } = author
&& let Some(state) = tracker.identity_taint(identity)
{
return Some(format!("session tainted by {}", state.source));
}
for evidence_ref in evidence {
if let Some(reason) = tracker.evidence_quarantine_reason(&evidence_ref.session_id) {
return Some(reason);
}
}
None
}
}
fn peer_projection_sender_identity(line: &str) -> Option<&str> {
let name = if let Some(rest) = line.strip_prefix("Peer message from ") {
match rest.split_once(": ") {
Some((name, _)) => name.trim(),
None => rest.strip_suffix(':').unwrap_or(rest).trim(),
}
} else {
let rest = line.strip_prefix("Peer response from ")?;
rest.split(" (to request:").next()?.trim()
};
if name.is_empty() {
return None;
}
Some(name.rsplit('/').next().unwrap_or(name))
}
pub trait MemberAgentEventSink: Send + Sync {
fn observe(&self, identity: &str, envelope: &meerkat_core::event::EventEnvelope<AgentEvent>);
}
impl MemberAgentEventSink for SessionTaintTracker {
fn observe(&self, identity: &str, envelope: &meerkat_core::event::EventEnvelope<AgentEvent>) {
self.observe_agent_event(identity, &envelope.payload);
}
}
pub struct CompactionResetSink {
on_compacted: Arc<dyn Fn(&str) + Send + Sync>,
}
impl CompactionResetSink {
pub fn new(on_compacted: Arc<dyn Fn(&str) + Send + Sync>) -> Self {
Self { on_compacted }
}
}
impl MemberAgentEventSink for CompactionResetSink {
fn observe(&self, _identity: &str, envelope: &meerkat_core::event::EventEnvelope<AgentEvent>) {
if matches!(envelope.payload, AgentEvent::CompactionCompleted { .. })
&& let meerkat_core::event::EventSourceIdentity::Session { session_id } =
&envelope.source
{
(self.on_compacted)(&session_id.to_string());
}
}
}
#[derive(Clone)]
pub struct TaintObserverGuard {
_abort: Arc<AbortOnDrop>,
}
struct AbortOnDrop(tokio::task::JoinHandle<()>);
impl Drop for AbortOnDrop {
fn drop(&mut self) {
self.0.abort();
}
}
pub fn spawn_taint_observer(
handle: meerkat_mob::MobHandle,
tracker: SessionTaintTracker,
) -> TaintObserverGuard {
spawn_member_event_observer(handle, vec![Arc::new(tracker)])
}
pub fn spawn_member_event_observer(
handle: meerkat_mob::MobHandle,
sinks: Vec<Arc<dyn MemberAgentEventSink>>,
) -> TaintObserverGuard {
let task = tokio::spawn(run_member_event_observer(handle, sinks));
TaintObserverGuard {
_abort: Arc::new(AbortOnDrop(task)),
}
}
async fn run_member_event_observer(
handle: meerkat_mob::MobHandle,
sinks: Vec<Arc<dyn MemberAgentEventSink>>,
) {
use futures::StreamExt;
use futures::stream::SelectAll;
enum Observed {
Event(String, Box<meerkat_core::event::EventEnvelope<AgentEvent>>),
Closed(String),
}
let mut streams: SelectAll<futures::stream::BoxStream<'static, Observed>> = SelectAll::new();
let mut subscribed: std::collections::HashSet<String> = std::collections::HashSet::new();
let mut warned: std::collections::HashSet<String> = std::collections::HashSet::new();
let mut reconcile = tokio::time::interval(std::time::Duration::from_secs(1));
reconcile.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
Some(observed) = streams.next() => match observed {
Observed::Event(identity, envelope) => {
for sink in &sinks {
sink.observe(&identity, &envelope);
}
}
Observed::Closed(identity) => {
subscribed.remove(&identity);
}
},
_ = reconcile.tick() => {
for entry in handle.list_members_including_retiring().await {
if entry.status != meerkat_mob::MobMemberStatus::Active {
continue;
}
let identity = entry.agent_identity.to_string();
if subscribed.contains(&identity) {
continue;
}
match handle.subscribe_agent_events(&entry.agent_identity).await {
Ok(stream) => {
warned.remove(&identity);
subscribed.insert(identity.clone());
let close_key = identity.clone();
streams.push(
stream
.map(move |envelope| {
Observed::Event(identity.clone(), Box::new(envelope))
})
.chain(futures::stream::once(async move {
Observed::Closed(close_key)
}))
.boxed(),
);
}
Err(error) => {
if warned.insert(identity.clone()) {
tracing::warn!(
identity = %identity,
error = %error,
"agent memory taint observer: failed to subscribe; will retry"
);
} else {
tracing::debug!(
identity = %identity,
error = %error,
"agent memory taint observer: subscribe still failing"
);
}
}
}
}
}
}
}
}
fn now_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|duration| duration.as_millis() as u64)
.unwrap_or(0)
}
#[cfg(test)]
#[allow(
clippy::expect_used,
clippy::panic,
clippy::redundant_clone,
clippy::unwrap_used
)]
mod tests {
use super::*;
use meerkat_core::types::{ContentBlock, ServerToolKind, SessionId};
use serde_json::json;
fn run_started(session: &SessionId) -> AgentEvent {
AgentEvent::RunStarted {
session_id: session.clone(),
input: meerkat_core::types::RunInput::Content {
content: meerkat_core::ContentInput::Text("hi".to_string()),
},
}
}
fn tool_result(name: &str) -> AgentEvent {
AgentEvent::ToolResultReceived {
id: "tool-1".to_string(),
name: name.to_string(),
content: vec![ContentBlock::Text {
text: "ok".to_string(),
}],
is_error: false,
}
}
fn peer_ingested(taint: Option<meerkat_core::comms::SenderContentTaint>) -> AgentEvent {
AgentEvent::PeerContentIngested {
kind: meerkat_core::types::CommsNoticeKind::Message,
peer: None,
request_id: None,
sender_taint: taint,
}
}
#[test]
fn declared_peer_taint_taints_receiver_but_none_or_clean_does_not() {
use meerkat_core::comms::SenderContentTaint;
let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
tracker.note_current_session("identity:b", "sess-b");
tracker.observe_agent_event("identity:b", &peer_ingested(None));
assert!(
tracker.session_taint("sess-b").is_none(),
"no declaration must not taint"
);
tracker.observe_agent_event(
"identity:b",
&peer_ingested(Some(SenderContentTaint::Clean)),
);
assert!(
tracker.session_taint("sess-b").is_none(),
"an affirmative Clean declaration must not taint"
);
tracker.observe_agent_event(
"identity:b",
&peer_ingested(Some(SenderContentTaint::Tainted)),
);
assert!(
tracker.session_taint("sess-b").is_some(),
"a declared-tainted peer delivery taints the receiving session"
);
}
#[test]
fn taint_transitions_emit_timeline_events_when_sink_wired() {
let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
let sink = std::sync::Arc::new(crate::memory::events::CollectingEventSink::new());
tracker.set_event_sink(sink.clone());
tracker.note_current_session("identity:a", "sess-1");
tracker.observe_agent_event("identity:a", &tool_result("web_fetch"));
tracker.mark_reset_boundary("sess-1");
tracker.mark_reset_boundary("sess-1");
tracker.note_current_session("identity:a", "sess-2");
let types = sink.types();
assert_eq!(
types,
vec![
"memory.taint.transition", "memory.taint.transition", "memory.taint.transition", ]
);
let events = sink.events.lock().unwrap();
let kinds: Vec<String> = events
.iter()
.map(|event| match event {
crate::memory::events::MemoryTimelineEvent::TaintTransition { kind, .. } => {
kind.clone()
}
other => panic!("unexpected event {other:?}"),
})
.collect();
assert_eq!(kinds, vec!["tainted", "reset_boundary", "rotated_clean"]);
}
#[test]
fn content_trust_parse_rejects_unknown_fields_and_bad_types() {
let err = ContentTrustConfig::from_json_value(&json!({"servers": []}))
.expect_err("unknown field must fail loud");
assert!(err.contains("unsupported content_trust fields"), "{err}");
let err = ContentTrustConfig::from_json_value(&json!({"trusted_mcp_servers": "kg"}))
.expect_err("non-array must fail loud");
assert!(err.contains("must be an array"), "{err}");
let err = ContentTrustConfig::from_json_value(&json!({"untrusted_tools": [1]}))
.expect_err("non-string entry must fail loud");
assert!(err.contains("non-empty strings"), "{err}");
let err =
ContentTrustConfig::from_json_value(&json!([])).expect_err("non-object must fail loud");
assert!(err.contains("must be an object"), "{err}");
}
#[test]
fn content_trust_parse_accepts_full_block() {
let config = ContentTrustConfig::from_json_value(&json!({
"trusted_mcp_servers": ["knowledge_graph"],
"untrusted_tools": ["scrape_page"],
"trusted_tools": ["mcp__scanner__lint"],
}))
.expect("valid block parses");
assert_eq!(config.trusted_mcp_servers, vec!["knowledge_graph"]);
assert_eq!(config.untrusted_tools, vec!["scrape_page"]);
assert_eq!(config.trusted_tools, vec!["mcp__scanner__lint"]);
}
#[test]
fn classification_precedence_holds() {
let config = ContentTrustConfig {
trusted_mcp_servers: vec!["kg".to_string()],
untrusted_tools: vec!["scrape_page".to_string()],
trusted_tools: vec!["web_search".to_string(), "mcp__evil__probe".to_string()],
};
assert!(matches!(
config.classify_tool("web_search"),
ToolContentTrust::Untrusted { .. }
));
assert!(matches!(
config.classify_tool("scrape_page"),
ToolContentTrust::Untrusted { .. }
));
assert_eq!(
config.classify_tool("mcp__evil__probe"),
ToolContentTrust::Trusted
);
assert!(matches!(
config.classify_tool("mcp__other__search"),
ToolContentTrust::Untrusted { .. }
));
assert_eq!(
config.classify_tool("mcp__kg__query"),
ToolContentTrust::Trusted
);
assert_eq!(config.classify_tool("shell"), ToolContentTrust::Trusted);
}
#[test]
fn tracker_taints_on_untrusted_tool_and_clears_on_rotation() {
let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
let session = SessionId::new();
tracker.observe_agent_event("identity:a", &run_started(&session));
assert!(tracker.identity_taint("identity:a").is_none());
tracker.observe_agent_event("identity:a", &tool_result("shell"));
assert!(tracker.identity_taint("identity:a").is_none());
tracker.observe_agent_event("identity:a", &tool_result("web_search"));
let taint = tracker
.identity_taint("identity:a")
.expect("web tool result taints the session");
assert!(taint.source.contains("web_search"), "{}", taint.source);
assert!(tracker.session_taint(&session.to_string()).is_some());
tracker.observe_agent_event("identity:a", &tool_result("shell"));
assert!(tracker.identity_taint("identity:a").is_some());
let fresh = SessionId::new();
tracker.observe_agent_event("identity:a", &run_started(&fresh));
assert!(tracker.identity_taint("identity:a").is_none());
assert!(tracker.session_taint(&session.to_string()).is_some());
}
#[test]
fn tracker_taints_on_server_tool_content() {
let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
let session = SessionId::new();
tracker.note_current_session("identity:a", &session.to_string());
tracker.observe_agent_event(
"identity:a",
&AgentEvent::ServerToolContent {
id: None,
kind: ServerToolKind::WebSearch,
content: json!({"results": []}),
},
);
let taint = tracker.identity_taint("identity:a").expect("taints");
assert!(taint.source.contains("web_search"), "{}", taint.source);
}
#[test]
fn pre_attribution_taint_holds_identity_sticky_then_transfers() {
let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
tracker.observe_agent_event("identity:a", &tool_result("fetch"));
assert!(tracker.identity_taint("identity:a").is_some());
let session = SessionId::new();
tracker.observe_agent_event("identity:a", &run_started(&session));
assert!(tracker.session_taint(&session.to_string()).is_some());
assert!(tracker.identity_taint("identity:a").is_some());
}
#[test]
fn clear_identity_drops_attribution_but_keeps_the_session_fact() {
let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
let session = SessionId::new();
tracker.observe_agent_event("identity:a", &run_started(&session));
tracker.observe_agent_event("identity:a", &tool_result("web_fetch"));
assert!(tracker.identity_taint("identity:a").is_some());
tracker.clear_identity("identity:a");
assert!(tracker.identity_taint("identity:a").is_none());
assert!(tracker.session_taint(&session.to_string()).is_some());
assert!(
tracker
.evidence_quarantine_reason(&session.to_string())
.is_some()
);
}
#[test]
fn outbound_declarer_stamps_tainted_on_ingestion_and_clears_on_boundaries() {
use std::sync::{Arc, Mutex};
type Recorded = Vec<(String, Option<meerkat_core::comms::SenderContentTaint>)>;
let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
let calls: Arc<Mutex<Recorded>> = Arc::new(Mutex::new(Vec::new()));
let sink = calls.clone();
tracker.set_outbound_taint_declarer(Arc::new(move |identity, taint| {
sink.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push((identity.to_string(), taint));
}));
let session = SessionId::new();
tracker.observe_agent_event("identity:a", &run_started(&session));
tracker.observe_agent_event("identity:a", &tool_result("shell"));
assert!(calls.lock().unwrap().is_empty());
tracker.observe_agent_event("identity:a", &tool_result("web_search"));
{
let recorded = calls.lock().unwrap();
assert_eq!(recorded.len(), 1, "one declaration expected: {recorded:?}");
assert_eq!(recorded[0].0, "identity:a");
assert_eq!(
recorded[0].1,
Some(meerkat_core::comms::SenderContentTaint::Tainted)
);
}
let fresh = SessionId::new();
tracker.note_current_session("identity:a", &fresh.to_string());
{
let recorded = calls.lock().unwrap();
assert_eq!(recorded.len(), 2, "rotation clears: {recorded:?}");
assert_eq!(recorded[1].1, None);
}
tracker.observe_agent_event("identity:a", &tool_result("web_fetch"));
tracker.clear_identity("identity:a");
{
let recorded = calls.lock().unwrap();
assert_eq!(
recorded.last().expect("a clear declaration").1,
None,
"reset clears the outbound declaration: {recorded:?}"
);
}
}
#[test]
fn gate_quarantines_tainted_agents_and_quarantined_policy() {
let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
let session = SessionId::new();
tracker.observe_agent_event("identity:a", &run_started(&session));
let gate = TaintLlmWriteGate::new(Some(tracker.clone()), AgentMemoryLlmWrites::Observed);
let agent = MemoryAuthor::Agent {
identity: "identity:a".to_string(),
};
assert!(
gate.quarantine_reason(&agent, StagedBatchKind::FreshWrite, &[])
.is_none()
);
assert!(
gate.quarantine_reason(&MemoryAuthor::Application, StagedBatchKind::FreshWrite, &[])
.is_none()
);
tracker.observe_agent_event("identity:a", &tool_result("web_search"));
let reason = gate
.quarantine_reason(&agent, StagedBatchKind::FreshWrite, &[])
.expect("tainted quarantines");
assert!(reason.contains("session tainted"), "{reason}");
assert!(
gate.quarantine_reason(&MemoryAuthor::Application, StagedBatchKind::FreshWrite, &[])
.is_none()
);
assert!(
gate.quarantine_reason(&MemoryAuthor::Operator, StagedBatchKind::FreshWrite, &[])
.is_none()
);
let strict = TaintLlmWriteGate::new(None, AgentMemoryLlmWrites::Quarantined);
let reason = strict
.quarantine_reason(
&MemoryAuthor::Agent {
identity: "identity:clean".to_string(),
},
StagedBatchKind::FreshWrite,
&[],
)
.expect("policy quarantines untainted writes");
assert!(reason.contains("llm_writes=quarantined"), "{reason}");
assert!(
strict
.quarantine_reason(
&MemoryAuthor::Distiller {
run_id: "run-1".to_string()
},
StagedBatchKind::FreshWrite,
&[]
)
.is_some()
);
let steward = MemoryAuthor::Steward {
run_id: "run-1".to_string(),
};
assert!(
strict
.quarantine_reason(&steward, StagedBatchKind::FreshWrite, &[])
.is_some(),
"fresh steward LLM output (consolidate/harvest/rank) respects the posture"
);
assert!(
strict
.quarantine_reason(&steward, StagedBatchKind::ReviewVerdict, &[])
.is_none()
);
assert!(
strict
.quarantine_reason(&MemoryAuthor::Operator, StagedBatchKind::FreshWrite, &[])
.is_none()
);
}
#[test]
fn quarantined_posture_still_gates_review_verdicts_on_tainted_evidence() {
let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
let session = SessionId::new();
tracker.observe_agent_event("identity:a", &run_started(&session));
tracker.observe_agent_event("identity:a", &tool_result("web_search"));
let gate = TaintLlmWriteGate::new(Some(tracker), AgentMemoryLlmWrites::Quarantined);
let steward = MemoryAuthor::Steward {
run_id: "run-1".to_string(),
};
assert!(
gate.quarantine_reason(&steward, StagedBatchKind::ReviewVerdict, &[])
.is_none()
);
let reason = gate
.quarantine_reason(
&steward,
StagedBatchKind::ReviewVerdict,
&evidence_for(&session),
)
.expect("tainted evidence still quarantines review verdicts");
assert!(reason.contains("evidence session tainted"), "{reason}");
}
fn evidence_for(session: &SessionId) -> Vec<EvidenceRef> {
vec![EvidenceRef {
session_id: session.to_string(),
generation: 0,
revision: None,
range: Some((0, 4)),
}]
}
#[test]
fn gate_quarantines_llm_writes_citing_tainted_evidence() {
let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
let session = SessionId::new();
tracker.observe_agent_event("identity:a", &run_started(&session));
tracker.observe_agent_event("identity:a", &tool_result("web_search"));
let fresh = SessionId::new();
tracker.observe_agent_event("identity:a", &run_started(&fresh));
let gate = TaintLlmWriteGate::new(Some(tracker), AgentMemoryLlmWrites::Observed);
let distiller = MemoryAuthor::Distiller {
run_id: "run-1".to_string(),
};
let reason = gate
.quarantine_reason(
&distiller,
StagedBatchKind::FreshWrite,
&evidence_for(&session),
)
.expect("tainted evidence range quarantines (session-tainted ⇒ range-tainted)");
assert!(reason.contains("evidence session tainted"), "{reason}");
assert!(
gate.quarantine_reason(
&distiller,
StagedBatchKind::FreshWrite,
&evidence_for(&fresh)
)
.is_none()
);
assert!(
gate.quarantine_reason(
&MemoryAuthor::Operator,
StagedBatchKind::FreshWrite,
&evidence_for(&session)
)
.is_none()
);
}
#[test]
fn reset_boundary_quarantines_evidence_without_content_taint() {
let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
let session = SessionId::new();
tracker.observe_agent_event("identity:a", &run_started(&session));
assert!(
tracker
.evidence_quarantine_reason(&session.to_string())
.is_none()
);
tracker.mark_reset_boundary(&session.to_string());
let gate = TaintLlmWriteGate::new(Some(tracker), AgentMemoryLlmWrites::Observed);
let reason = gate
.quarantine_reason(
&MemoryAuthor::Distiller {
run_id: "run-1".to_string(),
},
StagedBatchKind::FreshWrite,
&evidence_for(&session),
)
.expect("reset boundary quarantines distillates");
assert!(reason.contains("reset boundary"), "{reason}");
}
#[test]
fn peer_projection_sender_parses_message_and_response_shapes() {
assert_eq!(
peer_projection_sender_identity("Peer message from mob-1/worker/identity:bob:"),
Some("identity:bob")
);
assert_eq!(
peer_projection_sender_identity(
"Peer response from mob-1/worker/identity:bob (to request: req-9)"
),
Some("identity:bob")
);
assert_eq!(
peer_projection_sender_identity("Peer message from scout:"),
Some("scout")
);
assert_eq!(
peer_projection_sender_identity("Peer request from peer_id 018fabc (id: r-1)"),
None
);
assert_eq!(peer_projection_sender_identity("ordinary text"), None);
}
#[test]
fn comms_join_taints_receiver_of_message_from_tainted_sender() {
let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
let sender_session = SessionId::new();
tracker.observe_agent_event("identity:bob", &run_started(&sender_session));
tracker.observe_agent_event("identity:bob", &tool_result("web_search"));
let receiver_session = SessionId::new();
let delivery = AgentEvent::RunStarted {
session_id: receiver_session.clone(),
input: meerkat_core::types::RunInput::Content {
content: meerkat_core::ContentInput::Text(
"Peer message from mob-1/worker/identity:bob:\nplease remember X".to_string(),
),
},
};
tracker.observe_agent_event("identity:alice", &delivery);
let taint = tracker
.identity_taint("identity:alice")
.expect("receiver session taints (peer-laundering close, §10.1)");
assert!(taint.source.contains("identity:bob"), "{}", taint.source);
assert!(
tracker
.session_taint(&receiver_session.to_string())
.is_some()
);
let clean_session = SessionId::new();
tracker.observe_agent_event("identity:carol", &run_started(&clean_session));
let receiver2 = SessionId::new();
let clean_delivery = AgentEvent::RunStarted {
session_id: receiver2.clone(),
input: meerkat_core::types::RunInput::Content {
content: meerkat_core::ContentInput::Text(
"Peer message from mob-1/worker/identity:carol:\nhello".to_string(),
),
},
};
tracker.observe_agent_event("identity:dave", &clean_delivery);
assert!(tracker.identity_taint("identity:dave").is_none());
}
}