refactor: extract inline loops, optimize GC scheduling, and remove redundancies
This commit is contained in:
1 parent
0723811686
commit
a5011048b1
1 file changed
+79
-84
+79
-84
@@ -81,6 +81,83 @@ enum GateCommands {
|
|||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn garbage_collector_worker(state: Arc<MemoryState>) {
|
||||||
|
loop {
|
||||||
|
// Run every 6 hours
|
||||||
|
tokio::time::sleep(tokio::time::Duration::from_secs(6 * 3600)).await;
|
||||||
|
|
||||||
|
let now = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs();
|
||||||
|
|
||||||
|
// 1. Task GC (14 days)
|
||||||
|
let fourteen_days = 14 * 24 * 3600;
|
||||||
|
let task_cutoff = now.saturating_sub(fourteen_days);
|
||||||
|
state.tasks.modify(|tasks| {
|
||||||
|
let initial_len = tasks.len();
|
||||||
|
tasks.retain(|task| !(task.status.to_lowercase() == "completed" && task.created_at < task_cutoff));
|
||||||
|
if tasks.len() < initial_len {
|
||||||
|
eprintln!("GC: Removed {} old completed tasks", initial_len - tasks.len());
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
// 2. Ledger GC (7 days or max 1000 items)
|
||||||
|
state.ledger.modify(|ledger| {
|
||||||
|
let seven_days = now.saturating_sub(7 * 24 * 3600);
|
||||||
|
ledger.retain(|c| c.timestamp >= seven_days);
|
||||||
|
if ledger.len() > 1000 {
|
||||||
|
let excess = ledger.len() - 1000;
|
||||||
|
ledger.drain(0..excess);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
// 3. Sticky Notes GC (24 hours)
|
||||||
|
state.sticky.modify(|notes| {
|
||||||
|
notes.retain(|note| note.timestamp >= now.saturating_sub(24 * 3600));
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn git_sync_worker(state: Arc<MemoryState>) {
|
||||||
|
let repo_path = std::env::current_dir().unwrap_or_else(|_| ".".into());
|
||||||
|
let mut last_commit_id = String::new();
|
||||||
|
|
||||||
|
loop {
|
||||||
|
tokio::time::sleep(tokio::time::Duration::from_secs(30)).await;
|
||||||
|
|
||||||
|
if let Ok(repo) = git2::Repository::discover(&repo_path) {
|
||||||
|
if let Ok(head) = repo.head() {
|
||||||
|
if let Ok(commit) = head.peel_to_commit() {
|
||||||
|
let current_id = commit.id().to_string();
|
||||||
|
if current_id != last_commit_id && !last_commit_id.is_empty() {
|
||||||
|
let msg = commit.message().unwrap_or("").to_string();
|
||||||
|
let branch = head.shorthand().unwrap_or("unknown").to_string();
|
||||||
|
|
||||||
|
state.ledger.modify(|changes| {
|
||||||
|
changes.push(crate::models::CodeChange {
|
||||||
|
git_commit: Some(current_id.clone()),
|
||||||
|
git_branch: Some(branch),
|
||||||
|
description: format!("Auto-synced commit: {}", msg.trim()),
|
||||||
|
timestamp: std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs(),
|
||||||
|
file_path: "".to_string(),
|
||||||
|
});
|
||||||
|
});
|
||||||
|
eprintln!("Git Sync: Logged new commit {}", current_id);
|
||||||
|
|
||||||
|
state.tasks.modify(|tasks| {
|
||||||
|
for task in tasks.iter_mut() {
|
||||||
|
if task.status != "completed" && msg.to_lowercase().contains(&task.title.to_lowercase()) {
|
||||||
|
task.status = "completed".to_string();
|
||||||
|
eprintln!("Git Sync: Auto-completed task '{}'", task.title);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
last_commit_id = current_id;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
async fn reconcile_worker(state: Arc<MemoryState>) {
|
async fn reconcile_worker(state: Arc<MemoryState>) {
|
||||||
loop {
|
loop {
|
||||||
sleep(Duration::from_secs(5)).await;
|
sleep(Duration::from_secs(5)).await;
|
||||||
@@ -98,21 +175,6 @@ async fn reconcile_worker(state: Arc<MemoryState>) {
|
|||||||
}).await;
|
}).await;
|
||||||
}
|
}
|
||||||
|
|
||||||
let now = SystemTime::now()
|
|
||||||
.duration_since(UNIX_EPOCH)
|
|
||||||
.unwrap()
|
|
||||||
.as_secs();
|
|
||||||
state.ledger.modify(|ledger| {
|
|
||||||
let seven_days = now.saturating_sub(7 * 24 * 60 * 60);
|
|
||||||
ledger.retain(|c| c.timestamp >= seven_days);
|
|
||||||
if ledger.len() > 1000 {
|
|
||||||
let excess = ledger.len() - 1000;
|
|
||||||
ledger.drain(0..excess);
|
|
||||||
}
|
|
||||||
});
|
|
||||||
state.sticky.modify(|notes| {
|
|
||||||
notes.retain(|note| note.timestamp >= now.saturating_sub(24 * 60 * 60));
|
|
||||||
});
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -423,76 +485,9 @@ fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::Error>>
|
|||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
// Background Garbage Collection for old tasks
|
tokio::spawn(garbage_collector_worker(Arc::clone(&state)));
|
||||||
let state_gc = Arc::clone(&state);
|
|
||||||
tokio::spawn(async move {
|
|
||||||
loop {
|
|
||||||
// Run every 24 hours
|
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(24 * 3600)).await;
|
|
||||||
|
|
||||||
let now = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs();
|
tokio::spawn(git_sync_worker(Arc::clone(&state)));
|
||||||
let fourteen_days = 14 * 24 * 3600;
|
|
||||||
let cutoff = now.saturating_sub(fourteen_days);
|
|
||||||
|
|
||||||
state_gc.tasks.modify(|tasks| {
|
|
||||||
let initial_len = tasks.len();
|
|
||||||
tasks.retain(|task| {
|
|
||||||
if task.status.to_lowercase() == "completed" && task.created_at < cutoff {
|
|
||||||
false // remove
|
|
||||||
} else {
|
|
||||||
true // keep
|
|
||||||
}
|
|
||||||
});
|
|
||||||
if tasks.len() < initial_len {
|
|
||||||
eprintln!("GC: Removed {} old completed tasks", initial_len - tasks.len());
|
|
||||||
}
|
|
||||||
});
|
|
||||||
}
|
|
||||||
});
|
|
||||||
|
|
||||||
// Git Native Sync Background Task
|
|
||||||
let state_git = Arc::clone(&state);
|
|
||||||
tokio::spawn(async move {
|
|
||||||
let repo_path = std::env::current_dir().unwrap_or_else(|_| ".".into());
|
|
||||||
let mut last_commit_id = String::new();
|
|
||||||
|
|
||||||
loop {
|
|
||||||
tokio::time::sleep(tokio::time::Duration::from_secs(30)).await;
|
|
||||||
|
|
||||||
if let Ok(repo) = git2::Repository::discover(&repo_path) {
|
|
||||||
if let Ok(head) = repo.head() {
|
|
||||||
if let Ok(commit) = head.peel_to_commit() {
|
|
||||||
let current_id = commit.id().to_string();
|
|
||||||
if current_id != last_commit_id && !last_commit_id.is_empty() {
|
|
||||||
let msg = commit.message().unwrap_or("").to_string();
|
|
||||||
let branch = head.shorthand().unwrap_or("unknown").to_string();
|
|
||||||
|
|
||||||
state_git.ledger.modify(|changes| {
|
|
||||||
changes.push(crate::models::CodeChange {
|
|
||||||
git_commit: Some(current_id.clone()),
|
|
||||||
git_branch: Some(branch),
|
|
||||||
description: format!("Auto-synced commit: {}", msg.trim()),
|
|
||||||
timestamp: std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs(),
|
|
||||||
file_path: "".to_string(),
|
|
||||||
});
|
|
||||||
});
|
|
||||||
eprintln!("Git Sync: Logged new commit {}", current_id);
|
|
||||||
|
|
||||||
state_git.tasks.modify(|tasks| {
|
|
||||||
for task in tasks.iter_mut() {
|
|
||||||
if task.status != "completed" && msg.to_lowercase().contains(&task.title.to_lowercase()) {
|
|
||||||
task.status = "completed".to_string();
|
|
||||||
eprintln!("Git Sync: Auto-completed task '{}'", task.title);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
});
|
|
||||||
}
|
|
||||||
last_commit_id = current_id;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
});
|
|
||||||
eprintln!("MCP Memory Server running on http://127.0.0.1:3000/sse");
|
eprintln!("MCP Memory Server running on http://127.0.0.1:3000/sse");
|
||||||
if let Err(e) = axum::serve(listener, app).await {
|
if let Err(e) = axum::serve(listener, app).await {
|
||||||
let log_path = dirs::home_dir()
|
let log_path = dirs::home_dir()
|
||||||
|
|||||||
Reference in new issue
Block a user