chore: fix formatting and clippy lints
This commit is contained in:
1 parent
8952bd5399
commit
9d9e959744
35 files changed
+1091
-582
No files matched your search
+122
-38
@@ -4,6 +4,7 @@
|
||||
)]
|
||||
|
||||
pub mod api;
|
||||
pub mod config;
|
||||
pub mod db;
|
||||
pub mod embedding;
|
||||
pub mod error;
|
||||
@@ -105,11 +106,10 @@ pub async fn ttl_sweeper_worker(state: Arc<MemoryState>) {
|
||||
|
||||
state.project.tasks.read_with(|tasks| {
|
||||
for t in tasks.iter() {
|
||||
if let Some(exp) = t.expires_at {
|
||||
if t.is_active() {
|
||||
if let Some(exp) = t.expires_at
|
||||
&& t.is_active() {
|
||||
next_expiry = Some(next_expiry.map_or(exp, |curr| curr.min(exp)));
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
@@ -160,13 +160,12 @@ pub async fn ttl_sweeper_worker(state: Arc<MemoryState>) {
|
||||
let mut expired_tasks = Vec::new();
|
||||
state.project.tasks.modify(|tasks| {
|
||||
for t in tasks.iter_mut() {
|
||||
if let Some(exp) = t.expires_at {
|
||||
if exp <= now && t.is_active() {
|
||||
if let Some(exp) = t.expires_at
|
||||
&& exp <= now && t.is_active() {
|
||||
t.status = "expired".to_string();
|
||||
t.updated_at = now;
|
||||
expired_tasks.push(t.id.clone());
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
@@ -249,8 +248,8 @@ pub async fn condense_graph_worker(state: Arc<MemoryState>) {
|
||||
}
|
||||
});
|
||||
|
||||
if let Some((content, names)) = snippet_condensation {
|
||||
if !content.is_empty() {
|
||||
if let Some((content, names)) = snippet_condensation
|
||||
&& !content.is_empty() {
|
||||
let name = format!("Snippet History {}", now);
|
||||
state.modify_graph(|graph| {
|
||||
graph.entities.insert(
|
||||
@@ -271,6 +270,84 @@ pub async fn condense_graph_worker(state: Arc<MemoryState>) {
|
||||
});
|
||||
tracing::info!("Condensed snippets into Historical Summary.");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn memory_consolidation_worker(state: Arc<MemoryState>) {
|
||||
let mut interval = tokio::time::interval(std::time::Duration::from_secs(300));
|
||||
loop {
|
||||
tokio::select! {
|
||||
_ = state.shutdown_notify.notified() => break,
|
||||
_ = interval.tick() => {},
|
||||
}
|
||||
|
||||
let entities: Vec<_> = state.graph.read_with(|g| {
|
||||
g.entities
|
||||
.values()
|
||||
.map(|e| (e.name.clone(), e.entity_type.clone()))
|
||||
.collect()
|
||||
});
|
||||
|
||||
if entities.len() < 2 {
|
||||
continue;
|
||||
}
|
||||
|
||||
let mut entity_summaries = String::new();
|
||||
for (name, e_type) in entities.iter().take(50) {
|
||||
entity_summaries.push_str(&format!("- [{}] {}\n", e_type, name));
|
||||
}
|
||||
|
||||
let prompt = format!(
|
||||
"Analyze the following list of entities and identify exactly TWO that represent the exact same concept or item but have slightly different names (e.g. 'auth_service' and 'AuthService'). Return ONLY a valid JSON array containing exactly two strings: the two names to merge. If no obvious duplicates exist, return an empty array []. Do not output any markdown formatting or extra text.\n\nEntities:\n{}",
|
||||
entity_summaries
|
||||
);
|
||||
|
||||
if let Ok(response) = state
|
||||
.ollama
|
||||
.generate(
|
||||
&prompt,
|
||||
None,
|
||||
Some("You are a helpful JSON-only data deduplication assistant. Output only JSON."),
|
||||
)
|
||||
.await
|
||||
{
|
||||
let cleaned = response
|
||||
.trim()
|
||||
.trim_start_matches("```json")
|
||||
.trim_start_matches("```")
|
||||
.trim_end_matches("```")
|
||||
.trim();
|
||||
if let Ok(duplicates) = serde_json::from_str::<Vec<String>>(cleaned)
|
||||
&& duplicates.len() == 2 {
|
||||
let e1_name = &duplicates[0];
|
||||
let e2_name = &duplicates[1];
|
||||
|
||||
if e1_name != e2_name {
|
||||
tracing::info!(
|
||||
"Memory Consolidation Daemon: Merging '{}' into '{}'",
|
||||
e2_name,
|
||||
e1_name
|
||||
);
|
||||
state.modify_graph(|g| {
|
||||
if let Some(mut e2) = g.entities.remove(e2_name) {
|
||||
if let Some(e1) = g.entities.get_mut(e1_name) {
|
||||
e1.observations.append(&mut e2.observations);
|
||||
} else {
|
||||
g.entities.insert(e2_name.clone(), e2);
|
||||
}
|
||||
}
|
||||
|
||||
for rel in g.relations.iter_mut() {
|
||||
if rel.from == *e2_name {
|
||||
rel.from = e1_name.clone();
|
||||
}
|
||||
if rel.to == *e2_name {
|
||||
rel.to = e1_name.clone();
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -287,6 +364,7 @@ pub async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::erro
|
||||
tokio::spawn(index_committer_worker(Arc::clone(&state)));
|
||||
tokio::spawn(ttl_sweeper_worker(Arc::clone(&state)));
|
||||
tokio::spawn(condense_graph_worker(Arc::clone(&state)));
|
||||
tokio::spawn(memory_consolidation_worker(Arc::clone(&state)));
|
||||
crate::watcher::spawn_watcher(Arc::clone(&state));
|
||||
crate::handlers::vision::spawn_clipboard_listener(Arc::clone(&state));
|
||||
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
|
||||
@@ -337,7 +415,7 @@ pub async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::erro
|
||||
if let Ok(socket) = tokio::net::UdpSocket::bind(format!("127.0.0.1:{}", port1)).await {
|
||||
let socket = Arc::new(socket);
|
||||
let socket_rx = socket.clone();
|
||||
|
||||
|
||||
let mut subscribers: HashMap<(String, String), std::net::SocketAddr> = HashMap::new();
|
||||
let mut buf = vec![0u8; 65536];
|
||||
let mut event_rx = udp_state.handler.state.event_bus_tx.subscribe();
|
||||
@@ -371,32 +449,29 @@ pub async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::erro
|
||||
} else if let Ok(json_payload) = serde_json::from_slice::<serde_json::Value>(&buf[..len]) {
|
||||
if json_payload.get("type").and_then(|t| t.as_str()) == Some("ping") {
|
||||
let _ = socket.send_to(b"pong", addr).await;
|
||||
} else if json_payload.get("type").and_then(|t| t.as_str()) == Some("gate_wait") {
|
||||
if let (Some(action), Some(target)) = (
|
||||
} else if json_payload.get("type").and_then(|t| t.as_str()) == Some("gate_wait")
|
||||
&& let (Some(action), Some(target)) = (
|
||||
json_payload.get("action").and_then(|a| a.as_str()),
|
||||
json_payload.get("target").and_then(|t| t.as_str())
|
||||
) {
|
||||
subscribers.insert((action.to_string(), target.to_string()), addr);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(event) = event_rx.recv() => {
|
||||
if event.topic == "gate:event" {
|
||||
if let (Some(action), Some(target), Some(status)) = (
|
||||
if event.topic == "gate:event"
|
||||
&& let (Some(action), Some(target), Some(status)) = (
|
||||
event.payload.get("action").and_then(|a| a.as_str()),
|
||||
event.payload.get("target").and_then(|t| t.as_str()),
|
||||
event.payload.get("status").and_then(|s| s.as_str()),
|
||||
) {
|
||||
if status == "authorized" || status == "blocked" {
|
||||
if let Some(addr) = subscribers.remove(&(action.to_string(), target.to_string())) {
|
||||
if (status == "authorized" || status == "blocked")
|
||||
&& let Some(addr) = subscribers.remove(&(action.to_string(), target.to_string())) {
|
||||
let response = if status == "authorized" { b"APPROVED" } else { b"REJECTED" };
|
||||
let _ = socket.send_to(response, addr).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -476,34 +551,42 @@ pub async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::erro
|
||||
if payload.event.starts_with("agent_") || payload.event.starts_with("diff_") {
|
||||
let payload_val = serde_json::json!(&payload);
|
||||
// 1. General event topic (e.g. nvim:ui:agent_prompt_response, nvim:ui:agent_diff_accepted)
|
||||
let _ = nvim_udp_state.handler.state.event_bus_tx.send(crate::state::GenericEvent {
|
||||
topic: format!("nvim:ui:{}", payload.event),
|
||||
session_id: Some(payload.session_id.clone()),
|
||||
payload: payload_val.clone(),
|
||||
});
|
||||
let _ = nvim_udp_state.handler.state.event_bus_tx.send(
|
||||
crate::state::GenericEvent {
|
||||
topic: format!("nvim:ui:{}", payload.event),
|
||||
session_id: Some(payload.session_id.clone()),
|
||||
payload: payload_val.clone(),
|
||||
},
|
||||
);
|
||||
|
||||
// 2. Correlated request_id topic (e.g. nvim:ui:agent_prompt_response:REQ_ID)
|
||||
if let Some(ref req_id) = payload.request_id {
|
||||
let _ = nvim_udp_state.handler.state.event_bus_tx.send(crate::state::GenericEvent {
|
||||
topic: format!("nvim:ui:{}:{}", payload.event, req_id),
|
||||
session_id: Some(payload.session_id.clone()),
|
||||
payload: payload_val.clone(),
|
||||
});
|
||||
let _ = nvim_udp_state.handler.state.event_bus_tx.send(
|
||||
crate::state::GenericEvent {
|
||||
topic: format!("nvim:ui:{}:{}", payload.event, req_id),
|
||||
session_id: Some(payload.session_id.clone()),
|
||||
payload: payload_val.clone(),
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
// 3. Correlated diff_id topics
|
||||
if let Some(ref diff_id) = payload.diff_id {
|
||||
let _ = nvim_udp_state.handler.state.event_bus_tx.send(crate::state::GenericEvent {
|
||||
topic: format!("nvim:ui:{}:{}", payload.event, diff_id),
|
||||
session_id: Some(payload.session_id.clone()),
|
||||
payload: payload_val.clone(),
|
||||
});
|
||||
let _ = nvim_udp_state.handler.state.event_bus_tx.send(
|
||||
crate::state::GenericEvent {
|
||||
topic: format!("nvim:ui:{}:{}", payload.event, diff_id),
|
||||
session_id: Some(payload.session_id.clone()),
|
||||
payload: payload_val.clone(),
|
||||
},
|
||||
);
|
||||
// General diff decision topic
|
||||
let _ = nvim_udp_state.handler.state.event_bus_tx.send(crate::state::GenericEvent {
|
||||
topic: format!("nvim:ui:diff_decision:{}", diff_id),
|
||||
session_id: Some(payload.session_id.clone()),
|
||||
payload: payload_val.clone(),
|
||||
});
|
||||
let _ = nvim_udp_state.handler.state.event_bus_tx.send(
|
||||
crate::state::GenericEvent {
|
||||
topic: format!("nvim:ui:diff_decision:{}", diff_id),
|
||||
session_id: Some(payload.session_id.clone()),
|
||||
payload: payload_val.clone(),
|
||||
},
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -576,6 +659,7 @@ pub fn init_logging(app_name: &str) -> Option<tracing_appender::non_blocking::Wo
|
||||
}
|
||||
|
||||
pub fn run_cli() -> Result<(), Box<dyn std::error::Error>> {
|
||||
crate::config::load_mcp_config_env();
|
||||
let _guard = init_logging("mcp-memory-server");
|
||||
let cli = Cli::parse();
|
||||
|
||||
|
||||
Reference in new issue
Block a user