1use std::collections::HashMap;
54use std::sync::{Arc, Mutex};
55use std::time::{SystemTime, UNIX_EPOCH};
56
57use meerkat_core::event::AgentEvent;
58use serde::{Deserialize, Serialize};
59
60use crate::identity_first::agent_memory::AgentMemoryLlmWrites;
61use crate::memory::records::{EvidenceRef, MemoryAuthor};
62use crate::memory::staged::StagedBatchKind;
63
64const ALWAYS_UNTRUSTED_TOOL_NAMES: &[&str] = &["web_search", "web_fetch", "fetch", "http_request"];
70
71const MCP_QUALIFIED_PREFIX: &str = "mcp__";
75
76const MAX_TRACKED_TAINTED_SESSIONS: usize = 4096;
80
81#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
89pub struct ContentTrustConfig {
90 #[serde(default)]
93 pub trusted_mcp_servers: Vec<String>,
94 #[serde(default)]
97 pub untrusted_tools: Vec<String>,
98 #[serde(default)]
101 pub trusted_tools: Vec<String>,
102}
103
104#[derive(Debug, Clone, PartialEq, Eq)]
106pub enum ToolContentTrust {
107 Trusted,
108 Untrusted { source: String },
109}
110
111impl ContentTrustConfig {
112 pub fn from_json_value(value: &serde_json::Value) -> Result<Self, String> {
116 let object = value
117 .as_object()
118 .ok_or_else(|| "content_trust must be an object".to_string())?;
119 let supported = ["trusted_mcp_servers", "untrusted_tools", "trusted_tools"];
120 let unsupported = object
121 .keys()
122 .filter(|key| !supported.contains(&key.as_str()))
123 .map(String::as_str)
124 .collect::<Vec<_>>();
125 if !unsupported.is_empty() {
126 return Err(format!(
127 "unsupported content_trust fields: {}",
128 unsupported.join(", ")
129 ));
130 }
131 let parse_names = |key: &str| -> Result<Vec<String>, String> {
132 match object.get(key) {
133 None => Ok(Vec::new()),
134 Some(value) => {
135 let entries = value
136 .as_array()
137 .ok_or_else(|| format!("content_trust.{key} must be an array"))?;
138 entries
139 .iter()
140 .map(|entry| {
141 entry
142 .as_str()
143 .map(str::trim)
144 .filter(|name| !name.is_empty())
145 .map(ToString::to_string)
146 .ok_or_else(|| {
147 format!("content_trust.{key} entries must be non-empty strings")
148 })
149 })
150 .collect()
151 }
152 }
153 };
154 Ok(Self {
155 trusted_mcp_servers: parse_names("trusted_mcp_servers")?,
156 untrusted_tools: parse_names("untrusted_tools")?,
157 trusted_tools: parse_names("trusted_tools")?,
158 })
159 }
160
161 pub fn classify_tool(&self, name: &str) -> ToolContentTrust {
166 if ALWAYS_UNTRUSTED_TOOL_NAMES.contains(&name) {
167 return ToolContentTrust::Untrusted {
168 source: format!("web tool '{name}'"),
169 };
170 }
171 if self.untrusted_tools.iter().any(|tool| tool == name) {
172 return ToolContentTrust::Untrusted {
173 source: format!("configured untrusted tool '{name}'"),
174 };
175 }
176 if self.trusted_tools.iter().any(|tool| tool == name) {
177 return ToolContentTrust::Trusted;
178 }
179 if let Some(rest) = name.strip_prefix(MCP_QUALIFIED_PREFIX) {
180 let server = rest.split("__").next().unwrap_or(rest);
181 if self
182 .trusted_mcp_servers
183 .iter()
184 .any(|trusted| trusted == server)
185 {
186 return ToolContentTrust::Trusted;
187 }
188 return ToolContentTrust::Untrusted {
189 source: format!("MCP server '{server}' (tool '{name}')"),
190 };
191 }
192 ToolContentTrust::Trusted
193 }
194}
195
196#[derive(Debug, Clone, PartialEq, Eq)]
198pub struct TaintState {
199 pub tainted_at_ms: u64,
200 pub source: String,
201}
202
203#[derive(Default)]
204struct TaintInner {
205 current_session: HashMap<String, String>,
210 tainted: HashMap<String, TaintState>,
215 pending_identity_taint: HashMap<String, TaintState>,
219 reset_boundaries: HashMap<String, u64>,
222}
223
224pub type OutboundTaintDeclarer =
232 Arc<dyn Fn(&str, Option<meerkat_core::comms::SenderContentTaint>) + Send + Sync>;
233
234#[derive(Clone, Default)]
237pub struct SessionTaintTracker {
238 config: Arc<ContentTrustConfig>,
239 inner: Arc<Mutex<TaintInner>>,
240 event_sink: Arc<Mutex<Option<Arc<dyn crate::memory::events::MemoryEventSink>>>>,
242 outbound_declarer: Arc<Mutex<Option<OutboundTaintDeclarer>>>,
246}
247
248impl SessionTaintTracker {
249 pub fn new(config: ContentTrustConfig) -> Self {
250 Self {
251 config: Arc::new(config),
252 inner: Arc::new(Mutex::new(TaintInner::default())),
253 event_sink: Arc::new(Mutex::new(None)),
254 outbound_declarer: Arc::new(Mutex::new(None)),
255 }
256 }
257
258 pub fn set_outbound_taint_declarer(&self, declarer: OutboundTaintDeclarer) {
261 *self
262 .outbound_declarer
263 .lock()
264 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(declarer);
265 }
266
267 fn declare_outbound(
271 &self,
272 identity: &str,
273 taint: Option<meerkat_core::comms::SenderContentTaint>,
274 ) {
275 if let Some(declarer) = self
276 .outbound_declarer
277 .lock()
278 .unwrap_or_else(std::sync::PoisonError::into_inner)
279 .as_ref()
280 {
281 declarer(identity, taint);
282 }
283 }
284
285 pub fn set_event_sink(&self, sink: Arc<dyn crate::memory::events::MemoryEventSink>) {
288 *self
289 .event_sink
290 .lock()
291 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(sink);
292 }
293
294 fn emit_event(&self, event: crate::memory::events::MemoryTimelineEvent) {
295 if let Some(sink) = self
296 .event_sink
297 .lock()
298 .unwrap_or_else(std::sync::PoisonError::into_inner)
299 .as_ref()
300 {
301 sink.emit(event);
302 }
303 }
304
305 pub fn observe_agent_event(&self, identity: &str, event: &AgentEvent) {
308 match event {
309 AgentEvent::RunStarted { session_id, input } => {
310 let session_key = session_id.to_string();
311 self.note_current_session(identity, &session_key);
312 if let Some(text) = input.prompt_text() {
317 self.observe_inbound_peer_content(identity, &session_key, &text);
318 }
319 }
320 AgentEvent::PeerContentIngested {
327 peer, sender_taint, ..
328 } => {
329 if *sender_taint == Some(meerkat_core::comms::SenderContentTaint::Tainted) {
330 let sender = peer
331 .as_ref()
332 .and_then(|peer| peer.display_name.clone())
333 .unwrap_or_else(|| "peer".to_string());
334 self.mark_identity_tainted(
335 identity,
336 format!("peer content declared tainted by sender '{sender}'"),
337 );
338 }
339 }
340 AgentEvent::ToolResultReceived { name, .. }
344 | AgentEvent::ToolExecutionCompleted { name, .. } => {
345 if let ToolContentTrust::Untrusted { source } = self.config.classify_tool(name) {
346 self.mark_identity_tainted(identity, source);
347 }
348 }
349 AgentEvent::ServerToolContent { kind, .. } => {
352 self.mark_identity_tainted(
353 identity,
354 format!("provider server tool '{}'", kind.provider_name()),
355 );
356 }
357 _ => {}
358 }
359 }
360
361 pub fn note_current_session(&self, identity: &str, session_key: &str) {
365 let mut inner = self
366 .inner
367 .lock()
368 .unwrap_or_else(std::sync::PoisonError::into_inner);
369 let pending = inner.pending_identity_taint.remove(identity);
370 let previous = inner
371 .current_session
372 .insert(identity.to_string(), session_key.to_string());
373 let rotated_away_from_tainted = previous.as_deref() != Some(session_key)
374 && previous.is_some_and(|prior| inner.tainted.contains_key(&prior));
375 if rotated_away_from_tainted {
376 tracing::warn!(
377 identity,
378 session_key,
379 "agent memory taint: session rotated away from a tainted session; \
380 new session starts clean"
381 );
382 self.emit_event(
383 crate::memory::events::MemoryTimelineEvent::TaintTransition {
384 identity: Some(identity.to_string()),
385 session_key: session_key.to_string(),
386 kind: "rotated_clean".to_string(),
387 source: "session rotation away from tainted session".to_string(),
388 },
389 );
390 }
391 let pending_reapplied = pending.is_some();
394 if let Some(state) = pending {
395 self.insert_taint(&mut inner, session_key.to_string(), state);
396 }
397 if rotated_away_from_tainted && !pending_reapplied {
400 drop(inner);
401 self.declare_outbound(identity, None);
402 }
403 }
404
405 pub fn clear_identity(&self, identity: &str) {
413 let mut inner = self
414 .inner
415 .lock()
416 .unwrap_or_else(std::sync::PoisonError::into_inner);
417 inner.pending_identity_taint.remove(identity);
418 inner.current_session.remove(identity);
419 drop(inner);
422 self.declare_outbound(identity, None);
423 }
424
425 pub fn mark_reset_boundary(&self, session_key: &str) {
429 let mut inner = self
430 .inner
431 .lock()
432 .unwrap_or_else(std::sync::PoisonError::into_inner);
433 if inner
434 .reset_boundaries
435 .insert(session_key.to_string(), now_ms())
436 .is_none()
437 {
438 tracing::warn!(
439 session_key,
440 "agent memory taint: reset boundary marked; distillates over this \
441 session will land quarantined pending steward review (§8.4)"
442 );
443 self.emit_event(
444 crate::memory::events::MemoryTimelineEvent::TaintTransition {
445 identity: None,
446 session_key: session_key.to_string(),
447 kind: "reset_boundary".to_string(),
448 source: "reset() boundary (§8.4)".to_string(),
449 },
450 );
451 }
452 if inner.reset_boundaries.len() > MAX_TRACKED_TAINTED_SESSIONS
453 && let Some(oldest) = inner
454 .reset_boundaries
455 .iter()
456 .min_by_key(|(_, at_ms)| **at_ms)
457 .map(|(key, _)| key.clone())
458 {
459 inner.reset_boundaries.remove(&oldest);
460 }
461 }
462
463 pub fn evidence_quarantine_reason(&self, session_key: &str) -> Option<String> {
468 let inner = self
469 .inner
470 .lock()
471 .unwrap_or_else(std::sync::PoisonError::into_inner);
472 if let Some(state) = inner.tainted.get(session_key) {
473 return Some(format!(
474 "evidence session tainted by {} (session-tainted ⇒ range-tainted)",
475 state.source
476 ));
477 }
478 if inner.reset_boundaries.contains_key(session_key) {
479 return Some(
480 "evidence session closed at a reset boundary; distillates quarantine \
481 pending steward review (§8.4)"
482 .to_string(),
483 );
484 }
485 None
486 }
487
488 pub fn session_taint(&self, session_key: &str) -> Option<TaintState> {
490 let inner = self
491 .inner
492 .lock()
493 .unwrap_or_else(std::sync::PoisonError::into_inner);
494 inner.tainted.get(session_key).cloned()
495 }
496
497 pub fn identity_taint(&self, identity: &str) -> Option<TaintState> {
500 let inner = self
501 .inner
502 .lock()
503 .unwrap_or_else(std::sync::PoisonError::into_inner);
504 if let Some(state) = inner.pending_identity_taint.get(identity) {
505 return Some(state.clone());
506 }
507 inner
508 .current_session
509 .get(identity)
510 .and_then(|session| inner.tainted.get(session))
511 .cloned()
512 }
513
514 fn observe_inbound_peer_content(&self, identity: &str, session_key: &str, text: &str) {
539 for line in text.lines() {
540 let Some(sender_identity) = peer_projection_sender_identity(line) else {
541 continue;
542 };
543 if sender_identity == identity {
544 continue;
545 }
546 let source = {
547 let inner = self
548 .inner
549 .lock()
550 .unwrap_or_else(std::sync::PoisonError::into_inner);
551 inner
552 .pending_identity_taint
553 .get(sender_identity)
554 .or_else(|| {
555 inner
556 .current_session
557 .get(sender_identity)
558 .and_then(|session| inner.tainted.get(session))
559 })
560 .map(|state| state.source.clone())
561 };
562 if let Some(source) = source {
563 let state = TaintState {
564 tainted_at_ms: now_ms(),
565 source: format!(
566 "peer message from tainted sender '{sender_identity}' \
567 (sender session tainted by {source})"
568 ),
569 };
570 let mut inner = self
571 .inner
572 .lock()
573 .unwrap_or_else(std::sync::PoisonError::into_inner);
574 self.insert_taint(&mut inner, session_key.to_string(), state);
575 }
576 }
577 }
578
579 fn mark_identity_tainted(&self, identity: &str, source: String) {
580 let state = TaintState {
581 tainted_at_ms: now_ms(),
582 source,
583 };
584 let mut inner = self
585 .inner
586 .lock()
587 .unwrap_or_else(std::sync::PoisonError::into_inner);
588 match inner.current_session.get(identity).cloned() {
589 Some(session) => self.insert_taint(&mut inner, session, state),
590 None => {
591 match inner.pending_identity_taint.entry(identity.to_string()) {
594 std::collections::hash_map::Entry::Vacant(slot) => {
595 tracing::warn!(
596 identity,
597 source = %state.source,
598 "agent memory taint: untrusted ingestion observed before session \
599 attribution; holding identity-sticky taint"
600 );
601 slot.insert(state);
602 }
603 std::collections::hash_map::Entry::Occupied(mut slot) => {
604 slot.insert(state);
605 }
606 }
607 }
608 }
609 drop(inner);
613 self.declare_outbound(
614 identity,
615 Some(meerkat_core::comms::SenderContentTaint::Tainted),
616 );
617 }
618
619 fn insert_taint(&self, inner: &mut TaintInner, session: String, state: TaintState) {
620 match inner.tainted.entry(session) {
621 std::collections::hash_map::Entry::Occupied(_) => return,
622 std::collections::hash_map::Entry::Vacant(slot) => {
623 tracing::warn!(
624 session_key = %slot.key(),
625 source = %state.source,
626 "agent memory taint: session ingested untrusted content; LLM-authored \
627 memory writes from this session will quarantine until a fresh-context \
628 boundary (reset/respawn/fresh spawn)"
629 );
630 self.emit_event(
631 crate::memory::events::MemoryTimelineEvent::TaintTransition {
632 identity: None,
633 session_key: slot.key().clone(),
634 kind: "tainted".to_string(),
635 source: state.source.clone(),
636 },
637 );
638 slot.insert(state);
639 }
640 }
641 if inner.tainted.len() > MAX_TRACKED_TAINTED_SESSIONS
642 && let Some(oldest) = inner
643 .tainted
644 .iter()
645 .min_by_key(|(_, state)| state.tainted_at_ms)
646 .map(|(key, _)| key.clone())
647 {
648 inner.tainted.remove(&oldest);
649 }
650 }
651}
652
653pub trait LlmWriteGate: Send + Sync {
658 fn quarantine_reason(
665 &self,
666 author: &MemoryAuthor,
667 kind: StagedBatchKind,
668 evidence: &[EvidenceRef],
669 ) -> Option<String>;
670}
671
672pub struct TaintLlmWriteGate {
693 tracker: Option<SessionTaintTracker>,
694 llm_writes: AgentMemoryLlmWrites,
695}
696
697impl TaintLlmWriteGate {
698 pub fn new(tracker: Option<SessionTaintTracker>, llm_writes: AgentMemoryLlmWrites) -> Self {
699 Self {
700 tracker,
701 llm_writes,
702 }
703 }
704}
705
706impl LlmWriteGate for TaintLlmWriteGate {
707 fn quarantine_reason(
708 &self,
709 author: &MemoryAuthor,
710 kind: StagedBatchKind,
711 evidence: &[EvidenceRef],
712 ) -> Option<String> {
713 if !author.is_llm() {
714 return None;
715 }
716 if self.llm_writes == AgentMemoryLlmWrites::Quarantined
717 && kind != StagedBatchKind::ReviewVerdict
718 {
719 return Some("llm_writes=quarantined policy".to_string());
720 }
721 let tracker = self.tracker.as_ref()?;
722 if let MemoryAuthor::Agent { identity } = author
723 && let Some(state) = tracker.identity_taint(identity)
724 {
725 return Some(format!("session tainted by {}", state.source));
726 }
727 for evidence_ref in evidence {
728 if let Some(reason) = tracker.evidence_quarantine_reason(&evidence_ref.session_id) {
729 return Some(reason);
730 }
731 }
732 None
733 }
734}
735
736fn peer_projection_sender_identity(line: &str) -> Option<&str> {
745 let name = if let Some(rest) = line.strip_prefix("Peer message from ") {
746 match rest.split_once(": ") {
750 Some((name, _)) => name.trim(),
751 None => rest.strip_suffix(':').unwrap_or(rest).trim(),
752 }
753 } else if let Some(rest) = line.strip_prefix("Peer response from ") {
754 rest.split(" (to request:").next()?.trim()
755 } else {
756 return None;
757 };
758 if name.is_empty() {
759 return None;
760 }
761 Some(name.rsplit('/').next().unwrap_or(name))
762}
763
764pub trait MemberAgentEventSink: Send + Sync {
773 fn observe(&self, identity: &str, envelope: &meerkat_core::event::EventEnvelope<AgentEvent>);
774}
775
776impl MemberAgentEventSink for SessionTaintTracker {
777 fn observe(&self, identity: &str, envelope: &meerkat_core::event::EventEnvelope<AgentEvent>) {
778 self.observe_agent_event(identity, &envelope.payload);
779 }
780}
781
782pub struct CompactionResetSink {
789 on_compacted: Arc<dyn Fn(&str) + Send + Sync>,
790}
791
792impl CompactionResetSink {
793 pub fn new(on_compacted: Arc<dyn Fn(&str) + Send + Sync>) -> Self {
794 Self { on_compacted }
795 }
796}
797
798impl MemberAgentEventSink for CompactionResetSink {
799 fn observe(&self, _identity: &str, envelope: &meerkat_core::event::EventEnvelope<AgentEvent>) {
800 if matches!(envelope.payload, AgentEvent::CompactionCompleted { .. })
801 && let meerkat_core::event::EventSourceIdentity::Session { session_id } =
802 &envelope.source
803 {
804 (self.on_compacted)(&session_id.to_string());
805 }
806 }
807}
808
809#[derive(Clone)]
812pub struct TaintObserverGuard {
813 _abort: Arc<AbortOnDrop>,
814}
815
816struct AbortOnDrop(tokio::task::JoinHandle<()>);
817
818impl Drop for AbortOnDrop {
819 fn drop(&mut self) {
820 self.0.abort();
821 }
822}
823
824pub fn spawn_taint_observer(
828 handle: meerkat_mob::MobHandle,
829 tracker: SessionTaintTracker,
830) -> TaintObserverGuard {
831 spawn_member_event_observer(handle, vec![Arc::new(tracker)])
832}
833
834pub fn spawn_member_event_observer(
838 handle: meerkat_mob::MobHandle,
839 sinks: Vec<Arc<dyn MemberAgentEventSink>>,
840) -> TaintObserverGuard {
841 let task = tokio::spawn(run_member_event_observer(handle, sinks));
842 TaintObserverGuard {
843 _abort: Arc::new(AbortOnDrop(task)),
844 }
845}
846
847async fn run_member_event_observer(
848 handle: meerkat_mob::MobHandle,
849 sinks: Vec<Arc<dyn MemberAgentEventSink>>,
850) {
851 use futures::StreamExt;
852 use futures::stream::SelectAll;
853
854 enum Observed {
855 Event(String, Box<meerkat_core::event::EventEnvelope<AgentEvent>>),
856 Closed(String),
857 }
858
859 let mut streams: SelectAll<futures::stream::BoxStream<'static, Observed>> = SelectAll::new();
860 let mut subscribed: std::collections::HashSet<String> = std::collections::HashSet::new();
861 let mut warned: std::collections::HashSet<String> = std::collections::HashSet::new();
862 let mut reconcile = tokio::time::interval(std::time::Duration::from_secs(1));
863 reconcile.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
864
865 loop {
866 tokio::select! {
867 Some(observed) = streams.next() => match observed {
868 Observed::Event(identity, envelope) => {
869 for sink in &sinks {
870 sink.observe(&identity, &envelope);
871 }
872 }
873 Observed::Closed(identity) => {
874 subscribed.remove(&identity);
875 }
876 },
877 _ = reconcile.tick() => {
878 for entry in handle.list_members_including_retiring().await {
879 if entry.status != meerkat_mob::MobMemberStatus::Active {
883 continue;
884 }
885 let identity = entry.agent_identity.to_string();
886 if subscribed.contains(&identity) {
887 continue;
888 }
889 match handle.subscribe_agent_events(&entry.agent_identity).await {
890 Ok(stream) => {
891 warned.remove(&identity);
892 subscribed.insert(identity.clone());
893 let close_key = identity.clone();
894 streams.push(
895 stream
896 .map(move |envelope| {
897 Observed::Event(identity.clone(), Box::new(envelope))
898 })
899 .chain(futures::stream::once(async move {
900 Observed::Closed(close_key)
901 }))
902 .boxed(),
903 );
904 }
905 Err(error) => {
906 if warned.insert(identity.clone()) {
909 tracing::warn!(
910 identity = %identity,
911 error = %error,
912 "agent memory taint observer: failed to subscribe; will retry"
913 );
914 } else {
915 tracing::debug!(
916 identity = %identity,
917 error = %error,
918 "agent memory taint observer: subscribe still failing"
919 );
920 }
921 }
922 }
923 }
924 }
925 }
926 }
927}
928
929fn now_ms() -> u64 {
930 SystemTime::now()
931 .duration_since(UNIX_EPOCH)
932 .map(|duration| duration.as_millis() as u64)
933 .unwrap_or(0)
934}
935
936#[cfg(test)]
937#[allow(
938 clippy::expect_used,
939 clippy::panic,
940 clippy::redundant_clone,
941 clippy::unwrap_used
942)]
943mod tests {
944 use super::*;
945 use meerkat_core::types::{ContentBlock, ServerToolKind, SessionId};
946 use serde_json::json;
947
948 fn run_started(session: &SessionId) -> AgentEvent {
949 AgentEvent::RunStarted {
950 session_id: session.clone(),
951 input: meerkat_core::types::RunInput::Content {
952 content: meerkat_core::ContentInput::Text("hi".to_string()),
953 },
954 }
955 }
956
957 fn tool_result(name: &str) -> AgentEvent {
958 AgentEvent::ToolResultReceived {
959 id: "tool-1".to_string(),
960 name: name.to_string(),
961 content: vec![ContentBlock::Text {
962 text: "ok".to_string(),
963 }],
964 is_error: false,
965 }
966 }
967
968 fn peer_ingested(taint: Option<meerkat_core::comms::SenderContentTaint>) -> AgentEvent {
969 AgentEvent::PeerContentIngested {
970 kind: meerkat_core::types::CommsNoticeKind::Message,
971 peer: None,
972 request_id: None,
973 sender_taint: taint,
974 }
975 }
976
977 #[test]
981 fn declared_peer_taint_taints_receiver_but_none_or_clean_does_not() {
982 use meerkat_core::comms::SenderContentTaint;
983 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
984 tracker.note_current_session("identity:b", "sess-b");
985
986 tracker.observe_agent_event("identity:b", &peer_ingested(None));
987 assert!(
988 tracker.session_taint("sess-b").is_none(),
989 "no declaration must not taint"
990 );
991
992 tracker.observe_agent_event(
993 "identity:b",
994 &peer_ingested(Some(SenderContentTaint::Clean)),
995 );
996 assert!(
997 tracker.session_taint("sess-b").is_none(),
998 "an affirmative Clean declaration must not taint"
999 );
1000
1001 tracker.observe_agent_event(
1002 "identity:b",
1003 &peer_ingested(Some(SenderContentTaint::Tainted)),
1004 );
1005 assert!(
1006 tracker.session_taint("sess-b").is_some(),
1007 "a declared-tainted peer delivery taints the receiving session"
1008 );
1009 }
1010
1011 #[test]
1012 fn taint_transitions_emit_timeline_events_when_sink_wired() {
1013 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1014 let sink = std::sync::Arc::new(crate::memory::events::CollectingEventSink::new());
1015 tracker.set_event_sink(sink.clone());
1016
1017 tracker.note_current_session("identity:a", "sess-1");
1018 tracker.observe_agent_event("identity:a", &tool_result("web_fetch"));
1019 tracker.mark_reset_boundary("sess-1");
1020 tracker.mark_reset_boundary("sess-1");
1022 tracker.note_current_session("identity:a", "sess-2");
1023
1024 let types = sink.types();
1025 assert_eq!(
1026 types,
1027 vec![
1028 "memory.taint.transition", "memory.taint.transition", "memory.taint.transition", ]
1032 );
1033 let events = sink.events.lock().unwrap();
1034 let kinds: Vec<String> = events
1035 .iter()
1036 .map(|event| match event {
1037 crate::memory::events::MemoryTimelineEvent::TaintTransition { kind, .. } => {
1038 kind.clone()
1039 }
1040 other => panic!("unexpected event {other:?}"),
1041 })
1042 .collect();
1043 assert_eq!(kinds, vec!["tainted", "reset_boundary", "rotated_clean"]);
1044 }
1045
1046 #[test]
1047 fn content_trust_parse_rejects_unknown_fields_and_bad_types() {
1048 let err = ContentTrustConfig::from_json_value(&json!({"servers": []}))
1049 .expect_err("unknown field must fail loud");
1050 assert!(err.contains("unsupported content_trust fields"), "{err}");
1051 let err = ContentTrustConfig::from_json_value(&json!({"trusted_mcp_servers": "kg"}))
1052 .expect_err("non-array must fail loud");
1053 assert!(err.contains("must be an array"), "{err}");
1054 let err = ContentTrustConfig::from_json_value(&json!({"untrusted_tools": [1]}))
1055 .expect_err("non-string entry must fail loud");
1056 assert!(err.contains("non-empty strings"), "{err}");
1057 let err =
1058 ContentTrustConfig::from_json_value(&json!([])).expect_err("non-object must fail loud");
1059 assert!(err.contains("must be an object"), "{err}");
1060 }
1061
1062 #[test]
1063 fn content_trust_parse_accepts_full_block() {
1064 let config = ContentTrustConfig::from_json_value(&json!({
1065 "trusted_mcp_servers": ["knowledge_graph"],
1066 "untrusted_tools": ["scrape_page"],
1067 "trusted_tools": ["mcp__scanner__lint"],
1068 }))
1069 .expect("valid block parses");
1070 assert_eq!(config.trusted_mcp_servers, vec!["knowledge_graph"]);
1071 assert_eq!(config.untrusted_tools, vec!["scrape_page"]);
1072 assert_eq!(config.trusted_tools, vec!["mcp__scanner__lint"]);
1073 }
1074
1075 #[test]
1076 fn classification_precedence_holds() {
1077 let config = ContentTrustConfig {
1078 trusted_mcp_servers: vec!["kg".to_string()],
1079 untrusted_tools: vec!["scrape_page".to_string()],
1080 trusted_tools: vec!["web_search".to_string(), "mcp__evil__probe".to_string()],
1082 };
1083 assert!(matches!(
1084 config.classify_tool("web_search"),
1085 ToolContentTrust::Untrusted { .. }
1086 ));
1087 assert!(matches!(
1088 config.classify_tool("scrape_page"),
1089 ToolContentTrust::Untrusted { .. }
1090 ));
1091 assert_eq!(
1093 config.classify_tool("mcp__evil__probe"),
1094 ToolContentTrust::Trusted
1095 );
1096 assert!(matches!(
1098 config.classify_tool("mcp__other__search"),
1099 ToolContentTrust::Untrusted { .. }
1100 ));
1101 assert_eq!(
1102 config.classify_tool("mcp__kg__query"),
1103 ToolContentTrust::Trusted
1104 );
1105 assert_eq!(config.classify_tool("shell"), ToolContentTrust::Trusted);
1107 }
1108
1109 #[test]
1110 fn tracker_taints_on_untrusted_tool_and_clears_on_rotation() {
1111 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1112 let session = SessionId::new();
1113 tracker.observe_agent_event("identity:a", &run_started(&session));
1114 assert!(tracker.identity_taint("identity:a").is_none());
1115
1116 tracker.observe_agent_event("identity:a", &tool_result("shell"));
1117 assert!(tracker.identity_taint("identity:a").is_none());
1118
1119 tracker.observe_agent_event("identity:a", &tool_result("web_search"));
1120 let taint = tracker
1121 .identity_taint("identity:a")
1122 .expect("web tool result taints the session");
1123 assert!(taint.source.contains("web_search"), "{}", taint.source);
1124 assert!(tracker.session_taint(&session.to_string()).is_some());
1125
1126 tracker.observe_agent_event("identity:a", &tool_result("shell"));
1128 assert!(tracker.identity_taint("identity:a").is_some());
1129
1130 let fresh = SessionId::new();
1132 tracker.observe_agent_event("identity:a", &run_started(&fresh));
1133 assert!(tracker.identity_taint("identity:a").is_none());
1134 assert!(tracker.session_taint(&session.to_string()).is_some());
1136 }
1137
1138 #[test]
1139 fn tracker_taints_on_server_tool_content() {
1140 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1141 let session = SessionId::new();
1142 tracker.note_current_session("identity:a", &session.to_string());
1143 tracker.observe_agent_event(
1144 "identity:a",
1145 &AgentEvent::ServerToolContent {
1146 id: None,
1147 kind: ServerToolKind::WebSearch,
1148 content: json!({"results": []}),
1149 },
1150 );
1151 let taint = tracker.identity_taint("identity:a").expect("taints");
1152 assert!(taint.source.contains("web_search"), "{}", taint.source);
1153 }
1154
1155 #[test]
1156 fn pre_attribution_taint_holds_identity_sticky_then_transfers() {
1157 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1158 tracker.observe_agent_event("identity:a", &tool_result("fetch"));
1160 assert!(tracker.identity_taint("identity:a").is_some());
1161
1162 let session = SessionId::new();
1164 tracker.observe_agent_event("identity:a", &run_started(&session));
1165 assert!(tracker.session_taint(&session.to_string()).is_some());
1166 assert!(tracker.identity_taint("identity:a").is_some());
1167 }
1168
1169 #[test]
1170 fn clear_identity_drops_attribution_but_keeps_the_session_fact() {
1171 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1172 let session = SessionId::new();
1173 tracker.observe_agent_event("identity:a", &run_started(&session));
1174 tracker.observe_agent_event("identity:a", &tool_result("web_fetch"));
1175 assert!(tracker.identity_taint("identity:a").is_some());
1176 tracker.clear_identity("identity:a");
1177 assert!(tracker.identity_taint("identity:a").is_none());
1179 assert!(tracker.session_taint(&session.to_string()).is_some());
1182 assert!(
1183 tracker
1184 .evidence_quarantine_reason(&session.to_string())
1185 .is_some()
1186 );
1187 }
1188
1189 #[test]
1190 fn outbound_declarer_stamps_tainted_on_ingestion_and_clears_on_boundaries() {
1191 use std::sync::{Arc, Mutex};
1192 type Recorded = Vec<(String, Option<meerkat_core::comms::SenderContentTaint>)>;
1193 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1194 let calls: Arc<Mutex<Recorded>> = Arc::new(Mutex::new(Vec::new()));
1195 let sink = calls.clone();
1196 tracker.set_outbound_taint_declarer(Arc::new(move |identity, taint| {
1197 sink.lock()
1198 .unwrap_or_else(std::sync::PoisonError::into_inner)
1199 .push((identity.to_string(), taint));
1200 }));
1201
1202 let session = SessionId::new();
1203 tracker.observe_agent_event("identity:a", &run_started(&session));
1204 tracker.observe_agent_event("identity:a", &tool_result("shell"));
1206 assert!(calls.lock().unwrap().is_empty());
1207
1208 tracker.observe_agent_event("identity:a", &tool_result("web_search"));
1210 {
1211 let recorded = calls.lock().unwrap();
1212 assert_eq!(recorded.len(), 1, "one declaration expected: {recorded:?}");
1213 assert_eq!(recorded[0].0, "identity:a");
1214 assert_eq!(
1215 recorded[0].1,
1216 Some(meerkat_core::comms::SenderContentTaint::Tainted)
1217 );
1218 }
1219
1220 let fresh = SessionId::new();
1223 tracker.note_current_session("identity:a", &fresh.to_string());
1224 {
1225 let recorded = calls.lock().unwrap();
1226 assert_eq!(recorded.len(), 2, "rotation clears: {recorded:?}");
1227 assert_eq!(recorded[1].1, None);
1228 }
1229
1230 tracker.observe_agent_event("identity:a", &tool_result("web_fetch"));
1232 tracker.clear_identity("identity:a");
1233 {
1234 let recorded = calls.lock().unwrap();
1235 assert_eq!(
1236 recorded.last().expect("a clear declaration").1,
1237 None,
1238 "reset clears the outbound declaration: {recorded:?}"
1239 );
1240 }
1241 }
1242
1243 #[test]
1244 fn gate_quarantines_tainted_agents_and_quarantined_policy() {
1245 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1246 let session = SessionId::new();
1247 tracker.observe_agent_event("identity:a", &run_started(&session));
1248
1249 let gate = TaintLlmWriteGate::new(Some(tracker.clone()), AgentMemoryLlmWrites::Observed);
1250 let agent = MemoryAuthor::Agent {
1251 identity: "identity:a".to_string(),
1252 };
1253 assert!(
1254 gate.quarantine_reason(&agent, StagedBatchKind::FreshWrite, &[])
1255 .is_none()
1256 );
1257 assert!(
1258 gate.quarantine_reason(&MemoryAuthor::Application, StagedBatchKind::FreshWrite, &[])
1259 .is_none()
1260 );
1261
1262 tracker.observe_agent_event("identity:a", &tool_result("web_search"));
1263 let reason = gate
1264 .quarantine_reason(&agent, StagedBatchKind::FreshWrite, &[])
1265 .expect("tainted quarantines");
1266 assert!(reason.contains("session tainted"), "{reason}");
1267 assert!(
1269 gate.quarantine_reason(&MemoryAuthor::Application, StagedBatchKind::FreshWrite, &[])
1270 .is_none()
1271 );
1272 assert!(
1273 gate.quarantine_reason(&MemoryAuthor::Operator, StagedBatchKind::FreshWrite, &[])
1274 .is_none()
1275 );
1276
1277 let strict = TaintLlmWriteGate::new(None, AgentMemoryLlmWrites::Quarantined);
1281 let reason = strict
1282 .quarantine_reason(
1283 &MemoryAuthor::Agent {
1284 identity: "identity:clean".to_string(),
1285 },
1286 StagedBatchKind::FreshWrite,
1287 &[],
1288 )
1289 .expect("policy quarantines untainted writes");
1290 assert!(reason.contains("llm_writes=quarantined"), "{reason}");
1291 assert!(
1292 strict
1293 .quarantine_reason(
1294 &MemoryAuthor::Distiller {
1295 run_id: "run-1".to_string()
1296 },
1297 StagedBatchKind::FreshWrite,
1298 &[]
1299 )
1300 .is_some()
1301 );
1302 let steward = MemoryAuthor::Steward {
1303 run_id: "run-1".to_string(),
1304 };
1305 assert!(
1306 strict
1307 .quarantine_reason(&steward, StagedBatchKind::FreshWrite, &[])
1308 .is_some(),
1309 "fresh steward LLM output (consolidate/harvest/rank) respects the posture"
1310 );
1311 assert!(
1316 strict
1317 .quarantine_reason(&steward, StagedBatchKind::ReviewVerdict, &[])
1318 .is_none()
1319 );
1320 assert!(
1321 strict
1322 .quarantine_reason(&MemoryAuthor::Operator, StagedBatchKind::FreshWrite, &[])
1323 .is_none()
1324 );
1325 }
1326
1327 #[test]
1328 fn quarantined_posture_still_gates_review_verdicts_on_tainted_evidence() {
1329 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1333 let session = SessionId::new();
1334 tracker.observe_agent_event("identity:a", &run_started(&session));
1335 tracker.observe_agent_event("identity:a", &tool_result("web_search"));
1336
1337 let gate = TaintLlmWriteGate::new(Some(tracker), AgentMemoryLlmWrites::Quarantined);
1338 let steward = MemoryAuthor::Steward {
1339 run_id: "run-1".to_string(),
1340 };
1341 assert!(
1342 gate.quarantine_reason(&steward, StagedBatchKind::ReviewVerdict, &[])
1343 .is_none()
1344 );
1345 let reason = gate
1346 .quarantine_reason(
1347 &steward,
1348 StagedBatchKind::ReviewVerdict,
1349 &evidence_for(&session),
1350 )
1351 .expect("tainted evidence still quarantines review verdicts");
1352 assert!(reason.contains("evidence session tainted"), "{reason}");
1353 }
1354
1355 fn evidence_for(session: &SessionId) -> Vec<EvidenceRef> {
1356 vec![EvidenceRef {
1357 session_id: session.to_string(),
1358 generation: 0,
1359 revision: None,
1360 range: Some((0, 4)),
1361 }]
1362 }
1363
1364 #[test]
1365 fn gate_quarantines_llm_writes_citing_tainted_evidence() {
1366 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1367 let session = SessionId::new();
1368 tracker.observe_agent_event("identity:a", &run_started(&session));
1369 tracker.observe_agent_event("identity:a", &tool_result("web_search"));
1370 let fresh = SessionId::new();
1373 tracker.observe_agent_event("identity:a", &run_started(&fresh));
1374
1375 let gate = TaintLlmWriteGate::new(Some(tracker), AgentMemoryLlmWrites::Observed);
1376 let distiller = MemoryAuthor::Distiller {
1377 run_id: "run-1".to_string(),
1378 };
1379 let reason = gate
1380 .quarantine_reason(
1381 &distiller,
1382 StagedBatchKind::FreshWrite,
1383 &evidence_for(&session),
1384 )
1385 .expect("tainted evidence range quarantines (session-tainted ⇒ range-tainted)");
1386 assert!(reason.contains("evidence session tainted"), "{reason}");
1387 assert!(
1389 gate.quarantine_reason(
1390 &distiller,
1391 StagedBatchKind::FreshWrite,
1392 &evidence_for(&fresh)
1393 )
1394 .is_none()
1395 );
1396 assert!(
1398 gate.quarantine_reason(
1399 &MemoryAuthor::Operator,
1400 StagedBatchKind::FreshWrite,
1401 &evidence_for(&session)
1402 )
1403 .is_none()
1404 );
1405 }
1406
1407 #[test]
1408 fn reset_boundary_quarantines_evidence_without_content_taint() {
1409 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1410 let session = SessionId::new();
1411 tracker.observe_agent_event("identity:a", &run_started(&session));
1412 assert!(
1413 tracker
1414 .evidence_quarantine_reason(&session.to_string())
1415 .is_none()
1416 );
1417 tracker.mark_reset_boundary(&session.to_string());
1418 let gate = TaintLlmWriteGate::new(Some(tracker), AgentMemoryLlmWrites::Observed);
1419 let reason = gate
1420 .quarantine_reason(
1421 &MemoryAuthor::Distiller {
1422 run_id: "run-1".to_string(),
1423 },
1424 StagedBatchKind::FreshWrite,
1425 &evidence_for(&session),
1426 )
1427 .expect("reset boundary quarantines distillates");
1428 assert!(reason.contains("reset boundary"), "{reason}");
1429 }
1430
1431 #[test]
1432 fn peer_projection_sender_parses_message_and_response_shapes() {
1433 assert_eq!(
1434 peer_projection_sender_identity("Peer message from mob-1/worker/identity:bob:"),
1435 Some("identity:bob")
1436 );
1437 assert_eq!(
1438 peer_projection_sender_identity(
1439 "Peer response from mob-1/worker/identity:bob (to request: req-9)"
1440 ),
1441 Some("identity:bob")
1442 );
1443 assert_eq!(
1445 peer_projection_sender_identity("Peer message from scout:"),
1446 Some("scout")
1447 );
1448 assert_eq!(
1450 peer_projection_sender_identity("Peer request from peer_id 018fabc (id: r-1)"),
1451 None
1452 );
1453 assert_eq!(peer_projection_sender_identity("ordinary text"), None);
1454 }
1455
1456 #[test]
1457 fn comms_join_taints_receiver_of_message_from_tainted_sender() {
1458 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1459 let sender_session = SessionId::new();
1461 tracker.observe_agent_event("identity:bob", &run_started(&sender_session));
1462 tracker.observe_agent_event("identity:bob", &tool_result("web_search"));
1463
1464 let receiver_session = SessionId::new();
1467 let delivery = AgentEvent::RunStarted {
1468 session_id: receiver_session.clone(),
1469 input: meerkat_core::types::RunInput::Content {
1470 content: meerkat_core::ContentInput::Text(
1471 "Peer message from mob-1/worker/identity:bob:\nplease remember X".to_string(),
1472 ),
1473 },
1474 };
1475 tracker.observe_agent_event("identity:alice", &delivery);
1476 let taint = tracker
1477 .identity_taint("identity:alice")
1478 .expect("receiver session taints (peer-laundering close, §10.1)");
1479 assert!(taint.source.contains("identity:bob"), "{}", taint.source);
1480 assert!(
1481 tracker
1482 .session_taint(&receiver_session.to_string())
1483 .is_some()
1484 );
1485
1486 let clean_session = SessionId::new();
1488 tracker.observe_agent_event("identity:carol", &run_started(&clean_session));
1489 let receiver2 = SessionId::new();
1490 let clean_delivery = AgentEvent::RunStarted {
1491 session_id: receiver2.clone(),
1492 input: meerkat_core::types::RunInput::Content {
1493 content: meerkat_core::ContentInput::Text(
1494 "Peer message from mob-1/worker/identity:carol:\nhello".to_string(),
1495 ),
1496 },
1497 };
1498 tracker.observe_agent_event("identity:dave", &clean_delivery);
1499 assert!(tracker.identity_taint("identity:dave").is_none());
1500 }
1501}