chore: housekeeping

This commit is contained in:
Riz Ashraf committed 2026-10-09 09:33:54 +01:00
1 parent 4ea8f861cc
commit 0c2fc96ff8
29 files changed
+13 -1232

No files matched your search

-293
View File
@@ -1,293 +0,0 @@
if let Some(idx) = gates.iter().position(|g| {
g.action == q.action
&& g.target == q.target
&& g.namespace == q.namespace
&& g.params == q.params
}) {
found = Some(gates[idx].clone());
if q.consume {
to_remove = Some(idx);
}
}
if let Some(idx) = to_remove {
gates.remove(idx);
}
});
match found {
Some(record) => {
if record.status == "authorized" {
(axum::http::StatusCode::OK, "Authorized").into_response()
} else {
let msg = if let Some(r) = record.reason {
format!("Action blocked. Reason: {}", r)
} else {
"Action blocked.".to_string()
};
(axum::http::StatusCode::FORBIDDEN, msg).into_response()
}
}
None => (
axum::http::StatusCode::NOT_FOUND,
"Action not yet authorized (no gate record found).",
)
.into_response(),
}
}
async fn gate_set_handler(
State(app_state): State<Arc<AppState>>,
Json(body): Json<GateSetReq>,
) -> axum::response::Response {
let status = if body.block.unwrap_or(false) {
"blocked".to_string()
} else if body.authorize.unwrap_or(false) {
"authorized".to_string()
} else {
"pending".to_string()
};
let record = GateRecord {
id: uuid::Uuid::new_v4().to_string(),
action: body.action.clone(),
target: body.target.clone(),
namespace: body.namespace.clone(),
params: body.params.clone(),
status,
reason: body.reason.clone(),
timestamp: crate::handlers::utils::now_secs(),
};
app_state.handler.state.gates.modify(|gates| {
gates.retain(|g| !(g.action == record.action && g.target == record.target));
gates.push(record);
});
(axum::http::StatusCode::OK, "Gate state updated.").into_response()
}
async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::Error>> {
state.rebuild_index().await;
tokio::spawn(index_committer_worker(Arc::clone(&state)));
let app_state = Arc::new(AppState {
handler: Arc::new(MemoryHandler::new(Arc::clone(&state))),
clients: RwLock::new(HashMap::new()),
next_id: AtomicUsize::new(1),
});
let app_state_clone = Arc::clone(&app_state);
let mut rx = state.activity_tx.subscribe();
tokio::spawn(async move {
while let Ok(msg) = rx.recv().await {
let senders: Vec<_> = app_state_clone
.clients
.read()
.unwrap_or_else(|e| e.into_inner())
.values()
.cloned()
.collect();
for client_tx in senders {
let _ = client_tx.try_send(msg.clone());
}
}
});
let app = Router::new()
.route(
"/api/version",
get(|| async move {
axum::Json(serde_json::json!({
"version": env!("APP_VERSION"),
"git_hash": option_env!("GIT_HASH").unwrap_or("unknown")
}))
}),
)
.route("/ws", get(ws_handler))
.route("/health", get(health_handler))
.route("/nvim/telemetry", post(nvim_telemetry_handler))
.route("/gate/verify", get(gate_verify_handler))
.route("/gate/set", post(gate_set_handler))
.route(
"/shutdown",
post(
|headers: axum::http::HeaderMap, State(state): State<Arc<AppState>>| async move {
let token_path = state.handler.state.base_dir.join("admin.token");
let expected_token = tokio::fs::read_to_string(&token_path)
.await
.unwrap_or_default()
.trim()
.to_string();
let auth_header = headers
.get(axum::http::header::AUTHORIZATION)
.and_then(|h| h.to_str().ok())
.unwrap_or_default();
if expected_token.is_empty() || auth_header != format!("Bearer {}", expected_token) {
return (axum::http::StatusCode::UNAUTHORIZED, "Unauthorized").into_response();
}
std::thread::spawn(|| {
tracing::info!(
"Received shutdown request via /shutdown endpoint. Exiting process cleanly."
);
std::thread::sleep(std::time::Duration::from_millis(100));
std::process::exit(0);
});
(axum::http::StatusCode::OK, "Shutting down...").into_response()
},
),
)
.route(
"/",
get(|| async move { axum::response::Html(include_str!("dashboard.html")) }),
)
.route(
"/api/graph",
get({
let state_clone = app_state.handler.state.clone();
move || async move {
let graph_json = state_clone.read_graph(|g| serde_json::to_string(g).unwrap_or_else(|_| "{}".to_string()));
([(axum::http::header::CONTENT_TYPE, "application/json")], graph_json)
}
}),
)
.route(
"/api/tasks/{id}/complete",
post({
let state_clone = app_state.handler.state.clone();
move |axum::extract::Path(id): axum::extract::Path<String>| async move {
state_clone.tasks.modify(|tasks| {
for t in tasks.iter_mut() {
if t.id == id {
t.status = "completed".to_string();
break;
}
}
});
axum::Json(serde_json::json!({"status": "success"}))
}
}),
)
.route(
"/api/tasks",
get({
let state_clone = app_state.handler.state.clone();
move || async move {
let tasks_json = state_clone.tasks.read_with(|t| serde_json::to_string(t).unwrap_or_else(|_| "[]".to_string()));
([(axum::http::header::CONTENT_TYPE, "application/json")], tasks_json)
}
}),
)
.route(
"/api/sticky",
get({
let state_clone = app_state.handler.state.clone();
move || async move {
let sticky_json = state_clone.sticky.read_with(|s| serde_json::to_string(s).unwrap_or_else(|_| "[]".to_string()));
([(axum::http::header::CONTENT_TYPE, "application/json")], sticky_json)
}
}),
)
.route(
"/api/search",
get({
let state_clone = app_state.handler.state.clone();
move |axum::extract::Query(params): axum::extract::Query<
std::collections::HashMap<String, String>,
>| async move {
if let Some(q) = params.get("q")
&& let Ok(idx) = state_clone.search_index.read()
&& let Ok(results) = idx.search(q, None) {
let mut formatted_results = Vec::new();
for (id, doc_type, title, body, score) in results {
formatted_results.push(serde_json::json!({
"id": id,
"type_name": doc_type,
"title": title,
"content": body,
"score": score
}));
}
return axum::Json(
serde_json::json!({ "results": formatted_results }),
);
}
axum::Json(serde_json::json!({ "results": [] }))
}
}),
)
.route(
"/api/activity",
get({
let state_clone = app_state.handler.state.clone();
move || async move {
let activities_json = state_clone.recent_activities.read_with(|a| serde_json::to_string(a).unwrap_or_else(|_| "[]".to_string()));
([(axum::http::header::CONTENT_TYPE, "application/json")], activities_json)
}
}),
)
.route(
"/api/stats",
get({
let state_clone = app_state.handler.state.clone();
move || async move {
let (entities, relations) = state_clone.read_graph(|g| (g.entities.len(), g.relations.len()));
let tasks = state_clone.tasks.read_with(|items| items.len());
let snippets = state_clone.snippets.read_with(|items| items.len());
let tech_debts = state_clone.tech_debts.read_with(|items| items.len());
let adrs = state_clone.adrs.read_with(|items| items.len());
let ledger = state_clone.ledger.read_with(|items| items.len());
let error_fixes = state_clone.error_fixes.read_with(|items| items.len());
let session_summaries = state_clone.session_summaries.read_with(|items| items.len());
let handoff_memos = state_clone.handoff_memos.read_with(|items| items.len());
let env_fingerprints = state_clone.env_fingerprints.read_with(|items| items.len());
let env_requirements = state_clone.env_requirements.read_with(|items| items.len());
let milestones = state_clone.milestones.read_with(|items| items.len());
let environments = state_clone.environments.read_with(|items| items.len());
let gates = state_clone.gates.read_with(|items| items.len());
axum::Json(serde_json::json!({
"entities": entities,
"relations": relations,
"tasks": tasks,
"snippets": snippets,
"tech_debts": tech_debts,
"adrs": adrs,
"ledger": ledger,
"error_fixes": error_fixes,
"session_summaries": session_summaries,
"handoff_memos": handoff_memos,
"env_fingerprints": env_fingerprints,
"env_requirements": env_requirements,
"milestones": milestones,
"environments": environments,
"gates": gates
}))
}
}),
)
.with_state(app_state);
tracing::info!("MCP Memory Server running on http://127.0.0.1:3000/sse");
let addr = std::env::var("MCP_PORT").unwrap_or_else(|_| "3000".to_string());
let addr: std::net::SocketAddr = format!("127.0.0.1:{}", addr)
.parse()
.expect("Invalid bind address");
let listener = match tokio::net::TcpListener::bind(&addr).await {
Ok(l) => l,
Err(e) => {
let log_path = dirs::home_dir()
.unwrap_or_default()
.join(".gemini/mcp_memory/daemon_error.log");
let _ =
tokio::fs::write(&log_path, format!("Failed to bind to {}: {}\n", addr, e)).await;
return Ok(());
}
};
if let Err(e) = axum::serve(listener, app.into_make_service()).await {
let log_path = dirs::home_dir()
.unwrap_or_default()
.join(".gemini/mcp_memory/daemon_error.log");
let _ = tokio::fs::write(&log_path, format!("Server crashed: {}\n", e)).await;
}
+7 -60
View File
@@ -117,75 +117,22 @@ pub fn init_redb(base: &Path) -> Arc<Database> {
}
};
// Ensure the table exists and migrate legacy JSON files
// Ensure the table exists
match db.begin_write() {
Ok(write_txn) => {
let mut opened_ok = false;
if let Ok(mut table) = write_txn.open_table(STORE_TABLE) {
if let Ok(_) = write_txn.open_table(STORE_TABLE) {
opened_ok = true;
if !is_in_memory {
let stores = vec![
("knowledge_graph_master", "knowledge_graph_master.json"),
("audit_ledger", "audit_ledger.json"),
("tasks", "tasks.json"),
("snippets", "snippets.json"),
("adrs", "adrs.json"),
("error_fixes", "error_fixes.json"),
("session_summaries", "session_summaries.json"),
("handoff_memos", "handoff_memos.json"),
("env_fingerprints", "env_fingerprints.json"),
("env_requirements", "env_requirements.json"),
("milestones", "milestones.json"),
("environments", "environments.json"),
("tech_debts", "tech_debts.json"),
("gates", "gates.json"),
("state_snapshots", "state_snapshots.json"),
("hypotheses", "hypotheses.json"),
("agent_signals", "agent_signals.json"),
];
for (key, file_name) in stores.iter() {
let is_missing = match table.get(*key) {
Ok(res) => res.is_none(),
Err(e) => {
tracing::warn!("Failed to read key '{}' from redb: {}", key, e);
false
}
};
if is_missing {
let json_path = base.join(file_name);
if json_path.exists()
&& let Ok(data) = std::fs::read(&json_path)
&& serde_json::from_slice::<serde_json::Value>(&data).is_ok()
{
if let Err(e) = table.insert(*key, data.as_slice()) {
tracing::error!(
"Failed to insert migrated key '{}': {}",
key,
e
);
} else {
let migrated_path = json_path.with_extension("json.migrated");
if std::fs::rename(&json_path, &migrated_path).is_err()
&& migrated_path.exists()
{
let _ = std::fs::remove_file(&migrated_path);
let _ = std::fs::rename(&json_path, &migrated_path);
}
}
}
}
}
}
}
if opened_ok && let Err(e) = write_txn.commit() {
tracing::error!("Failed to commit database migration transaction: {}", e);
if opened_ok {
if let Err(e) = write_txn.commit() {
tracing::error!("Failed to commit database initialization transaction: {}", e);
}
}
}
Err(e) => {
tracing::error!(
"Failed to begin write transaction for redb migration: {}",
"Failed to begin write transaction for redb initialization: {}",
e
);
}
+3 -3
View File
@@ -977,8 +977,8 @@ impl McpTool for GetSubgraphHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: GetSubgraphTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let root = req.root_entity.or(req.root_node).ok_or_else(|| {
crate::error::AppError::Internal("root_entity or root_node is required".to_string())
let root = req.root_entity.ok_or_else(|| {
crate::error::AppError::Internal("root_entity is required".to_string())
})?;
let depth = req.depth.unwrap_or(2);
let format = req.format.unwrap_or(SubgraphFormat::Json);
@@ -1043,7 +1043,7 @@ impl McpTool for GetSubgraphHandler {
}
let result = serde_json::json!({
"root_node": root,
"root_entity": root,
"depth": depth,
"entities": matched_entities,
"relations": matched_relations,
-93
View File
@@ -843,69 +843,7 @@ impl McpTool for ClipboardHandler {
})
.to_string())
}
ClipboardAction::Read => {
let engine = ensure_ocr_engine().await;
let out = tokio::task::spawn_blocking(
move || -> crate::error::Result<serde_json::Map<String, Value>> {
let mut out = serde_json::Map::new();
if let Some(text) = get_native_clipboard_text() {
out.insert("text".into(), json!(text));
}
if let Some(dynamic_img) = get_native_clipboard_image() {
let mut img = dynamic_img.clone();
let max_dim = 1440;
if img.width() > max_dim || img.height() > max_dim {
img = img.resize(max_dim, max_dim, FilterType::Lanczos3);
}
let rgb_img = img.into_rgb8();
let mut jpeg_bytes = std::io::Cursor::new(Vec::new());
let mut encoder = image::codecs::jpeg::JpegEncoder::new_with_quality(
&mut jpeg_bytes,
88,
);
if encoder
.encode(
&rgb_img,
rgb_img.width(),
rgb_img.height(),
image::ExtendedColorType::Rgb8,
)
.is_ok()
{
let bytes = jpeg_bytes.into_inner();
let cache_dir = dirs::home_dir()
.unwrap_or_default()
.join(".gemini/mcp_memory/clipboard");
let _ = std::fs::create_dir_all(&cache_dir);
let file_path = cache_dir.join("clipboard_latest_image.jpg");
if std::fs::write(&file_path, &bytes).is_ok() {
let path_str = file_path.to_string_lossy().to_string();
out.insert("image_path".into(), json!(path_str));
let wsl_path = to_wsl_path(&path_str);
out.insert("image_path_wsl".into(), json!(wsl_path));
}
}
if let Some(eng) = engine
&& let Some(ocr_text) = perform_ocrs_ocr(eng, &dynamic_img)
{
out.insert("image_analysis".to_string(), json!(ocr_text.trim()));
}
}
Ok(out)
},
)
.await
.map_err(|e| crate::error::AppError::Internal(format!("Task panic: {}", e)))??;
state.record_activity("clipboard", "Read contents from OS clipboard", None);
Ok::<String, crate::error::AppError>(serde_yaml::to_string(&Value::Object(out))?)
}
ClipboardAction::Write => {
let text_opt = req.text;
let image_path_opt = req.image_path;
@@ -1023,38 +961,7 @@ mod tests {
);
}
#[tokio::test]
async fn test_read_clipboard() {
let dir = tempdir().unwrap();
let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap()));
let handler = ClipboardHandler;
let result = handler
.execute(json!({"action": "read"}), state)
.await
.map_err(|e| format!("Failed to read clipboard: {}", e))
.unwrap();
// Returns a JSON string, possibly {}
let parsed: serde_json::Value = serde_yaml::from_str(&result).unwrap();
assert!(parsed.is_object());
}
#[tokio::test]
async fn test_read_clipboard_empty() {
let dir = tempfile::tempdir().unwrap();
let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap()));
let handler = ClipboardHandler;
let result = handler
.execute(serde_json::json!({"action": "read"}), state)
.await
.map_err(|e| format!("Failed to read clipboard: {}", e))
.unwrap();
let parsed: serde_json::Value = serde_yaml::from_str(&result).unwrap();
assert!(parsed.is_object());
}
#[test]
fn test_no_subprocess_clipboard_regression() {
+1 -1
View File
@@ -59,7 +59,7 @@ The server registers 5 high-signal workflow prompts to initiate standardized age
## 5. Consolidated Smart Tools Architecture (11 Primary Tools)
The server consolidates granular single-purpose tools into domain-named smart tools. Always prefer the consolidated tools over legacy aliases:
The server consolidates granular single-purpose tools into domain-named smart tools.
* **`tasks`**: Complete task lifecycle management.
- `action: "add"`: Create a new task (requires `title`, optional `description`, `git_branch`, `repo_name`, `priority: "low" | "medium" | "high" | "urgent"`, `assigned_agent`, `verification_command`, `parent_id`, `dependencies`).
+1 -1
View File
@@ -1428,7 +1428,7 @@ mod tests {
("decisions", json!({"action": "query"})),
("tech_debt", json!({"action": "list"})),
("snippets", json!({"action": "search", "query": "test"})),
("clipboard", json!({"action": "read"})),
("clipboard", json!({"action": "history"})),
("environment", json!({"action": "read_fingerprint"})),
("omni_search", json!({"query": "test"})),
("get_project_health", json!({})),
-5
View File
@@ -172,8 +172,6 @@ pub enum SubgraphFormat {
pub struct GetSubgraphTool {
/// The root entity name to start the subgraph search from.
pub root_entity: Option<String>,
/// Legacy alias for root_entity.
pub root_node: Option<String>,
/// Maximum search depth (hops). Defaults to 2.
pub depth: Option<u32>,
/// Output format: 'json' (raw entities and relations) or 'markdown_tree' (compact topology tree). Defaults to 'json'.
@@ -1014,8 +1012,6 @@ pub enum ClipboardAction {
History,
#[serde(alias = "clear", alias = "CLEAR", alias = "Clear")]
Clear,
#[serde(alias = "read", alias = "READ", alias = "Read")]
Read,
#[serde(alias = "write", alias = "WRITE", alias = "Write")]
Write,
}
@@ -1026,7 +1022,6 @@ pub enum ClipboardAction {
/// - 'text': Get latest clipboard text (or normalized Markdown if HTML was copied).
/// - 'history': View recent clipboard history ring buffer (images and text with timestamps).
/// - 'clear': Clear OS clipboard and memory cache.
/// - 'read': Read current clipboard contents (legacy alias).
/// - 'write': Write content to OS clipboard. Optional: text, html, files, image_path.
///
/// Triggers: Call 'image' immediately when user says "look at image in clipboard", "see screenshot",
-1
View File
@@ -1 +0,0 @@
test
-1
View File
@@ -1 +0,0 @@
test