1use std::collections::HashMap;
56use std::path::{Path, PathBuf};
57use std::sync::{Arc, Mutex};
58use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
59
60use async_trait::async_trait;
61use futures::StreamExt;
62use serde::Deserialize;
63
64use meerkat_client::{LlmClient, LlmDoneOutcome, LlmError, LlmEvent, LlmRequest};
65use meerkat_core::event::AgentEvent;
66use meerkat_core::{Message, Provider, UserMessage};
67
68use crate::identity_first::agent_memory::{
69 AgentMemoryError, AgentMemoryProvider, MEMORY_TOOL_NAME, compact_whitespace,
70 truncate_utf8_boundary,
71};
72use crate::memory::guards::{BackgroundBudget, BackgroundBudgetConfig};
73use crate::memory::records::{
74 EvidenceRef, ManifestTier, MemoryAuthor, MemoryKind, MemoryScope, NewMemoryRecord, RecordMeta,
75};
76use crate::memory::selector::FactorySelectorHandle;
77use crate::memory::taint::{MemberAgentEventSink, SessionTaintTracker};
78
79pub const EMBEDDED_PROMPT_V0: &str = include_str!("distiller_prompt_v0.md");
84
85const MANIFEST_PLACEHOLDER: &str = "{{existing_manifest}}";
86const TOMBSTONES_PLACEHOLDER: &str = "{{recent_tombstones}}";
87const TRANSCRIPT_PLACEHOLDER: &str = "{{transcript}}";
88
89pub const PRE_ROTATION_TIMEOUT: Duration = Duration::from_secs(15);
93pub const MIN_SECONDS_BETWEEN_RUNS: u64 = 120;
95const TOMBSTONE_LOOKBACK_MS: u64 = 7 * 24 * 60 * 60 * 1000;
97const COMPACTION_HARVEST_LIMIT: usize = 64;
99const MAX_TRANSCRIPT_MESSAGE_BYTES: usize = 4 * 1024;
101const MAX_TRANSCRIPT_TOTAL_BYTES: usize = 48 * 1024;
102const MAX_TRACKED_WINDOWS: usize = 4096;
104const DEFAULT_MAX_OUTPUT_TOKENS: u32 = 2048;
106
107#[derive(Debug)]
112pub enum DistillerError {
113 Profile(String),
114 Auth(String),
115 Client(String),
116 Parse(String),
117 Store(String),
118 Transcript(String),
119}
120
121impl std::fmt::Display for DistillerError {
122 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
123 match self {
124 Self::Profile(msg) => write!(f, "distiller profile error: {msg}"),
125 Self::Auth(msg) => write!(f, "distiller auth error: {msg}"),
126 Self::Client(msg) => write!(f, "distiller client error: {msg}"),
127 Self::Parse(msg) => write!(f, "distiller parse error: {msg}"),
128 Self::Store(msg) => write!(f, "distiller store error: {msg}"),
129 Self::Transcript(msg) => write!(f, "distiller transcript error: {msg}"),
130 }
131 }
132}
133
134impl std::error::Error for DistillerError {}
135
136#[derive(Debug, Clone, Deserialize)]
141pub struct DistillerParams {
142 #[serde(default = "default_temperature")]
143 pub temperature: f32,
144 #[serde(default = "default_max_output_tokens")]
145 pub max_output_tokens: u32,
146 #[serde(default = "default_max_manifest_records")]
148 pub max_manifest_records: usize,
149 #[serde(default = "default_max_tombstones")]
150 pub max_tombstones: usize,
151}
152
153fn default_temperature() -> f32 {
154 0.0
155}
156fn default_max_output_tokens() -> u32 {
157 DEFAULT_MAX_OUTPUT_TOKENS
158}
159fn default_max_manifest_records() -> usize {
160 200
161}
162fn default_max_tombstones() -> usize {
163 32
164}
165
166impl Default for DistillerParams {
167 fn default() -> Self {
168 Self {
169 temperature: default_temperature(),
170 max_output_tokens: default_max_output_tokens(),
171 max_manifest_records: default_max_manifest_records(),
172 max_tombstones: default_max_tombstones(),
173 }
174 }
175}
176
177#[derive(Debug, Clone)]
179pub struct DistillerProfile {
180 pub stage: String,
181 pub version: String,
182 pub model: String,
183 pub provider: Provider,
184 pub prompt_bundle: String,
185 pub prompt_template: String,
186 pub params: DistillerParams,
187}
188
189#[derive(Debug, Deserialize)]
190struct RawProfile {
191 stage: String,
192 version: String,
193 model: String,
194 #[serde(default)]
195 provider: Option<String>,
196 prompt_bundle: String,
197 #[serde(default)]
198 params: Option<DistillerParams>,
199}
200
201impl DistillerProfile {
202 pub fn embedded_default() -> Self {
207 Self {
208 stage: "distiller".to_string(),
209 version: "0".to_string(),
210 model: "claude-haiku-4-5".to_string(),
211 provider: Provider::Anthropic,
212 prompt_bundle: "prompts/distiller-v0.md".to_string(),
213 prompt_template: EMBEDDED_PROMPT_V0.to_string(),
214 params: DistillerParams::default(),
215 }
216 }
217
218 pub fn with_model_override(mut self, model: &str) -> Result<Self, DistillerError> {
222 let model = model.trim();
223 if model.is_empty() {
224 return Err(DistillerError::Profile(
225 "distiller model override must not be empty".to_string(),
226 ));
227 }
228 self.provider = meerkat_models::infer_provider(model).ok_or_else(|| {
229 DistillerError::Profile(format!(
230 "distiller model override '{model}' is not in the model catalog"
231 ))
232 })?;
233 self.model = model.to_string();
234 Ok(self)
235 }
236
237 pub fn load(path: &Path) -> Result<Self, DistillerError> {
240 let text = std::fs::read_to_string(path).map_err(|err| {
241 DistillerError::Profile(format!("cannot read profile '{}': {err}", path.display()))
242 })?;
243 let raw: RawProfile = toml::from_str(&text).map_err(|err| {
244 DistillerError::Profile(format!("invalid profile '{}': {err}", path.display()))
245 })?;
246 if raw.stage != "distiller" {
247 return Err(DistillerError::Profile(format!(
248 "profile '{}' is for stage '{}', not 'distiller'",
249 path.display(),
250 raw.stage
251 )));
252 }
253 if raw.model.trim().is_empty() || raw.model == "PLACEHOLDER" {
254 return Err(DistillerError::Profile(format!(
255 "profile '{}' does not name a model",
256 path.display()
257 )));
258 }
259 let provider = match raw.provider.as_deref() {
260 Some(name) => Provider::parse_strict(name).ok_or_else(|| {
261 DistillerError::Profile(format!(
262 "profile '{}': unknown provider '{name}'",
263 path.display()
264 ))
265 })?,
266 None => meerkat_models::infer_provider(&raw.model).ok_or_else(|| {
267 DistillerError::Profile(format!(
268 "profile '{}': model '{}' is not in the catalog; set `provider` explicitly",
269 path.display(),
270 raw.model
271 ))
272 })?,
273 };
274 let base = path.parent().unwrap_or_else(|| Path::new("."));
275 let candidates = [
276 base.join(&raw.prompt_bundle),
277 base.parent()
278 .unwrap_or_else(|| Path::new("."))
279 .join(&raw.prompt_bundle),
280 ];
281 let bundle_path = candidates.iter().find(|p| p.is_file()).ok_or_else(|| {
282 DistillerError::Profile(format!(
283 "profile '{}': prompt_bundle '{}' does not resolve",
284 path.display(),
285 raw.prompt_bundle
286 ))
287 })?;
288 let prompt_template = std::fs::read_to_string(bundle_path).map_err(|err| {
289 DistillerError::Profile(format!(
290 "cannot read prompt bundle '{}': {err}",
291 bundle_path.display()
292 ))
293 })?;
294 let profile = Self {
295 stage: raw.stage,
296 version: raw.version,
297 model: raw.model,
298 provider,
299 prompt_bundle: raw.prompt_bundle,
300 prompt_template,
301 params: raw.params.unwrap_or_default(),
302 };
303 profile.validate()?;
304 Ok(profile)
305 }
306
307 fn validate(&self) -> Result<(), DistillerError> {
308 for placeholder in [
309 MANIFEST_PLACEHOLDER,
310 TOMBSTONES_PLACEHOLDER,
311 TRANSCRIPT_PLACEHOLDER,
312 ] {
313 if !self.prompt_template.contains(placeholder) {
314 return Err(DistillerError::Profile(format!(
315 "prompt bundle '{}' is missing placeholder `{placeholder}`",
316 self.prompt_bundle
317 )));
318 }
319 }
320 Ok(())
321 }
322}
323
324#[derive(Debug, Clone, PartialEq, Eq)]
332pub struct DistillerConfig {
333 pub enabled: bool,
334 pub runs_per_hour: u32,
336 pub min_interactions: u32,
339 pub model: Option<String>,
341}
342
343impl Default for DistillerConfig {
344 fn default() -> Self {
345 Self {
346 enabled: false,
347 runs_per_hour: crate::memory::guards::DEFAULT_RUNS_PER_HOUR,
348 min_interactions: 3,
349 model: None,
350 }
351 }
352}
353
354#[derive(Debug, Clone, Copy, PartialEq, Eq)]
362pub enum Epistemic {
363 OperatorSaid,
364 Observed,
365}
366
367impl Epistemic {
368 fn parse(value: &str) -> Option<Self> {
369 match value {
370 "operator_said" => Some(Self::OperatorSaid),
371 "observed" => Some(Self::Observed),
372 _ => None,
373 }
374 }
375}
376
377#[derive(Debug, Clone, PartialEq, Eq)]
378pub enum ProposedAction {
379 Remember,
380 Update { target_id: String },
381}
382
383#[derive(Debug, Clone, PartialEq)]
385pub struct ProposedOp {
386 pub action: ProposedAction,
387 pub kind: MemoryKind,
388 pub title: String,
389 pub description: String,
390 pub body: String,
391 pub tags: Vec<String>,
392 pub epistemic: Epistemic,
393 pub evidence_range: Option<(u64, u64)>,
394}
395
396#[derive(Debug, Deserialize)]
397struct RawOp {
398 action: String,
399 #[serde(default)]
400 target_id: Option<String>,
401 #[serde(default)]
402 kind: Option<String>,
403 title: String,
404 #[serde(default)]
405 description: String,
406 body: String,
407 #[serde(default)]
408 tags: Vec<String>,
409 epistemic: String,
410 #[serde(default)]
411 evidence_range: Option<(u64, u64)>,
412}
413
414pub fn parse_ops(reply: &str) -> Result<Vec<RawParsedOp>, String> {
417 let trimmed = reply.trim();
418 let raw: Vec<RawOp> = match serde_json::from_str(trimmed) {
419 Ok(raw) => raw,
420 Err(first_err) => {
421 let (Some(start), Some(end)) = (trimmed.find('['), trimmed.rfind(']')) else {
422 return Err(format!("no JSON array in reply: {first_err}"));
423 };
424 if start >= end {
425 return Err(format!("no JSON array in reply: {first_err}"));
426 }
427 serde_json::from_str(&trimmed[start..=end]).map_err(|err| err.to_string())?
428 }
429 };
430 Ok(raw.into_iter().map(RawParsedOp).collect())
431}
432
433#[derive(Debug)]
436pub struct RawParsedOp(RawOp);
437
438pub fn validate_op(op: RawParsedOp, manifest_ids: &[String]) -> Result<ProposedOp, String> {
442 let raw = op.0;
443 let action = match raw.action.as_str() {
444 "remember" => ProposedAction::Remember,
445 "update" => {
446 let target = raw
447 .target_id
448 .as_deref()
449 .map(str::trim)
450 .filter(|id| !id.is_empty())
451 .ok_or_else(|| "update op is missing target_id".to_string())?;
452 if !manifest_ids.iter().any(|id| id == target) {
453 return Err(format!(
454 "update op targets '{target}', which is not in the manifest"
455 ));
456 }
457 ProposedAction::Update {
458 target_id: target.to_string(),
459 }
460 }
461 other => return Err(format!("unknown action '{other}'")),
462 };
463 let kind = match raw.kind.as_deref() {
464 None => MemoryKind::Fact,
465 Some(kind) => {
466 MemoryKind::parse(kind).ok_or_else(|| format!("unknown record kind '{kind}'"))?
467 }
468 };
469 let epistemic = Epistemic::parse(&raw.epistemic)
470 .ok_or_else(|| format!("unknown epistemic status '{}'", raw.epistemic))?;
471 let title = compact_whitespace(&raw.title);
472 if title.is_empty() {
473 return Err("op has an empty title".to_string());
474 }
475 let body = raw.body.trim().to_string();
476 if body.is_empty() {
477 return Err("op has an empty body".to_string());
478 }
479 if let Some((start, end)) = raw.evidence_range
480 && start > end
481 {
482 return Err(format!("evidence_range [{start}, {end}] is inverted"));
483 }
484 Ok(ProposedOp {
485 action,
486 kind,
487 title,
488 description: compact_whitespace(&raw.description),
489 body,
490 tags: raw.tags,
491 epistemic,
492 evidence_range: raw.evidence_range,
493 })
494}
495
496#[derive(Debug, Clone, PartialEq, Eq)]
504pub struct TranscriptMessage {
505 pub index: u64,
506 pub role: &'static str,
507 pub text: String,
508}
509
510#[derive(Debug, Clone, PartialEq, Eq)]
513pub struct TranscriptSlice {
514 pub session_key: String,
515 pub start_index: u64,
516 pub end_index: u64,
517 pub messages: Vec<TranscriptMessage>,
518 pub head_revision: Option<String>,
523}
524
525#[async_trait]
528pub trait TranscriptSource: Send + Sync {
529 async fn read(
532 &self,
533 session_key: &str,
534 from_index: u64,
535 ) -> Result<Option<TranscriptSlice>, DistillerError>;
536}
537
538pub struct SessionStoreTranscriptSource {
543 store: Arc<dyn meerkat::SessionStore>,
544}
545
546impl SessionStoreTranscriptSource {
547 pub fn new(store: Arc<dyn meerkat::SessionStore>) -> Self {
548 Self { store }
549 }
550}
551
552#[async_trait]
553impl TranscriptSource for SessionStoreTranscriptSource {
554 async fn read(
555 &self,
556 session_key: &str,
557 from_index: u64,
558 ) -> Result<Option<TranscriptSlice>, DistillerError> {
559 let session_id = meerkat_core::types::SessionId::parse(session_key).map_err(|err| {
560 DistillerError::Transcript(format!("invalid session key '{session_key}': {err}"))
561 })?;
562 let session = self
563 .store
564 .load(&session_id)
565 .await
566 .map_err(|err| DistillerError::Transcript(err.to_string()))?;
567 let Some(session) = session else {
568 return Ok(None);
569 };
570 let head_revision = session.transcript_revision().ok();
575 let all = session.messages();
576 let end_index = all.len() as u64;
577 let start_index = from_index.min(end_index);
578 let messages = all[start_index as usize..]
579 .iter()
580 .enumerate()
581 .filter_map(|(offset, message)| {
582 let index = start_index + offset as u64;
583 project_message(message).map(|(role, text)| TranscriptMessage {
584 index,
585 role,
586 text: truncate_utf8_boundary(
587 &compact_whitespace(&text),
588 MAX_TRANSCRIPT_MESSAGE_BYTES,
589 ),
590 })
591 })
592 .collect();
593 Ok(Some(TranscriptSlice {
594 session_key: session_key.to_string(),
595 start_index,
596 end_index,
597 messages,
598 head_revision,
599 }))
600 }
601}
602
603fn project_message(message: &Message) -> Option<(&'static str, String)> {
608 match message {
609 Message::System(_) => None,
610 Message::SystemNotice(notice) => notice
611 .body
612 .as_deref()
613 .map(|body| ("system notice", body.to_string())),
614 Message::User(user) => Some(("user", user.text_content())),
615 Message::BlockAssistant(assistant) => Some((
616 "assistant",
617 assistant.text_blocks().collect::<Vec<_>>().join("\n"),
618 )),
619 Message::ToolResults { results, .. } => {
620 let text = results
621 .iter()
622 .map(|result| meerkat_core::types::text_content(&result.content))
623 .collect::<Vec<_>>()
624 .join("\n");
625 Some(("tool results", text))
626 }
627 }
628}
629
630#[derive(Debug, Clone, PartialEq)]
633pub struct DiscardEntry {
634 pub content: String,
635 pub range: Option<(u64, u64)>,
638}
639
640#[async_trait]
642pub trait CompactionDiscardSource: Send + Sync {
643 async fn read_discards(
644 &self,
645 session_key: &str,
646 limit: usize,
647 ) -> Result<Vec<DiscardEntry>, DistillerError>;
648
649 async fn drop_scope(&self, session_key: &str) -> Result<usize, DistillerError> {
659 let _ = session_key;
660 Ok(0)
661 }
662}
663
664pub struct HnswDiscardSource {
682 dir: PathBuf,
683}
684
685const HARVEST_PAGE_ROWS: usize = 256;
687
688impl HnswDiscardSource {
689 pub fn new(dir: impl Into<PathBuf>) -> Self {
690 Self { dir: dir.into() }
691 }
692}
693
694#[async_trait]
695impl CompactionDiscardSource for HnswDiscardSource {
696 async fn read_discards(
697 &self,
698 session_key: &str,
699 limit: usize,
700 ) -> Result<Vec<DiscardEntry>, DistillerError> {
701 use meerkat_core::memory::{MemoryEnumerationRequest, MemorySearchScope, MemoryStore};
702
703 if !self.dir.is_dir() {
704 return Ok(Vec::new());
707 }
708 let session_id = meerkat_core::types::SessionId::parse(session_key).map_err(|err| {
709 DistillerError::Transcript(format!("invalid session key '{session_key}': {err}"))
710 })?;
711 let dir = self.dir.clone();
712 let store = tokio::task::spawn_blocking(move || meerkat_memory::HnswMemoryStore::open(dir))
713 .await
714 .map_err(|err| DistillerError::Store(err.to_string()))?
715 .map_err(|err| DistillerError::Store(err.to_string()))?;
716 let scope = MemorySearchScope::for_session(session_id);
717
718 let mut entries = Vec::new();
722 let mut offset = 0usize;
723 loop {
724 let remaining = limit.saturating_sub(entries.len());
725 if remaining == 0 {
726 break;
727 }
728 let page = store
729 .enumerate_scoped(
730 &scope,
731 MemoryEnumerationRequest {
732 limit: HARVEST_PAGE_ROWS.min(remaining),
733 offset,
734 source_overlap: None,
735 indexed_after: None,
736 },
737 )
738 .await
739 .map_err(|err| DistillerError::Store(err.to_string()))?;
740 for record in page.records {
741 entries.push(DiscardEntry {
742 range: record
743 .metadata
744 .source
745 .source_range()
746 .map(|range| (range.start(), range.end())),
747 content: record.content,
748 });
749 }
750 match page.next_offset {
751 Some(next) if entries.len() < limit => offset = next,
752 Some(_) => {
753 tracing::warn!(
754 session_key,
755 limit,
756 "compaction discard harvest hit its row ceiling; \
757 scope has more rows than were read"
758 );
759 break;
760 }
761 None => break,
762 }
763 }
764 Ok(entries)
765 }
766
767 async fn drop_scope(&self, session_key: &str) -> Result<usize, DistillerError> {
768 use meerkat_core::memory::{MemoryOwner, MemoryStore};
769
770 if !self.dir.is_dir() {
771 return Ok(0);
773 }
774 let session_id = meerkat_core::types::SessionId::parse(session_key).map_err(|err| {
775 DistillerError::Transcript(format!("invalid session key '{session_key}': {err}"))
776 })?;
777 let dir = self.dir.clone();
778 let store = tokio::task::spawn_blocking(move || meerkat_memory::HnswMemoryStore::open(dir))
779 .await
780 .map_err(|err| DistillerError::Store(err.to_string()))?
781 .map_err(|err| DistillerError::Store(err.to_string()))?;
782 let receipt = store
783 .drop_scope(&MemoryOwner::canonical_session(session_id))
784 .await
785 .map_err(|err| DistillerError::Store(err.to_string()))?;
786 Ok(receipt.dropped_entries)
787 }
788}
789
790#[derive(Debug, Clone, PartialEq, Eq)]
797pub struct TombstoneMeta {
798 pub title: String,
799 pub kind: MemoryKind,
800 pub tombstoned_at_ms: u64,
801}
802
803#[async_trait]
804pub trait TombstoneSource: Send + Sync {
805 async fn recent_tombstones(
806 &self,
807 scope: &MemoryScope,
808 since_ms: u64,
809 limit: usize,
810 ) -> Result<Vec<TombstoneMeta>, AgentMemoryError>;
811}
812
813#[async_trait]
818pub trait DistillerClientHandle: Send + Sync {
819 async fn client(&self) -> Result<Arc<dyn LlmClient>, DistillerError>;
820 fn invalidate(&self);
821}
822
823pub struct FactoryDistillerHandle {
826 inner: FactorySelectorHandle,
827}
828
829impl FactoryDistillerHandle {
830 pub fn new(
831 store_path: impl Into<PathBuf>,
832 config: meerkat::Config,
833 realm: impl Into<String>,
834 profile: &DistillerProfile,
835 ) -> Self {
836 Self {
837 inner: FactorySelectorHandle::for_model(
838 store_path,
839 config,
840 realm,
841 &profile.model,
842 profile.provider,
843 ),
844 }
845 }
846}
847
848#[async_trait]
849impl DistillerClientHandle for FactoryDistillerHandle {
850 async fn client(&self) -> Result<Arc<dyn LlmClient>, DistillerError> {
851 use crate::memory::selector::{SelectorError, SelectorHandle};
852 self.inner.client().await.map_err(|err| match err {
853 SelectorError::Auth(msg) => DistillerError::Auth(msg),
854 other => DistillerError::Client(other.to_string()),
855 })
856 }
857
858 fn invalidate(&self) {
859 use crate::memory::selector::SelectorHandle;
860 self.inner.invalidate();
861 }
862}
863
864pub fn render_prompt(
869 profile: &DistillerProfile,
870 manifest: &[RecordMeta],
871 tombstones: &[TombstoneMeta],
872 transcript_text: &str,
873) -> String {
874 let manifest_text = if manifest.is_empty() {
875 "(no records)".to_string()
876 } else {
877 manifest
878 .iter()
879 .take(profile.params.max_manifest_records)
880 .map(crate::memory::selector::render_manifest_row)
881 .collect::<Vec<_>>()
882 .join("\n")
883 };
884 let tombstones_text = if tombstones.is_empty() {
885 "(none)".to_string()
886 } else {
887 tombstones
888 .iter()
889 .take(profile.params.max_tombstones)
890 .map(|tombstone| {
891 format!(
892 "- [{}] {}",
893 tombstone.kind.as_str(),
894 compact_whitespace(&tombstone.title)
895 )
896 })
897 .collect::<Vec<_>>()
898 .join("\n")
899 };
900 profile
901 .prompt_template
902 .replace(MANIFEST_PLACEHOLDER, &manifest_text)
903 .replace(TOMBSTONES_PLACEHOLDER, &tombstones_text)
904 .replace(TRANSCRIPT_PLACEHOLDER, transcript_text)
905}
906
907pub fn render_transcript(slice: &TranscriptSlice) -> String {
911 let mut lines: Vec<String> = Vec::new();
912 let mut total = 0usize;
913 for message in slice.messages.iter().rev() {
914 let line = format!("[{}] {}: {}", message.index, message.role, message.text);
915 if total + line.len() + 1 > MAX_TRANSCRIPT_TOTAL_BYTES && !lines.is_empty() {
916 lines.push("(earlier messages omitted for budget)".to_string());
917 break;
918 }
919 total += line.len() + 1;
920 lines.push(line);
921 }
922 lines.reverse();
923 lines.join("\n")
924}
925
926fn render_discards(entries: &[DiscardEntry]) -> String {
927 let mut lines: Vec<String> = vec![
928 "(compaction-discarded content recovered from session semantic memory; \
929 [N-M] are pre-compaction message offsets)"
930 .to_string(),
931 ];
932 let mut total = 0usize;
933 for entry in entries {
934 let prefix = match entry.range {
935 Some((start, end)) => format!("[{start}-{end}]"),
936 None => "[?]".to_string(),
937 };
938 let line = format!(
939 "{prefix} {}",
940 truncate_utf8_boundary(
941 &compact_whitespace(&entry.content),
942 MAX_TRANSCRIPT_MESSAGE_BYTES
943 )
944 );
945 if total + line.len() + 1 > MAX_TRANSCRIPT_TOTAL_BYTES && lines.len() > 1 {
946 lines.push("(further discards omitted for budget)".to_string());
947 break;
948 }
949 total += line.len() + 1;
950 lines.push(line);
951 }
952 lines.join("\n")
953}
954
955async fn complete_text(
956 client: &dyn LlmClient,
957 profile: &DistillerProfile,
958 prompt: String,
959) -> Result<String, DistillerError> {
960 let request = LlmRequest::new(
961 &profile.model,
962 vec![Message::User(UserMessage::text(prompt))],
963 )
964 .with_max_tokens(profile.params.max_output_tokens)
965 .with_temperature(profile.params.temperature);
966 let mut stream = client.stream(&request);
967 let mut text = String::new();
968 while let Some(event) = stream.next().await {
969 match event.map_err(classify_llm_error)? {
970 LlmEvent::TextDelta { delta, .. } => text.push_str(&delta),
971 LlmEvent::Done { outcome } => match outcome {
972 LlmDoneOutcome::Success { .. } => break,
973 LlmDoneOutcome::Error { error } => return Err(classify_llm_error(error)),
974 },
975 _ => {}
976 }
977 }
978 Ok(text)
979}
980
981fn classify_llm_error(error: LlmError) -> DistillerError {
982 match error {
983 LlmError::AuthenticationFailed { .. } | LlmError::InvalidApiKey => {
984 DistillerError::Auth(error.to_string())
985 }
986 other => DistillerError::Client(other.to_string()),
987 }
988}
989
990pub async fn extract(
993 profile: &DistillerProfile,
994 client: &dyn LlmClient,
995 manifest: &[RecordMeta],
996 tombstones: &[TombstoneMeta],
997 transcript_text: &str,
998) -> Result<Vec<RawParsedOp>, DistillerError> {
999 let prompt = render_prompt(profile, manifest, tombstones, transcript_text);
1000 let reply = complete_text(client, profile, prompt).await?;
1001 match parse_ops(&reply) {
1002 Ok(ops) => Ok(ops),
1003 Err(first_err) => {
1004 let repair_prompt = format!(
1005 "The following reply was supposed to be exactly one JSON array of memory ops \
1006 (each {{\"action\": \"remember\" | \"update\", \"target_id\"?, \"kind\", \
1007 \"title\", \"description\", \"body\", \"tags\", \"epistemic\", \
1008 \"evidence_range\"?}}) but did not parse ({first_err}). Reply with ONLY the \
1009 corrected JSON array, no other text.\n\n{reply}"
1010 );
1011 let repaired = complete_text(client, profile, repair_prompt).await?;
1012 parse_ops(&repaired).map_err(DistillerError::Parse)
1013 }
1014 }
1015}
1016
1017#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1026pub enum DistillCause {
1027 Interactions,
1028 Respawn,
1029 Retire,
1030 Delete,
1031 Reset,
1032 ResumeFallback,
1033 Compaction,
1034}
1035
1036impl DistillCause {
1037 pub fn as_str(&self) -> &'static str {
1038 match self {
1039 Self::Interactions => "interactions",
1040 Self::Respawn => "respawn",
1041 Self::Retire => "retire",
1042 Self::Delete => "delete",
1043 Self::Reset => "reset",
1044 Self::ResumeFallback => "resume_fallback",
1045 Self::Compaction => "compaction",
1046 }
1047 }
1048
1049 pub fn quarantines_output(&self) -> bool {
1052 matches!(self, Self::Reset)
1053 }
1054}
1055
1056#[derive(Default, Clone)]
1061struct WindowState {
1062 cursor: u64,
1063 completed_runs: u32,
1064 recorder_wrote: bool,
1070 generation: u64,
1071 last_run_at: Option<Instant>,
1072 last_activity_at: Option<Instant>,
1073 in_flight: bool,
1074}
1075
1076#[derive(Debug, Clone, PartialEq, Eq)]
1078pub enum DistillOutcome {
1079 Skipped {
1080 reason: String,
1081 },
1082 Completed {
1083 run_id: String,
1084 written: usize,
1085 quarantined: usize,
1086 },
1087}
1088
1089pub(crate) const SKIP_NO_DISCARD_SOURCE: &str = "no compaction discard source wired";
1093pub(crate) const SKIP_NO_DISCARDS: &str = "no compaction discards to harvest";
1094
1095impl DistillOutcome {
1096 pub fn compaction_harvest_satisfied(&self) -> bool {
1102 match self {
1103 Self::Completed { .. } => true,
1104 Self::Skipped { reason } => {
1105 reason == SKIP_NO_DISCARDS || reason == SKIP_NO_DISCARD_SOURCE
1106 }
1107 }
1108 }
1109}
1110
1111pub type CompactionFollowUp = Arc<dyn Fn(&str, &str, &DistillOutcome) + Send + Sync>;
1116
1117pub type CompactionObserved = Arc<dyn Fn(&str, &str) + Send + Sync>;
1126
1127pub struct DistillerEngine {
1128 profile: DistillerProfile,
1129 config: DistillerConfig,
1130 handle: Arc<dyn DistillerClientHandle>,
1131 provider: Arc<dyn AgentMemoryProvider>,
1132 tombstones: Arc<dyn TombstoneSource>,
1133 transcripts: Arc<dyn TranscriptSource>,
1134 compaction: Option<Arc<dyn CompactionDiscardSource>>,
1135 tracker: Option<SessionTaintTracker>,
1136 budget: BackgroundBudget,
1137 realm: String,
1138 events: Mutex<Option<Arc<dyn crate::memory::events::MemoryEventSink>>>,
1140 windows: Mutex<HashMap<(String, String), WindowState>>,
1142 pre_rotation_timeout: Duration,
1144 run_counter: std::sync::atomic::AtomicU64,
1145 compaction_follow_up: Mutex<Option<CompactionFollowUp>>,
1148 compaction_observed: Mutex<Option<CompactionObserved>>,
1150}
1151
1152impl DistillerEngine {
1153 #[allow(clippy::too_many_arguments)]
1154 pub fn new(
1155 profile: DistillerProfile,
1156 config: DistillerConfig,
1157 handle: Arc<dyn DistillerClientHandle>,
1158 provider: Arc<dyn AgentMemoryProvider>,
1159 tombstones: Arc<dyn TombstoneSource>,
1160 transcripts: Arc<dyn TranscriptSource>,
1161 compaction: Option<Arc<dyn CompactionDiscardSource>>,
1162 tracker: Option<SessionTaintTracker>,
1163 realm: impl Into<String>,
1164 ) -> Self {
1165 let budget = BackgroundBudget::new(BackgroundBudgetConfig {
1166 runs_per_window: config.runs_per_hour,
1167 window: Duration::from_hours(1),
1168 max_concurrent: crate::memory::guards::DEFAULT_MAX_CONCURRENT,
1169 });
1170 Self {
1171 profile,
1172 config,
1173 handle,
1174 provider,
1175 tombstones,
1176 transcripts,
1177 compaction,
1178 tracker,
1179 budget,
1180 realm: realm.into(),
1181 events: Mutex::new(None),
1182 windows: Mutex::new(HashMap::new()),
1183 pre_rotation_timeout: PRE_ROTATION_TIMEOUT,
1184 run_counter: std::sync::atomic::AtomicU64::new(0),
1185 compaction_follow_up: Mutex::new(None),
1186 compaction_observed: Mutex::new(None),
1187 }
1188 }
1189
1190 #[cfg(test)]
1191 fn with_pre_rotation_timeout(mut self, timeout: Duration) -> Self {
1192 self.pre_rotation_timeout = timeout;
1193 self
1194 }
1195
1196 pub fn set_event_sink(&self, sink: Arc<dyn crate::memory::events::MemoryEventSink>) {
1199 self.budget.set_event_sink(sink.clone());
1200 *self
1201 .events
1202 .lock()
1203 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(sink);
1204 }
1205
1206 pub fn pre_rotation_timeout(&self) -> Duration {
1207 self.pre_rotation_timeout
1208 }
1209
1210 pub fn set_compaction_follow_up(&self, hook: CompactionFollowUp) {
1213 *self
1214 .compaction_follow_up
1215 .lock()
1216 .unwrap_or_else(|err| err.into_inner()) = Some(hook);
1217 }
1218
1219 pub fn set_compaction_observed(&self, hook: CompactionObserved) {
1222 *self
1223 .compaction_observed
1224 .lock()
1225 .unwrap_or_else(|err| err.into_inner()) = Some(hook);
1226 }
1227
1228 fn note_compaction_observed(&self, identity: &str, session_key: &str) {
1232 let hook = self
1233 .compaction_observed
1234 .lock()
1235 .unwrap_or_else(|err| err.into_inner())
1236 .clone();
1237 if let Some(hook) = hook {
1238 hook(identity, session_key);
1239 }
1240 }
1241
1242 pub fn distilled_cursor(&self, identity: &str, session_key: &str) -> u64 {
1246 self.windows
1247 .lock()
1248 .unwrap_or_else(std::sync::PoisonError::into_inner)
1249 .get(&(identity.to_string(), session_key.to_string()))
1250 .map(|state| state.cursor)
1251 .unwrap_or(0)
1252 }
1253
1254 fn with_window<T>(
1255 &self,
1256 identity: &str,
1257 session_key: &str,
1258 f: impl FnOnce(&mut WindowState) -> T,
1259 ) -> T {
1260 let mut windows = self
1261 .windows
1262 .lock()
1263 .unwrap_or_else(std::sync::PoisonError::into_inner);
1264 if windows.len() >= MAX_TRACKED_WINDOWS
1265 && !windows.contains_key(&(identity.to_string(), session_key.to_string()))
1266 && let Some(oldest) = windows
1267 .iter()
1268 .min_by_key(|(_, state)| state.last_activity_at)
1269 .map(|(key, _)| key.clone())
1270 {
1271 windows.remove(&oldest);
1272 }
1273 let state = windows
1274 .entry((identity.to_string(), session_key.to_string()))
1275 .or_default();
1276 state.last_activity_at = Some(Instant::now());
1277 f(state)
1278 }
1279
1280 pub fn note_session_generation(&self, identity: &str, session_key: &str, generation: u64) {
1286 self.with_window(identity, session_key, |state| {
1287 state.generation = generation;
1288 });
1289 }
1290
1291 fn note_recorder_write(&self, identity: &str, session_key: &str) {
1292 self.with_window(identity, session_key, |state| {
1293 state.recorder_wrote = true;
1294 });
1295 }
1296
1297 fn note_run_completed(&self, identity: &str, session_key: &str) -> bool {
1300 let min_interactions = self.config.min_interactions;
1301 self.with_window(identity, session_key, |state| {
1302 state.completed_runs += 1;
1303 if state.in_flight {
1304 return false;
1305 }
1306 if state.completed_runs < min_interactions {
1307 return false;
1308 }
1309 if let Some(last) = state.last_run_at
1310 && last.elapsed() < Duration::from_secs(MIN_SECONDS_BETWEEN_RUNS)
1311 {
1312 return false;
1313 }
1314 state.in_flight = true;
1315 true
1316 })
1317 }
1318
1319 fn identity_scope(&self, identity: &str) -> MemoryScope {
1320 MemoryScope::Identity {
1321 realm: self.realm.clone(),
1322 identity: identity.to_string(),
1323 }
1324 }
1325
1326 fn mint_run_id(&self) -> String {
1327 let seq = self
1328 .run_counter
1329 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1330 format!("distill-{}-{seq}", now_ms())
1331 }
1332
1333 pub async fn distill_now(
1338 self: &Arc<Self>,
1339 identity: &str,
1340 session_key: &str,
1341 cause: DistillCause,
1342 ) -> DistillOutcome {
1343 let outcome = self.distill_inner(identity, session_key, cause).await;
1344 self.with_window(identity, session_key, |state| {
1345 state.in_flight = false;
1346 });
1347 if cause == DistillCause::Compaction {
1348 let follow_up = self
1349 .compaction_follow_up
1350 .lock()
1351 .unwrap_or_else(|err| err.into_inner())
1352 .clone();
1353 if let Some(follow_up) = follow_up {
1354 follow_up(identity, session_key, &outcome);
1355 }
1356 }
1357 match &outcome {
1358 DistillOutcome::Skipped { reason } => {
1359 tracing::debug!(
1360 identity,
1361 session_key,
1362 cause = cause.as_str(),
1363 reason,
1364 "agent memory distiller: run skipped"
1365 );
1366 }
1367 DistillOutcome::Completed {
1368 run_id,
1369 written,
1370 quarantined,
1371 } => {
1372 tracing::info!(
1373 identity,
1374 session_key,
1375 cause = cause.as_str(),
1376 run_id,
1377 written,
1378 quarantined,
1379 "agent memory distiller: run completed"
1380 );
1381 }
1382 }
1383 outcome
1384 }
1385
1386 async fn distill_inner(
1387 self: &Arc<Self>,
1388 identity: &str,
1389 session_key: &str,
1390 cause: DistillCause,
1391 ) -> DistillOutcome {
1392 let (cursor, recorder_wrote, generation) =
1393 self.with_window(identity, session_key, |state| {
1394 (state.cursor, state.recorder_wrote, state.generation)
1395 });
1396
1397 let (evidence_text, evidence_range, window_end, head_revision) = match cause {
1400 DistillCause::Compaction => {
1401 let Some(compaction) = self.compaction.as_ref() else {
1402 return DistillOutcome::Skipped {
1403 reason: SKIP_NO_DISCARD_SOURCE.to_string(),
1404 };
1405 };
1406 let entries = match compaction
1407 .read_discards(session_key, COMPACTION_HARVEST_LIMIT)
1408 .await
1409 {
1410 Ok(entries) => entries,
1411 Err(err) => {
1412 tracing::warn!(
1413 identity,
1414 session_key,
1415 error = %err,
1416 "agent memory distiller: compaction harvest failed"
1417 );
1418 return DistillOutcome::Skipped {
1419 reason: format!("compaction harvest failed: {err}"),
1420 };
1421 }
1422 };
1423 if entries.is_empty() {
1424 return DistillOutcome::Skipped {
1425 reason: SKIP_NO_DISCARDS.to_string(),
1426 };
1427 }
1428 let range = discard_evidence_range(&entries);
1429 (render_discards(&entries), range, None, None)
1432 }
1433 _ => {
1434 let slice = match self.transcripts.read(session_key, cursor).await {
1435 Ok(Some(slice)) => slice,
1436 Ok(None) => {
1437 return DistillOutcome::Skipped {
1438 reason: "session not found in the session store".to_string(),
1439 };
1440 }
1441 Err(err) => {
1442 tracing::warn!(
1443 identity,
1444 session_key,
1445 error = %err,
1446 "agent memory distiller: transcript read failed"
1447 );
1448 return DistillOutcome::Skipped {
1449 reason: format!("transcript read failed: {err}"),
1450 };
1451 }
1452 };
1453 if slice.messages.is_empty() {
1454 return DistillOutcome::Skipped {
1455 reason: "empty evidence window".to_string(),
1456 };
1457 }
1458 if cause == DistillCause::Interactions && recorder_wrote {
1459 self.with_window(identity, session_key, |state| {
1462 state.cursor = slice.end_index;
1463 state.recorder_wrote = false;
1464 state.completed_runs = 0;
1465 });
1466 return DistillOutcome::Skipped {
1467 reason: "recorder wrote in window (mutual exclusion)".to_string(),
1468 };
1469 }
1470 let range = Some((slice.start_index, slice.end_index.saturating_sub(1)));
1471 let head_revision = slice.head_revision.clone();
1472 (
1473 render_transcript(&slice),
1474 range,
1475 Some(slice.end_index),
1476 head_revision,
1477 )
1478 }
1479 };
1480
1481 let _permit = match self.budget.try_acquire(&self.realm, "distiller") {
1483 Ok(permit) => permit,
1484 Err(denied) => {
1485 return DistillOutcome::Skipped {
1486 reason: format!("budget denied: {denied}"),
1487 };
1488 }
1489 };
1490
1491 if cause.quarantines_output()
1495 && let Some(tracker) = self.tracker.as_ref()
1496 {
1497 tracker.mark_reset_boundary(session_key);
1498 }
1499
1500 let scope = self.identity_scope(identity);
1501 let manifest = match self
1502 .provider
1503 .manifest(std::slice::from_ref(&scope), ManifestTier::Full)
1504 .await
1505 {
1506 Ok(manifest) => manifest,
1507 Err(err) => {
1508 tracing::warn!(identity, error = %err, "agent memory distiller: manifest read failed");
1509 return DistillOutcome::Skipped {
1510 reason: format!("manifest read failed: {err}"),
1511 };
1512 }
1513 };
1514 let tombstones = match self
1515 .tombstones
1516 .recent_tombstones(
1517 &scope,
1518 now_ms().saturating_sub(TOMBSTONE_LOOKBACK_MS),
1519 self.profile.params.max_tombstones,
1520 )
1521 .await
1522 {
1523 Ok(tombstones) => tombstones,
1524 Err(err) => {
1525 tracing::warn!(identity, error = %err, "agent memory distiller: tombstone read failed");
1526 return DistillOutcome::Skipped {
1527 reason: format!("tombstone read failed: {err}"),
1528 };
1529 }
1530 };
1531
1532 let client = match self.handle.client().await {
1533 Ok(client) => client,
1534 Err(err) => {
1535 tracing::warn!(identity, error = %err, "agent memory distiller: client acquisition failed");
1536 return DistillOutcome::Skipped {
1537 reason: format!("client acquisition failed: {err}"),
1538 };
1539 }
1540 };
1541 let raw_ops = match extract(
1542 &self.profile,
1543 &*client,
1544 &manifest,
1545 &tombstones,
1546 &evidence_text,
1547 )
1548 .await
1549 {
1550 Ok(ops) => ops,
1551 Err(DistillerError::Auth(message)) => {
1552 tracing::warn!(error = %message, "distiller auth failure; re-resolving client");
1554 self.handle.invalidate();
1555 let retried = match self.handle.client().await {
1556 Ok(client) => {
1557 extract(
1558 &self.profile,
1559 &*client,
1560 &manifest,
1561 &tombstones,
1562 &evidence_text,
1563 )
1564 .await
1565 }
1566 Err(err) => Err(err),
1567 };
1568 match retried {
1569 Ok(ops) => ops,
1570 Err(err) => {
1571 tracing::warn!(identity, error = %err, "agent memory distiller: extraction failed");
1572 return DistillOutcome::Skipped {
1573 reason: format!("extraction failed: {err}"),
1574 };
1575 }
1576 }
1577 }
1578 Err(err) => {
1579 tracing::warn!(identity, error = %err, "agent memory distiller: extraction failed");
1580 return DistillOutcome::Skipped {
1581 reason: format!("extraction failed: {err}"),
1582 };
1583 }
1584 };
1585
1586 let manifest_ids: Vec<String> = manifest.iter().map(|meta| meta.id.clone()).collect();
1587 let run_id = self.mint_run_id();
1588 let author = MemoryAuthor::Distiller {
1589 run_id: run_id.clone(),
1590 };
1591 let mut written = 0usize;
1592 let mut quarantined = 0usize;
1593 for raw in raw_ops {
1594 let op = match validate_op(raw, &manifest_ids) {
1595 Ok(op) => op,
1596 Err(reason) => {
1597 tracing::warn!(run_id, reason, "agent memory distiller: op dropped");
1598 continue;
1599 }
1600 };
1601 let record = self.build_record(
1602 &op,
1603 session_key,
1604 generation,
1605 evidence_range,
1606 head_revision.as_deref(),
1607 );
1608 let result = match &op.action {
1609 ProposedAction::Remember => {
1610 self.provider
1611 .remember_authored(&scope, record, author.clone())
1612 .await
1613 }
1614 ProposedAction::Update { target_id } => {
1615 self.provider
1616 .supersede_authored(&scope, target_id, record, author.clone())
1617 .await
1618 }
1619 };
1620 match result {
1621 Ok(receipt) => {
1622 written += 1;
1623 if matches!(
1624 receipt.status,
1625 crate::memory::records::RecordStatus::Quarantined { .. }
1626 ) {
1627 quarantined += 1;
1628 }
1629 }
1630 Err(err) => {
1631 tracing::warn!(run_id, error = %err, "agent memory distiller: write rejected");
1635 }
1636 }
1637 }
1638
1639 self.with_window(identity, session_key, |state| {
1640 if let Some(end) = window_end {
1647 state.cursor = end;
1648 state.recorder_wrote = false;
1649 state.completed_runs = 0;
1650 }
1651 state.last_run_at = Some(Instant::now());
1652 });
1653 DistillOutcome::Completed {
1654 run_id,
1655 written,
1656 quarantined,
1657 }
1658 }
1659
1660 fn build_record(
1661 &self,
1662 op: &ProposedOp,
1663 session_key: &str,
1664 generation: u64,
1665 window_range: Option<(u64, u64)>,
1666 revision: Option<&str>,
1667 ) -> NewMemoryRecord {
1668 let mut tags = op.tags.clone();
1669 if op.epistemic == Epistemic::OperatorSaid
1670 && !tags.iter().any(|tag| tag == "epistemic:operator_said")
1671 {
1672 tags.push("epistemic:operator_said".to_string());
1675 }
1676 let range = match (op.evidence_range, window_range) {
1680 (Some((start, end)), Some((window_start, window_end)))
1681 if start >= window_start && end <= window_end =>
1682 {
1683 Some((start, end))
1684 }
1685 (_, window) => window,
1686 };
1687 NewMemoryRecord {
1688 kind: op.kind,
1689 title: op.title.clone(),
1690 description: op.description.clone(),
1691 body: op.body.clone(),
1692 tags,
1693 evidence: vec![EvidenceRef {
1694 session_id: session_key.to_string(),
1695 generation,
1696 revision: revision.map(str::to_string),
1700 range,
1701 }],
1702 verification: None,
1703 }
1704 }
1705
1706 pub async fn distill_before_rotation(
1709 self: &Arc<Self>,
1710 identity: &str,
1711 session_key: &str,
1712 cause: DistillCause,
1713 ) {
1714 let timeout = self.pre_rotation_timeout;
1715 match tokio::time::timeout(timeout, self.distill_now(identity, session_key, cause)).await {
1716 Ok(_) => {}
1717 Err(_) => {
1718 tracing::warn!(
1719 identity,
1720 session_key,
1721 cause = cause.as_str(),
1722 timeout_ms = timeout.as_millis() as u64,
1723 "agent memory distiller: pre-rotation distillation timed out; \
1724 rotation proceeds without it"
1725 );
1726 if let Some(sink) = self
1727 .events
1728 .lock()
1729 .unwrap_or_else(std::sync::PoisonError::into_inner)
1730 .as_ref()
1731 {
1732 sink.emit(
1733 crate::memory::events::MemoryTimelineEvent::DistillationTimedOut {
1734 identity: identity.to_string(),
1735 session_key: session_key.to_string(),
1736 cause: cause.as_str().to_string(),
1737 },
1738 );
1739 }
1740 }
1741 }
1742 }
1743
1744 pub async fn drop_orphaned_session_scope(
1751 self: &Arc<Self>,
1752 session_key: &str,
1753 cause: DistillCause,
1754 ) {
1755 let Some(compaction) = self.compaction.as_ref() else {
1756 return;
1757 };
1758 match compaction.drop_scope(session_key).await {
1759 Ok(0) => {}
1760 Ok(dropped) => tracing::debug!(
1761 session_key,
1762 cause = cause.as_str(),
1763 dropped,
1764 "agent memory GC: reclaimed orphaned session semantic-memory rows"
1765 ),
1766 Err(err) => tracing::warn!(
1767 session_key,
1768 cause = cause.as_str(),
1769 error = %err,
1770 "agent memory GC: dropping orphaned session scope failed; \
1771 rows remain (re-embed tax persists), rotation proceeds"
1772 ),
1773 }
1774 }
1775
1776 pub fn spawn_detached(
1780 self: &Arc<Self>,
1781 identity: &str,
1782 session_key: &str,
1783 cause: DistillCause,
1784 ) {
1785 let engine = self.clone();
1786 let identity = identity.to_string();
1787 let session_key = session_key.to_string();
1788 tokio::spawn(async move {
1789 engine.distill_now(&identity, &session_key, cause).await;
1790 });
1791 }
1792}
1793
1794fn discard_evidence_range(entries: &[DiscardEntry]) -> Option<(u64, u64)> {
1795 let mut bounds: Option<(u64, u64)> = None;
1796 for entry in entries {
1797 if let Some((start, end)) = entry.range {
1798 bounds = Some(match bounds {
1799 None => (start, end),
1800 Some((lo, hi)) => (lo.min(start), hi.max(end)),
1801 });
1802 }
1803 }
1804 bounds
1805}
1806
1807pub struct DistillerTriggers {
1813 engine: Arc<DistillerEngine>,
1814}
1815
1816impl DistillerTriggers {
1817 pub fn new(engine: Arc<DistillerEngine>) -> Self {
1818 Self { engine }
1819 }
1820}
1821
1822impl MemberAgentEventSink for DistillerTriggers {
1823 fn observe(&self, identity: &str, envelope: &meerkat_core::event::EventEnvelope<AgentEvent>) {
1824 let identity = crate::member_comms_id::logical_memory_identity(identity);
1831 let identity = identity.as_str();
1832 match &envelope.payload {
1833 AgentEvent::ToolCallRequested { name, args, .. } if name == MEMORY_TOOL_NAME => {
1834 let is_write = args
1835 .as_value()
1836 .get("action")
1837 .and_then(serde_json::Value::as_str)
1838 .is_some_and(|action| {
1839 matches!(action, "remember" | "update" | "forget" | "propose_to_mob")
1840 });
1841 if is_write && let Some(session) = current_session_of(envelope) {
1842 self.engine.note_recorder_write(identity, &session);
1843 }
1844 }
1845 AgentEvent::RunCompleted { session_id, .. } => {
1846 let session = session_id.to_string();
1847 if self.engine.note_run_completed(identity, &session) {
1848 self.engine
1849 .spawn_detached(identity, &session, DistillCause::Interactions);
1850 }
1851 }
1852 AgentEvent::CompactionCompleted { .. } => {
1853 if let Some(session) = current_session_of(envelope) {
1854 self.engine.note_compaction_observed(identity, &session);
1858 self.engine
1859 .spawn_detached(identity, &session, DistillCause::Compaction);
1860 } else {
1861 tracing::warn!(
1862 identity,
1863 "agent memory distiller: compaction event without session \
1864 attribution; harvest skipped"
1865 );
1866 }
1867 }
1868 _ => {}
1869 }
1870 }
1871}
1872
1873fn current_session_of(envelope: &meerkat_core::event::EventEnvelope<AgentEvent>) -> Option<String> {
1876 match &envelope.source {
1877 meerkat_core::event::EventSourceIdentity::Session { session_id } => {
1878 Some(session_id.to_string())
1879 }
1880 _ => None,
1881 }
1882}
1883
1884fn now_ms() -> u64 {
1885 SystemTime::now()
1886 .duration_since(UNIX_EPOCH)
1887 .map(|duration| duration.as_millis() as u64)
1888 .unwrap_or(0)
1889}
1890
1891#[cfg(test)]
1892#[allow(clippy::expect_used, clippy::redundant_clone, clippy::unwrap_used)]
1893mod tests {
1894 use super::*;
1895 use crate::identity_first::agent_memory::{
1896 AgentMemoryRecallRequest, AgentMemoryRecord, AuthoredWriteReceipt,
1897 };
1898 use crate::memory::records::RecordStatus;
1899 use futures::stream;
1900 use std::sync::Mutex as StdMutex;
1901
1902 #[tokio::test]
1908 async fn hnsw_discard_source_enumerates_then_drops_scope()
1909 -> Result<(), Box<dyn std::error::Error>> {
1910 use meerkat_core::memory::{
1911 MemoryIndexRequest, MemoryIndexScope, MemoryMetadata, MemorySource, MemoryStore,
1912 MessageRange,
1913 };
1914 use meerkat_core::types::{MemoryIndexableContent, SessionId};
1915
1916 let dir = tempfile::tempdir()?;
1917 let session = SessionId::new();
1918
1919 {
1921 let store = meerkat_memory::HnswMemoryStore::open(dir.path())?;
1922 let scope = MemoryIndexScope::for_session(session.clone());
1923 for (start, end, text) in [
1924 (0u64, 2u64, "operator prefers terse status updates"),
1925 (2u64, 5u64, "root cause was an expired upstream TLS cert"),
1926 ] {
1927 let request = MemoryIndexRequest::new(
1928 scope.clone(),
1929 MemoryIndexableContent::Indexable(text.to_string()),
1930 MemoryMetadata {
1931 session_id: session.clone(),
1932 source: MemorySource::Compaction {
1933 source_range: MessageRange::new(start, end)?,
1934 },
1935 indexed_at: meerkat_core::time_compat::SystemTime::now(),
1936 },
1937 )?;
1938 store.index_scoped(request).await?;
1939 }
1940 }
1941
1942 let source = HnswDiscardSource::new(dir.path());
1943
1944 let discards = source
1946 .read_discards(&session.to_string(), COMPACTION_HARVEST_LIMIT)
1947 .await?;
1948 assert_eq!(discards.len(), 2, "enumeration must return every scope row");
1949 assert!(discards.iter().any(|d| d.range == Some((0, 2))));
1950 assert!(
1951 discards
1952 .iter()
1953 .any(|d| d.content.contains("expired upstream TLS cert"))
1954 );
1955
1956 let dropped = source.drop_scope(&session.to_string()).await?;
1958 assert_eq!(dropped, 2, "drop_scope must reclaim every seeded row");
1959 let after = source
1960 .read_discards(&session.to_string(), COMPACTION_HARVEST_LIMIT)
1961 .await?;
1962 assert!(after.is_empty(), "dropped scope must enumerate empty");
1963 Ok(())
1964 }
1965
1966 struct ScriptedLlm {
1969 replies: StdMutex<Vec<String>>,
1970 prompts: StdMutex<Vec<String>>,
1971 }
1972
1973 impl ScriptedLlm {
1974 fn new(replies: Vec<&str>) -> Self {
1975 Self {
1976 replies: StdMutex::new(replies.into_iter().map(str::to_string).collect()),
1977 prompts: StdMutex::new(Vec::new()),
1978 }
1979 }
1980
1981 fn prompts(&self) -> Vec<String> {
1982 self.prompts
1983 .lock()
1984 .unwrap_or_else(std::sync::PoisonError::into_inner)
1985 .clone()
1986 }
1987 }
1988
1989 #[async_trait]
1990 impl LlmClient for ScriptedLlm {
1991 fn stream<'a>(&'a self, request: &'a LlmRequest) -> meerkat_client::types::LlmStream<'a> {
1992 let prompt = request
1993 .messages
1994 .iter()
1995 .map(|message| match message {
1996 Message::User(user) => user.text_content(),
1997 _ => String::new(),
1998 })
1999 .collect::<Vec<_>>()
2000 .join("\n");
2001 self.prompts
2002 .lock()
2003 .unwrap_or_else(std::sync::PoisonError::into_inner)
2004 .push(prompt);
2005 let reply = {
2006 let mut replies = self
2007 .replies
2008 .lock()
2009 .unwrap_or_else(std::sync::PoisonError::into_inner);
2010 if replies.is_empty() {
2011 String::new()
2012 } else {
2013 replies.remove(0)
2014 }
2015 };
2016 Box::pin(stream::iter(vec![
2017 Ok(LlmEvent::TextDelta {
2018 delta: reply,
2019 meta: None,
2020 }),
2021 Ok(LlmEvent::Done {
2022 outcome: LlmDoneOutcome::Success {
2023 stop_reason: meerkat_core::StopReason::EndTurn,
2024 },
2025 }),
2026 ]))
2027 }
2028
2029 fn provider(&self) -> Provider {
2030 Provider::Other
2031 }
2032
2033 async fn health_check(&self) -> Result<(), LlmError> {
2034 Ok(())
2035 }
2036 }
2037
2038 struct ScriptedHandle {
2039 client: Arc<ScriptedLlm>,
2040 }
2041
2042 #[async_trait]
2043 impl DistillerClientHandle for ScriptedHandle {
2044 async fn client(&self) -> Result<Arc<dyn LlmClient>, DistillerError> {
2045 Ok(self.client.clone())
2046 }
2047 fn invalidate(&self) {}
2048 }
2049
2050 struct HangingHandle;
2053
2054 #[async_trait]
2055 impl DistillerClientHandle for HangingHandle {
2056 async fn client(&self) -> Result<Arc<dyn LlmClient>, DistillerError> {
2057 futures::future::pending().await
2058 }
2059 fn invalidate(&self) {}
2060 }
2061
2062 #[derive(Default)]
2065 struct CapturingProvider {
2066 remembers: StdMutex<Vec<(MemoryScope, NewMemoryRecord, MemoryAuthor)>>,
2067 supersedes: StdMutex<Vec<(MemoryScope, String, NewMemoryRecord, MemoryAuthor)>>,
2068 manifest: StdMutex<Vec<RecordMeta>>,
2069 quarantine_all: bool,
2070 }
2071
2072 #[async_trait]
2073 impl AgentMemoryProvider for CapturingProvider {
2074 async fn recall(
2075 &self,
2076 _request: AgentMemoryRecallRequest,
2077 ) -> Result<Vec<AgentMemoryRecord>, AgentMemoryError> {
2078 Ok(Vec::new())
2079 }
2080
2081 async fn manifest(
2082 &self,
2083 _scopes: &[MemoryScope],
2084 _tier: ManifestTier,
2085 ) -> Result<Vec<RecordMeta>, AgentMemoryError> {
2086 Ok(self
2087 .manifest
2088 .lock()
2089 .unwrap_or_else(std::sync::PoisonError::into_inner)
2090 .clone())
2091 }
2092
2093 fn supports_manifest(&self) -> bool {
2094 true
2095 }
2096
2097 async fn remember_authored(
2098 &self,
2099 scope: &MemoryScope,
2100 record: NewMemoryRecord,
2101 author: MemoryAuthor,
2102 ) -> Result<AuthoredWriteReceipt, AgentMemoryError> {
2103 self.remembers
2104 .lock()
2105 .unwrap_or_else(std::sync::PoisonError::into_inner)
2106 .push((scope.clone(), record, author));
2107 Ok(AuthoredWriteReceipt {
2108 memory_id: "mem-new".to_string(),
2109 status: if self.quarantine_all {
2110 RecordStatus::Quarantined {
2111 reason: "test".to_string(),
2112 }
2113 } else {
2114 RecordStatus::Active
2115 },
2116 })
2117 }
2118
2119 async fn supersede_authored(
2120 &self,
2121 scope: &MemoryScope,
2122 prior: &str,
2123 record: NewMemoryRecord,
2124 author: MemoryAuthor,
2125 ) -> Result<AuthoredWriteReceipt, AgentMemoryError> {
2126 self.supersedes
2127 .lock()
2128 .unwrap_or_else(std::sync::PoisonError::into_inner)
2129 .push((scope.clone(), prior.to_string(), record, author));
2130 Ok(AuthoredWriteReceipt {
2131 memory_id: "mem-updated".to_string(),
2132 status: RecordStatus::Active,
2133 })
2134 }
2135
2136 fn supports_authored_writes(&self) -> bool {
2137 true
2138 }
2139 }
2140
2141 struct StaticTombstones(Vec<TombstoneMeta>);
2142
2143 #[async_trait]
2144 impl TombstoneSource for StaticTombstones {
2145 async fn recent_tombstones(
2146 &self,
2147 _scope: &MemoryScope,
2148 _since_ms: u64,
2149 _limit: usize,
2150 ) -> Result<Vec<TombstoneMeta>, AgentMemoryError> {
2151 Ok(self.0.clone())
2152 }
2153 }
2154
2155 struct StaticTranscript(Option<TranscriptSlice>);
2156
2157 #[async_trait]
2158 impl TranscriptSource for StaticTranscript {
2159 async fn read(
2160 &self,
2161 _session_key: &str,
2162 from_index: u64,
2163 ) -> Result<Option<TranscriptSlice>, DistillerError> {
2164 Ok(self.0.clone().map(|mut slice| {
2165 slice.messages.retain(|message| message.index >= from_index);
2166 slice.start_index = from_index.max(slice.start_index);
2167 slice
2168 }))
2169 }
2170 }
2171
2172 fn slice(messages: &[(&'static str, &str)]) -> TranscriptSlice {
2173 TranscriptSlice {
2174 session_key: "sess-1".to_string(),
2175 start_index: 0,
2176 end_index: messages.len() as u64,
2177 messages: messages
2178 .iter()
2179 .enumerate()
2180 .map(|(index, (role, text))| TranscriptMessage {
2181 index: index as u64,
2182 role,
2183 text: text.to_string(),
2184 })
2185 .collect(),
2186 head_revision: Some("rev-head".to_string()),
2187 }
2188 }
2189
2190 fn engine_with(
2191 replies: Vec<&str>,
2192 provider: Arc<CapturingProvider>,
2193 transcript: Option<TranscriptSlice>,
2194 tombstones: Vec<TombstoneMeta>,
2195 tracker: Option<SessionTaintTracker>,
2196 config: DistillerConfig,
2197 ) -> (Arc<DistillerEngine>, Arc<ScriptedLlm>) {
2198 let client = Arc::new(ScriptedLlm::new(replies));
2199 let engine = Arc::new(DistillerEngine::new(
2200 DistillerProfile::embedded_default(),
2201 config,
2202 Arc::new(ScriptedHandle {
2203 client: client.clone(),
2204 }),
2205 provider,
2206 Arc::new(StaticTombstones(tombstones)),
2207 Arc::new(StaticTranscript(transcript)),
2208 None,
2209 tracker,
2210 "family",
2211 ));
2212 (engine, client)
2213 }
2214
2215 fn enabled_config() -> DistillerConfig {
2216 DistillerConfig {
2217 enabled: true,
2218 ..DistillerConfig::default()
2219 }
2220 }
2221
2222 const NOOP_REPLY: &str = "[]";
2223 const ONE_OP_REPLY: &str = r#"[{"action": "remember", "kind": "gotcha",
2224 "title": "Cargo goes through the wrapper",
2225 "description": "When running cargo commands in this repo",
2226 "body": "Operator said: \"always use ./scripts/repo-cargo, never raw cargo\".",
2227 "tags": [], "epistemic": "operator_said", "evidence_range": [1, 2]}]"#;
2228
2229 #[test]
2232 fn prompt_renders_manifest_tombstones_and_transcript() {
2233 let profile = DistillerProfile::embedded_default();
2234 let manifest = vec![RecordMeta {
2235 id: "mem-1".to_string(),
2236 kind: MemoryKind::Gotcha,
2237 title: "Wrapper cargo".to_string(),
2238 description: "When building".to_string(),
2239 age_days: 2,
2240 rank: Some(1),
2241 }];
2242 let tombstones = vec![TombstoneMeta {
2243 title: "Operator phone number".to_string(),
2244 kind: MemoryKind::Fact,
2245 tombstoned_at_ms: 5,
2246 }];
2247 let transcript = render_transcript(&slice(&[
2248 ("user", "please run the build"),
2249 ("assistant", "running it now"),
2250 ]));
2251 let prompt = render_prompt(&profile, &manifest, &tombstones, &transcript);
2252 assert!(prompt.contains("- mem-1 [gotcha, saved 2 days ago, rank 1] Wrapper cargo"));
2253 assert!(prompt.contains("- [fact] Operator phone number"));
2254 assert!(prompt.contains("[0] user: please run the build"));
2255 assert!(prompt.contains("[1] assistant: running it now"));
2256 assert!(!prompt.contains("{{existing_manifest}}"));
2257 assert!(!prompt.contains("{{recent_tombstones}}"));
2258 assert!(!prompt.contains("{{transcript}}"));
2259 }
2260
2261 #[test]
2262 fn prompt_renders_empty_sections_honestly() {
2263 let profile = DistillerProfile::embedded_default();
2264 let prompt = render_prompt(&profile, &[], &[], "(empty)");
2265 assert!(prompt.contains("(no records)"));
2266 assert!(prompt.contains("(none)"));
2267 }
2268
2269 #[tokio::test]
2272 async fn extract_parses_ops_and_tolerates_fences() -> Result<(), Box<dyn std::error::Error>> {
2273 let fenced = format!("```json\n{ONE_OP_REPLY}\n```");
2274 let client = ScriptedLlm::new(vec![&fenced]);
2275 let profile = DistillerProfile::embedded_default();
2276 let ops = extract(&profile, &client, &[], &[], "[0] user: hi").await?;
2277 assert_eq!(ops.len(), 1);
2278 assert_eq!(client.prompts().len(), 1, "fenced JSON needs no repair");
2279 let op = validate_op(ops.into_iter().next().unwrap(), &[]).expect("valid op");
2280 assert_eq!(op.kind, MemoryKind::Gotcha);
2281 assert_eq!(op.epistemic, Epistemic::OperatorSaid);
2282 assert_eq!(op.evidence_range, Some((1, 2)));
2283 Ok(())
2284 }
2285
2286 #[tokio::test]
2287 async fn extract_repairs_malformed_output_once() -> Result<(), Box<dyn std::error::Error>> {
2288 let client = ScriptedLlm::new(vec!["Here are my thoughts, no JSON", NOOP_REPLY]);
2289 let profile = DistillerProfile::embedded_default();
2290 let ops = extract(&profile, &client, &[], &[], "[0] user: hi").await?;
2291 assert!(ops.is_empty());
2292 let prompts = client.prompts();
2293 assert_eq!(prompts.len(), 2, "exactly one repair round-trip");
2294 assert!(prompts[1].contains("ONLY the corrected JSON array"));
2295 Ok(())
2296 }
2297
2298 #[tokio::test]
2299 async fn extract_errors_after_failed_repair() {
2300 let client = ScriptedLlm::new(vec!["nope", "still nope"]);
2301 let profile = DistillerProfile::embedded_default();
2302 let result = extract(&profile, &client, &[], &[], "[0] user: hi").await;
2303 assert!(
2304 matches!(result, Err(DistillerError::Parse(_))),
2305 "{result:?}"
2306 );
2307 }
2308
2309 #[test]
2310 fn validate_op_enforces_action_kind_epistemic_and_targets() {
2311 let parse_one = |json: &str| -> Result<ProposedOp, String> {
2312 let ops = parse_ops(json).expect("parses");
2313 validate_op(
2314 ops.into_iter().next().expect("one op"),
2315 &["mem-1".to_string()],
2316 )
2317 };
2318 let err =
2320 parse_one(r#"[{"action":"remember","title":"t","body":"b","epistemic":"vibes"}]"#)
2321 .expect_err("unknown epistemic");
2322 assert!(err.contains("epistemic"), "{err}");
2323 let err = parse_one(
2325 r#"[{"action":"update","target_id":"mem-9","title":"t","body":"b","epistemic":"observed"}]"#,
2326 )
2327 .expect_err("unknown target");
2328 assert!(err.contains("not in the manifest"), "{err}");
2329 let ok = parse_one(
2330 r#"[{"action":"update","target_id":"mem-1","title":"t","body":"b","epistemic":"observed"}]"#,
2331 )
2332 .expect("valid update");
2333 assert_eq!(
2334 ok.action,
2335 ProposedAction::Update {
2336 target_id: "mem-1".to_string()
2337 }
2338 );
2339 let err = parse_one(
2341 r#"[{"action":"remember","title":"t","body":"b","epistemic":"observed","evidence_range":[9,2]}]"#,
2342 )
2343 .expect_err("inverted range");
2344 assert!(err.contains("inverted"), "{err}");
2345 }
2346
2347 #[tokio::test]
2350 async fn distill_writes_through_authored_seam_with_distiller_author_and_evidence() {
2351 let provider = Arc::new(CapturingProvider::default());
2352 let (engine, _client) = engine_with(
2353 vec![ONE_OP_REPLY],
2354 provider.clone(),
2355 Some(slice(&[
2356 ("user", "use the wrapper"),
2357 ("assistant", "noted"),
2358 ("user", "always ./scripts/repo-cargo, never raw cargo"),
2359 ])),
2360 Vec::new(),
2361 None,
2362 enabled_config(),
2363 );
2364 engine.note_session_generation("identity:a", "sess-1", 3);
2365 let outcome = engine
2366 .distill_now("identity:a", "sess-1", DistillCause::Retire)
2367 .await;
2368 assert!(
2369 matches!(outcome, DistillOutcome::Completed { written: 1, .. }),
2370 "{outcome:?}"
2371 );
2372 let remembers = provider.remembers.lock().unwrap();
2373 assert_eq!(remembers.len(), 1);
2374 let (scope, record, author) = &remembers[0];
2375 assert_eq!(
2376 *scope,
2377 MemoryScope::Identity {
2378 realm: "family".to_string(),
2379 identity: "identity:a".to_string()
2380 }
2381 );
2382 assert!(matches!(author, MemoryAuthor::Distiller { .. }));
2383 assert!(author.is_llm(), "Distiller must classify as an LLM author");
2384 assert_eq!(record.evidence.len(), 1);
2385 let evidence = &record.evidence[0];
2386 assert_eq!(evidence.session_id, "sess-1");
2387 assert_eq!(evidence.generation, 3);
2388 assert_eq!(evidence.revision.as_deref(), Some("rev-head"));
2391 assert_eq!(
2392 evidence.range,
2393 Some((1, 2)),
2394 "model-cited range within window"
2395 );
2396 assert!(record.tags.contains(&"epistemic:operator_said".to_string()));
2397 }
2398
2399 #[tokio::test]
2400 async fn hallucinated_evidence_range_falls_back_to_window_bounds() {
2401 let reply = r#"[{"action": "remember", "kind": "fact", "title": "t",
2402 "description": "", "body": "b", "tags": [],
2403 "epistemic": "observed", "evidence_range": [90, 95]}]"#;
2404 let provider = Arc::new(CapturingProvider::default());
2405 let (engine, _client) = engine_with(
2406 vec![reply],
2407 provider.clone(),
2408 Some(slice(&[("user", "a"), ("assistant", "b")])),
2409 Vec::new(),
2410 None,
2411 enabled_config(),
2412 );
2413 engine
2414 .distill_now("identity:a", "sess-1", DistillCause::Retire)
2415 .await;
2416 let remembers = provider.remembers.lock().unwrap();
2417 assert_eq!(remembers[0].1.evidence[0].range, Some((0, 1)));
2418 }
2419
2420 #[tokio::test]
2421 async fn update_ops_supersede_the_target_record() {
2422 let reply = r#"[{"action": "update", "target_id": "mem-1", "kind": "gotcha",
2423 "title": "t", "description": "", "body": "b", "tags": [],
2424 "epistemic": "observed"}]"#;
2425 let provider = Arc::new(CapturingProvider::default());
2426 provider.manifest.lock().unwrap().push(RecordMeta {
2427 id: "mem-1".to_string(),
2428 kind: MemoryKind::Gotcha,
2429 title: "old".to_string(),
2430 description: String::new(),
2431 age_days: 1,
2432 rank: None,
2433 });
2434 let (engine, _client) = engine_with(
2435 vec![reply],
2436 provider.clone(),
2437 Some(slice(&[("user", "a")])),
2438 Vec::new(),
2439 None,
2440 enabled_config(),
2441 );
2442 engine
2443 .distill_now("identity:a", "sess-1", DistillCause::Retire)
2444 .await;
2445 let supersedes = provider.supersedes.lock().unwrap();
2446 assert_eq!(supersedes.len(), 1);
2447 assert_eq!(supersedes[0].1, "mem-1");
2448 assert!(provider.remembers.lock().unwrap().is_empty());
2449 }
2450
2451 #[tokio::test]
2454 async fn recorder_write_in_window_skips_extraction_and_advances_cursor() {
2455 let provider = Arc::new(CapturingProvider::default());
2456 let (engine, client) = engine_with(
2457 vec![ONE_OP_REPLY, ONE_OP_REPLY],
2458 provider.clone(),
2459 Some(slice(&[("user", "a"), ("assistant", "b")])),
2460 Vec::new(),
2461 None,
2462 enabled_config(),
2463 );
2464 engine.note_recorder_write("identity:a", "sess-1");
2465 let outcome = engine
2466 .distill_now("identity:a", "sess-1", DistillCause::Interactions)
2467 .await;
2468 assert!(
2469 matches!(&outcome, DistillOutcome::Skipped { reason } if reason.contains("mutual exclusion")),
2470 "{outcome:?}"
2471 );
2472 assert!(
2473 client.prompts().is_empty(),
2474 "no LLM call in a skipped window"
2475 );
2476 assert!(provider.remembers.lock().unwrap().is_empty());
2477
2478 let outcome = engine
2481 .distill_now("identity:a", "sess-1", DistillCause::Interactions)
2482 .await;
2483 assert!(
2484 matches!(&outcome, DistillOutcome::Skipped { reason } if reason.contains("empty")),
2485 "{outcome:?}"
2486 );
2487
2488 let (engine, _client) = engine_with(
2491 vec![ONE_OP_REPLY],
2492 provider.clone(),
2493 Some(slice(&[("user", "a")])),
2494 Vec::new(),
2495 None,
2496 enabled_config(),
2497 );
2498 engine.note_recorder_write("identity:a", "sess-1");
2499 let outcome = engine
2500 .distill_now("identity:a", "sess-1", DistillCause::Retire)
2501 .await;
2502 assert!(
2503 matches!(outcome, DistillOutcome::Completed { .. }),
2504 "{outcome:?}"
2505 );
2506 }
2507
2508 #[test]
2511 fn interaction_trigger_respects_min_interactions_and_in_flight() {
2512 let provider = Arc::new(CapturingProvider::default());
2513 let (engine, _client) = engine_with(
2514 vec![],
2515 provider,
2516 None,
2517 Vec::new(),
2518 None,
2519 DistillerConfig {
2520 enabled: true,
2521 min_interactions: 3,
2522 ..DistillerConfig::default()
2523 },
2524 );
2525 assert!(!engine.note_run_completed("identity:a", "sess-1"));
2526 assert!(!engine.note_run_completed("identity:a", "sess-1"));
2527 assert!(
2528 engine.note_run_completed("identity:a", "sess-1"),
2529 "third run trips"
2530 );
2531 assert!(!engine.note_run_completed("identity:a", "sess-1"));
2533 }
2534
2535 #[tokio::test]
2536 async fn budget_guard_skips_runs_beyond_the_window_cap() {
2537 let provider = Arc::new(CapturingProvider::default());
2538 let (engine, _client) = engine_with(
2539 vec![NOOP_REPLY, NOOP_REPLY],
2540 provider,
2541 Some(slice(&[("user", "a")])),
2542 Vec::new(),
2543 None,
2544 DistillerConfig {
2545 enabled: true,
2546 runs_per_hour: 1,
2547 ..DistillerConfig::default()
2548 },
2549 );
2550 let first = engine
2551 .distill_now("identity:a", "sess-1", DistillCause::Retire)
2552 .await;
2553 assert!(
2554 matches!(first, DistillOutcome::Completed { .. }),
2555 "{first:?}"
2556 );
2557 engine.with_window("identity:a", "sess-1", |state| state.cursor = 0);
2559 let second = engine
2560 .distill_now("identity:a", "sess-1", DistillCause::Retire)
2561 .await;
2562 assert!(
2563 matches!(&second, DistillOutcome::Skipped { reason } if reason.contains("budget denied")),
2564 "{second:?}"
2565 );
2566 }
2567
2568 fn envelope(
2571 session: &meerkat_core::types::SessionId,
2572 payload: AgentEvent,
2573 ) -> meerkat_core::event::EventEnvelope<AgentEvent> {
2574 meerkat_core::event::EventEnvelope {
2575 event_id: Default::default(),
2576 source: meerkat_core::event::EventSourceIdentity::Session {
2577 session_id: session.clone(),
2578 },
2579 seq: 0,
2580 mob_id: None,
2581 timestamp_ms: 0,
2582 payload,
2583 }
2584 }
2585
2586 #[tokio::test]
2587 async fn trigger_sink_marks_memory_tool_writes_but_not_reads() {
2588 let provider = Arc::new(CapturingProvider::default());
2589 let session = meerkat_core::types::SessionId::new();
2590 let session_key = session.to_string();
2591 let make = |transcript: TranscriptSlice| {
2592 engine_with(
2593 vec![NOOP_REPLY],
2594 provider.clone(),
2595 Some(transcript),
2596 Vec::new(),
2597 None,
2598 enabled_config(),
2599 )
2600 };
2601 let tool_call = |action: &str| AgentEvent::ToolCallRequested {
2602 id: "t-1".to_string(),
2603 name: MEMORY_TOOL_NAME.to_string(),
2604 args: meerkat_core::event::ToolCallArguments::from_value(serde_json::json!({
2605 "action": action, "title": "t", "body": "b"
2606 }))
2607 .expect("object args"),
2608 };
2609
2610 let mut transcript = slice(&[("user", "a")]);
2612 transcript.session_key = session_key.clone();
2613 let (engine, _client) = make(transcript.clone());
2614 let sink = DistillerTriggers::new(engine.clone());
2615 sink.observe("identity:a", &envelope(&session, tool_call("remember")));
2616 let outcome = engine
2617 .distill_now("identity:a", &session_key, DistillCause::Interactions)
2618 .await;
2619 assert!(
2620 matches!(&outcome, DistillOutcome::Skipped { reason } if reason.contains("mutual exclusion")),
2621 "{outcome:?}"
2622 );
2623
2624 let (engine, _client) = make(transcript);
2626 let sink = DistillerTriggers::new(engine.clone());
2627 sink.observe("identity:a", &envelope(&session, tool_call("recall")));
2628 let outcome = engine
2629 .distill_now("identity:a", &session_key, DistillCause::Interactions)
2630 .await;
2631 assert!(
2632 matches!(outcome, DistillOutcome::Completed { .. }),
2633 "{outcome:?}"
2634 );
2635 }
2636
2637 #[tokio::test]
2640 async fn reset_cause_marks_boundary_so_evidence_quarantines() {
2641 let tracker = SessionTaintTracker::new(Default::default());
2642 let provider = Arc::new(CapturingProvider::default());
2643 let (engine, _client) = engine_with(
2644 vec![ONE_OP_REPLY],
2645 provider,
2646 Some(slice(&[("user", "a"), ("assistant", "b"), ("user", "c")])),
2647 Vec::new(),
2648 Some(tracker.clone()),
2649 enabled_config(),
2650 );
2651 assert!(tracker.evidence_quarantine_reason("sess-1").is_none());
2652 engine
2653 .distill_now("identity:a", "sess-1", DistillCause::Reset)
2654 .await;
2655 let reason = tracker
2656 .evidence_quarantine_reason("sess-1")
2657 .expect("reset boundary marked before writes");
2658 assert!(reason.contains("reset"), "{reason}");
2659 }
2660
2661 #[tokio::test]
2664 async fn pre_rotation_distillation_never_blocks_rotation() {
2665 let provider = Arc::new(CapturingProvider::default());
2666 let engine = Arc::new(
2667 DistillerEngine::new(
2668 DistillerProfile::embedded_default(),
2669 enabled_config(),
2670 Arc::new(HangingHandle),
2671 provider,
2672 Arc::new(StaticTombstones(Vec::new())),
2673 Arc::new(StaticTranscript(Some(slice(&[("user", "a")])))),
2674 None,
2675 None,
2676 "family",
2677 )
2678 .with_pre_rotation_timeout(Duration::from_millis(50)),
2679 );
2680 let started = Instant::now();
2681 engine
2682 .distill_before_rotation("identity:a", "sess-1", DistillCause::Respawn)
2683 .await;
2684 assert!(
2685 started.elapsed() < Duration::from_secs(2),
2686 "pre-rotation hook must return at the timeout, not hang"
2687 );
2688 }
2689
2690 #[test]
2693 fn embedded_prompt_matches_calibration_bundle() -> Result<(), Box<dyn std::error::Error>> {
2694 let bundle =
2695 Path::new(env!("CARGO_MANIFEST_DIR")).join("../memory-evals/prompts/distiller-v0.md");
2696 if !bundle.is_file() {
2697 return Ok(());
2698 }
2699 let text = std::fs::read_to_string(bundle)?;
2700 assert_eq!(
2701 text, EMBEDDED_PROMPT_V0,
2702 "memory-evals/prompts/distiller-v0.md and \
2703 src/memory/distiller_prompt_v0.md have drifted"
2704 );
2705 Ok(())
2706 }
2707
2708 #[test]
2709 fn embedded_default_profile_validates_and_names_a_catalog_model() {
2710 let profile = DistillerProfile::embedded_default();
2711 profile.validate().expect("embedded profile must validate");
2712 assert_eq!(
2713 meerkat_models::infer_provider(&profile.model),
2714 Some(profile.provider),
2715 "embedded default model must resolve in the catalog"
2716 );
2717 }
2718
2719 #[test]
2720 fn external_profile_loads_from_evals_layout() -> Result<(), Box<dyn std::error::Error>> {
2721 let path = Path::new(env!("CARGO_MANIFEST_DIR"))
2722 .join("../memory-evals/profiles/distiller-v0.toml");
2723 if !path.is_file() {
2724 return Ok(());
2725 }
2726 let profile = DistillerProfile::load(&path)?;
2727 assert_eq!(profile.stage, "distiller");
2728 assert_eq!(profile.prompt_template, EMBEDDED_PROMPT_V0);
2729 Ok(())
2730 }
2731
2732 #[test]
2733 fn model_override_is_fail_loud() {
2734 let profile = DistillerProfile::embedded_default();
2735 assert!(profile.clone().with_model_override("not-a-model").is_err());
2736 assert!(profile.clone().with_model_override(" ").is_err());
2737 let overridden = profile
2738 .with_model_override("claude-haiku-4-5")
2739 .expect("catalog model accepted");
2740 assert_eq!(overridden.model, "claude-haiku-4-5");
2741 }
2742
2743 #[test]
2746 fn compaction_harvest_satisfaction_classifies_skip_reasons() {
2747 assert!(
2748 DistillOutcome::Completed {
2749 run_id: "r".to_string(),
2750 written: 0,
2751 quarantined: 0,
2752 }
2753 .compaction_harvest_satisfied()
2754 );
2755 assert!(
2756 DistillOutcome::Skipped {
2757 reason: SKIP_NO_DISCARDS.to_string(),
2758 }
2759 .compaction_harvest_satisfied()
2760 );
2761 assert!(
2762 DistillOutcome::Skipped {
2763 reason: SKIP_NO_DISCARD_SOURCE.to_string(),
2764 }
2765 .compaction_harvest_satisfied()
2766 );
2767 for reason in [
2768 "budget denied: window budget exhausted (2/2 runs)",
2769 "compaction harvest failed: io",
2770 "extraction failed: parse",
2771 ] {
2772 assert!(
2773 !DistillOutcome::Skipped {
2774 reason: reason.to_string(),
2775 }
2776 .compaction_harvest_satisfied(),
2777 "{reason}"
2778 );
2779 }
2780 }
2781
2782 #[tokio::test]
2783 async fn compaction_follow_up_fires_with_outcome_and_cursor_is_readable() {
2784 let provider = Arc::new(CapturingProvider::default());
2785 let (engine, _client) = engine_with(
2788 vec![NOOP_REPLY],
2789 provider,
2790 Some(slice(&[("user", "hello")])),
2791 Vec::new(),
2792 None,
2793 enabled_config(),
2794 );
2795 let seen: Arc<StdMutex<Vec<(String, String, bool)>>> = Arc::new(StdMutex::new(Vec::new()));
2796 let sink = seen.clone();
2797 engine.set_compaction_follow_up(Arc::new(move |identity, session, outcome| {
2798 sink.lock()
2799 .unwrap_or_else(std::sync::PoisonError::into_inner)
2800 .push((
2801 identity.to_string(),
2802 session.to_string(),
2803 outcome.compaction_harvest_satisfied(),
2804 ));
2805 }));
2806 engine
2807 .distill_now("identity:a", "sess-1", DistillCause::Compaction)
2808 .await;
2809 engine
2811 .distill_now("identity:a", "sess-1", DistillCause::Interactions)
2812 .await;
2813 let seen = seen
2814 .lock()
2815 .unwrap_or_else(std::sync::PoisonError::into_inner)
2816 .clone();
2817 assert_eq!(
2818 seen,
2819 vec![("identity:a".to_string(), "sess-1".to_string(), true)]
2820 );
2821 assert_eq!(engine.distilled_cursor("identity:never", "sess-x"), 0);
2823 engine.with_window("identity:a", "sess-1", |state| state.cursor = 7);
2824 assert_eq!(engine.distilled_cursor("identity:a", "sess-1"), 7);
2825 }
2826
2827 struct StaticDiscards(Vec<DiscardEntry>);
2828
2829 #[async_trait]
2830 impl CompactionDiscardSource for StaticDiscards {
2831 async fn read_discards(
2832 &self,
2833 _session_key: &str,
2834 _limit: usize,
2835 ) -> Result<Vec<DiscardEntry>, DistillerError> {
2836 Ok(self.0.clone())
2837 }
2838 }
2839
2840 #[tokio::test]
2841 async fn compaction_harvest_preserves_interaction_window_state() {
2842 let provider = Arc::new(CapturingProvider::default());
2843 let client = Arc::new(ScriptedLlm::new(vec![ONE_OP_REPLY]));
2844 let engine = Arc::new(DistillerEngine::new(
2845 DistillerProfile::embedded_default(),
2846 enabled_config(),
2847 Arc::new(ScriptedHandle { client }),
2848 provider,
2849 Arc::new(StaticTombstones(Vec::new())),
2850 Arc::new(StaticTranscript(Some(slice(&[
2851 ("user", "use the wrapper"),
2852 ("assistant", "noted"),
2853 ])))),
2854 Some(Arc::new(StaticDiscards(vec![DiscardEntry {
2855 content: "discarded: wrapper reminder".to_string(),
2856 range: Some((0, 5)),
2857 }]))),
2858 None,
2859 "family",
2860 ));
2861 engine.note_recorder_write("identity:a", "sess-1");
2864 assert!(!engine.note_run_completed("identity:a", "sess-1"));
2865
2866 let outcome = engine
2868 .distill_now("identity:a", "sess-1", DistillCause::Compaction)
2869 .await;
2870 assert!(
2871 matches!(outcome, DistillOutcome::Completed { written: 1, .. }),
2872 "{outcome:?}"
2873 );
2874
2875 let (cursor, recorder_wrote, completed_runs) =
2879 engine.with_window("identity:a", "sess-1", |state| {
2880 (state.cursor, state.recorder_wrote, state.completed_runs)
2881 });
2882 assert_eq!(cursor, 0, "compaction never advances the transcript cursor");
2883 assert!(
2884 recorder_wrote,
2885 "compaction-cause harvests must not clear the recorder \
2886 mutual-exclusion flag for a window they did not distill"
2887 );
2888 assert_eq!(
2889 completed_runs, 1,
2890 "compaction-cause harvests must not zero the interaction counter"
2891 );
2892
2893 let outcome = engine
2896 .distill_now("identity:a", "sess-1", DistillCause::Interactions)
2897 .await;
2898 assert!(
2899 matches!(&outcome, DistillOutcome::Skipped { reason } if reason.contains("mutual exclusion")),
2900 "{outcome:?}"
2901 );
2902 }
2903
2904 #[tokio::test]
2911 async fn trigger_sink_normalizes_roster_ids_to_the_logical_identity_across_generations() {
2912 let provider = Arc::new(CapturingProvider::default());
2913 let (engine, _client) = engine_with(
2914 vec![NOOP_REPLY],
2915 provider,
2916 Some(slice(&[("user", "hello")])),
2917 Vec::new(),
2918 None,
2919 enabled_config(),
2920 );
2921 let seen: Arc<StdMutex<Vec<(String, String)>>> = Arc::new(StdMutex::new(Vec::new()));
2922 let observed = seen.clone();
2923 engine.set_compaction_observed(Arc::new(move |identity, session| {
2924 observed
2925 .lock()
2926 .unwrap_or_else(std::sync::PoisonError::into_inner)
2927 .push((identity.to_string(), session.to_string()));
2928 }));
2929
2930 let session = meerkat_core::types::SessionId::new();
2931 let sink = DistillerTriggers::new(engine.clone());
2932 for generation in ["rt:identity:parent-1:0", "rt:identity:parent-1:1"] {
2935 let roster_id = crate::member_comms_id::mob_member_id_str(generation).into_owned();
2936 assert!(roster_id.starts_with("mk--"), "{roster_id}");
2937 sink.observe(
2938 &roster_id,
2939 &envelope(
2940 &session,
2941 AgentEvent::CompactionCompleted {
2942 summary_tokens: 10,
2943 messages_before: 20,
2944 messages_after: 2,
2945 },
2946 ),
2947 );
2948 }
2949 let identities: Vec<String> = seen
2950 .lock()
2951 .unwrap_or_else(std::sync::PoisonError::into_inner)
2952 .iter()
2953 .map(|(identity, _)| identity.clone())
2954 .collect();
2955 assert_eq!(
2956 identities,
2957 vec![
2958 "identity:parent-1".to_string(),
2959 "identity:parent-1".to_string()
2960 ],
2961 "both generations must key the ONE logical scope"
2962 );
2963 }
2964
2965 #[tokio::test]
2966 async fn compaction_observed_hook_fires_at_observation_even_without_discards() {
2967 let provider = Arc::new(CapturingProvider::default());
2968 let (engine, _client) = engine_with(
2972 vec![NOOP_REPLY],
2973 provider,
2974 Some(slice(&[("user", "hello")])),
2975 Vec::new(),
2976 None,
2977 enabled_config(),
2978 );
2979 let seen: Arc<StdMutex<Vec<(String, String)>>> = Arc::new(StdMutex::new(Vec::new()));
2980 let observed = seen.clone();
2981 engine.set_compaction_observed(Arc::new(move |identity, session| {
2982 observed
2983 .lock()
2984 .unwrap_or_else(std::sync::PoisonError::into_inner)
2985 .push((identity.to_string(), session.to_string()));
2986 }));
2987
2988 let session = meerkat_core::types::SessionId::new();
2989 let sink = DistillerTriggers::new(engine.clone());
2990 sink.observe(
2991 "identity:a",
2992 &envelope(
2993 &session,
2994 AgentEvent::CompactionCompleted {
2995 summary_tokens: 10,
2996 messages_before: 20,
2997 messages_after: 2,
2998 },
2999 ),
3000 );
3001 assert_eq!(
3004 seen.lock()
3005 .unwrap_or_else(std::sync::PoisonError::into_inner)
3006 .clone(),
3007 vec![("identity:a".to_string(), session.to_string())]
3008 );
3009
3010 sink.observe(
3012 "identity:a",
3013 &envelope(
3014 &session,
3015 AgentEvent::RunCompleted {
3016 session_id: session.clone(),
3017 result: "done".to_string(),
3018 structured_output: None,
3019 extraction_required: false,
3020 usage: Default::default(),
3021 terminal_cause_kind: None,
3022 },
3023 ),
3024 );
3025 assert_eq!(
3026 seen.lock()
3027 .unwrap_or_else(std::sync::PoisonError::into_inner)
3028 .len(),
3029 1
3030 );
3031 }
3032}