1use std::sync::Arc;
19use std::time::{Duration, Instant};
20
21use chrono::{DateTime, Utc};
22use serde::Deserialize;
23
24use bamboo_domain::ledger::{LedgerRecord, LedgerScope, RecordStatus};
25use bamboo_memory::ledger_store::{LedgerScheduleBridge, LedgerStore, RecordFilter};
26use bamboo_memory::memory_store::{DurableMemoryType, MemoryStore};
27
28use crate::auto_dream::AutoDreamContext;
29use crate::gardener::{collect_model_json, resolve_background_model};
30
31const LEDGER_GARDENER_TRACING_TARGET: &str = "bamboo.ledger_gardener";
32const EXPIRY_GRACE_HOURS: i64 = 48;
35const DISTILL_MAX_RECORDS_PER_RUN: usize = 8;
37const DISTILLED_TAG: &str = "distilled";
39
40const DISTILL_SYSTEM_INSTRUCTION: &str = "You are Bamboo's background ledger gardener. From the user's completed ledger records, extract only durable, long-term facts worth remembering (habits, recurring obligations, stable preferences, notable life events). Return only the specified JSON array. No prose, no markdown fences.";
41
42pub struct LedgerGardenerContext {
43 pub dream: AutoDreamContext,
44 pub schedule_bridge: Option<Arc<dyn LedgerScheduleBridge>>,
45}
46
47#[derive(Debug, Clone, PartialEq, Eq, Default)]
48pub struct LedgerGardenerRunResult {
49 pub scanned: usize,
50 pub expired: usize,
51 pub schedules_released: usize,
52 pub schedules_synced: usize,
53 pub distilled_records: usize,
54 pub memories_written: usize,
55 pub failed: usize,
56}
57
58#[derive(Debug, Clone, Deserialize, PartialEq, Eq)]
59pub struct DistilledMemoryCandidate {
60 pub title: String,
61 pub content: String,
62 #[serde(default)]
63 pub r#type: Option<String>,
64 #[serde(default)]
65 pub tags: Vec<String>,
66}
67
68pub fn should_expire(record: &LedgerRecord, now: DateTime<Utc>) -> bool {
74 if record.status.is_terminal() {
75 return false;
76 }
77 if record.time.due_at.is_some() || record.time.recurrence.is_some() {
78 return false;
79 }
80 let cutoff = now - chrono::Duration::hours(EXPIRY_GRACE_HOURS);
81 if let Some(end) = record.time.ends_at.or(record.time.starts_at) {
82 return end < cutoff;
83 }
84 if !record.time.remind_at.is_empty() {
85 return record.time.remind_at.iter().all(|at| *at < cutoff);
86 }
87 false
88}
89
90pub fn parse_distilled_candidates(raw: &str) -> Vec<DistilledMemoryCandidate> {
94 let Some(start) = raw.find('[') else {
95 return Vec::new();
96 };
97 let Some(end) = raw.rfind(']') else {
98 return Vec::new();
99 };
100 if end < start {
101 return Vec::new();
102 }
103 serde_json::from_str::<Vec<DistilledMemoryCandidate>>(&raw[start..=end])
104 .unwrap_or_default()
105 .into_iter()
106 .filter(|candidate| {
107 !candidate.title.trim().is_empty() && !candidate.content.trim().is_empty()
108 })
109 .collect()
110}
111
112fn build_distillation_prompt(docs: &[(String, String)]) -> String {
113 let mut prompt = String::from(
114 "The user completed the following ledger records (todos/events/reminders).\n\
115 Extract durable long-term memories ONLY where a record reveals a lasting fact:\n\
116 a habit or routine, a recurring obligation, a stable preference, or a notable\n\
117 life event worth recalling months from now. Routine one-off chores yield nothing.\n\n\
118 Records:\n",
119 );
120 for (title, detail) in docs {
121 prompt.push_str(&format!("- {title}\n {detail}\n"));
122 }
123 prompt.push_str(
124 "\nReturn a JSON array (possibly empty), each item:\n\
125 {\"title\": \"specific descriptive title\", \"content\": \"one atomic fact\", \
126 \"type\": \"user|reference\", \"tags\": [\"...\"]}\n",
127 );
128 prompt
129}
130
131fn all_scopes(project_keys: &[String]) -> Vec<(LedgerScope, Option<String>)> {
132 let mut scopes: Vec<(LedgerScope, Option<String>)> = vec![(LedgerScope::Global, None)];
133 scopes.extend(
134 project_keys
135 .iter()
136 .map(|key| (LedgerScope::Project, Some(key.clone()))),
137 );
138 scopes
139}
140
141pub async fn run_ledger_gardener_once(
142 ctx: &LedgerGardenerContext,
143) -> Result<Option<LedgerGardenerRunResult>, String> {
144 let ledger = LedgerStore::new(ctx.dream.session_store.bamboo_home_dir());
145 run_ledger_gardener_once_with_store(ctx, &ledger).await
146}
147
148async fn run_ledger_gardener_once_with_store(
149 ctx: &LedgerGardenerContext,
150 ledger: &LedgerStore,
151) -> Result<Option<LedgerGardenerRunResult>, String> {
152 let config_snapshot = ctx.dream.config.read().await.clone();
153 let memory_cfg = config_snapshot.memory().clone().unwrap_or_default();
154 if !memory_cfg.ledger_gardener_enabled {
155 return Ok(None);
156 }
157
158 let now = Utc::now();
159 let mut result = LedgerGardenerRunResult::default();
160 let project_keys = ledger
161 .list_project_keys()
162 .await
163 .map_err(|error| format!("ledger gardener scope scan failed: {error}"))?;
164 let scopes = all_scopes(&project_keys);
165
166 for (scope, project_key) in &scopes {
168 let docs = ledger
169 .list_records(
170 *scope,
171 project_key.as_deref(),
172 &RecordFilter {
173 include_terminal: true,
174 ..RecordFilter::default()
175 },
176 )
177 .await
178 .map_err(|error| format!("ledger gardener list failed: {error}"))?;
179 result.scanned += docs.len();
180
181 for doc in docs {
182 let record = &doc.record;
183
184 if should_expire(record, now) {
186 match ledger
187 .transition_record(
188 *scope,
189 project_key.as_deref(),
190 &record.id,
191 RecordStatus::Expired,
192 Some("auto-expired by the ledger gardener"),
193 )
194 .await
195 {
196 Ok(Some(expired_doc)) => {
197 result.expired += 1;
198 if let (Some(bridge), false) =
203 (&ctx.schedule_bridge, record.schedule_ids.is_empty())
204 {
205 if bridge.release_schedules(&record.schedule_ids).await.is_ok() {
206 let mut cleared = expired_doc.record.clone();
207 cleared.schedule_ids.clear();
208 let _ = ledger
209 .write_record(cleared, Some(expired_doc.body.clone()))
210 .await;
211 result.schedules_released += record.schedule_ids.len();
212 }
213 }
214 continue;
215 }
216 Ok(None) => {}
217 Err(error) => {
218 result.failed += 1;
219 tracing::warn!(
220 target: LEDGER_GARDENER_TRACING_TARGET,
221 record_id = %record.id,
222 "[ledger-gardener] expiry transition failed: {error}"
223 );
224 continue;
225 }
226 }
227 }
228
229 let Some(bridge) = &ctx.schedule_bridge else {
231 continue;
232 };
233 if record.status.is_terminal() && !record.schedule_ids.is_empty() {
234 match bridge.release_schedules(&record.schedule_ids).await {
235 Ok(()) => {
236 let mut cleared = record.clone();
237 cleared.schedule_ids.clear();
238 if ledger
239 .write_record(cleared, Some(doc.body.clone()))
240 .await
241 .is_ok()
242 {
243 result.schedules_released += record.schedule_ids.len();
244 }
245 }
246 Err(error) => {
247 result.failed += 1;
248 tracing::warn!(
249 target: LEDGER_GARDENER_TRACING_TARGET,
250 record_id = %record.id,
251 "[ledger-gardener] schedule release failed: {error}"
252 );
253 }
254 }
255 } else if !record.status.is_terminal()
256 && record.schedule_ids.is_empty()
257 && (record.time.recurrence.is_some()
258 || record.time.remind_at.iter().any(|at| *at > now))
259 {
260 match bridge.sync_record_schedules(record).await {
261 Ok(ids) if !ids.is_empty() => {
262 let mut updated = record.clone();
263 updated.schedule_ids = ids.clone();
264 if ledger
265 .write_record(updated, Some(doc.body.clone()))
266 .await
267 .is_ok()
268 {
269 result.schedules_synced += ids.len();
270 }
271 }
272 Ok(_) => {}
273 Err(error) => {
274 result.failed += 1;
275 tracing::warn!(
276 target: LEDGER_GARDENER_TRACING_TARGET,
277 record_id = %record.id,
278 "[ledger-gardener] schedule sync failed: {error}"
279 );
280 }
281 }
282 }
283 }
284 }
285
286 if memory_cfg.ledger_distillation_enabled {
288 if let Err(error) = run_distillation_pass(ctx, ledger, &scopes, &mut result).await {
289 result.failed += 1;
290 tracing::warn!(
291 target: LEDGER_GARDENER_TRACING_TARGET,
292 "[ledger-gardener] distillation failed: {error}"
293 );
294 }
295 }
296
297 tracing::info!(
298 target: LEDGER_GARDENER_TRACING_TARGET,
299 event = "run_complete",
300 scanned = result.scanned,
301 expired = result.expired,
302 schedules_released = result.schedules_released,
303 schedules_synced = result.schedules_synced,
304 distilled_records = result.distilled_records,
305 memories_written = result.memories_written,
306 failed = result.failed,
307 "[ledger-gardener] run complete"
308 );
309 Ok(Some(result))
310}
311
312async fn run_distillation_pass(
313 ctx: &LedgerGardenerContext,
314 ledger: &LedgerStore,
315 scopes: &[(LedgerScope, Option<String>)],
316 result: &mut LedgerGardenerRunResult,
317) -> Result<(), String> {
318 let mut pending = Vec::new();
320 for (scope, project_key) in scopes {
321 let docs = ledger
322 .list_records(
323 *scope,
324 project_key.as_deref(),
325 &RecordFilter {
326 statuses: Some([RecordStatus::Done].into_iter().collect()),
327 include_terminal: true,
328 ..RecordFilter::default()
329 },
330 )
331 .await
332 .map_err(|error| format!("distillation list failed: {error}"))?;
333 pending.extend(
334 docs.into_iter()
335 .filter(|doc| !doc.record.tags.iter().any(|tag| tag == DISTILLED_TAG)),
336 );
337 if pending.len() >= DISTILL_MAX_RECORDS_PER_RUN {
338 break;
339 }
340 }
341 pending.truncate(DISTILL_MAX_RECORDS_PER_RUN);
342 if pending.is_empty() {
343 return Ok(());
344 }
345
346 let config_snapshot = ctx.dream.config.read().await.clone();
347 let Some((provider, model)) = resolve_background_model(&ctx.dream, &config_snapshot) else {
348 return Ok(());
351 };
352
353 let lines: Vec<(String, String)> = pending
354 .iter()
355 .map(|doc| {
356 let record = &doc.record;
357 let detail = format!(
358 "kind={}, completed={}, notes: {}",
359 record.kind.as_str(),
360 record.updated_at.format("%Y-%m-%d"),
361 doc.body
362 .chars()
363 .take(200)
364 .collect::<String>()
365 .replace('\n', " "),
366 );
367 (record.title.clone(), detail)
368 })
369 .collect();
370 let raw = collect_model_json(
371 provider,
372 &model,
373 DISTILL_SYSTEM_INSTRUCTION,
374 build_distillation_prompt(&lines),
375 )
376 .await?;
377 let candidates = parse_distilled_candidates(&raw);
378
379 let memory = MemoryStore::new(ctx.dream.session_store.bamboo_home_dir());
380 for candidate in &candidates {
381 let r#type = match candidate.r#type.as_deref() {
382 Some("reference") => DurableMemoryType::Reference,
383 _ => DurableMemoryType::User,
384 };
385 match memory
386 .write_memory(
387 bamboo_memory::memory_store::MemoryScope::Global,
388 None,
389 r#type,
390 candidate.title.trim(),
391 candidate.content.trim(),
392 &candidate.tags,
393 None,
394 "ledger-gardener",
395 true, None,
397 )
398 .await
399 {
400 Ok(_) => result.memories_written += 1,
401 Err(error) => {
402 result.failed += 1;
403 tracing::warn!(
404 target: LEDGER_GARDENER_TRACING_TARGET,
405 "[ledger-gardener] distilled memory write failed: {error}"
406 );
407 }
408 }
409 }
410
411 for doc in pending {
414 let mut record = doc.record.clone();
415 record.tags.push(DISTILLED_TAG.to_string());
416 if ledger
417 .write_record(record, Some(doc.body.clone()))
418 .await
419 .is_ok()
420 {
421 result.distilled_records += 1;
422 }
423 }
424 Ok(())
425}
426
427pub fn spawn_ledger_gardener_task(ctx: LedgerGardenerContext) {
430 tokio::spawn(async move {
431 let interval_secs = {
432 let guard = ctx.dream.config.read().await;
433 guard
434 .memory()
435 .as_ref()
436 .map(|memory| memory.ledger_gardener_interval_secs)
437 .filter(|secs| *secs > 0)
438 .unwrap_or_else(|| {
439 bamboo_config::MemoryConfig::default().ledger_gardener_interval_secs
440 })
441 };
442 let poll_secs = 300.min(interval_secs.max(1));
445 let mut ticker = tokio::time::interval(Duration::from_secs(poll_secs));
446 let mut last_run: Option<Instant> = None;
447
448 loop {
449 ticker.tick().await;
450 let due = last_run
451 .map(|at| at.elapsed().as_secs() >= interval_secs)
452 .unwrap_or(true);
453 if !due {
454 continue;
455 }
456 if let Err(error) = run_ledger_gardener_once(&ctx).await {
457 tracing::warn!(
458 target: LEDGER_GARDENER_TRACING_TARGET,
459 event = "run_failed",
460 "[ledger-gardener] run failed: {error}"
461 );
462 }
463 last_run = Some(Instant::now());
464 }
465 });
466}
467
468#[cfg(test)]
469mod tests {
470 use super::*;
471 use std::collections::HashMap;
472 use std::sync::Mutex;
473
474 use async_trait::async_trait;
475 use bamboo_agent_core::Message;
476 use bamboo_domain::ledger::RecordKind;
477 use bamboo_domain::schedule::ScheduleTrigger;
478 use bamboo_llm::{Config, LLMChunk, LLMError, LLMProvider, LLMStream, ProviderRegistry};
479 use bamboo_storage::SessionStoreV2;
480 use chrono::Duration as ChronoDuration;
481 use futures::stream;
482 use tokio::sync::{Mutex as AsyncMutex, RwLock};
483
484 fn config_with_memory(memory: bamboo_config::MemoryConfig) -> Config {
485 let mut config = Config::default();
486 *config.memory_mut() = Some(memory);
487 config
488 }
489
490 fn record(kind: RecordKind, status: RecordStatus) -> LedgerRecord {
491 let mut record = LedgerRecord::new("rec_test", kind, "Test record");
492 record.status = status;
493 record
494 }
495
496 #[test]
497 fn expiry_policy_expires_only_fully_past_events_and_reminders() {
498 let now = Utc::now();
499 let past = now - ChronoDuration::hours(EXPIRY_GRACE_HOURS + 1);
500 let recent_past = now - ChronoDuration::hours(1);
501
502 let mut event = record(RecordKind::Event, RecordStatus::Open);
504 event.time.starts_at = Some(past);
505 assert!(should_expire(&event, now));
506
507 let mut recent_event = record(RecordKind::Event, RecordStatus::Open);
509 recent_event.time.starts_at = Some(recent_past);
510 assert!(!should_expire(&recent_event, now));
511
512 let mut reminder = record(RecordKind::Reminder, RecordStatus::Open);
514 reminder.time.remind_at = vec![past, past - ChronoDuration::hours(2)];
515 assert!(should_expire(&reminder, now));
516
517 let mut live_reminder = record(RecordKind::Reminder, RecordStatus::Open);
519 live_reminder.time.remind_at = vec![past, now + ChronoDuration::hours(1)];
520 assert!(!should_expire(&live_reminder, now));
521
522 let mut todo = record(RecordKind::Todo, RecordStatus::Open);
524 todo.time.due_at = Some(past);
525 assert!(!should_expire(&todo, now));
526
527 let mut habit = record(RecordKind::Habit, RecordStatus::Open);
529 habit.time.remind_at = vec![past];
530 habit.time.recurrence = Some(ScheduleTrigger::Daily {
531 hour: 9,
532 minute: 0,
533 second: 0,
534 });
535 assert!(!should_expire(&habit, now));
536
537 assert!(!should_expire(
539 &record(RecordKind::Todo, RecordStatus::Done),
540 now
541 ));
542 assert!(!should_expire(
543 &record(RecordKind::Todo, RecordStatus::Open),
544 now
545 ));
546 }
547
548 #[test]
549 fn parse_distilled_candidates_is_lenient() {
550 let wrapped = "Sure! Here you go:\n```json\n[{\"title\": \"Takes medication daily\", \"content\": \"User takes medication every morning at 9.\", \"type\": \"user\", \"tags\": [\"health\"]}]\n```";
551 let parsed = parse_distilled_candidates(wrapped);
552 assert_eq!(parsed.len(), 1);
553 assert_eq!(parsed[0].title, "Takes medication daily");
554
555 assert!(parse_distilled_candidates("no json here").is_empty());
556 assert!(parse_distilled_candidates("[]").is_empty());
557 assert!(parse_distilled_candidates("[{\"title\": \" \", \"content\": \"x\"}]").is_empty());
558 assert!(parse_distilled_candidates("[{broken").is_empty());
559 }
560
561 #[test]
562 fn distillation_prompt_lists_records_and_requests_json() {
563 let prompt = build_distillation_prompt(&[(
564 "Renew passport".to_string(),
565 "kind=todo, completed=2026-07-13".to_string(),
566 )]);
567 assert!(prompt.contains("Renew passport"));
568 assert!(prompt.contains("JSON array"));
569 }
570
571 #[derive(Clone)]
572 struct CannedProvider {
573 responses: Arc<Mutex<Vec<String>>>,
574 }
575
576 #[async_trait]
577 impl LLMProvider for CannedProvider {
578 async fn chat_stream(
579 &self,
580 _messages: &[Message],
581 _tools: &[bamboo_agent_core::tools::ToolSchema],
582 _max_output_tokens: Option<u32>,
583 _model: &str,
584 ) -> Result<LLMStream, LLMError> {
585 let mut responses = self.responses.lock().expect("lock poisoned");
586 let text = if responses.is_empty() {
587 "[]".to_string()
588 } else {
589 responses.remove(0)
590 };
591 Ok(Box::pin(stream::iter(vec![
592 Ok(LLMChunk::Token(text)),
593 Ok(LLMChunk::Done),
594 ])))
595 }
596 }
597
598 #[derive(Default)]
599 struct RecordingBridge {
600 synced: AsyncMutex<Vec<String>>,
601 released: AsyncMutex<Vec<String>>,
602 }
603
604 #[async_trait]
605 impl LedgerScheduleBridge for RecordingBridge {
606 async fn sync_record_schedules(
607 &self,
608 record: &LedgerRecord,
609 ) -> Result<Vec<String>, String> {
610 self.synced.lock().await.push(record.id.clone());
611 Ok(vec![format!("sched_for_{}", record.id)])
612 }
613
614 async fn release_schedules(&self, schedule_ids: &[String]) -> Result<(), String> {
615 self.released
616 .lock()
617 .await
618 .extend(schedule_ids.iter().cloned());
619 Ok(())
620 }
621 }
622
623 #[tokio::test]
624 async fn full_run_expires_reconciles_and_distills() {
625 let temp = tempfile::tempdir().expect("tempdir");
626 let session_store = Arc::new(
627 SessionStoreV2::new(temp.path().to_path_buf())
628 .await
629 .unwrap(),
630 );
631 let ledger = LedgerStore::new(session_store.bamboo_home_dir());
632 let now = Utc::now();
633
634 let mut past_event = LedgerRecord::new("rec_event", RecordKind::Event, "Old meetup");
636 past_event.time.starts_at = Some(now - ChronoDuration::hours(EXPIRY_GRACE_HOURS + 10));
637 past_event.schedule_ids = vec!["sched_old".to_string()];
638 ledger.write_record(past_event, None).await.unwrap();
639
640 let mut drifted = LedgerRecord::new("rec_drift", RecordKind::Reminder, "Call the bank");
643 drifted.time.remind_at = vec![now + ChronoDuration::hours(5)];
644 ledger.write_record(drifted, None).await.unwrap();
645
646 let mut done = LedgerRecord::new("rec_done", RecordKind::Habit, "Morning run streak");
649 done.status = RecordStatus::Done;
650 ledger.write_record(done, None).await.unwrap();
651
652 let provider: Arc<dyn LLMProvider> = Arc::new(CannedProvider {
653 responses: Arc::new(Mutex::new(vec![
654 "[{\"title\": \"Runs every morning\", \"content\": \"User keeps a morning run habit.\", \"type\": \"user\", \"tags\": [\"health\"]}]".to_string(),
655 ])),
656 });
657 let config = Arc::new(RwLock::new(config_with_memory(
658 bamboo_config::MemoryConfig {
659 background_model: Some("fast-model".to_string()),
660 ..bamboo_config::MemoryConfig::default()
661 },
662 )));
663 let bridge = Arc::new(RecordingBridge::default());
664 let ctx = LedgerGardenerContext {
665 dream: AutoDreamContext {
666 session_store: session_store.clone(),
667 storage: session_store.clone(),
668 provider,
669 config,
670 provider_registry: Arc::new(ProviderRegistry::new(
671 HashMap::new(),
672 "test".to_string(),
673 )),
674 },
675 schedule_bridge: Some(bridge.clone()),
676 };
677
678 let result = run_ledger_gardener_once_with_store(&ctx, &ledger)
679 .await
680 .unwrap()
681 .expect("gardener enabled by default");
682
683 assert_eq!(result.expired, 1);
684 assert_eq!(result.schedules_released, 1);
685 assert_eq!(result.schedules_synced, 1);
686 assert_eq!(result.memories_written, 1);
687 assert_eq!(result.distilled_records, 1);
688 assert_eq!(result.failed, 0);
689
690 let expired = ledger
691 .get_record(LedgerScope::Global, None, "rec_event")
692 .await
693 .unwrap()
694 .unwrap();
695 assert_eq!(expired.record.status, RecordStatus::Expired);
696 assert!(expired.record.schedule_ids.is_empty());
697 assert!(
698 !expired.record.transitions.is_empty(),
699 "history must survive"
700 );
701 assert_eq!(*bridge.released.lock().await, vec!["sched_old".to_string()]);
702
703 let drifted = ledger
704 .get_record(LedgerScope::Global, None, "rec_drift")
705 .await
706 .unwrap()
707 .unwrap();
708 assert_eq!(drifted.record.schedule_ids, vec!["sched_for_rec_drift"]);
709
710 let done = ledger
711 .get_record(LedgerScope::Global, None, "rec_done")
712 .await
713 .unwrap()
714 .unwrap();
715 assert!(done.record.tags.iter().any(|tag| tag == DISTILLED_TAG));
716
717 let second = run_ledger_gardener_once_with_store(&ctx, &ledger)
719 .await
720 .unwrap()
721 .unwrap();
722 assert_eq!(second.expired, 0);
723 assert_eq!(second.distilled_records, 0);
724 assert_eq!(second.memories_written, 0);
725 }
726}