Skip to main content

navi_core/memory/
maintenance.rs

1use crate::memory::history_store::HistoryStore;
2use crate::memory::{AutoMemoryStore, GlobalMemoryStore};
3use crate::model::{ModelMessage, ModelProvider, ModelRequest, ThinkingConfig};
4use anyhow::{Context, Result};
5use serde::Deserialize;
6use std::path::PathBuf;
7use std::time::{SystemTime, UNIX_EPOCH};
8
9pub const DREAM_SYSTEM: &str = r#"You are a memory maintenance subagent named NAVI Dream.
10Reflect on persistent memory and recent sessions, then produce a consolidated memory store.
11
12Offline synthesis rules:
13- Do not merely append a transcript summary.
14- Merge duplicates.
15- Resolve contradictions by preferring the newest verified session evidence.
16- Drop stale, temporary, speculative, or one-off debugging notes.
17- Preserve stable project architecture, commands, conventions, user preferences, and reusable lessons.
18- Surface new durable insights for future sessions.
19
20Output ONLY these XML blocks:
21<updated_project_memory>...</updated_project_memory>
22<updated_global_memory>...</updated_global_memory>
23<dream_report>Briefly list what changed, what was removed, and notable unresolved contradictions.</dream_report>
24"#;
25
26pub const DREAM_PROMPT: &str = r#"Existing project memory index:
27{project_memory}
28
29Existing global memory index:
30{global_memory}
31
32Current checkpoint:
33{checkpoint}
34
35Current notes:
36{notes}
37
38Recent sessions:
39{recent_sessions}
40
41Additional dream instructions:
42{instructions}
43"#;
44
45pub const DISTILL_SYSTEM: &str = r#"You are a process distillation subagent named NAVI Distill.
46Analyze recent conversation histories and extract reusable processes (SOPs, skills, checklists).
47Identify repeated successful patterns, workflows, checklists, or setups.
48Generate a reusable SOP in Markdown.
49Output ONLY inside a `<sop_artifact filename="name.md">...</sop_artifact>` block."#;
50
51pub const DISTILL_PROMPT: &str = r#"Recent Session History:
52{recent_history}
53"#;
54
55pub const MEMORY_CONSOLIDATION_SYSTEM: &str = r#"You are a memory consolidation subagent for NAVI.
56Review SQLite-stored memories and produce consolidation actions:
571. obsolete — contradicted, irrelevant, or superseded
582. merge — duplicates; surviving id + combined body
593. update — adjust confidence (0.0–1.0)
60
61Each action: { "action": "obsolete"|"merge"|"update", "id": "...", "merged_body"?: "...", "confidence"?: 0.0-1.0 }
62If none needed, return []. Output ONLY the JSON array — no markdown fences."#;
63
64pub const MEMORY_CONSOLIDATION_PROMPT: &str = r#"Current memories (JSON array):
65{memories_json}
66"#;
67
68/// A consolidation action from the model.
69#[derive(Debug, Clone, Deserialize)]
70struct ConsolidationAction {
71    action: String,
72    id: String,
73    merged_body: Option<String>,
74    confidence: Option<f64>,
75}
76
77fn sanitize_input(text: &str) -> String {
78    text.replace("<updated_project_memory>", "[updated_project_memory]")
79        .replace("</updated_project_memory>", "[/updated_project_memory]")
80        .replace("<updated_global_memory>", "[updated_global_memory]")
81        .replace("</updated_global_memory>", "[/updated_global_memory]")
82        .replace("<dream_report>", "[dream_report]")
83        .replace("</dream_report>", "[/dream_report]")
84        .replace("<sop_artifact", "[sop_artifact")
85        .replace("</sop_artifact>", "[/sop_artifact]")
86}
87
88#[derive(Debug, Clone)]
89pub struct DreamOptions {
90    pub session_limit: usize,
91    pub instructions: Option<String>,
92    pub apply: bool,
93}
94
95impl Default for DreamOptions {
96    fn default() -> Self {
97        Self {
98            session_limit: 10,
99            instructions: None,
100            apply: false,
101        }
102    }
103}
104
105#[derive(Debug, Clone)]
106pub struct DreamResult {
107    pub output_dir: PathBuf,
108    pub project_memory_path: PathBuf,
109    pub global_memory_path: PathBuf,
110    pub report_path: PathBuf,
111    pub applied: bool,
112    pub auto_memory_report: Option<crate::memory::ConsolidationReport>,
113}
114
115pub async fn run_dream_maintenance(
116    auto_memory: &AutoMemoryStore,
117    global_memory: &GlobalMemoryStore,
118    history_store: &HistoryStore,
119    model_provider: &dyn ModelProvider,
120    model_name: &str,
121) -> Result<DreamResult> {
122    run_dream_maintenance_with_options(
123        auto_memory,
124        global_memory,
125        history_store,
126        model_provider,
127        model_name,
128        DreamOptions::default(),
129    )
130    .await
131}
132
133pub async fn run_dream_maintenance_with_options(
134    auto_memory: &AutoMemoryStore,
135    global_memory: &GlobalMemoryStore,
136    history_store: &HistoryStore,
137    model_provider: &dyn ModelProvider,
138    model_name: &str,
139    options: DreamOptions,
140) -> Result<DreamResult> {
141    let project_memory = sanitize_input(&auto_memory.render_index());
142    let global_memory_text = sanitize_input(&global_memory.read_index().unwrap_or_default());
143    let checkpoint = sanitize_input(&auto_memory.read_checkpoint().unwrap_or_default());
144    let notes = sanitize_input(&auto_memory.read_notes().unwrap_or_default());
145    let recent_sessions = sanitize_input(&format_recent_sessions(
146        history_store,
147        options.session_limit.clamp(1, 100),
148    )?);
149    let instructions = sanitize_input(options.instructions.as_deref().unwrap_or(
150        "Focus on stable coding workflow, project architecture, commands, and user preferences.",
151    ));
152
153    let prompt = DREAM_PROMPT
154        .replace("{project_memory}", &project_memory)
155        .replace("{global_memory}", &global_memory_text)
156        .replace("{checkpoint}", &checkpoint)
157        .replace("{notes}", &notes)
158        .replace("{recent_sessions}", &recent_sessions)
159        .replace("{instructions}", &instructions);
160
161    let request = ModelRequest {
162        model: model_name.to_string(),
163        instructions: None,
164        messages: vec![
165            ModelMessage::system(DREAM_SYSTEM),
166            ModelMessage::user(prompt),
167        ],
168        thinking: ThinkingConfig::Off,
169        tools: vec![],
170        session_id: None,
171    };
172
173    let response = model_provider.complete(request).await?;
174    let text = response.text;
175
176    let updated_pm = extract_block(
177        &text,
178        "<updated_project_memory>",
179        "</updated_project_memory>",
180    )
181    .context("dream response did not include <updated_project_memory>")?;
182    let updated_gm = extract_block(&text, "<updated_global_memory>", "</updated_global_memory>")
183        .context("dream response did not include <updated_global_memory>")?;
184    let dream_report = extract_block(&text, "<dream_report>", "</dream_report>")
185        .unwrap_or_else(|| "Dream completed without a report.".to_string());
186
187    if updated_pm.trim().is_empty() {
188        anyhow::bail!("dream response produced empty project memory");
189    }
190    if updated_gm.trim().is_empty() {
191        anyhow::bail!("dream response produced empty global memory");
192    }
193
194    let output_dir = dream_output_dir(auto_memory)?;
195    std::fs::create_dir_all(&output_dir)
196        .with_context(|| format!("failed to create {}", output_dir.display()))?;
197
198    let project_memory_path = output_dir.join("project-memory.md");
199    let global_memory_path = output_dir.join("global-memory.md");
200    let report_path = output_dir.join("dream-report.md");
201    crate::memory::memory_store::write_atomic(&project_memory_path, updated_pm.trim())?;
202    crate::memory::memory_store::write_atomic(&global_memory_path, updated_gm.trim())?;
203    crate::memory::memory_store::write_atomic(&report_path, dream_report.trim())?;
204
205    if options.apply {
206        global_memory.write_from_markdown(updated_gm.trim())?;
207
208        // Model-based SQLite memory consolidation
209        if let Err(e) = run_model_based_consolidation(auto_memory, model_provider, model_name).await
210        {
211            tracing::warn!("model-based memory consolidation failed: {}", e);
212        }
213    }
214
215    // Consolidate auto-memory SQLite store (mechanical: mark stale, deduplicate, backfill embeddings)
216    let auto_memory_report = {
217        match auto_memory.consolidate(30) {
218            Ok(report) => {
219                tracing::info!(
220                    "auto-memory consolidation: {} stale, {} duplicates, {} active",
221                    report.marked_stale,
222                    report.duplicates_merged,
223                    report.remaining_active
224                );
225
226                // Backfill embeddings for memories without them
227                if crate::memory::embeddings_available() {
228                    let db_path = &auto_memory.db_path;
229                    let models_dir = db_path
230                        .parent()
231                        .unwrap_or(std::path::Path::new("."))
232                        .join("models");
233                    let model_path = models_dir.join(crate::memory::DEFAULT_MODEL_FILE);
234                    let tokenizer_path = models_dir.join(crate::memory::DEFAULT_TOKENIZER_FILE);
235
236                    if let Some(embedder) =
237                        crate::memory::embedding::get_cached_embedder(&model_path, &tokenizer_path)
238                    {
239                        let missing = auto_memory.list_without_embeddings().unwrap_or_default();
240
241                        if !missing.is_empty() {
242                            tracing::info!("backfilling embeddings for {} memories", missing.len());
243                            for m in &missing {
244                                if let Some(text) =
245                                    auto_memory.get_memory_text(&m.id).unwrap_or(None)
246                                {
247                                    match embedder.embed(&text) {
248                                        Ok(emb) => {
249                                            let _ = auto_memory.set_embedding(&m.id, &emb);
250                                        }
251                                        Err(e) => {
252                                            tracing::debug!(
253                                                "embedding backfill failed for {}: {}",
254                                                m.id,
255                                                e
256                                            );
257                                        }
258                                    }
259                                }
260                            }
261                        }
262                    }
263                }
264
265                Some(report)
266            }
267            Err(e) => {
268                tracing::warn!("auto-memory consolidation failed: {}", e);
269                None
270            }
271        }
272    };
273
274    Ok(DreamResult {
275        output_dir,
276        project_memory_path,
277        global_memory_path,
278        report_path,
279        applied: options.apply,
280        auto_memory_report,
281    })
282}
283
284/// Runs model-based consolidation on the SQLite auto-memory store.
285///
286/// Sends all active memories (with full body text) to the model, which returns
287/// consolidation actions: mark obsolete, merge duplicates, update confidence.
288/// Actions are applied directly to the SQLite store.
289async fn run_model_based_consolidation(
290    auto_memory: &AutoMemoryStore,
291    model_provider: &dyn ModelProvider,
292    model_name: &str,
293) -> Result<()> {
294    let entries = auto_memory.list_full_entries()?;
295    if entries.is_empty() {
296        return Ok(());
297    }
298
299    let memories_json = serde_json::to_string_pretty(
300        &entries
301            .iter()
302            .map(|e| {
303                serde_json::json!({
304                    "id": e.id,
305                    "type": e.memory_type.as_str(),
306                    "name": e.name,
307                    "description": e.description,
308                    "body": e.body,
309                    "confidence": e.confidence,
310                })
311            })
312            .collect::<Vec<_>>(),
313    )?;
314
315    let prompt = MEMORY_CONSOLIDATION_PROMPT.replace("{memories_json}", &memories_json);
316
317    let request = ModelRequest {
318        model: model_name.to_string(),
319        instructions: None,
320        messages: vec![
321            ModelMessage::system(MEMORY_CONSOLIDATION_SYSTEM),
322            ModelMessage::user(prompt),
323        ],
324        thinking: ThinkingConfig::Off,
325        tools: vec![],
326        session_id: None,
327    };
328
329    let response = model_provider.complete(request).await?;
330    let text = response.text.trim();
331
332    let actions: Vec<ConsolidationAction> = if text.starts_with('[') {
333        serde_json::from_str(text).unwrap_or_default()
334    } else if let Some(start) = text.find('[') {
335        if let Some(end) = text.rfind(']') {
336            serde_json::from_str(&text[start..=end]).unwrap_or_default()
337        } else {
338            Vec::new()
339        }
340    } else {
341        Vec::new()
342    };
343
344    let mut applied = 0;
345    for action in &actions {
346        let result = match action.action.as_str() {
347            "obsolete" => auto_memory.mark_obsolete(&action.id),
348            "merge" => {
349                if let Some(ref body) = action.merged_body {
350                    auto_memory.update_consolidated(&action.id, Some(body), None)
351                } else {
352                    Ok(())
353                }
354            }
355            "update" => auto_memory.update_consolidated(&action.id, None, action.confidence),
356            _ => Ok(()),
357        };
358        if result.is_ok() {
359            applied += 1;
360        }
361    }
362
363    if applied > 0 {
364        tracing::info!("model-based consolidation: applied {} actions", applied);
365    }
366
367    Ok(())
368}
369
370pub async fn run_distill_maintenance(
371    auto_memory: &AutoMemoryStore,
372    history_store: &HistoryStore,
373    model_provider: &dyn ModelProvider,
374    model_name: &str,
375) -> Result<()> {
376    let sessions = history_store.list_sessions()?;
377    let mut history_text = String::new();
378    for session in sessions.iter().take(3) {
379        let events = history_store.get_recent_events(&session.id, Some(50))?;
380        for e in events {
381            if let Some(ref content) = e.content {
382                history_text.push_str(&format!("[Session {}]: {}\n", session.id, content));
383            }
384        }
385    }
386    let history_text = sanitize_input(&history_text);
387
388    let prompt = DISTILL_PROMPT.replace("{recent_history}", &history_text);
389
390    let request = ModelRequest {
391        model: model_name.to_string(),
392        instructions: None,
393        messages: vec![
394            ModelMessage::system(DISTILL_SYSTEM),
395            ModelMessage::user(prompt),
396        ],
397        thinking: ThinkingConfig::Off,
398        tools: vec![],
399        session_id: None,
400    };
401
402    let response = model_provider.complete(request).await?;
403    let text = response.text;
404
405    if let Some(sop_block) = extract_block_with_attr(&text, "<sop_artifact", "</sop_artifact>") {
406        let (filename, content) = sop_block;
407        if !content.trim().is_empty() {
408            let sops_dir = auto_memory
409                .db_path
410                .parent()
411                .unwrap_or(std::path::Path::new("."))
412                .join("sops");
413            if !sops_dir.exists() {
414                std::fs::create_dir_all(&sops_dir).with_context(|| {
415                    format!("Failed to create SOPs directory: {}", sops_dir.display())
416                })?;
417            }
418            // Append timestamp to filename to avoid concurrent write collisions
419            let timestamp = SystemTime::now()
420                .duration_since(UNIX_EPOCH)
421                .map(|d| d.as_secs())
422                .unwrap_or(0);
423            let safe_name = if let Some(stem) = filename.rsplit_once('.') {
424                format!("{}_{}.{}", stem.0, timestamp, stem.1)
425            } else {
426                format!("{}_{}", filename, timestamp)
427            };
428            let sop_path = sops_dir.join(safe_name);
429            crate::memory::memory_store::write_atomic(&sop_path, &content)
430                .with_context(|| format!("Failed to write SOP artifact: {}", sop_path.display()))?;
431        }
432    }
433
434    Ok(())
435}
436
437fn extract_block(text: &str, start_tag: &str, end_tag: &str) -> Option<String> {
438    let start_idx = text.find(start_tag)?;
439    let end_idx = text.find(end_tag)?;
440    if start_idx < end_idx {
441        Some(text[start_idx + start_tag.len()..end_idx].to_string())
442    } else {
443        None
444    }
445}
446
447fn dream_output_dir(auto_memory: &AutoMemoryStore) -> Result<PathBuf> {
448    let timestamp = SystemTime::now()
449        .duration_since(UNIX_EPOCH)
450        .map(|duration| duration.as_secs())
451        .unwrap_or_default();
452    let memory_root = auto_memory
453        .db_path
454        .parent()
455        .unwrap_or(std::path::Path::new("."))
456        .to_path_buf();
457    let dreams_dir = memory_root.join("dreams");
458    let mut candidate = dreams_dir.join(format!("dream-{timestamp}"));
459    let mut suffix = 1;
460    while candidate.exists() {
461        candidate = dreams_dir.join(format!("dream-{timestamp}-{suffix}"));
462        suffix += 1;
463    }
464    Ok(candidate)
465}
466
467fn format_recent_sessions(history_store: &HistoryStore, session_limit: usize) -> Result<String> {
468    let sessions = history_store.list_sessions()?;
469    if sessions.is_empty() {
470        return Ok("No recorded sessions yet.".to_string());
471    }
472
473    let mut rendered = String::new();
474    for session in sessions.iter().take(session_limit) {
475        rendered.push_str(&format!(
476            "## Session {}\nStarted: {}\nProject: {}\n",
477            session.id, session.started_at, session.project_id
478        ));
479        let events = history_store.get_recent_events(&session.id, Some(80))?;
480        for event in events {
481            if event.event_type != "message" {
482                continue;
483            }
484            let role = event.role.as_deref().unwrap_or("unknown");
485            if let Some(content) = event.content {
486                rendered.push_str(&format!(
487                    "[{}] {}\n",
488                    role,
489                    truncate_chars(content.trim(), 2_000)
490                ));
491            }
492            if let Some(tool_output) = event.tool_output {
493                rendered.push_str(&format!(
494                    "[tool-output] {}\n",
495                    truncate_chars(tool_output.trim(), 1_000)
496                ));
497            }
498        }
499        rendered.push_str("\n---\n");
500    }
501
502    Ok(truncate_chars(&rendered, 80_000))
503}
504
505fn truncate_chars(text: &str, max_chars: usize) -> String {
506    if text.chars().count() <= max_chars {
507        return text.to_string();
508    }
509
510    let mut truncated: String = text.chars().take(max_chars).collect();
511    truncated.push_str("\n[truncated]");
512    truncated
513}
514
515fn extract_block_with_attr(
516    text: &str,
517    start_tag_prefix: &str,
518    end_tag: &str,
519) -> Option<(String, String)> {
520    let start_idx = text.find(start_tag_prefix)?;
521    let end_tag_start = text[start_idx..].find('>')?;
522    let start_tag_full_len = end_tag_start + 1;
523    let start_tag_content = &text[start_idx..start_idx + start_tag_full_len];
524
525    let filename = if let Some(fn_start) = start_tag_content.find("filename=\"") {
526        let fn_sub = &start_tag_content[fn_start + "filename=\"".len()..];
527        if let Some(fn_end) = fn_sub.find('"') {
528            fn_sub[..fn_end].to_string()
529        } else {
530            "sop.md".to_string()
531        }
532    } else {
533        "sop.md".to_string()
534    };
535
536    let content_start = start_idx + start_tag_full_len;
537    let end_idx = text[content_start..].find(end_tag)?;
538    let content = text[content_start..content_start + end_idx].to_string();
539
540    Some((filename, content))
541}