1use std::collections::HashMap;
67use std::sync::{Arc, Mutex};
68use std::time::{SystemTime, UNIX_EPOCH};
69
70use meerkat_core::event::AgentEvent;
71use meerkat_core::types::{ServerToolKind, ToolProvenance, ToolSourceKind};
72use serde::{Deserialize, Serialize};
73
74use crate::identity_first::agent_memory::AgentMemoryLlmWrites;
75use crate::memory::records::{EvidenceRef, MemoryAuthor};
76use crate::memory::staged::StagedBatchKind;
77
78const ALWAYS_UNTRUSTED_TOOL_NAMES: &[&str] = &["web_search", "web_fetch", "fetch", "http_request"];
84
85const MCP_QUALIFIED_PREFIX: &str = "mcp__";
89
90const MAX_TRACKED_TAINTED_SESSIONS: usize = 4096;
94
95#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
103pub struct ContentTrustConfig {
104 #[serde(default)]
107 pub trusted_mcp_servers: Vec<String>,
108 #[serde(default)]
111 pub untrusted_tools: Vec<String>,
112 #[serde(default)]
115 pub trusted_tools: Vec<String>,
116}
117
118#[derive(Debug, Clone, PartialEq, Eq)]
120pub enum ToolContentTrust {
121 Trusted,
122 Untrusted { source: String },
123}
124
125impl ContentTrustConfig {
126 pub fn from_json_value(value: &serde_json::Value) -> Result<Self, String> {
130 let object = value
131 .as_object()
132 .ok_or_else(|| "content_trust must be an object".to_string())?;
133 let supported = ["trusted_mcp_servers", "untrusted_tools", "trusted_tools"];
134 let unsupported = object
135 .keys()
136 .filter(|key| !supported.contains(&key.as_str()))
137 .map(String::as_str)
138 .collect::<Vec<_>>();
139 if !unsupported.is_empty() {
140 return Err(format!(
141 "unsupported content_trust fields: {}",
142 unsupported.join(", ")
143 ));
144 }
145 let parse_names = |key: &str| -> Result<Vec<String>, String> {
146 match object.get(key) {
147 None => Ok(Vec::new()),
148 Some(value) => {
149 let entries = value
150 .as_array()
151 .ok_or_else(|| format!("content_trust.{key} must be an array"))?;
152 entries
153 .iter()
154 .map(|entry| {
155 entry
156 .as_str()
157 .map(str::trim)
158 .filter(|name| !name.is_empty())
159 .map(ToString::to_string)
160 .ok_or_else(|| {
161 format!("content_trust.{key} entries must be non-empty strings")
162 })
163 })
164 .collect()
165 }
166 }
167 };
168 Ok(Self {
169 trusted_mcp_servers: parse_names("trusted_mcp_servers")?,
170 untrusted_tools: parse_names("untrusted_tools")?,
171 trusted_tools: parse_names("trusted_tools")?,
172 })
173 }
174
175 pub fn classify_tool(&self, name: &str) -> ToolContentTrust {
180 self.classify_tool_with_provenance(name, None)
181 }
182
183 pub fn classify_tool_with_provenance(
194 &self,
195 name: &str,
196 provenance: Option<&ToolProvenance>,
197 ) -> ToolContentTrust {
198 if ALWAYS_UNTRUSTED_TOOL_NAMES.contains(&name) {
199 return ToolContentTrust::Untrusted {
200 source: format!("web tool '{name}'"),
201 };
202 }
203 if self.untrusted_tools.iter().any(|tool| tool == name) {
204 return ToolContentTrust::Untrusted {
205 source: format!("configured untrusted tool '{name}'"),
206 };
207 }
208 if self.trusted_tools.iter().any(|tool| tool == name) {
209 return ToolContentTrust::Trusted;
210 }
211 if let Some(provenance) = provenance
212 && provenance.kind == ToolSourceKind::Mcp
213 {
214 return self.classify_mcp_server(provenance.source_id.as_str(), name);
215 }
216 if let Some(rest) = name.strip_prefix(MCP_QUALIFIED_PREFIX) {
217 let server = rest.split("__").next().unwrap_or(rest);
218 return self.classify_mcp_server(server, name);
219 }
220 ToolContentTrust::Trusted
221 }
222
223 fn classify_mcp_server(&self, server: &str, name: &str) -> ToolContentTrust {
224 if self
225 .trusted_mcp_servers
226 .iter()
227 .any(|trusted| trusted == server)
228 {
229 return ToolContentTrust::Trusted;
230 }
231 ToolContentTrust::Untrusted {
232 source: format!("MCP server '{server}' (tool '{name}')"),
233 }
234 }
235}
236
237#[derive(Debug, Clone, PartialEq, Eq)]
239pub struct TaintState {
240 pub tainted_at_ms: u64,
241 pub source: String,
242}
243
244#[derive(Default)]
245struct TaintInner {
246 current_session: HashMap<String, String>,
251 tainted: HashMap<String, TaintState>,
256 pending_identity_taint: HashMap<String, TaintState>,
260 reset_boundaries: HashMap<String, u64>,
263}
264
265pub type OutboundTaintDeclarer =
273 Arc<dyn Fn(&str, Option<meerkat_core::comms::SenderContentTaint>) + Send + Sync>;
274
275#[derive(Clone, Default)]
278pub struct SessionTaintTracker {
279 config: Arc<ContentTrustConfig>,
280 inner: Arc<Mutex<TaintInner>>,
281 event_sink: Arc<Mutex<Option<Arc<dyn crate::memory::events::MemoryEventSink>>>>,
283 outbound_declarer: Arc<Mutex<Option<OutboundTaintDeclarer>>>,
287}
288
289impl SessionTaintTracker {
290 pub fn new(config: ContentTrustConfig) -> Self {
291 Self {
292 config: Arc::new(config),
293 inner: Arc::new(Mutex::new(TaintInner::default())),
294 event_sink: Arc::new(Mutex::new(None)),
295 outbound_declarer: Arc::new(Mutex::new(None)),
296 }
297 }
298
299 pub fn set_outbound_taint_declarer(&self, declarer: OutboundTaintDeclarer) {
302 *self
303 .outbound_declarer
304 .lock()
305 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(declarer);
306 }
307
308 fn declare_outbound(
312 &self,
313 identity: &str,
314 taint: Option<meerkat_core::comms::SenderContentTaint>,
315 ) {
316 if let Some(declarer) = self
317 .outbound_declarer
318 .lock()
319 .unwrap_or_else(std::sync::PoisonError::into_inner)
320 .as_ref()
321 {
322 declarer(identity, taint);
323 }
324 }
325
326 pub fn set_event_sink(&self, sink: Arc<dyn crate::memory::events::MemoryEventSink>) {
329 *self
330 .event_sink
331 .lock()
332 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(sink);
333 }
334
335 fn emit_event(&self, event: crate::memory::events::MemoryTimelineEvent) {
336 if let Some(sink) = self
337 .event_sink
338 .lock()
339 .unwrap_or_else(std::sync::PoisonError::into_inner)
340 .as_ref()
341 {
342 sink.emit(event);
343 }
344 }
345
346 pub fn observe_agent_event(&self, identity: &str, event: &AgentEvent) {
349 match event {
350 AgentEvent::RunStarted { session_id, input } => {
351 let session_key = session_id.to_string();
352 self.note_current_session(identity, &session_key);
353 if let Some(text) = input.prompt_text() {
358 self.observe_inbound_peer_content(identity, &session_key, &text);
359 }
360 }
361 AgentEvent::PeerContentIngested {
368 peer, sender_taint, ..
369 } => {
370 if *sender_taint == Some(meerkat_core::comms::SenderContentTaint::Tainted) {
371 let sender = peer
372 .as_ref()
373 .and_then(|peer| peer.display_name.clone())
374 .unwrap_or_else(|| "peer".to_string());
375 self.mark_identity_tainted(
376 identity,
377 format!("peer content declared tainted by sender '{sender}'"),
378 );
379 }
380 }
381 AgentEvent::ToolResultReceived { name, .. }
385 | AgentEvent::ToolExecutionCompleted { name, .. } => {
386 if let ToolContentTrust::Untrusted { source } = self.config.classify_tool(name) {
387 self.mark_identity_tainted(identity, source);
388 }
389 }
390 AgentEvent::ServerToolContent { kind, .. } => {
393 self.mark_identity_tainted(
394 identity,
395 format!("provider server tool '{}'", kind.provider_name()),
396 );
397 }
398 _ => {}
399 }
400 }
401
402 pub fn observe_dispatched_tool_result(
409 &self,
410 identity: &str,
411 name: &str,
412 provenance: Option<&ToolProvenance>,
413 ) {
414 if let ToolContentTrust::Untrusted { source } =
415 self.config.classify_tool_with_provenance(name, provenance)
416 {
417 self.mark_identity_tainted(identity, source);
418 }
419 }
420
421 pub fn observe_dispatched_server_tool(&self, identity: &str, kind: &ServerToolKind) {
426 self.mark_identity_tainted(
427 identity,
428 format!("provider server tool '{}'", kind.provider_name()),
429 );
430 }
431
432 pub fn note_current_session(&self, identity: &str, session_key: &str) {
436 let mut inner = self
437 .inner
438 .lock()
439 .unwrap_or_else(std::sync::PoisonError::into_inner);
440 let pending = inner.pending_identity_taint.remove(identity);
441 let previous = inner
442 .current_session
443 .insert(identity.to_string(), session_key.to_string());
444 let rotated_away_from_tainted = previous.as_deref() != Some(session_key)
445 && previous.is_some_and(|prior| inner.tainted.contains_key(&prior));
446 if rotated_away_from_tainted {
447 tracing::warn!(
448 identity,
449 session_key,
450 "agent memory taint: session rotated away from a tainted session; \
451 new session starts clean"
452 );
453 self.emit_event(
454 crate::memory::events::MemoryTimelineEvent::TaintTransition {
455 identity: Some(identity.to_string()),
456 session_key: session_key.to_string(),
457 kind: "rotated_clean".to_string(),
458 source: "session rotation away from tainted session".to_string(),
459 },
460 );
461 }
462 let pending_reapplied = pending.is_some();
465 if let Some(state) = pending {
466 self.insert_taint(&mut inner, session_key.to_string(), state);
467 }
468 if rotated_away_from_tainted && !pending_reapplied {
471 drop(inner);
472 self.declare_outbound(identity, None);
473 }
474 }
475
476 pub fn clear_identity(&self, identity: &str) {
484 let mut inner = self
485 .inner
486 .lock()
487 .unwrap_or_else(std::sync::PoisonError::into_inner);
488 inner.pending_identity_taint.remove(identity);
489 inner.current_session.remove(identity);
490 drop(inner);
493 self.declare_outbound(identity, None);
494 }
495
496 pub fn mark_reset_boundary(&self, session_key: &str) {
500 let mut inner = self
501 .inner
502 .lock()
503 .unwrap_or_else(std::sync::PoisonError::into_inner);
504 if inner
505 .reset_boundaries
506 .insert(session_key.to_string(), now_ms())
507 .is_none()
508 {
509 tracing::warn!(
510 session_key,
511 "agent memory taint: reset boundary marked; distillates over this \
512 session will land quarantined pending steward review (§8.4)"
513 );
514 self.emit_event(
515 crate::memory::events::MemoryTimelineEvent::TaintTransition {
516 identity: None,
517 session_key: session_key.to_string(),
518 kind: "reset_boundary".to_string(),
519 source: "reset() boundary (§8.4)".to_string(),
520 },
521 );
522 }
523 if inner.reset_boundaries.len() > MAX_TRACKED_TAINTED_SESSIONS
524 && let Some(oldest) = inner
525 .reset_boundaries
526 .iter()
527 .min_by_key(|(_, at_ms)| **at_ms)
528 .map(|(key, _)| key.clone())
529 {
530 inner.reset_boundaries.remove(&oldest);
531 }
532 }
533
534 pub fn evidence_quarantine_reason(&self, session_key: &str) -> Option<String> {
539 let inner = self
540 .inner
541 .lock()
542 .unwrap_or_else(std::sync::PoisonError::into_inner);
543 if let Some(state) = inner.tainted.get(session_key) {
544 return Some(format!(
545 "evidence session tainted by {} (session-tainted ⇒ range-tainted)",
546 state.source
547 ));
548 }
549 if inner.reset_boundaries.contains_key(session_key) {
550 return Some(
551 "evidence session closed at a reset boundary; distillates quarantine \
552 pending steward review (§8.4)"
553 .to_string(),
554 );
555 }
556 None
557 }
558
559 pub fn session_taint(&self, session_key: &str) -> Option<TaintState> {
561 let inner = self
562 .inner
563 .lock()
564 .unwrap_or_else(std::sync::PoisonError::into_inner);
565 inner.tainted.get(session_key).cloned()
566 }
567
568 pub fn identity_taint(&self, identity: &str) -> Option<TaintState> {
571 let inner = self
572 .inner
573 .lock()
574 .unwrap_or_else(std::sync::PoisonError::into_inner);
575 if let Some(state) = inner.pending_identity_taint.get(identity) {
576 return Some(state.clone());
577 }
578 inner
579 .current_session
580 .get(identity)
581 .and_then(|session| inner.tainted.get(session))
582 .cloned()
583 }
584
585 fn observe_inbound_peer_content(&self, identity: &str, session_key: &str, text: &str) {
610 for line in text.lines() {
611 let Some(sender_identity) = peer_projection_sender_identity(line) else {
612 continue;
613 };
614 if sender_identity == identity {
615 continue;
616 }
617 let source = {
618 let inner = self
619 .inner
620 .lock()
621 .unwrap_or_else(std::sync::PoisonError::into_inner);
622 inner
623 .pending_identity_taint
624 .get(sender_identity)
625 .or_else(|| {
626 inner
627 .current_session
628 .get(sender_identity)
629 .and_then(|session| inner.tainted.get(session))
630 })
631 .map(|state| state.source.clone())
632 };
633 if let Some(source) = source {
634 let state = TaintState {
635 tainted_at_ms: now_ms(),
636 source: format!(
637 "peer message from tainted sender '{sender_identity}' \
638 (sender session tainted by {source})"
639 ),
640 };
641 let mut inner = self
642 .inner
643 .lock()
644 .unwrap_or_else(std::sync::PoisonError::into_inner);
645 self.insert_taint(&mut inner, session_key.to_string(), state);
646 }
647 }
648 }
649
650 fn mark_identity_tainted(&self, identity: &str, source: String) {
651 let state = TaintState {
652 tainted_at_ms: now_ms(),
653 source,
654 };
655 let mut inner = self
656 .inner
657 .lock()
658 .unwrap_or_else(std::sync::PoisonError::into_inner);
659 match inner.current_session.get(identity).cloned() {
660 Some(session) => self.insert_taint(&mut inner, session, state),
661 None => {
662 match inner.pending_identity_taint.entry(identity.to_string()) {
665 std::collections::hash_map::Entry::Vacant(slot) => {
666 tracing::warn!(
667 identity,
668 source = %state.source,
669 "agent memory taint: untrusted ingestion observed before session \
670 attribution; holding identity-sticky taint"
671 );
672 slot.insert(state);
673 }
674 std::collections::hash_map::Entry::Occupied(mut slot) => {
675 slot.insert(state);
676 }
677 }
678 }
679 }
680 drop(inner);
684 self.declare_outbound(
685 identity,
686 Some(meerkat_core::comms::SenderContentTaint::Tainted),
687 );
688 }
689
690 fn insert_taint(&self, inner: &mut TaintInner, session: String, state: TaintState) {
691 match inner.tainted.entry(session) {
692 std::collections::hash_map::Entry::Occupied(_) => return,
693 std::collections::hash_map::Entry::Vacant(slot) => {
694 tracing::warn!(
695 session_key = %slot.key(),
696 source = %state.source,
697 "agent memory taint: session ingested untrusted content; LLM-authored \
698 memory writes from this session will quarantine until a fresh-context \
699 boundary (reset/respawn/fresh spawn)"
700 );
701 self.emit_event(
702 crate::memory::events::MemoryTimelineEvent::TaintTransition {
703 identity: None,
704 session_key: slot.key().clone(),
705 kind: "tainted".to_string(),
706 source: state.source.clone(),
707 },
708 );
709 slot.insert(state);
710 }
711 }
712 if inner.tainted.len() > MAX_TRACKED_TAINTED_SESSIONS
713 && let Some(oldest) = inner
714 .tainted
715 .iter()
716 .min_by_key(|(_, state)| state.tainted_at_ms)
717 .map(|(key, _)| key.clone())
718 {
719 inner.tainted.remove(&oldest);
720 }
721 }
722}
723
724pub trait LlmWriteGate: Send + Sync {
729 fn quarantine_reason(
736 &self,
737 author: &MemoryAuthor,
738 kind: StagedBatchKind,
739 evidence: &[EvidenceRef],
740 ) -> Option<String>;
741}
742
743pub struct TaintLlmWriteGate {
764 tracker: Option<SessionTaintTracker>,
765 llm_writes: AgentMemoryLlmWrites,
766}
767
768impl TaintLlmWriteGate {
769 pub fn new(tracker: Option<SessionTaintTracker>, llm_writes: AgentMemoryLlmWrites) -> Self {
770 Self {
771 tracker,
772 llm_writes,
773 }
774 }
775}
776
777impl LlmWriteGate for TaintLlmWriteGate {
778 fn quarantine_reason(
779 &self,
780 author: &MemoryAuthor,
781 kind: StagedBatchKind,
782 evidence: &[EvidenceRef],
783 ) -> Option<String> {
784 if !author.is_llm() {
785 return None;
786 }
787 if self.llm_writes == AgentMemoryLlmWrites::Quarantined
788 && kind != StagedBatchKind::ReviewVerdict
789 {
790 return Some("llm_writes=quarantined policy".to_string());
791 }
792 let tracker = self.tracker.as_ref()?;
793 if let MemoryAuthor::Agent { identity } = author
794 && let Some(state) = tracker.identity_taint(identity)
795 {
796 return Some(format!("session tainted by {}", state.source));
797 }
798 for evidence_ref in evidence {
799 if let Some(reason) = tracker.evidence_quarantine_reason(&evidence_ref.session_id) {
800 return Some(reason);
801 }
802 }
803 None
804 }
805}
806
807fn peer_projection_sender_identity(line: &str) -> Option<&str> {
816 let name = if let Some(rest) = line.strip_prefix("Peer message from ") {
817 match rest.split_once(": ") {
821 Some((name, _)) => name.trim(),
822 None => rest.strip_suffix(':').unwrap_or(rest).trim(),
823 }
824 } else {
825 let rest = line.strip_prefix("Peer response from ")?;
826 rest.split(" (to request:").next()?.trim()
827 };
828 if name.is_empty() {
829 return None;
830 }
831 Some(name.rsplit('/').next().unwrap_or(name))
832}
833
834pub trait MemberAgentEventSink: Send + Sync {
843 fn observe(&self, identity: &str, envelope: &meerkat_core::event::EventEnvelope<AgentEvent>);
844}
845
846impl MemberAgentEventSink for SessionTaintTracker {
847 fn observe(&self, identity: &str, envelope: &meerkat_core::event::EventEnvelope<AgentEvent>) {
848 self.observe_agent_event(identity, &envelope.payload);
849 }
850}
851
852pub struct CompactionResetSink {
859 on_compacted: Arc<dyn Fn(&str) + Send + Sync>,
860}
861
862impl CompactionResetSink {
863 pub fn new(on_compacted: Arc<dyn Fn(&str) + Send + Sync>) -> Self {
864 Self { on_compacted }
865 }
866}
867
868impl MemberAgentEventSink for CompactionResetSink {
869 fn observe(&self, _identity: &str, envelope: &meerkat_core::event::EventEnvelope<AgentEvent>) {
870 if matches!(envelope.payload, AgentEvent::CompactionCompleted { .. })
871 && let meerkat_core::event::EventSourceIdentity::Session { session_id } =
872 &envelope.source
873 {
874 (self.on_compacted)(&session_id.to_string());
875 }
876 }
877}
878
879#[derive(Clone)]
882pub struct TaintObserverGuard {
883 _abort: Arc<AbortOnDrop>,
884}
885
886struct AbortOnDrop(std::sync::Mutex<Option<tokio::task::JoinHandle<()>>>);
887
888impl Drop for AbortOnDrop {
889 fn drop(&mut self) {
890 if let Some(task) = self
891 .0
892 .get_mut()
893 .unwrap_or_else(std::sync::PoisonError::into_inner)
894 .take()
895 {
896 task.abort();
897 }
898 }
899}
900
901impl TaintObserverGuard {
902 pub async fn abort_and_join(self) {
906 let task = self
907 ._abort
908 .0
909 .lock()
910 .unwrap_or_else(std::sync::PoisonError::into_inner)
911 .take();
912 if let Some(task) = task {
913 task.abort();
914 let _ = task.await;
915 }
916 }
917}
918
919pub fn spawn_taint_observer(
923 handle: meerkat_mob::MobHandle,
924 tracker: SessionTaintTracker,
925) -> TaintObserverGuard {
926 spawn_member_event_observer(handle, vec![Arc::new(tracker)])
927}
928
929pub fn spawn_member_event_observer(
933 handle: meerkat_mob::MobHandle,
934 sinks: Vec<Arc<dyn MemberAgentEventSink>>,
935) -> TaintObserverGuard {
936 let task = tokio::spawn(run_member_event_observer(handle, sinks));
937 TaintObserverGuard {
938 _abort: Arc::new(AbortOnDrop(std::sync::Mutex::new(Some(task)))),
939 }
940}
941
942async fn run_member_event_observer(
943 handle: meerkat_mob::MobHandle,
944 sinks: Vec<Arc<dyn MemberAgentEventSink>>,
945) {
946 use futures::StreamExt;
947 use futures::stream::SelectAll;
948
949 enum Observed {
950 Event(String, Box<meerkat_core::event::EventEnvelope<AgentEvent>>),
951 Closed(String),
952 }
953
954 let mut streams: SelectAll<futures::stream::BoxStream<'static, Observed>> = SelectAll::new();
955 let mut subscribed: std::collections::HashSet<String> = std::collections::HashSet::new();
956 let mut warned: std::collections::HashSet<String> = std::collections::HashSet::new();
957 let mut reconcile = tokio::time::interval(std::time::Duration::from_secs(1));
958 reconcile.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
959
960 loop {
961 tokio::select! {
962 Some(observed) = streams.next() => match observed {
963 Observed::Event(identity, envelope) => {
964 for sink in &sinks {
965 sink.observe(&identity, &envelope);
966 }
967 }
968 Observed::Closed(identity) => {
969 subscribed.remove(&identity);
970 }
971 },
972 _ = reconcile.tick() => {
973 for entry in handle.list_members_including_retiring().await {
974 if entry.status != meerkat_mob::MobMemberStatus::Active {
978 continue;
979 }
980 let identity = entry.agent_identity.to_string();
981 if subscribed.contains(&identity) {
982 continue;
983 }
984 let sink_identity =
996 crate::member_comms_id::logical_memory_identity(&identity);
997 match handle.subscribe_agent_events(&entry.agent_identity).await {
998 Ok(stream) => {
999 warned.remove(&identity);
1000 subscribed.insert(identity.clone());
1001 let close_key = identity.clone();
1002 streams.push(
1003 stream
1004 .map(move |envelope| {
1005 Observed::Event(
1006 sink_identity.clone(),
1007 Box::new(envelope),
1008 )
1009 })
1010 .chain(futures::stream::once(async move {
1011 Observed::Closed(close_key)
1012 }))
1013 .boxed(),
1014 );
1015 }
1016 Err(error) => {
1017 if warned.insert(identity.clone()) {
1020 tracing::warn!(
1021 identity = %identity,
1022 error = %error,
1023 "agent memory taint observer: failed to subscribe; will retry"
1024 );
1025 } else {
1026 tracing::debug!(
1027 identity = %identity,
1028 error = %error,
1029 "agent memory taint observer: subscribe still failing"
1030 );
1031 }
1032 }
1033 }
1034 }
1035 }
1036 }
1037 }
1038}
1039
1040fn now_ms() -> u64 {
1041 SystemTime::now()
1042 .duration_since(UNIX_EPOCH)
1043 .map(|duration| duration.as_millis() as u64)
1044 .unwrap_or(0)
1045}
1046
1047#[cfg(test)]
1048#[allow(
1049 clippy::expect_used,
1050 clippy::panic,
1051 clippy::redundant_clone,
1052 clippy::unwrap_used
1053)]
1054mod tests {
1055 use super::*;
1056 use meerkat_core::types::{ContentBlock, ServerToolKind, SessionId};
1057 use serde_json::json;
1058
1059 fn run_started(session: &SessionId) -> AgentEvent {
1060 AgentEvent::RunStarted {
1061 session_id: session.clone(),
1062 input: meerkat_core::types::RunInput::Content {
1063 content: meerkat_core::ContentInput::Text("hi".to_string()),
1064 },
1065 }
1066 }
1067
1068 fn tool_result(name: &str) -> AgentEvent {
1069 AgentEvent::ToolResultReceived {
1070 id: "tool-1".to_string(),
1071 name: name.to_string(),
1072 content: vec![ContentBlock::Text {
1073 text: "ok".to_string(),
1074 }],
1075 is_error: false,
1076 }
1077 }
1078
1079 fn peer_ingested(taint: Option<meerkat_core::comms::SenderContentTaint>) -> AgentEvent {
1080 AgentEvent::PeerContentIngested {
1081 kind: meerkat_core::types::CommsNoticeKind::Message,
1082 peer: None,
1083 request_id: None,
1084 sender_taint: taint,
1085 }
1086 }
1087
1088 #[test]
1092 fn declared_peer_taint_taints_receiver_but_none_or_clean_does_not() {
1093 use meerkat_core::comms::SenderContentTaint;
1094 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1095 tracker.note_current_session("identity:b", "sess-b");
1096
1097 tracker.observe_agent_event("identity:b", &peer_ingested(None));
1098 assert!(
1099 tracker.session_taint("sess-b").is_none(),
1100 "no declaration must not taint"
1101 );
1102
1103 tracker.observe_agent_event(
1104 "identity:b",
1105 &peer_ingested(Some(SenderContentTaint::Clean)),
1106 );
1107 assert!(
1108 tracker.session_taint("sess-b").is_none(),
1109 "an affirmative Clean declaration must not taint"
1110 );
1111
1112 tracker.observe_agent_event(
1113 "identity:b",
1114 &peer_ingested(Some(SenderContentTaint::Tainted)),
1115 );
1116 assert!(
1117 tracker.session_taint("sess-b").is_some(),
1118 "a declared-tainted peer delivery taints the receiving session"
1119 );
1120 }
1121
1122 #[test]
1123 fn taint_transitions_emit_timeline_events_when_sink_wired() {
1124 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1125 let sink = std::sync::Arc::new(crate::memory::events::CollectingEventSink::new());
1126 tracker.set_event_sink(sink.clone());
1127
1128 tracker.note_current_session("identity:a", "sess-1");
1129 tracker.observe_agent_event("identity:a", &tool_result("web_fetch"));
1130 tracker.mark_reset_boundary("sess-1");
1131 tracker.mark_reset_boundary("sess-1");
1133 tracker.note_current_session("identity:a", "sess-2");
1134
1135 let types = sink.types();
1136 assert_eq!(
1137 types,
1138 vec![
1139 "memory.taint.transition", "memory.taint.transition", "memory.taint.transition", ]
1143 );
1144 let events = sink.events.lock().unwrap();
1145 let kinds: Vec<String> = events
1146 .iter()
1147 .map(|event| match event {
1148 crate::memory::events::MemoryTimelineEvent::TaintTransition { kind, .. } => {
1149 kind.clone()
1150 }
1151 other => panic!("unexpected event {other:?}"),
1152 })
1153 .collect();
1154 assert_eq!(kinds, vec!["tainted", "reset_boundary", "rotated_clean"]);
1155 }
1156
1157 #[test]
1158 fn content_trust_parse_rejects_unknown_fields_and_bad_types() {
1159 let err = ContentTrustConfig::from_json_value(&json!({"servers": []}))
1160 .expect_err("unknown field must fail loud");
1161 assert!(err.contains("unsupported content_trust fields"), "{err}");
1162 let err = ContentTrustConfig::from_json_value(&json!({"trusted_mcp_servers": "kg"}))
1163 .expect_err("non-array must fail loud");
1164 assert!(err.contains("must be an array"), "{err}");
1165 let err = ContentTrustConfig::from_json_value(&json!({"untrusted_tools": [1]}))
1166 .expect_err("non-string entry must fail loud");
1167 assert!(err.contains("non-empty strings"), "{err}");
1168 let err =
1169 ContentTrustConfig::from_json_value(&json!([])).expect_err("non-object must fail loud");
1170 assert!(err.contains("must be an object"), "{err}");
1171 }
1172
1173 #[test]
1174 fn content_trust_parse_accepts_full_block() {
1175 let config = ContentTrustConfig::from_json_value(&json!({
1176 "trusted_mcp_servers": ["knowledge_graph"],
1177 "untrusted_tools": ["scrape_page"],
1178 "trusted_tools": ["mcp__scanner__lint"],
1179 }))
1180 .expect("valid block parses");
1181 assert_eq!(config.trusted_mcp_servers, vec!["knowledge_graph"]);
1182 assert_eq!(config.untrusted_tools, vec!["scrape_page"]);
1183 assert_eq!(config.trusted_tools, vec!["mcp__scanner__lint"]);
1184 }
1185
1186 #[test]
1187 fn classification_precedence_holds() {
1188 let config = ContentTrustConfig {
1189 trusted_mcp_servers: vec!["kg".to_string()],
1190 untrusted_tools: vec!["scrape_page".to_string()],
1191 trusted_tools: vec!["web_search".to_string(), "mcp__evil__probe".to_string()],
1193 };
1194 assert!(matches!(
1195 config.classify_tool("web_search"),
1196 ToolContentTrust::Untrusted { .. }
1197 ));
1198 assert!(matches!(
1199 config.classify_tool("scrape_page"),
1200 ToolContentTrust::Untrusted { .. }
1201 ));
1202 assert_eq!(
1204 config.classify_tool("mcp__evil__probe"),
1205 ToolContentTrust::Trusted
1206 );
1207 assert!(matches!(
1209 config.classify_tool("mcp__other__search"),
1210 ToolContentTrust::Untrusted { .. }
1211 ));
1212 assert_eq!(
1213 config.classify_tool("mcp__kg__query"),
1214 ToolContentTrust::Trusted
1215 );
1216 assert_eq!(config.classify_tool("shell"), ToolContentTrust::Trusted);
1218 }
1219
1220 #[test]
1224 fn typed_mcp_provenance_attributes_unqualified_tool_names() {
1225 use meerkat_core::types::{ToolProvenance, ToolSourceKind};
1226 let config = ContentTrustConfig {
1227 trusted_mcp_servers: vec!["kg".to_string()],
1228 ..ContentTrustConfig::default()
1229 };
1230 let kg = ToolProvenance {
1231 kind: ToolSourceKind::Mcp,
1232 source_id: "kg".into(),
1233 };
1234 let scraper = ToolProvenance {
1235 kind: ToolSourceKind::Mcp,
1236 source_id: "scraper".into(),
1237 };
1238 let verdict = config.classify_tool_with_provenance("scrape_page", Some(&scraper));
1240 match verdict {
1241 ToolContentTrust::Untrusted { source } => {
1242 assert!(source.contains("MCP server 'scraper'"), "{source}");
1243 assert!(source.contains("scrape_page"), "{source}");
1244 }
1245 ToolContentTrust::Trusted => panic!("untrusted-server MCP tool must taint"),
1246 }
1247 assert_eq!(
1249 config.classify_tool_with_provenance("query", Some(&kg)),
1250 ToolContentTrust::Trusted
1251 );
1252 assert_eq!(
1255 config.classify_tool_with_provenance("mcp__evil__query", Some(&kg)),
1256 ToolContentTrust::Trusted
1257 );
1258 assert!(matches!(
1261 config.classify_tool_with_provenance("web_search", Some(&kg)),
1262 ToolContentTrust::Untrusted { .. }
1263 ));
1264 let listed = ContentTrustConfig {
1265 trusted_tools: vec!["scrape_page".to_string()],
1266 ..ContentTrustConfig::default()
1267 };
1268 assert_eq!(
1269 listed.classify_tool_with_provenance("scrape_page", Some(&scraper)),
1270 ToolContentTrust::Trusted,
1271 "explicit trusted_tools overrides server-level distrust, as on the name path"
1272 );
1273 }
1274
1275 #[test]
1278 fn absent_or_non_mcp_provenance_falls_back_to_name_classification() {
1279 use meerkat_core::types::{ToolProvenance, ToolSourceKind};
1280 let config = ContentTrustConfig {
1281 trusted_mcp_servers: vec!["kg".to_string()],
1282 untrusted_tools: vec!["scrape_page".to_string()],
1283 trusted_tools: vec!["mcp__evil__probe".to_string()],
1284 };
1285 for name in [
1286 "web_search",
1287 "scrape_page",
1288 "mcp__evil__probe",
1289 "mcp__other__search",
1290 "mcp__kg__query",
1291 "shell",
1292 ] {
1293 assert_eq!(
1294 config.classify_tool_with_provenance(name, None),
1295 config.classify_tool(name),
1296 "provenance-absent classification must match the name path for '{name}'"
1297 );
1298 }
1299 let builtin = ToolProvenance {
1301 kind: ToolSourceKind::Builtin,
1302 source_id: "builtin".into(),
1303 };
1304 assert_eq!(
1305 config.classify_tool_with_provenance("shell", Some(&builtin)),
1306 ToolContentTrust::Trusted
1307 );
1308 assert_eq!(
1309 config.classify_tool_with_provenance("mcp__other__search", Some(&builtin)),
1310 config.classify_tool("mcp__other__search")
1311 );
1312 }
1313
1314 #[test]
1315 fn dispatched_tool_result_marks_identity_before_any_session_attribution() {
1316 use meerkat_core::types::{ToolProvenance, ToolSourceKind};
1317 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1318 let provenance = ToolProvenance {
1319 kind: ToolSourceKind::Mcp,
1320 source_id: "scraper".into(),
1321 };
1322 tracker.observe_dispatched_tool_result("identity:a", "scrape_page", Some(&provenance));
1325 let taint = tracker
1326 .identity_taint("identity:a")
1327 .expect("dispatch-time mark must be visible to the gate immediately");
1328 assert!(
1329 taint.source.contains("MCP server 'scraper'"),
1330 "{}",
1331 taint.source
1332 );
1333
1334 tracker.observe_dispatched_tool_result("identity:b", "shell", None);
1336 assert!(tracker.identity_taint("identity:b").is_none());
1337
1338 tracker.observe_dispatched_server_tool("identity:c", &ServerToolKind::WebSearch);
1340 let taint = tracker
1341 .identity_taint("identity:c")
1342 .expect("server tool marks");
1343 assert!(taint.source.contains("web_search"), "{}", taint.source);
1344 }
1345
1346 #[test]
1347 fn tracker_taints_on_untrusted_tool_and_clears_on_rotation() {
1348 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1349 let session = SessionId::new();
1350 tracker.observe_agent_event("identity:a", &run_started(&session));
1351 assert!(tracker.identity_taint("identity:a").is_none());
1352
1353 tracker.observe_agent_event("identity:a", &tool_result("shell"));
1354 assert!(tracker.identity_taint("identity:a").is_none());
1355
1356 tracker.observe_agent_event("identity:a", &tool_result("web_search"));
1357 let taint = tracker
1358 .identity_taint("identity:a")
1359 .expect("web tool result taints the session");
1360 assert!(taint.source.contains("web_search"), "{}", taint.source);
1361 assert!(tracker.session_taint(&session.to_string()).is_some());
1362
1363 tracker.observe_agent_event("identity:a", &tool_result("shell"));
1365 assert!(tracker.identity_taint("identity:a").is_some());
1366
1367 let fresh = SessionId::new();
1369 tracker.observe_agent_event("identity:a", &run_started(&fresh));
1370 assert!(tracker.identity_taint("identity:a").is_none());
1371 assert!(tracker.session_taint(&session.to_string()).is_some());
1373 }
1374
1375 #[test]
1376 fn tracker_taints_on_server_tool_content() {
1377 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1378 let session = SessionId::new();
1379 tracker.note_current_session("identity:a", &session.to_string());
1380 tracker.observe_agent_event(
1381 "identity:a",
1382 &AgentEvent::ServerToolContent {
1383 id: None,
1384 kind: ServerToolKind::WebSearch,
1385 content: json!({"results": []}),
1386 },
1387 );
1388 let taint = tracker.identity_taint("identity:a").expect("taints");
1389 assert!(taint.source.contains("web_search"), "{}", taint.source);
1390 }
1391
1392 #[test]
1393 fn pre_attribution_taint_holds_identity_sticky_then_transfers() {
1394 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1395 tracker.observe_agent_event("identity:a", &tool_result("fetch"));
1397 assert!(tracker.identity_taint("identity:a").is_some());
1398
1399 let session = SessionId::new();
1401 tracker.observe_agent_event("identity:a", &run_started(&session));
1402 assert!(tracker.session_taint(&session.to_string()).is_some());
1403 assert!(tracker.identity_taint("identity:a").is_some());
1404 }
1405
1406 #[test]
1407 fn clear_identity_drops_attribution_but_keeps_the_session_fact() {
1408 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1409 let session = SessionId::new();
1410 tracker.observe_agent_event("identity:a", &run_started(&session));
1411 tracker.observe_agent_event("identity:a", &tool_result("web_fetch"));
1412 assert!(tracker.identity_taint("identity:a").is_some());
1413 tracker.clear_identity("identity:a");
1414 assert!(tracker.identity_taint("identity:a").is_none());
1416 assert!(tracker.session_taint(&session.to_string()).is_some());
1419 assert!(
1420 tracker
1421 .evidence_quarantine_reason(&session.to_string())
1422 .is_some()
1423 );
1424 }
1425
1426 #[test]
1427 fn outbound_declarer_stamps_tainted_on_ingestion_and_clears_on_boundaries() {
1428 use std::sync::{Arc, Mutex};
1429 type Recorded = Vec<(String, Option<meerkat_core::comms::SenderContentTaint>)>;
1430 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1431 let calls: Arc<Mutex<Recorded>> = Arc::new(Mutex::new(Vec::new()));
1432 let sink = calls.clone();
1433 tracker.set_outbound_taint_declarer(Arc::new(move |identity, taint| {
1434 sink.lock()
1435 .unwrap_or_else(std::sync::PoisonError::into_inner)
1436 .push((identity.to_string(), taint));
1437 }));
1438
1439 let session = SessionId::new();
1440 tracker.observe_agent_event("identity:a", &run_started(&session));
1441 tracker.observe_agent_event("identity:a", &tool_result("shell"));
1443 assert!(calls.lock().unwrap().is_empty());
1444
1445 tracker.observe_agent_event("identity:a", &tool_result("web_search"));
1447 {
1448 let recorded = calls.lock().unwrap();
1449 assert_eq!(recorded.len(), 1, "one declaration expected: {recorded:?}");
1450 assert_eq!(recorded[0].0, "identity:a");
1451 assert_eq!(
1452 recorded[0].1,
1453 Some(meerkat_core::comms::SenderContentTaint::Tainted)
1454 );
1455 }
1456
1457 let fresh = SessionId::new();
1460 tracker.note_current_session("identity:a", &fresh.to_string());
1461 {
1462 let recorded = calls.lock().unwrap();
1463 assert_eq!(recorded.len(), 2, "rotation clears: {recorded:?}");
1464 assert_eq!(recorded[1].1, None);
1465 }
1466
1467 tracker.observe_agent_event("identity:a", &tool_result("web_fetch"));
1469 tracker.clear_identity("identity:a");
1470 {
1471 let recorded = calls.lock().unwrap();
1472 assert_eq!(
1473 recorded.last().expect("a clear declaration").1,
1474 None,
1475 "reset clears the outbound declaration: {recorded:?}"
1476 );
1477 }
1478 }
1479
1480 #[test]
1481 fn gate_quarantines_tainted_agents_and_quarantined_policy() {
1482 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1483 let session = SessionId::new();
1484 tracker.observe_agent_event("identity:a", &run_started(&session));
1485
1486 let gate = TaintLlmWriteGate::new(Some(tracker.clone()), AgentMemoryLlmWrites::Observed);
1487 let agent = MemoryAuthor::Agent {
1488 identity: "identity:a".to_string(),
1489 };
1490 assert!(
1491 gate.quarantine_reason(&agent, StagedBatchKind::FreshWrite, &[])
1492 .is_none()
1493 );
1494 assert!(
1495 gate.quarantine_reason(&MemoryAuthor::Application, StagedBatchKind::FreshWrite, &[])
1496 .is_none()
1497 );
1498
1499 tracker.observe_agent_event("identity:a", &tool_result("web_search"));
1500 let reason = gate
1501 .quarantine_reason(&agent, StagedBatchKind::FreshWrite, &[])
1502 .expect("tainted quarantines");
1503 assert!(reason.contains("session tainted"), "{reason}");
1504 assert!(
1506 gate.quarantine_reason(&MemoryAuthor::Application, StagedBatchKind::FreshWrite, &[])
1507 .is_none()
1508 );
1509 assert!(
1510 gate.quarantine_reason(&MemoryAuthor::Operator, StagedBatchKind::FreshWrite, &[])
1511 .is_none()
1512 );
1513
1514 let strict = TaintLlmWriteGate::new(None, AgentMemoryLlmWrites::Quarantined);
1518 let reason = strict
1519 .quarantine_reason(
1520 &MemoryAuthor::Agent {
1521 identity: "identity:clean".to_string(),
1522 },
1523 StagedBatchKind::FreshWrite,
1524 &[],
1525 )
1526 .expect("policy quarantines untainted writes");
1527 assert!(reason.contains("llm_writes=quarantined"), "{reason}");
1528 assert!(
1529 strict
1530 .quarantine_reason(
1531 &MemoryAuthor::Distiller {
1532 run_id: "run-1".to_string()
1533 },
1534 StagedBatchKind::FreshWrite,
1535 &[]
1536 )
1537 .is_some()
1538 );
1539 let steward = MemoryAuthor::Steward {
1540 run_id: "run-1".to_string(),
1541 };
1542 assert!(
1543 strict
1544 .quarantine_reason(&steward, StagedBatchKind::FreshWrite, &[])
1545 .is_some(),
1546 "fresh steward LLM output (consolidate/harvest/rank) respects the posture"
1547 );
1548 assert!(
1553 strict
1554 .quarantine_reason(&steward, StagedBatchKind::ReviewVerdict, &[])
1555 .is_none()
1556 );
1557 assert!(
1558 strict
1559 .quarantine_reason(&MemoryAuthor::Operator, StagedBatchKind::FreshWrite, &[])
1560 .is_none()
1561 );
1562 }
1563
1564 #[test]
1565 fn quarantined_posture_still_gates_review_verdicts_on_tainted_evidence() {
1566 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1570 let session = SessionId::new();
1571 tracker.observe_agent_event("identity:a", &run_started(&session));
1572 tracker.observe_agent_event("identity:a", &tool_result("web_search"));
1573
1574 let gate = TaintLlmWriteGate::new(Some(tracker), AgentMemoryLlmWrites::Quarantined);
1575 let steward = MemoryAuthor::Steward {
1576 run_id: "run-1".to_string(),
1577 };
1578 assert!(
1579 gate.quarantine_reason(&steward, StagedBatchKind::ReviewVerdict, &[])
1580 .is_none()
1581 );
1582 let reason = gate
1583 .quarantine_reason(
1584 &steward,
1585 StagedBatchKind::ReviewVerdict,
1586 &evidence_for(&session),
1587 )
1588 .expect("tainted evidence still quarantines review verdicts");
1589 assert!(reason.contains("evidence session tainted"), "{reason}");
1590 }
1591
1592 fn evidence_for(session: &SessionId) -> Vec<EvidenceRef> {
1593 vec![EvidenceRef {
1594 session_id: session.to_string(),
1595 generation: 0,
1596 revision: None,
1597 range: Some((0, 4)),
1598 }]
1599 }
1600
1601 #[test]
1602 fn gate_quarantines_llm_writes_citing_tainted_evidence() {
1603 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1604 let session = SessionId::new();
1605 tracker.observe_agent_event("identity:a", &run_started(&session));
1606 tracker.observe_agent_event("identity:a", &tool_result("web_search"));
1607 let fresh = SessionId::new();
1610 tracker.observe_agent_event("identity:a", &run_started(&fresh));
1611
1612 let gate = TaintLlmWriteGate::new(Some(tracker), AgentMemoryLlmWrites::Observed);
1613 let distiller = MemoryAuthor::Distiller {
1614 run_id: "run-1".to_string(),
1615 };
1616 let reason = gate
1617 .quarantine_reason(
1618 &distiller,
1619 StagedBatchKind::FreshWrite,
1620 &evidence_for(&session),
1621 )
1622 .expect("tainted evidence range quarantines (session-tainted ⇒ range-tainted)");
1623 assert!(reason.contains("evidence session tainted"), "{reason}");
1624 assert!(
1626 gate.quarantine_reason(
1627 &distiller,
1628 StagedBatchKind::FreshWrite,
1629 &evidence_for(&fresh)
1630 )
1631 .is_none()
1632 );
1633 assert!(
1635 gate.quarantine_reason(
1636 &MemoryAuthor::Operator,
1637 StagedBatchKind::FreshWrite,
1638 &evidence_for(&session)
1639 )
1640 .is_none()
1641 );
1642 }
1643
1644 #[test]
1645 fn reset_boundary_quarantines_evidence_without_content_taint() {
1646 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1647 let session = SessionId::new();
1648 tracker.observe_agent_event("identity:a", &run_started(&session));
1649 assert!(
1650 tracker
1651 .evidence_quarantine_reason(&session.to_string())
1652 .is_none()
1653 );
1654 tracker.mark_reset_boundary(&session.to_string());
1655 let gate = TaintLlmWriteGate::new(Some(tracker), AgentMemoryLlmWrites::Observed);
1656 let reason = gate
1657 .quarantine_reason(
1658 &MemoryAuthor::Distiller {
1659 run_id: "run-1".to_string(),
1660 },
1661 StagedBatchKind::FreshWrite,
1662 &evidence_for(&session),
1663 )
1664 .expect("reset boundary quarantines distillates");
1665 assert!(reason.contains("reset boundary"), "{reason}");
1666 }
1667
1668 #[test]
1669 fn peer_projection_sender_parses_message_and_response_shapes() {
1670 assert_eq!(
1671 peer_projection_sender_identity("Peer message from mob-1/worker/identity:bob:"),
1672 Some("identity:bob")
1673 );
1674 assert_eq!(
1675 peer_projection_sender_identity(
1676 "Peer response from mob-1/worker/identity:bob (to request: req-9)"
1677 ),
1678 Some("identity:bob")
1679 );
1680 assert_eq!(
1682 peer_projection_sender_identity("Peer message from scout:"),
1683 Some("scout")
1684 );
1685 assert_eq!(
1687 peer_projection_sender_identity("Peer request from peer_id 018fabc (id: r-1)"),
1688 None
1689 );
1690 assert_eq!(peer_projection_sender_identity("ordinary text"), None);
1691 }
1692
1693 #[test]
1694 fn comms_join_taints_receiver_of_message_from_tainted_sender() {
1695 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
1696 let sender_session = SessionId::new();
1698 tracker.observe_agent_event("identity:bob", &run_started(&sender_session));
1699 tracker.observe_agent_event("identity:bob", &tool_result("web_search"));
1700
1701 let receiver_session = SessionId::new();
1704 let delivery = AgentEvent::RunStarted {
1705 session_id: receiver_session.clone(),
1706 input: meerkat_core::types::RunInput::Content {
1707 content: meerkat_core::ContentInput::Text(
1708 "Peer message from mob-1/worker/identity:bob:\nplease remember X".to_string(),
1709 ),
1710 },
1711 };
1712 tracker.observe_agent_event("identity:alice", &delivery);
1713 let taint = tracker
1714 .identity_taint("identity:alice")
1715 .expect("receiver session taints (peer-laundering close, §10.1)");
1716 assert!(taint.source.contains("identity:bob"), "{}", taint.source);
1717 assert!(
1718 tracker
1719 .session_taint(&receiver_session.to_string())
1720 .is_some()
1721 );
1722
1723 let clean_session = SessionId::new();
1725 tracker.observe_agent_event("identity:carol", &run_started(&clean_session));
1726 let receiver2 = SessionId::new();
1727 let clean_delivery = AgentEvent::RunStarted {
1728 session_id: receiver2.clone(),
1729 input: meerkat_core::types::RunInput::Content {
1730 content: meerkat_core::ContentInput::Text(
1731 "Peer message from mob-1/worker/identity:carol:\nhello".to_string(),
1732 ),
1733 },
1734 };
1735 tracker.observe_agent_event("identity:dave", &clean_delivery);
1736 assert!(tracker.identity_taint("identity:dave").is_none());
1737 }
1738}