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 record(kind: RecordKind, status: RecordStatus) -> LedgerRecord {
485 let mut record = LedgerRecord::new("rec_test", kind, "Test record");
486 record.status = status;
487 record
488 }
489
490 #[test]
491 fn expiry_policy_expires_only_fully_past_events_and_reminders() {
492 let now = Utc::now();
493 let past = now - ChronoDuration::hours(EXPIRY_GRACE_HOURS + 1);
494 let recent_past = now - ChronoDuration::hours(1);
495
496 let mut event = record(RecordKind::Event, RecordStatus::Open);
498 event.time.starts_at = Some(past);
499 assert!(should_expire(&event, now));
500
501 let mut recent_event = record(RecordKind::Event, RecordStatus::Open);
503 recent_event.time.starts_at = Some(recent_past);
504 assert!(!should_expire(&recent_event, now));
505
506 let mut reminder = record(RecordKind::Reminder, RecordStatus::Open);
508 reminder.time.remind_at = vec![past, past - ChronoDuration::hours(2)];
509 assert!(should_expire(&reminder, now));
510
511 let mut live_reminder = record(RecordKind::Reminder, RecordStatus::Open);
513 live_reminder.time.remind_at = vec![past, now + ChronoDuration::hours(1)];
514 assert!(!should_expire(&live_reminder, now));
515
516 let mut todo = record(RecordKind::Todo, RecordStatus::Open);
518 todo.time.due_at = Some(past);
519 assert!(!should_expire(&todo, now));
520
521 let mut habit = record(RecordKind::Habit, RecordStatus::Open);
523 habit.time.remind_at = vec![past];
524 habit.time.recurrence = Some(ScheduleTrigger::Daily {
525 hour: 9,
526 minute: 0,
527 second: 0,
528 });
529 assert!(!should_expire(&habit, now));
530
531 assert!(!should_expire(
533 &record(RecordKind::Todo, RecordStatus::Done),
534 now
535 ));
536 assert!(!should_expire(
537 &record(RecordKind::Todo, RecordStatus::Open),
538 now
539 ));
540 }
541
542 #[test]
543 fn parse_distilled_candidates_is_lenient() {
544 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```";
545 let parsed = parse_distilled_candidates(wrapped);
546 assert_eq!(parsed.len(), 1);
547 assert_eq!(parsed[0].title, "Takes medication daily");
548
549 assert!(parse_distilled_candidates("no json here").is_empty());
550 assert!(parse_distilled_candidates("[]").is_empty());
551 assert!(parse_distilled_candidates("[{\"title\": \" \", \"content\": \"x\"}]").is_empty());
552 assert!(parse_distilled_candidates("[{broken").is_empty());
553 }
554
555 #[test]
556 fn distillation_prompt_lists_records_and_requests_json() {
557 let prompt = build_distillation_prompt(&[(
558 "Renew passport".to_string(),
559 "kind=todo, completed=2026-07-13".to_string(),
560 )]);
561 assert!(prompt.contains("Renew passport"));
562 assert!(prompt.contains("JSON array"));
563 }
564
565 #[derive(Clone)]
566 struct CannedProvider {
567 responses: Arc<Mutex<Vec<String>>>,
568 }
569
570 #[async_trait]
571 impl LLMProvider for CannedProvider {
572 async fn chat_stream(
573 &self,
574 _messages: &[Message],
575 _tools: &[bamboo_agent_core::tools::ToolSchema],
576 _max_output_tokens: Option<u32>,
577 _model: &str,
578 ) -> Result<LLMStream, LLMError> {
579 let mut responses = self.responses.lock().expect("lock poisoned");
580 let text = if responses.is_empty() {
581 "[]".to_string()
582 } else {
583 responses.remove(0)
584 };
585 Ok(Box::pin(stream::iter(vec![
586 Ok(LLMChunk::Token(text)),
587 Ok(LLMChunk::Done),
588 ])))
589 }
590 }
591
592 #[derive(Default)]
593 struct RecordingBridge {
594 synced: AsyncMutex<Vec<String>>,
595 released: AsyncMutex<Vec<String>>,
596 }
597
598 #[async_trait]
599 impl LedgerScheduleBridge for RecordingBridge {
600 async fn sync_record_schedules(
601 &self,
602 record: &LedgerRecord,
603 ) -> Result<Vec<String>, String> {
604 self.synced.lock().await.push(record.id.clone());
605 Ok(vec![format!("sched_for_{}", record.id)])
606 }
607
608 async fn release_schedules(&self, schedule_ids: &[String]) -> Result<(), String> {
609 self.released
610 .lock()
611 .await
612 .extend(schedule_ids.iter().cloned());
613 Ok(())
614 }
615 }
616
617 #[tokio::test]
618 async fn full_run_expires_reconciles_and_distills() {
619 let temp = tempfile::tempdir().expect("tempdir");
620 let session_store = Arc::new(
621 SessionStoreV2::new(temp.path().to_path_buf())
622 .await
623 .unwrap(),
624 );
625 let ledger = LedgerStore::new(session_store.bamboo_home_dir());
626 let now = Utc::now();
627
628 let mut past_event = LedgerRecord::new("rec_event", RecordKind::Event, "Old meetup");
630 past_event.time.starts_at = Some(now - ChronoDuration::hours(EXPIRY_GRACE_HOURS + 10));
631 past_event.schedule_ids = vec!["sched_old".to_string()];
632 ledger.write_record(past_event, None).await.unwrap();
633
634 let mut drifted = LedgerRecord::new("rec_drift", RecordKind::Reminder, "Call the bank");
637 drifted.time.remind_at = vec![now + ChronoDuration::hours(5)];
638 ledger.write_record(drifted, None).await.unwrap();
639
640 let mut done = LedgerRecord::new("rec_done", RecordKind::Habit, "Morning run streak");
643 done.status = RecordStatus::Done;
644 ledger.write_record(done, None).await.unwrap();
645
646 let provider: Arc<dyn LLMProvider> = Arc::new(CannedProvider {
647 responses: Arc::new(Mutex::new(vec![
648 "[{\"title\": \"Runs every morning\", \"content\": \"User keeps a morning run habit.\", \"type\": \"user\", \"tags\": [\"health\"]}]".to_string(),
649 ])),
650 });
651 let config = Arc::new(RwLock::new(Config {
652 memory: Some(bamboo_config::MemoryConfig {
653 background_model: Some("fast-model".to_string()),
654 ..bamboo_config::MemoryConfig::default()
655 }),
656 ..Config::default()
657 }));
658 let bridge = Arc::new(RecordingBridge::default());
659 let ctx = LedgerGardenerContext {
660 dream: AutoDreamContext {
661 session_store: session_store.clone(),
662 storage: session_store.clone(),
663 provider,
664 config,
665 provider_registry: Arc::new(ProviderRegistry::new(
666 HashMap::new(),
667 "test".to_string(),
668 )),
669 },
670 schedule_bridge: Some(bridge.clone()),
671 };
672
673 let result = run_ledger_gardener_once_with_store(&ctx, &ledger)
674 .await
675 .unwrap()
676 .expect("gardener enabled by default");
677
678 assert_eq!(result.expired, 1);
679 assert_eq!(result.schedules_released, 1);
680 assert_eq!(result.schedules_synced, 1);
681 assert_eq!(result.memories_written, 1);
682 assert_eq!(result.distilled_records, 1);
683 assert_eq!(result.failed, 0);
684
685 let expired = ledger
686 .get_record(LedgerScope::Global, None, "rec_event")
687 .await
688 .unwrap()
689 .unwrap();
690 assert_eq!(expired.record.status, RecordStatus::Expired);
691 assert!(expired.record.schedule_ids.is_empty());
692 assert!(
693 !expired.record.transitions.is_empty(),
694 "history must survive"
695 );
696 assert_eq!(*bridge.released.lock().await, vec!["sched_old".to_string()]);
697
698 let drifted = ledger
699 .get_record(LedgerScope::Global, None, "rec_drift")
700 .await
701 .unwrap()
702 .unwrap();
703 assert_eq!(drifted.record.schedule_ids, vec!["sched_for_rec_drift"]);
704
705 let done = ledger
706 .get_record(LedgerScope::Global, None, "rec_done")
707 .await
708 .unwrap()
709 .unwrap();
710 assert!(done.record.tags.iter().any(|tag| tag == DISTILLED_TAG));
711
712 let second = run_ledger_gardener_once_with_store(&ctx, &ledger)
714 .await
715 .unwrap()
716 .unwrap();
717 assert_eq!(second.expired, 0);
718 assert_eq!(second.distilled_records, 0);
719 assert_eq!(second.memories_written, 0);
720 }
721}