pub mod collector;
pub mod compactor;
pub mod config;
pub mod multimodal;
pub mod store;
pub mod temporal;
pub use collector::{
CollectedBatch, CollectedItem, EpisodeData, MemoryEntryData, ShortTermCollector, SourceType,
};
pub use compactor::MemoryCompactor;
pub use config::ConsolidationConfig;
pub use multimodal::MultimodalRef;
pub use store::LongTermStore;
pub use temporal::{CompactedContent, ConsolidationReport, RecordImportance, TemporalRecord};
use anyhow::Result;
use std::path::PathBuf;
use std::time::Duration;
use tracing::info;
fn record_to_collected_item(record: &TemporalRecord) -> CollectedItem {
use crate::consolidation::temporal::RecordImportance;
let importance: u8 = match record.importance {
RecordImportance::Low => 1,
RecordImportance::Normal => 2,
RecordImportance::High => 3,
RecordImportance::Critical => 4,
};
CollectedItem {
source_id: record.id.clone(),
source_type: SourceType::Episode,
content: record.content.summary.clone(),
timestamp: record.created_at,
importance,
tags: record.tags.clone(),
metadata: std::collections::HashMap::new(),
related_ids: record.causal_parents.clone(),
session_id: record.session_id.clone(),
file_refs: record
.multimodal_refs
.iter()
.filter_map(|r| match r {
MultimodalRef::Screenshot { path, .. } => path.to_str().map(|s| s.to_string()),
_ => None,
})
.collect(),
}
}
pub struct ConsolidationEngine {
#[allow(dead_code)]
config: ConsolidationConfig,
#[allow(dead_code)]
collector: ShortTermCollector,
compactor: MemoryCompactor,
store: LongTermStore,
}
impl ConsolidationEngine {
pub fn new(config: ConsolidationConfig) -> Result<Self> {
let collector = ShortTermCollector::new(
(config.max_episode_age_hours * 3600.0) as u64,
config.min_importance,
);
let compactor = MemoryCompactor::new(config.clone())?;
let store = LongTermStore::new(PathBuf::from("consolidated_memory"));
Ok(Self {
config,
collector,
compactor,
store,
})
}
pub fn with_storage_dir(mut self, dir: PathBuf) -> Self {
self.store = LongTermStore::new(dir);
self
}
pub async fn consolidate_episodes(
&mut self,
episodes: Vec<EpisodeData>,
) -> Result<ConsolidationReport> {
info!(
"Starting consolidation cycle with {} episodes",
episodes.len()
);
let items = self.collector.collect_episodes(&episodes);
if items.is_empty() {
info!("No items to consolidate");
let now = chrono::Utc::now();
return Ok(ConsolidationReport {
started_at: now,
ended_at: now,
episodes_processed: 0,
records_produced: 0,
duplicates_removed: 0,
tokens_used: 0,
causal_links_created: 0,
multimodal_refs_count: 0,
errors: vec![],
});
}
let batch = self.collector.assemble_batch(items);
let (records, mut report) = self.compactor.compact(batch).await?;
let store_result = self.store.store(&records).await?;
if !store_result.errors.is_empty() {
report.errors.extend(store_result.errors);
}
info!(
"Consolidation complete: {} episodes -> {} records, {} tokens used",
report.episodes_processed, report.records_produced, report.tokens_used,
);
Ok(report)
}
pub fn start_periodic(self, interval: Duration) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let engine = self;
let mut tick = tokio::time::interval(interval);
loop {
tick.tick().await;
info!("Periodic consolidation triggered");
match engine.store.load_all() {
Ok(records) => {
if records.is_empty() {
info!("Periodic check: no records to re-consolidate");
continue;
}
info!(
"Periodic check: re-consolidating {} existing records",
records.len()
);
let items: Vec<CollectedItem> =
records.iter().map(record_to_collected_item).collect();
let batch = engine.collector.assemble_batch(items);
match engine.compactor.compact(batch).await {
Ok((new_records, report)) => {
info!(
"Periodic re-consolidation: {} records -> {} records, {} errors",
records.len(),
new_records.len(),
report.errors.len(),
);
if let Err(e) = engine.store.store(&new_records).await {
tracing::error!("Periodic consolidation store error: {e}");
}
}
Err(e) => {
tracing::error!("Periodic consolidation compaction error: {e}");
}
}
}
Err(e) => {
tracing::error!("Periodic consolidation load error: {e}");
}
}
}
})
}
pub fn store(&self) -> &LongTermStore {
&self.store
}
}
#[cfg(test)]
#[path = "../../tests/unit/consolidation/mod_test.rs"]
mod tests;