Implement Knowledge Graph Consolidation worker in server
This commit is contained in:
1 parent
d2ab8c89b6
commit
63f8ec6281
1 file changed
+77
@@ -127,6 +127,82 @@ async fn index_committer_worker(state: Arc<MemoryState>) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn condense_graph_worker(state: Arc<MemoryState>) {
|
||||||
|
loop {
|
||||||
|
tokio::time::sleep(Duration::from_secs(3600)).await;
|
||||||
|
|
||||||
|
let threshold: usize = std::env::var("MCP_MEMORY_CONDENSE_THRESHOLD")
|
||||||
|
.unwrap_or_else(|_| "100".to_string())
|
||||||
|
.parse()
|
||||||
|
.unwrap_or(100);
|
||||||
|
|
||||||
|
let now = std::time::SystemTime::now()
|
||||||
|
.duration_since(std::time::UNIX_EPOCH)
|
||||||
|
.unwrap_or_default()
|
||||||
|
.as_secs();
|
||||||
|
|
||||||
|
// Condense sticky notes
|
||||||
|
let mut condensed_sticky_content = String::new();
|
||||||
|
state.sticky.modify(|notes| {
|
||||||
|
if notes.len() > threshold {
|
||||||
|
notes.sort_by_key(|n| n.timestamp);
|
||||||
|
let to_remove = notes.len() - (threshold / 2);
|
||||||
|
let removed: Vec<_> = notes.drain(0..to_remove).collect();
|
||||||
|
for r in removed {
|
||||||
|
condensed_sticky_content.push_str(&format!("{}\n", r.content));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
if !condensed_sticky_content.is_empty() {
|
||||||
|
state.modify_graph(|graph| {
|
||||||
|
let name = format!("StickyNote History {}", now);
|
||||||
|
graph.entities.insert(
|
||||||
|
name.clone(),
|
||||||
|
crate::models::Entity {
|
||||||
|
name: name.clone(),
|
||||||
|
entity_type: "Historical Summary".to_string(),
|
||||||
|
observations: vec![condensed_sticky_content],
|
||||||
|
namespace: crate::models::default_namespace(),
|
||||||
|
git_branch: None,
|
||||||
|
},
|
||||||
|
);
|
||||||
|
});
|
||||||
|
tracing::info!("Condensed sticky notes into Historical Summary.");
|
||||||
|
}
|
||||||
|
|
||||||
|
// Condense snippets
|
||||||
|
let mut condensed_snippet_content = String::new();
|
||||||
|
state.snippets.modify(|snippets| {
|
||||||
|
if snippets.len() > threshold {
|
||||||
|
snippets.sort_by_key(|s| s.updated_at);
|
||||||
|
let to_remove = snippets.len() - (threshold / 2);
|
||||||
|
let removed: Vec<_> = snippets.drain(0..to_remove).collect();
|
||||||
|
for r in removed {
|
||||||
|
condensed_snippet_content.push_str(&format!("Name: {}\nDesc: {}\nCode: {}\n", r.name, r.description, r.code));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
if !condensed_snippet_content.is_empty() {
|
||||||
|
state.modify_graph(|graph| {
|
||||||
|
let name = format!("Snippet History {}", now);
|
||||||
|
graph.entities.insert(
|
||||||
|
name.clone(),
|
||||||
|
crate::models::Entity {
|
||||||
|
name: name.clone(),
|
||||||
|
entity_type: "Historical Summary".to_string(),
|
||||||
|
observations: vec![condensed_snippet_content],
|
||||||
|
namespace: crate::models::default_namespace(),
|
||||||
|
git_branch: None,
|
||||||
|
},
|
||||||
|
);
|
||||||
|
});
|
||||||
|
tracing::info!("Condensed snippets into Historical Summary.");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::Error>> {
|
async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::Error>> {
|
||||||
let state_for_index = Arc::clone(&state);
|
let state_for_index = Arc::clone(&state);
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
@@ -139,6 +215,7 @@ async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::E
|
|||||||
|
|
||||||
tokio::spawn(index_committer_worker(Arc::clone(&state)));
|
tokio::spawn(index_committer_worker(Arc::clone(&state)));
|
||||||
tokio::spawn(ttl_sweeper_worker(Arc::clone(&state)));
|
tokio::spawn(ttl_sweeper_worker(Arc::clone(&state)));
|
||||||
|
tokio::spawn(condense_graph_worker(Arc::clone(&state)));
|
||||||
crate::clipboard_watcher::spawn_watcher(Arc::clone(&state));
|
crate::clipboard_watcher::spawn_watcher(Arc::clone(&state));
|
||||||
crate::watcher::spawn_watcher(Arc::clone(&state));
|
crate::watcher::spawn_watcher(Arc::clone(&state));
|
||||||
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
|
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
|
||||||
|
|||||||
Reference in new issue
Block a user