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#[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}", ¬es)
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 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 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 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
284async 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 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}