diff --git a/agent-rules/global_subagents.md b/agent-rules/global_subagents.md index 9b194c0..56941c1 100644 --- a/agent-rules/global_subagents.md +++ b/agent-rules/global_subagents.md @@ -30,7 +30,7 @@ Invoke these dynamically using the `invoke_subagent` tool. Use `send_message` to * **Role:** Project & Task Orchestrator * **Permissions:** You MUST define this subagent with `enable_mcp_tools: true` and `enable_write_tools: true`. * **Trigger:** When we start a new feature or finish a task. -* **Action:** Passively monitors the `memory://tasks/active` and `memory://milestones` MCP resources. Creates tasks, organizes milestones, and updates statuses when criteria are met using `add_task`, `update_task_status`, `add_milestone`, and `verify_acceptance_criteria`. +* **Action:** Passively monitors the `memory://tasks/active` and `memory://milestones` MCP resources. Creates tasks, organizes milestones, and updates statuses when criteria are met using `tasks` and `milestones`. ## 5. DevOpsSRE * **Role:** Environment & Handoff Manager diff --git a/agent-rules/mcp_and_build_constraints.md b/agent-rules/mcp_and_build_constraints.md index 939b88f..531e059 100644 --- a/agent-rules/mcp_and_build_constraints.md +++ b/agent-rules/mcp_and_build_constraints.md @@ -12,5 +12,5 @@ description: Strict constraints for Model Context Protocol stdio transport and p - **Rule:** ALWAYS use the project's native `justfile` targets (e.g., `just build-wsl`, `just deploy-stub-wsl`, `just deploy-all`) to manage building and deploying binaries across Windows and WSL boundaries. # 3. Clipboard & Vision Tooling (Strict MCP Enforcement) -- **Rule (CRITICAL):** You MUST ALWAYS execute the `read_clipboard` MCP tool (`mcp-memory`) when reading or inspecting text or images from the OS clipboard. -- **Rule:** NEVER run ad-hoc PowerShell (`System.Windows.Forms.Clipboard`), bash, or shell scripts to extract clipboard images or saved files. `read_clipboard` handles native clipboard image extraction, auto-resizing, and compressed JPEG (`.jpg`) encoding to minimize LLM token usage. +- **Rule (CRITICAL):** You MUST ALWAYS execute the `clipboard` MCP tool with action="read" (`mcp-memory`) when reading or inspecting text or images from the OS clipboard. +- **Rule:** NEVER run ad-hoc PowerShell (`System.Windows.Forms.Clipboard`), bash, or shell scripts to extract clipboard images or saved files. `clipboard` (read) handles native clipboard image extraction, auto-resizing, and compressed JPEG (`.jpg`) encoding to minimize LLM token usage. diff --git a/agent-rules/mcp_memory_integrity.md b/agent-rules/mcp_memory_integrity.md index 076410b..5fe35a5 100644 --- a/agent-rules/mcp_memory_integrity.md +++ b/agent-rules/mcp_memory_integrity.md @@ -2,6 +2,6 @@ **CRITICAL RULE:** Under NO CIRCUMSTANCES should any agent or subagent directly manipulate, edit, or write raw JSON/data to the MCP Memory Graph files (e.g., `knowledge_graph_master.json`, `.db` files, or any files inside `~/.gemini/mcp_memory/`). -All interactions, additions, creations, and updates to the knowledge graph or task tracker MUST go through the officially exposed MCP tools (such as `add_task`, `set_acceptance_criteria`, `verify_acceptance_criteria`, `log_code_change`, `create_entities`, etc.). +All interactions, additions, creations, and updates to the knowledge graph or task tracker MUST go through the officially exposed MCP tools (such as `tasks` [add, set_criteria, verify], `log_code_change`, `create_entities`, etc.). Bypassing the MCP API by using file-editing tools (like `replace_file_content` or `write_to_file`) on the database files corrupts the state, bypasses indexing, and destroys the Redb Write-Ahead Log (WAL). If an agent attempts to do this, immediately stop them and route the request correctly through the provided MCP API tools. diff --git a/justfile b/justfile index 3654e7d..64d0ee8 100644 --- a/justfile +++ b/justfile @@ -16,7 +16,7 @@ deploy-all: deploy-win deploy-wsl @Write-Host "Deployment complete across both OS boundaries." -ForegroundColor Green # Build Windows-native binaries -build-win: build-server build-stub-win build-nvim-win +build-win: build-ui build-server build-stub-win build-nvim-win # Deploy Windows-native binaries and rules deploy-win: deploy-server deploy-stub-win deploy-nvim-win deploy-rules-win @@ -57,8 +57,18 @@ run-server: # 2. BUILD (Compile Binaries) # ========================================================= -# Build Windows-native server binary in release mode -build-server: +# Transpile Dashboard TypeScript (server/src/dashboard.ts) into JavaScript +build-ui: + @Write-Host "Building Dashboard JavaScript from server/src/dashboard.ts..." -ForegroundColor Cyan + bun build server/src/dashboard.ts --outfile=server/src/dashboard.js --target=browser + +# Type-check Dashboard TypeScript source code +check-ui: + @Write-Host "Type-checking Dashboard TypeScript..." -ForegroundColor Cyan + deno check server/src/dashboard.ts + +# Build Windows-native server binary in release mode (ensuring fresh UI build) +build-server: build-ui @Write-Host "Building Windows Server..." -ForegroundColor Cyan cargo build --release -p mcp-memory-server @@ -138,11 +148,16 @@ deploy-rules-wsl: # 4. TESTING & COVERAGE # ========================================================= -# Run fast parallel unit tests via cargo-nextest -test: +# Run fast parallel unit tests via cargo-nextest (and type-check UI) +test: check-ui @Write-Host "Running fast parallel tests via cargo-nextest..." -ForegroundColor Cyan cargo nextest run --workspace +# Run dashboard UI endpoint and integrity unit tests +test-ui: + @Write-Host "Running dashboard UI endpoint and integrity unit tests..." -ForegroundColor Cyan + cargo test -p mcp-memory-server --lib api::setup::tests::test_dashboard_endpoint_and_html_integrity -- --nocapture + # Run standard sequential cargo unit tests test-cargo: @Write-Host "Running standard unit tests across workspace..." -ForegroundColor Cyan diff --git a/nvim-core/src/lib.rs b/nvim-core/src/lib.rs index a31721c..74b365b 100644 --- a/nvim-core/src/lib.rs +++ b/nvim-core/src/lib.rs @@ -385,19 +385,6 @@ async fn get_nvim_connection() -> Result, String> { } }); - // Cleanup task for timed-out requests - let pending_clone3 = Arc::clone(&pending_requests); - tokio::spawn(async move { - let mut interval = tokio::time::interval(tokio::time::Duration::from_secs(10)); - loop { - interval.tick().await; - if Arc::strong_count(&pending_clone3) <= 1 { - break; // Socket closed and other tasks finished, no need to keep cleaning up - } - pending_clone3.retain(|_, sender| !sender.is_closed()); - } - }); - let mut conn_lock = NVIM_CONN.lock().unwrap_or_else(|e| e.into_inner()); if let Some(existing_sender) = conn_lock.as_ref() && !existing_sender.is_closed() diff --git a/server/src/api/setup.rs b/server/src/api/setup.rs index 3c81f7f..c509813 100644 --- a/server/src/api/setup.rs +++ b/server/src/api/setup.rs @@ -161,6 +161,12 @@ pub fn create_router(app_state: Arc) -> Router { "/", get(|| async move { axum::response::Html(include_str!("../dashboard.html")) }), ) + .route( + "/dashboard.js", + get(|| async move { + ([(axum::http::header::CONTENT_TYPE, "application/javascript")], include_str!("../dashboard.js")) + }), + ) .route( "/api/graph", get({ @@ -675,6 +681,38 @@ mod tests { let shut_resp = app_shut.oneshot(shut_req).await.unwrap(); assert_eq!(shut_resp.status(), StatusCode::UNAUTHORIZED); } + + #[tokio::test] + async fn test_dashboard_endpoint_and_html_integrity() { + let (_app, app_state, _dir) = setup_app().await; + + // Test GET / (dashboard.html) + let app_html = create_router(app_state.clone()); + let req_html = Request::builder().uri("/").body(Body::empty()).unwrap(); + let resp_html = app_html.oneshot(req_html).await.unwrap(); + assert_eq!(resp_html.status(), StatusCode::OK); + let content_type = resp_html.headers().get(axum::http::header::CONTENT_TYPE).unwrap().to_str().unwrap(); + assert!(content_type.contains("text/html")); + + let body_bytes = axum::body::to_bytes(resp_html.into_body(), usize::MAX).await.unwrap(); + let html_str = String::from_utf8(body_bytes.to_vec()).unwrap(); + assert!(html_str.contains("")); + + // Test GET /dashboard.js + let app_js = create_router(app_state.clone()); + let req_js = Request::builder().uri("/dashboard.js").body(Body::empty()).unwrap(); + let resp_js = app_js.oneshot(req_js).await.unwrap(); + assert_eq!(resp_js.status(), StatusCode::OK); + let content_type_js = resp_js.headers().get(axum::http::header::CONTENT_TYPE).unwrap().to_str().unwrap(); + assert!(content_type_js.contains("application/javascript")); + + let js_bytes = axum::body::to_bytes(resp_js.into_body(), usize::MAX).await.unwrap(); + let js_str = String::from_utf8(js_bytes.to_vec()).unwrap(); + assert!(js_str.contains("escapeHtml")); + assert!(js_str.contains("setupWS")); + assert!(js_str.contains("refreshActiveTab")); + assert!(js_str.contains("parseActivityPayload")); + } } diff --git a/server/src/api/ws.rs b/server/src/api/ws.rs index f3d29b2..4a2c454 100644 --- a/server/src/api/ws.rs +++ b/server/src/api/ws.rs @@ -54,6 +54,40 @@ pub async fn handle_socket(socket: WebSocket, state: Arc, _client_type .unwrap_or_else(|e| e.into_inner()) .insert(session_id.clone(), tx.clone()); + // Reconnection Catch-Up: Replay recent TASK_EVENT notifications so client receives missed Futures + let recent_task_notifications: Vec = state + .handler + .state + .telemetry + .recent_activities + .read_with(|activities| { + activities + .iter() + .filter_map(|act_val| { + if act_val["category"] == "TASK_EVENT" { + if let Some(details_str) = act_val["details"].as_str() { + if let Ok(event_val) = serde_json::from_str::(details_str) { + return Some( + serde_json::json!({ + "jsonrpc": "2.0", + "method": "notifications/task/completed", + "params": event_val + }) + .to_string(), + ); + } + } + } + None + }) + .take(5) + .collect() + }); + + for notif in recent_task_notifications.into_iter().rev() { + let _ = tx.try_send(notif); + } + let (mut sender, mut receiver) = socket.split(); let send_task = tokio::spawn(async move { @@ -185,4 +219,83 @@ mod tests { assert!(app_state.clients.read().unwrap().is_empty()); } + + #[tokio::test] + async fn test_task_event_broadcast_and_reconnection_replay() { + let dir = tempdir().unwrap(); + let mem_state = Arc::new(MemoryState::new(dir.path().to_str().unwrap())); + + let task_event = crate::models::TaskEvent { + task_id: "task-999".to_string(), + status: "completed".to_string(), + action: Some("update".to_string()), + result: Some(serde_json::json!({"status": "completed"})), + error: None, + timestamp: 1728129000, + session_id: None, + }; + + // Broadcast task event + mem_state.broadcast_task_event(task_event.clone()); + + // Verify recent activities recorded the event + let recorded = mem_state.telemetry.recent_activities.read_with(|act| act.clone()); + assert!(!recorded.is_empty()); + assert_eq!(recorded[0]["category"], "TASK_EVENT"); + + // Verify reconnection catch-up replay fetches the notification + let (shutdown_tx, _) = tokio::sync::oneshot::channel(); + let app_state = Arc::new(AppState { + handler: Arc::new(MemoryHandler::new(mem_state)), + clients: RwLock::new(HashMap::new()), + next_id: AtomicUsize::new(1), + shutdown_tx: std::sync::Mutex::new(Some(shutdown_tx)), + }); + + let (tx, mut rx) = mpsc::channel::(10); + app_state + .clients + .write() + .unwrap() + .insert("session-1".to_string(), tx.clone()); + + let recent_notifications: Vec = app_state + .handler + .state + .telemetry + .recent_activities + .read_with(|activities| { + activities + .iter() + .filter_map(|act_val| { + if act_val["category"] == "TASK_EVENT" { + if let Some(details_str) = act_val["details"].as_str() { + if let Ok(event_val) = serde_json::from_str::(details_str) { + return Some( + serde_json::json!({ + "jsonrpc": "2.0", + "method": "notifications/task/completed", + "params": event_val + }) + .to_string(), + ); + } + } + } + None + }) + .take(5) + .collect() + }); + + for notif in recent_notifications { + let _ = tx.try_send(notif); + } + + let replayed_msg = rx.recv().await.expect("Expected replayed task event notification"); + let parsed: serde_json::Value = serde_json::from_str(&replayed_msg).unwrap(); + assert_eq!(parsed["method"], "notifications/task/completed"); + assert_eq!(parsed["params"]["task_id"], "task-999"); + assert_eq!(parsed["params"]["status"], "completed"); + } } diff --git a/server/src/clipboard_watcher.rs b/server/src/clipboard_watcher.rs index 7ddcb55..87b0408 100644 --- a/server/src/clipboard_watcher.rs +++ b/server/src/clipboard_watcher.rs @@ -7,17 +7,18 @@ pub fn spawn_watcher(state: Arc) { tokio::spawn(async move { let mut last_text = String::new(); loop { - sleep(Duration::from_millis(1000)).await; - let is_enabled = { let watch = state.clipboard_watch_mode.read().await; *watch }; if !is_enabled { + state.clipboard_notify.notified().await; continue; } + sleep(Duration::from_millis(1000)).await; + if let Ok(mut clipboard) = Clipboard::new() && let Ok(text) = clipboard.get_text() && text != last_text @@ -69,6 +70,7 @@ mod tests { // Enable watch mode *state.clipboard_watch_mode.write().await = true; + state.clipboard_notify.notify_waiters(); tokio::time::sleep(Duration::from_millis(150)).await; } } diff --git a/server/src/dashboard.html b/server/src/dashboard.html index 9e93589..abced54 100644 --- a/server/src/dashboard.html +++ b/server/src/dashboard.html @@ -631,868 +631,6 @@ - + diff --git a/server/src/dashboard.js b/server/src/dashboard.js new file mode 100644 index 0000000..2e3979c --- /dev/null +++ b/server/src/dashboard.js @@ -0,0 +1,706 @@ +// server/src/dashboard.ts +var currentTabId = "graph-tab"; +function switchTab(tabId, btn) { + currentTabId = tabId; + document.querySelectorAll(".tab-content").forEach((el) => el.classList.remove("active")); + document.querySelectorAll(".tab-button").forEach((el) => el.classList.remove("active")); + const targetTab = document.getElementById(tabId); + if (targetTab) + targetTab.classList.add("active"); + if (btn) + btn.classList.add("active"); + switch (tabId) { + case "graph-tab": + if (network) { + network.redraw(); + } else { + loadGraph(); + } + break; + case "activity-tab": + loadActivityHistory(); + break; + case "task-tab": + loadTasks(); + break; + case "preferences-tab": + loadPreferences(); + break; + case "sticky-tab": + loadStickyNotes(); + break; + case "techdebt-tab": + loadTechDebt(); + break; + case "adrs-tab": + loadADRs(); + break; + case "workspaces-tab": + loadWorkspaces(); + break; + case "pinned-tab": + loadPinned(); + break; + case "memos-tab": + loadMemos(); + break; + case "snippets-tab": + loadSnippets(); + break; + case "pr-tab": + loadPRs(); + break; + case "terminal-tab": + loadTerminal(); + break; + } +} +if (localStorage.getItem("theme") === "dark" || !localStorage.getItem("theme") && window.matchMedia && window.matchMedia("(prefers-color-scheme: dark)").matches) { + document.documentElement.setAttribute("data-theme", "dark"); +} +var network = null; +var nodesData = new vis.DataSet; +var edgesData = new vis.DataSet; +var rawEntities = {}; +var rawRelations = []; +var activeFilters = new Set; +var allTypes = new Set; +var typeColors = { + Concept: { bg: "#9b59b6", border: "#8e44ad" }, + Technology: { bg: "#3498db", border: "#2980b9" }, + Infrastructure: { bg: "#e67e22", border: "#d35400" }, + Database_Table: { bg: "#f1c40f", border: "#f39c12" }, + Service: { bg: "#1abc9c", border: "#f39c12" }, + Tool: { bg: "#e74c3c", border: "#c0392b" }, + Integration: { bg: "#2ecc71", border: "#27ae60" }, + System: { bg: "#34495e", border: "#2c3e50" }, + Rule: { bg: "#fd79a8", border: "#e84393" } +}; +function getColorForType(type) { + if (typeColors[type]) + return typeColors[type]; + let hash = 0; + for (let i = 0;i < type.length; i++) + hash = type.charCodeAt(i) + ((hash << 5) - hash); + const hue = Math.abs(hash) % 360; + return { background: `hsl(${hue}, 70%, 60%)`, border: `hsl(${hue}, 70%, 40%)` }; +} +function zoomGraph(step) { + if (!network) + return; + const currentScale = network.getScale(); + network.moveTo({ scale: currentScale * (1 + step) }); +} +function closeInspector() { + const inspector = document.getElementById("inspector-panel"); + if (inspector) + inspector.classList.remove("open"); + if (network) + network.unselectAll(); +} +function showInspector(nodeId) { + const entity = rawEntities[nodeId]; + if (!entity) + return; + const titleEl = document.getElementById("inspector-title"); + if (titleEl) + titleEl.innerText = entity.name; + const typeEl = document.getElementById("inspector-type"); + if (typeEl) + typeEl.innerText = entity.entity_type; + const nsEl = document.getElementById("inspector-namespace"); + if (nsEl) + nsEl.innerText = entity.namespace || "global"; + const obsHtml = (entity.observations || []).map((o) => `
  • ${o}
  • `).join(""); + const obsEl = document.getElementById("inspector-observations"); + if (obsEl) + obsEl.innerHTML = obsHtml || "
  • No observations recorded.
  • "; + const rels = rawRelations.filter((r) => r.from === nodeId || r.to === nodeId); + const relHtml = rels.map((r) => { + if (r.from === nodeId) + return `
  • → ${r.relation_type} → ${r.to}
  • `; + return `
  • ← ${r.relation_type} ← ${r.from}
  • `; + }).join(""); + const relEl = document.getElementById("inspector-relations"); + if (relEl) + relEl.innerHTML = relHtml || "
  • No direct relations.
  • "; + const inspector = document.getElementById("inspector-panel"); + if (inspector) + inspector.classList.add("open"); +} +function renderFilters() { + const container = document.getElementById("graph-filters"); + if (!container) + return; + container.innerHTML = Array.from(allTypes).sort().map((type) => { + const isActive = activeFilters.has(type); + const color = getColorForType(type).bg || getColorForType(type).background; + return ``; + }).join(""); +} +function updateGraphData() { + const newNodes = []; + const newEdges = []; + const nodeIds = new Set; + for (const [name, entity] of Object.entries(rawEntities)) { + allTypes.add(entity.entity_type); + if (activeFilters.size > 0 && !activeFilters.has(entity.entity_type)) + continue; + const color = getColorForType(entity.entity_type); + newNodes.push({ + id: name, + label: name, + title: `${name} +Type: ${entity.entity_type}`, + color: { background: color.bg || color.background, border: color.border }, + font: { color: document.documentElement.getAttribute("data-theme") === "dark" ? "#eee" : "#333" } + }); + nodeIds.add(name); + } + rawRelations.forEach((r) => { + if (nodeIds.has(r.from) && nodeIds.has(r.to)) { + newEdges.push({ + id: r.from + "-" + r.to + "-" + r.relation_type, + from: r.from, + to: r.to, + label: r.relation_type, + arrows: "to", + font: { color: document.documentElement.getAttribute("data-theme") === "dark" ? "#aaa" : "#666", strokeWidth: 0 } + }); + } + }); + const currentNodes = nodesData.getIds(); + nodesData.remove(currentNodes.filter((id) => !nodeIds.has(String(id)))); + nodesData.update(newNodes); + const currentEdges = edgesData.getIds(); + const newEdgeIds = new Set(newEdges.map((e) => e.id)); + edgesData.remove(currentEdges.filter((id) => !newEdgeIds.has(String(id)))); + edgesData.update(newEdges); + renderFilters(); +} +async function loadGraph() { + try { + const res = await fetch("/api/graph"); + const data = await res.json(); + rawEntities = data.entities || {}; + rawRelations = data.relations || []; + updateGraphData(); + if (!network) { + const container = document.getElementById("network-container"); + if (!container) + return; + const options = { + nodes: { shape: "dot", size: 16, font: { size: 12 } }, + edges: { color: { inherit: "from", opacity: 0.6 }, font: { size: 10, align: "middle" }, smooth: { type: "continuous" } }, + physics: { barnesHut: { gravitationalConstant: -2000, centralGravity: 0.3, springLength: 95 } }, + interaction: { hover: true, tooltipDelay: 100, zoomView: false } + }; + network = new vis.Network(container, { nodes: nodesData, edges: edgesData }, options); + container.addEventListener("wheel", function(event) { + event.preventDefault(); + const direction = event.deltaY > 0 ? -0.15 : 0.15; + zoomGraph(direction); + }, { passive: false }); + network.on("selectNode", function(params) { + if (params.nodes.length > 0) + showInspector(params.nodes[0]); + }); + network.on("deselectNode", function(params) { + if (params.nodes.length === 0) + closeInspector(); + }); + } + } catch (err) { + console.error("Failed to load graph", err); + const container = document.getElementById("network-container"); + if (container) { + container.innerHTML = '
    Error loading graph: ' + err.message + "
    "; + } + } +} +function buildTaskTreeHTML(tasks, parentId, depth = 0) { + let html = ""; + const children = tasks.filter((t) => { + const pid = t.parent_id || t.parentId; + if (!parentId) + return !pid; + return pid === parentId; + }); + if (children.length === 0) + return html; + children.forEach((t) => { + const isCompleted = t.status === "completed" || t.status === "done"; + const isCancelled = t.status === "cancelled" || t.status === "abandoned"; + let cardClass = "task-card"; + if (isCompleted) + cardClass += " completed"; + if (isCancelled) + cardClass += " cancelled"; + let isBlocked = false; + let blockers = []; + const deps = t.dependencies || []; + deps.forEach((depId) => { + const depTask = tasks.find((dt) => dt.id === depId); + if (depTask && depTask.status !== "completed" && depTask.status !== "done") { + isBlocked = true; + blockers.push(depTask.title); + } + }); + const allChildren = tasks.filter((ct) => (ct.parentId || ct.parent_id) === t.id); + const completedChildren = allChildren.filter((ct) => ct.status === "completed" || ct.status === "done"); + let progressHtml = ""; + if (allChildren.length > 0) { + const pct = Math.round(completedChildren.length / allChildren.length * 100); + progressHtml = ` +
    +
    +
    +
    ${pct}% (${completedChildren.length}/${allChildren.length} child tasks)
    + `; + if (completedChildren.length < allChildren.length) { + isBlocked = true; + } + } + html += `
    `; + if (isBlocked && !isCompleted && !isCancelled) { + html += `
    [BLOCKED]
    `; + if (blockers.length > 0) { + html += `
    Waiting on: ${blockers.join(", ")}
    `; + } + } + if (isCancelled) { + html += `
    [CANCELLED]
    `; + } + html += `${t.title}${t.description}`; + const criteria = t.acceptance_criteria || []; + if (criteria.length > 0) { + html += `
      `; + let unmetCriteria = false; + criteria.forEach((c) => { + const isMet = c.is_met; + if (!isMet) + unmetCriteria = true; + const check = isMet ? "☑" : "☐"; + const strike = isMet ? "text-decoration: line-through;" : ""; + html += `
    • ${check} ${c.description}
    • `; + }); + html += `
    `; + if (unmetCriteria && !isCompleted && !isCancelled) + isBlocked = true; + } + html += progressHtml; + if (!isCompleted && !isCancelled && !isBlocked) { + html += ``; + } + if (allChildren.length > 0) { + html += `
    `; + html += buildTaskTreeHTML(tasks, t.id, 0); + html += `
    `; + } + html += `
    `; + }); + return html; +} +async function loadTasks() { + try { + const res = await fetch("/api/tasks"); + const tasks = await res.json(); + const taskContainer = document.getElementById("task-tree-container"); + if (!taskContainer) + return; + const rootHtml = buildTaskTreeHTML(tasks, null, 0); + if (!rootHtml) { + taskContainer.innerHTML = '
    No active tasks.
    '; + } else { + taskContainer.innerHTML = rootHtml; + } + } catch (err) { + console.error("Failed to load tasks", err); + } +} +var MAX_ACTIVITY_HISTORY = 100; +function escapeHtml(str) { + if (str === null || str === undefined) + return ""; + return String(str).replace(/&/g, "&").replace(//g, ">").replace(/"/g, """).replace(/'/g, "'"); +} +function parseActivityPayload(item) { + let category = "TOOL"; + let summary = ""; + let details = ""; + let timestamp = Date.now(); + if (typeof item === "string") { + try { + const parsed = JSON.parse(item); + return parseActivityPayload(parsed); + } catch (e) { + summary = item; + } + } else if (typeof item === "object" && item !== null) { + if (item.method === "notifications/activity" && item.params) { + return parseActivityPayload(item.params); + } + category = (item.category || item.type || "TOOL").toUpperCase(); + summary = item.summary || item.message || item.description || item.data || ""; + details = item.details || ""; + timestamp = item.timestamp || item.time || item.updated_at || timestamp; + if (typeof timestamp === "number" && timestamp < 10000000000) { + timestamp = timestamp * 1000; + } + } + const catUpper = category.toUpperCase(); + const colorMap = { + GRAPH: "#2ecc71", + DECISION: "#f39c12", + CODE: "#3498db", + TASK: "#1abc9c", + CLIPBOARD: "#9b59b6", + STICKY_NOTE: "#e67e22", + CHECKPOINT: "#e74c3c", + ERROR_FIX: "#e74c3c", + TECH_DEBT: "#d35400", + SUBAGENT: "#8e44ad", + SNIPPET: "#16a085", + TOOL: "#3498db", + SYSTEM: "#95a5a6" + }; + const badgeColor = colorMap[catUpper] || "#3498db"; + const dateObj = new Date(timestamp); + const validDate = isNaN(dateObj.getTime()) ? new Date : dateObj; + const timeStr = validDate.toLocaleString([], { + month: "2-digit", + day: "2-digit", + hour: "2-digit", + minute: "2-digit", + second: "2-digit" + }); + let html = `[${timeStr}] ${escapeHtml(catUpper)} ${escapeHtml(summary)}`; + if (details) { + html += `
    ${escapeHtml(details)}
    `; + } + return html; +} +async function loadActivityHistory() { + try { + const response = await fetch("/api/activity"); + const history = await response.json(); + const feed = document.getElementById("activity-feed"); + if (!feed) + return; + feed.innerHTML = ""; + const getMillis = (item) => { + if (!item) + return 0; + if (typeof item === "string") { + try { + item = JSON.parse(item); + } catch (e) { + return 0; + } + } + if (item.method === "notifications/activity" && item.params) { + item = item.params; + } + let val = item.timestamp || item.time || item.updated_at || 0; + if (typeof val === "number" && val < 10000000000) { + val = val * 1000; + } + return val; + }; + const sorted = [...history].sort((a, b) => getMillis(b) - getMillis(a)); + sorted.forEach((item) => { + const div = document.createElement("div"); + div.className = "feed-entry"; + div.innerHTML = parseActivityPayload(item); + feed.appendChild(div); + }); + if (sorted.length > 0) { + feed.scrollTop = 0; + } + } catch (e) { + console.error("Failed to load activity history", e); + } +} +function setupWS() { + const protocol = location.protocol === "https:" ? "wss:" : "ws:"; + const ws = new WebSocket(`${protocol}//${location.host}/ws?client=ui`); + const feed = document.getElementById("activity-feed"); + ws.onmessage = function(event) { + try { + const data = JSON.parse(event.data); + if (data.type === "activity" || data.method === "notifications/activity" || data.method === "notifications/task/completed" || data.method === "notifications/resources/updated") { + refreshActiveTab(); + const isScrolledToTop = feed ? feed.scrollTop <= 20 : true; + const div = document.createElement("div"); + div.className = "feed-entry"; + const payload = data.params || data.data || data; + div.innerHTML = parseActivityPayload(payload); + if (feed) { + feed.prepend(div); + while (feed.children.length > MAX_ACTIVITY_HISTORY && feed.lastChild) { + feed.removeChild(feed.lastChild); + } + if (isScrolledToTop) + feed.scrollTop = 0; + } + } + } catch (e) { + console.error("WebSocket message parse error", e); + } + }; + ws.onclose = function() { + console.log("WebSocket closed, attempting to reconnect in 3s..."); + setTimeout(setupWS, 3000); + }; +} +async function loadPreferences() { + try { + const res = await fetch("/api/preferences"); + const data = await res.json(); + const container = document.getElementById("preferences-container"); + if (!container) + return; + container.innerHTML = ""; + if (!data || Object.keys(data).length === 0) { + container.innerHTML = '
    No global preferences found.
    '; + return; + } + for (const [key, pref] of Object.entries(data)) { + const date = new Date(pref.updated_at * 1000).toLocaleString(); + container.innerHTML += ` +
    + ${key} +
    ${pref.value}
    +
    Last Updated: ${date}
    +
    `; + } + } catch (e) { + console.error("Error loading preferences:", e); + } +} +async function loadStickyNotes() { + try { + const res = await fetch("/api/sticky"); + const sticky = await res.json(); + const container = document.getElementById("sticky-notes-container"); + if (!container) + return; + container.innerHTML = ""; + sticky.forEach((note) => { + const card = document.createElement("div"); + card.className = "sticky-note"; + const date = new Date(note.timestamp * 1000).toLocaleString(); + card.innerHTML = `
    ${date}
    +
    ${note.content}
    `; + container.appendChild(card); + }); + } catch (err) { + console.error("Failed to load sticky notes", err); + } +} +async function loadGenericList(endpoint, containerId, formatter) { + try { + const res = await fetch(endpoint); + const data = await res.json(); + const container = document.getElementById(containerId); + if (!container) + return; + if (!data || data.length === 0) { + container.innerHTML = '
    No items recorded.
    '; + return; + } + container.innerHTML = data.map((item) => { + const date = item.timestamp ? new Date(item.timestamp * 1000).toLocaleString() : ""; + return `
    +
    ${date}
    + ${formatter(item)} +
    `; + }).join(""); + } catch (err) { + console.error(`Failed to load ${endpoint}`, err); + } +} +function loadTerminal() { + loadGenericList("/api/terminal/history", "terminal-container", (item) => ` + ${item.command} + Exit Code: ${item.exit_code} +
    + OS: ${item.os || "Unknown"} + CWD: ${item.cwd || "Unknown"} +
    + `); +} +function loadTechDebt() { + loadGenericList("/api/tech_debts", "techdebt-container", (item) => ` + ${item.id} ${item.is_resolved ? '(Resolved)' : '(Open)'} +
    Description: ${item.description}
    +
    Ideal Solution: ${item.ideal_solution}
    +
    + Git Commit: ${item.git_commit || "None"} + Branch: ${item.git_branch || "None"} +
    + `); + loadGenericList("/api/error_fixes", "errorfixes-container", (item) => ` + ${item.signature || "Error"} (Resolved) +
    Solution: ${item.solution || ""}
    +
    Commit: ${item.git_commit || "None"}
    + `); +} +function loadADRs() { + loadGenericList("/api/adrs", "adrs-container", (item) => ` + ${item.id} | ${item.title} + ${item.status} +
    Context: ${item.context || ""}
    +
    Decision: ${item.decision || ""}
    +
    Consequence: ${item.consequence || ""}
    + ${item.supersedes ? `
    Supersedes: ${item.supersedes}
    ` : ""} + `); +} +function loadWorkspaces() { + loadGenericList("/api/context_workspaces", "workspaces-container", (item) => ` + ${item.name} +
    ${item.description || ""}
    +
    ${(item.paths || []).join(", ")}
    + `); +} +function loadPinned() { + loadGenericList("/api/pinned_files", "pinned-container", (item) => ` + ${item.path || item.file_path || item.id} +
    ${item.reason || item.description || "Pinned"}
    + `); +} +function loadMemos() { + loadGenericList("/api/handoff_memos", "memos-container", (item) => ` + Memo from ${item.author || "System"} +
    ${item.content || item.summary || ""}
    + `); + loadGenericList("/api/milestones", "milestones-container", (item) => ` + ${item.title || item.name} +
    ${item.description || ""}
    + `); +} +function loadSnippets() { + loadGenericList("/api/snippets", "snippets-container", (item) => ` + ${item.description || "Snippet"} +
    Language: ${item.language || "txt"} | Tags: ${(item.tags || []).join(", ")}
    +
    ${item.content || item.code || ""}
    + `); +} +function loadPRs() { + loadGenericList("/api/pr_checklists", "pr-container", (item) => ` + ${item.name || "Checklist"} +
      + ${(item.items || []).map((i) => { + const check = i.is_completed ? "☑" : "☐"; + const strike = i.is_completed ? "text-decoration:line-through; color:var(--text-secondary);" : ""; + return `
    • ${check} ${i.description}
    • `; + }).join("")} +
    + `); +} +async function loadVersion() { + try { + const res = await fetch("/api/version"); + const data = await res.json(); + const verEl = document.getElementById("app-version"); + if (verEl) { + verEl.innerHTML = `v${data.version}`; + } + } catch (err) { + console.error("Failed to load version", err); + } +} +function setupSSE() { + try { + const sse = new EventSource("/api/activity/stream"); + sse.onmessage = function(event) { + if (event.data) { + try { + refreshActiveTab(); + const feed = document.getElementById("activity-feed"); + if (feed) { + const isScrolledToTop = feed.scrollTop <= 20; + const div = document.createElement("div"); + div.className = "feed-entry"; + div.innerHTML = parseActivityPayload(event.data); + feed.prepend(div); + while (feed.children.length > MAX_ACTIVITY_HISTORY && feed.lastChild) { + feed.removeChild(feed.lastChild); + } + if (isScrolledToTop) + feed.scrollTop = 0; + } + } catch (e) {} + } + }; + } catch (e) { + console.error("SSE initialization error", e); + } +} +document.addEventListener("keydown", function(e) { + if ((e.ctrlKey || e.metaKey) && e.key.toLowerCase() === "k") { + e.preventDefault(); + const searchTabBtn = document.querySelectorAll(".tab-button")[1]; + if (searchTabBtn) { + switchTab("search-tab", searchTabBtn); + } + const searchInput = document.getElementById("search-box"); + if (searchInput) { + searchInput.focus(); + searchInput.select(); + } + } +}); +function refreshActiveTab() { + switch (currentTabId) { + case "graph-tab": + loadGraph(); + break; + case "task-tab": + loadTasks(); + break; + case "sticky-tab": + loadStickyNotes(); + break; + case "techdebt-tab": + loadTechDebt(); + break; + case "adrs-tab": + loadADRs(); + break; + case "workspaces-tab": + loadWorkspaces(); + break; + case "pinned-tab": + loadPinned(); + break; + case "memos-tab": + loadMemos(); + break; + case "snippets-tab": + loadSnippets(); + break; + case "pr-tab": + loadPRs(); + break; + case "terminal-tab": + loadTerminal(); + break; + case "preferences-tab": + loadPreferences(); + break; + case "activity-tab": + loadActivityHistory(); + break; + } +} +loadVersion(); +loadGraph(); +loadActivityHistory(); +setupWS(); +setupSSE(); +var observer = new MutationObserver(() => updateGraphData()); +observer.observe(document.documentElement, { attributes: true, attributeFilter: ["data-theme"] }); diff --git a/server/src/dashboard.ts b/server/src/dashboard.ts new file mode 100644 index 0000000..f7610f4 --- /dev/null +++ b/server/src/dashboard.ts @@ -0,0 +1,934 @@ +/// +/// + +// --- Type Definitions for Dashboard --- +declare const vis: { + DataSet: new (data?: T[]) => { + getIds(): (string | number)[]; + remove(ids: (string | number)[]): void; + update(items: T[]): void; + }; + Network: new (container: HTMLElement, data: { nodes: any; edges: any }, options: any) => { + getScale(): number; + moveTo(options: { scale?: number }): void; + fit(options?: { animation?: { duration: number; easingFunction: string } }): void; + unselectAll(): void; + selectNodes(nodeIds: string[]): void; + focus(nodeId: string, options?: { scale?: number; animation?: boolean }): void; + redraw(): void; + on(event: string, callback: (params: any) => void): void; + }; +}; + +interface Entity { + name: string; + entity_type: string; + namespace?: string; + observations?: string[]; +} + +interface Relation { + from: string; + to: string; + relation_type: string; +} + +interface TaskItem { + id: string; + title: string; + description: string; + status: string; + parent_id?: string; + parentId?: string; + dependencies?: string[]; + acceptance_criteria?: { description: string; is_met: boolean }[]; +} + +interface ActivityItem { + category?: string; + type?: string; + summary?: string; + message?: string; + description?: string; + data?: any; + details?: string; + timestamp?: number; + time?: number; + updated_at?: number; + method?: string; + params?: any; +} + +interface SearchResultItem { + id: string; + title: string; + type_name: string; + score: number; + content: string; +} + +// --- Tabs --- +let currentTabId: string = 'graph-tab'; + +function switchTab(tabId: string, btn?: HTMLElement | null): void { + currentTabId = tabId; + + document.querySelectorAll('.tab-content').forEach(el => el.classList.remove('active')); + document.querySelectorAll('.tab-button').forEach(el => el.classList.remove('active')); + + const targetTab = document.getElementById(tabId); + if (targetTab) targetTab.classList.add('active'); + if (btn) btn.classList.add('active'); + + switch (tabId) { + case 'graph-tab': + if (network) { + network.redraw(); + } else { + loadGraph(); + } + break; + case 'activity-tab': + loadActivityHistory(); + break; + case 'task-tab': + loadTasks(); + break; + case 'preferences-tab': + loadPreferences(); + break; + case 'sticky-tab': + loadStickyNotes(); + break; + case 'techdebt-tab': + loadTechDebt(); + break; + case 'adrs-tab': + loadADRs(); + break; + case 'workspaces-tab': + loadWorkspaces(); + break; + case 'pinned-tab': + loadPinned(); + break; + case 'memos-tab': + loadMemos(); + break; + case 'snippets-tab': + loadSnippets(); + break; + case 'pr-tab': + loadPRs(); + break; + case 'terminal-tab': + loadTerminal(); + break; + } +} + +// --- Theme Toggle --- +function toggleTheme(): void { + const currentTheme = document.documentElement.getAttribute('data-theme'); + const newTheme = currentTheme === 'dark' ? 'light' : 'dark'; + document.documentElement.setAttribute('data-theme', newTheme); + localStorage.setItem('theme', newTheme); +} + +if (localStorage.getItem('theme') === 'dark' || (!localStorage.getItem('theme') && window.matchMedia && window.matchMedia('(prefers-color-scheme: dark)').matches)) { + document.documentElement.setAttribute('data-theme', 'dark'); +} + +// --- Global Graph Data --- +let network: any = null; +let nodesData = new vis.DataSet(); +let edgesData = new vis.DataSet(); +let rawEntities: Record = {}; +let rawRelations: Relation[] = []; +let activeFilters: Set = new Set(); +let allTypes: Set = new Set(); + +const typeColors: Record = { + 'Concept': { bg: '#9b59b6', border: '#8e44ad' }, + 'Technology': { bg: '#3498db', border: '#2980b9' }, + 'Infrastructure': { bg: '#e67e22', border: '#d35400' }, + 'Database_Table': { bg: '#f1c40f', border: '#f39c12' }, + 'Service': { bg: '#1abc9c', border: '#f39c12' }, + 'Tool': { bg: '#e74c3c', border: '#c0392b' }, + 'Integration': { bg: '#2ecc71', border: '#27ae60' }, + 'System': { bg: '#34495e', border: '#2c3e50' }, + 'Rule': { bg: '#fd79a8', border: '#e84393' } +}; + +function getColorForType(type: string): { bg?: string; border?: string; background?: string } { + if (typeColors[type]) return typeColors[type]; + let hash = 0; + for (let i = 0; i < type.length; i++) hash = type.charCodeAt(i) + ((hash << 5) - hash); + const hue = Math.abs(hash) % 360; + return { background: `hsl(${hue}, 70%, 60%)`, border: `hsl(${hue}, 70%, 40%)` }; +} + +// --- Graph Controls --- +function zoomGraph(step: number): void { + if (!network) return; + const currentScale = network.getScale(); + network.moveTo({ scale: currentScale * (1 + step) }); +} + +function resetGraph(): void { + if (!network) return; + network.fit({ animation: { duration: 500, easingFunction: 'easeInOutQuad' } }); +} + +// --- Graph Inspector --- +function closeInspector(): void { + const inspector = document.getElementById('inspector-panel'); + if (inspector) inspector.classList.remove('open'); + if (network) network.unselectAll(); +} + +function showInspector(nodeId: string): void { + const entity = rawEntities[nodeId]; + if (!entity) return; + + const titleEl = document.getElementById('inspector-title'); + if (titleEl) titleEl.innerText = entity.name; + + const typeEl = document.getElementById('inspector-type'); + if (typeEl) typeEl.innerText = entity.entity_type; + + const nsEl = document.getElementById('inspector-namespace'); + if (nsEl) nsEl.innerText = entity.namespace || 'global'; + + const obsHtml = (entity.observations || []).map(o => `
  • ${o}
  • `).join(''); + const obsEl = document.getElementById('inspector-observations'); + if (obsEl) obsEl.innerHTML = obsHtml || '
  • No observations recorded.
  • '; + + const rels = rawRelations.filter(r => r.from === nodeId || r.to === nodeId); + const relHtml = rels.map(r => { + if (r.from === nodeId) return `
  • → ${r.relation_type} → ${r.to}
  • `; + return `
  • ← ${r.relation_type} ← ${r.from}
  • `; + }).join(''); + + const relEl = document.getElementById('inspector-relations'); + if (relEl) relEl.innerHTML = relHtml || '
  • No direct relations.
  • '; + + const inspector = document.getElementById('inspector-panel'); + if (inspector) inspector.classList.add('open'); +} + +function toggleFilter(type: string): void { + if (activeFilters.has(type)) activeFilters.delete(type); + else activeFilters.add(type); + renderFilters(); + updateGraphData(); +} + +function renderFilters(): void { + const container = document.getElementById('graph-filters'); + if (!container) return; + container.innerHTML = Array.from(allTypes).sort().map(type => { + const isActive = activeFilters.has(type); + const color = getColorForType(type).bg || getColorForType(type).background; + return ``; + }).join(''); +} + +function updateGraphData(): void { + const newNodes: any[] = []; + const newEdges: any[] = []; + const nodeIds = new Set(); + + for (const [name, entity] of Object.entries(rawEntities)) { + allTypes.add(entity.entity_type); + if (activeFilters.size > 0 && !activeFilters.has(entity.entity_type)) continue; + + const color = getColorForType(entity.entity_type); + + newNodes.push({ + id: name, + label: name, + title: `${name}\nType: ${entity.entity_type}`, + color: { background: color.bg || color.background, border: color.border }, + font: { color: document.documentElement.getAttribute('data-theme') === 'dark' ? '#eee' : '#333' } + }); + nodeIds.add(name); + } + + rawRelations.forEach(r => { + if (nodeIds.has(r.from) && nodeIds.has(r.to)) { + newEdges.push({ + id: r.from + '-' + r.to + '-' + r.relation_type, + from: r.from, + to: r.to, + label: r.relation_type, + arrows: 'to', + font: { color: document.documentElement.getAttribute('data-theme') === 'dark' ? '#aaa' : '#666', strokeWidth: 0 } + }); + } + }); + + const currentNodes = nodesData.getIds(); + nodesData.remove(currentNodes.filter(id => !nodeIds.has(String(id)))); + nodesData.update(newNodes); + + const currentEdges = edgesData.getIds(); + const newEdgeIds = new Set(newEdges.map(e => e.id)); + edgesData.remove(currentEdges.filter(id => !newEdgeIds.has(String(id)))); + edgesData.update(newEdges); + + renderFilters(); +} + +async function loadGraph(): Promise { + try { + const res = await fetch('/api/graph'); + const data = await res.json(); + + rawEntities = data.entities || {}; + rawRelations = data.relations || []; + + updateGraphData(); + + if (!network) { + const container = document.getElementById('network-container'); + if (!container) return; + const options = { + nodes: { shape: 'dot', size: 16, font: { size: 12 } }, + edges: { color: { inherit: 'from', opacity: 0.6 }, font: { size: 10, align: 'middle' }, smooth: { type: 'continuous' } }, + physics: { barnesHut: { gravitationalConstant: -2000, centralGravity: 0.3, springLength: 95 } }, + interaction: { hover: true, tooltipDelay: 100, zoomView: false } + }; + network = new vis.Network(container, { nodes: nodesData, edges: edgesData }, options); + + container.addEventListener('wheel', function(event: WheelEvent) { + event.preventDefault(); + const direction = event.deltaY > 0 ? -0.15 : 0.15; + zoomGraph(direction); + }, { passive: false }); + + network.on("selectNode", function(params: any) { + if (params.nodes.length > 0) showInspector(params.nodes[0]); + }); + network.on("deselectNode", function(params: any) { + if (params.nodes.length === 0) closeInspector(); + }); + } + } catch (err: any) { + console.error("Failed to load graph", err); + const container = document.getElementById("network-container"); + if (container) { + container.innerHTML = '
    Error loading graph: ' + err.message + '
    '; + } + } +} + +// --- Task Tree (HTN/DAG) --- +async function completeTask(id: string): Promise { + try { + await fetch(`/api/tasks/${id}/complete`, { method: 'POST' }); + loadTasks(); + } catch(e) { console.error("Failed to complete task", e); } +} + +function buildTaskTreeHTML(tasks: TaskItem[], parentId: string | null, depth: number = 0): string { + let html = ''; + const children = tasks.filter(t => { + const pid = t.parent_id || t.parentId; + if (!parentId) return !pid; + return pid === parentId; + }); + + if (children.length === 0) return html; + + children.forEach(t => { + const isCompleted = t.status === 'completed' || t.status === 'done'; + const isCancelled = t.status === 'cancelled' || t.status === 'abandoned'; + let cardClass = 'task-card'; + if (isCompleted) cardClass += ' completed'; + if (isCancelled) cardClass += ' cancelled'; + + let isBlocked = false; + let blockers: string[] = []; + const deps = t.dependencies || []; + deps.forEach(depId => { + const depTask = tasks.find(dt => dt.id === depId); + if (depTask && depTask.status !== 'completed' && depTask.status !== 'done') { + isBlocked = true; + blockers.push(depTask.title); + } + }); + + const allChildren = tasks.filter(ct => (ct.parentId || ct.parent_id) === t.id); + const completedChildren = allChildren.filter(ct => ct.status === 'completed' || ct.status === 'done'); + let progressHtml = ''; + if (allChildren.length > 0) { + const pct = Math.round((completedChildren.length / allChildren.length) * 100); + progressHtml = ` +
    +
    +
    +
    ${pct}% (${completedChildren.length}/${allChildren.length} child tasks)
    + `; + if (completedChildren.length < allChildren.length) { + isBlocked = true; + } + } + + html += `
    `; + + if (isBlocked && !isCompleted && !isCancelled) { + html += `
    [BLOCKED]
    `; + if (blockers.length > 0) { + html += `
    Waiting on: ${blockers.join(', ')}
    `; + } + } + if (isCancelled) { + html += `
    [CANCELLED]
    `; + } + + html += `${t.title}${t.description}`; + + const criteria = t.acceptance_criteria || []; + if (criteria.length > 0) { + html += `
      `; + let unmetCriteria = false; + criteria.forEach(c => { + const isMet = c.is_met; + if (!isMet) unmetCriteria = true; + const check = isMet ? '☑' : '☐'; + const strike = isMet ? 'text-decoration: line-through;' : ''; + html += `
    • ${check} ${c.description}
    • `; + }); + html += `
    `; + if (unmetCriteria && !isCompleted && !isCancelled) isBlocked = true; + } + + html += progressHtml; + + if (!isCompleted && !isCancelled && !isBlocked) { + html += ``; + } + + if (allChildren.length > 0) { + html += `
    `; + html += buildTaskTreeHTML(tasks, t.id, 0); + html += `
    `; + } + + html += `
    `; + }); + return html; +} + +async function loadTasks(): Promise { + try { + const res = await fetch('/api/tasks'); + const tasks: TaskItem[] = await res.json(); + + const taskContainer = document.getElementById('task-tree-container'); + if (!taskContainer) return; + + const rootHtml = buildTaskTreeHTML(tasks, null, 0); + + if (!rootHtml) { + taskContainer.innerHTML = '
    No active tasks.
    '; + } else { + taskContainer.innerHTML = rootHtml; + } + } catch (err) { + console.error("Failed to load tasks", err); + } +} + +// --- Search --- +function highlightText(text: string, query: string): string { + if (!query) return text; + const regex = new RegExp(`(${query})`, 'gi'); + return text.replace(regex, '$1'); +} + +let searchDebounceTimer: any = null; +let activeSearchAbortController: AbortController | null = null; + +async function handleSearch(e: KeyboardEvent): Promise { + const input = e.target as HTMLInputElement; + const q = input.value.trim(); + + if (searchDebounceTimer) clearTimeout(searchDebounceTimer); + + if (!q) { + if (activeSearchAbortController) activeSearchAbortController.abort(); + const resultsEl = document.getElementById('search-results'); + if (resultsEl) resultsEl.innerHTML = ''; + return; + } + + const delay = (e.key === 'Enter') ? 0 : 250; + + searchDebounceTimer = setTimeout(async () => { + if (activeSearchAbortController) activeSearchAbortController.abort(); + activeSearchAbortController = new AbortController(); + + try { + const res = await fetch(`/api/search?q=${encodeURIComponent(q)}`, { + signal: activeSearchAbortController.signal + }); + const data = await res.json(); + const container = document.getElementById('search-results'); + if (!container) return; + + if (!data.results || data.results.length === 0) { + container.innerHTML = '
    No results found.
    '; + return; + } + + container.innerHTML = data.results.map((r: SearchResultItem) => ` +
    +
    +
    + ${highlightText(r.title, q)} + ${r.type_name} +
    +
    ${r.score.toFixed(2)}
    +
    +
    ${highlightText(r.content.substring(0, 150), q)}${r.content.length > 150 ? '...' : ''}
    +
    ID: ${r.id}
    +
    + `).join(''); + } catch (err: any) { + if (err.name !== 'AbortError') { + console.error("Search error:", err); + } + } + }, delay); +} + +// --- WebSocket Activity Feed --- +const MAX_ACTIVITY_HISTORY = 100; + +function escapeHtml(str: any): string { + if (str === null || str === undefined) return ''; + return String(str) + .replace(/&/g, '&') + .replace(//g, '>') + .replace(/"/g, '"') + .replace(/'/g, '''); +} + +function parseActivityPayload(item: any): string { + let category = 'TOOL'; + let summary = ''; + let details = ''; + let timestamp = Date.now(); + + if (typeof item === 'string') { + try { + const parsed = JSON.parse(item); + return parseActivityPayload(parsed); + } catch (e) { + summary = item; + } + } else if (typeof item === 'object' && item !== null) { + if (item.method === 'notifications/activity' && item.params) { + return parseActivityPayload(item.params); + } + category = (item.category || item.type || 'TOOL').toUpperCase(); + summary = item.summary || item.message || item.description || item.data || ''; + details = item.details || ''; + timestamp = item.timestamp || item.time || item.updated_at || timestamp; + + if (typeof timestamp === 'number' && timestamp < 10000000000) { + timestamp = timestamp * 1000; + } + } + + const catUpper = category.toUpperCase(); + const colorMap: Record = { + 'GRAPH': '#2ecc71', + 'DECISION': '#f39c12', + 'CODE': '#3498db', + 'TASK': '#1abc9c', + 'CLIPBOARD': '#9b59b6', + 'STICKY_NOTE': '#e67e22', + 'CHECKPOINT': '#e74c3c', + 'ERROR_FIX': '#e74c3c', + 'TECH_DEBT': '#d35400', + 'SUBAGENT': '#8e44ad', + 'SNIPPET': '#16a085', + 'TOOL': '#3498db', + 'SYSTEM': '#95a5a6' + }; + const badgeColor = colorMap[catUpper] || '#3498db'; + + const dateObj = new Date(timestamp); + const validDate = isNaN(dateObj.getTime()) ? new Date() : dateObj; + const timeStr = validDate.toLocaleString([], { + month: '2-digit', day: '2-digit', hour: '2-digit', minute:'2-digit', second:'2-digit' + }); + + let html = `[${timeStr}] ${escapeHtml(catUpper)} ${escapeHtml(summary)}`; + if (details) { + html += `
    ${escapeHtml(details)}
    `; + } + return html; +} + +async function loadActivityHistory(): Promise { + try { + const response = await fetch('/api/activity'); + const history: ActivityItem[] = await response.json(); + const feed = document.getElementById('activity-feed'); + if (!feed) return; + feed.innerHTML = ''; + + const getMillis = (item: any): number => { + if (!item) return 0; + if (typeof item === 'string') { + try { item = JSON.parse(item); } catch(e) { return 0; } + } + if (item.method === 'notifications/activity' && item.params) { + item = item.params; + } + let val = item.timestamp || item.time || item.updated_at || 0; + if (typeof val === 'number' && val < 10000000000) { + val = val * 1000; + } + return val; + }; + + const sorted = [...history].sort((a, b) => getMillis(b) - getMillis(a)); + + sorted.forEach(item => { + const div = document.createElement('div'); + div.className = 'feed-entry'; + div.innerHTML = parseActivityPayload(item); + feed.appendChild(div); + }); + if (sorted.length > 0) { + feed.scrollTop = 0; + } + } catch(e) { + console.error("Failed to load activity history", e); + } +} + +function setupWS(): void { + const protocol = location.protocol === 'https:' ? 'wss:' : 'ws:'; + const ws = new WebSocket(`${protocol}//${location.host}/ws?client=ui`); + const feed = document.getElementById('activity-feed'); + + ws.onmessage = function(event: MessageEvent) { + try { + const data = JSON.parse(event.data); + if (data.type === 'activity' || data.method === 'notifications/activity' || data.method === 'notifications/task/completed' || data.method === 'notifications/resources/updated') { + refreshActiveTab(); + const isScrolledToTop = feed ? feed.scrollTop <= 20 : true; + + const div = document.createElement('div'); + div.className = 'feed-entry'; + const payload = data.params || data.data || data; + div.innerHTML = parseActivityPayload(payload); + if (feed) { + feed.prepend(div); + while (feed.children.length > MAX_ACTIVITY_HISTORY && feed.lastChild) { + feed.removeChild(feed.lastChild); + } + if (isScrolledToTop) feed.scrollTop = 0; + } + } + } catch (e) { + console.error("WebSocket message parse error", e); + } + }; + + ws.onclose = function() { + console.log("WebSocket closed, attempting to reconnect in 3s..."); + setTimeout(setupWS, 3000); + }; +} + +async function loadPreferences(): Promise { + try { + const res = await fetch('/api/preferences'); + const data = await res.json(); + + const container = document.getElementById('preferences-container'); + if (!container) return; + container.innerHTML = ''; + + if (!data || Object.keys(data).length === 0) { + container.innerHTML = '
    No global preferences found.
    '; + return; + } + + for (const [key, pref] of Object.entries(data)) { + const date = new Date(pref.updated_at * 1000).toLocaleString(); + container.innerHTML += ` +
    + ${key} +
    ${pref.value}
    +
    Last Updated: ${date}
    +
    `; + } + } catch (e) { + console.error('Error loading preferences:', e); + } +} + +async function loadStickyNotes(): Promise { + try { + const res = await fetch('/api/sticky'); + const sticky = await res.json(); + + const container = document.getElementById('sticky-notes-container'); + if (!container) return; + container.innerHTML = ''; + + sticky.forEach((note: any) => { + const card = document.createElement('div'); + card.className = 'sticky-note'; + + const date = new Date(note.timestamp * 1000).toLocaleString(); + + card.innerHTML = `
    ${date}
    +
    ${note.content}
    `; + + container.appendChild(card); + }); + } catch (err) { + console.error("Failed to load sticky notes", err); + } +} + +async function loadGenericList(endpoint: string, containerId: string, formatter: (item: any) => string): Promise { + try { + const res = await fetch(endpoint); + const data = await res.json(); + const container = document.getElementById(containerId); + if (!container) return; + + if (!data || data.length === 0) { + container.innerHTML = '
    No items recorded.
    '; + return; + } + + container.innerHTML = data.map((item: any) => { + const date = item.timestamp ? new Date(item.timestamp * 1000).toLocaleString() : ''; + return `
    +
    ${date}
    + ${formatter(item)} +
    `; + }).join(''); + } catch (err) { + console.error(`Failed to load ${endpoint}`, err); + } +} + +function loadTerminal(): void { + loadGenericList('/api/terminal/history', 'terminal-container', item => ` + ${item.command} + Exit Code: ${item.exit_code} +
    + OS: ${item.os || 'Unknown'} + CWD: ${item.cwd || 'Unknown'} +
    + `); +} + +function loadTechDebt(): void { + loadGenericList('/api/tech_debts', 'techdebt-container', item => ` + ${item.id} ${item.is_resolved ? '(Resolved)' : '(Open)'} +
    Description: ${item.description}
    +
    Ideal Solution: ${item.ideal_solution}
    +
    + Git Commit: ${item.git_commit || 'None'} + Branch: ${item.git_branch || 'None'} +
    + `); + + loadGenericList('/api/error_fixes', 'errorfixes-container', item => ` + ${item.signature || 'Error'} (Resolved) +
    Solution: ${item.solution || ''}
    +
    Commit: ${item.git_commit || 'None'}
    + `); +} + +function loadADRs(): void { + loadGenericList('/api/adrs', 'adrs-container', item => ` + ${item.id} | ${item.title} + ${item.status} +
    Context: ${item.context || ''}
    +
    Decision: ${item.decision || ''}
    +
    Consequence: ${item.consequence || ''}
    + ${item.supersedes ? `
    Supersedes: ${item.supersedes}
    ` : ''} + `); +} + +function loadWorkspaces(): void { + loadGenericList('/api/context_workspaces', 'workspaces-container', item => ` + ${item.name} +
    ${item.description || ''}
    +
    ${(item.paths || []).join(', ')}
    + `); +} + +function loadPinned(): void { + loadGenericList('/api/pinned_files', 'pinned-container', item => ` + ${item.path || item.file_path || item.id} +
    ${item.reason || item.description || 'Pinned'}
    + `); +} + +function loadMemos(): void { + loadGenericList('/api/handoff_memos', 'memos-container', item => ` + Memo from ${item.author || 'System'} +
    ${item.content || item.summary || ''}
    + `); + + loadGenericList('/api/milestones', 'milestones-container', item => ` + ${item.title || item.name} +
    ${item.description || ''}
    + `); +} + +function loadSnippets(): void { + loadGenericList('/api/snippets', 'snippets-container', item => ` + ${item.description || 'Snippet'} +
    Language: ${item.language || 'txt'} | Tags: ${(item.tags || []).join(', ')}
    +
    ${item.content || item.code || ''}
    + `); +} + +function loadPRs(): void { + loadGenericList('/api/pr_checklists', 'pr-container', item => ` + ${item.name || 'Checklist'} +
      + ${(item.items || []).map((i: any) => { + const check = i.is_completed ? '☑' : '☐'; + const strike = i.is_completed ? 'text-decoration:line-through; color:var(--text-secondary);' : ''; + return `
    • ${check} ${i.description}
    • `; + }).join('')} +
    + `); +} + +function loadAllExtras(): void { + loadTerminal(); + loadTechDebt(); + loadADRs(); + loadWorkspaces(); + loadPinned(); + loadMemos(); + loadSnippets(); + loadPRs(); +} + +async function testClipboard(): Promise { + const modal = document.getElementById('clipboard-modal'); + const resultDiv = document.getElementById('clipboard-result'); + if (modal) modal.style.display = 'flex'; + if (resultDiv) resultDiv.innerHTML = '

    Analyzing your clipboard...

    '; + + try { + const res = await fetch('/api/clipboard/capture', { method: 'POST' }); + const data = await res.json(); + + if (resultDiv) { + if (data.success) { + resultDiv.innerHTML = ``; + } else { + resultDiv.innerHTML = `

    Failed to read clipboard: ${data.error}

    `; + } + } + } catch (err: any) { + if (resultDiv) { + resultDiv.innerHTML = `

    Error calling endpoint: ${err.message}

    `; + } + } +} + +async function loadVersion(): Promise { + try { + const res = await fetch('/api/version'); + const data = await res.json(); + const verEl = document.getElementById('app-version'); + if (verEl) { + verEl.innerHTML = `v${data.version}`; + } + } catch (err) { + console.error("Failed to load version", err); + } +} + +function setupSSE(): void { + try { + const sse = new EventSource('/api/activity/stream'); + sse.onmessage = function(event: MessageEvent) { + if (event.data) { + try { + refreshActiveTab(); + const feed = document.getElementById('activity-feed'); + if (feed) { + const isScrolledToTop = feed.scrollTop <= 20; + const div = document.createElement('div'); + div.className = 'feed-entry'; + div.innerHTML = parseActivityPayload(event.data); + feed.prepend(div); + while (feed.children.length > MAX_ACTIVITY_HISTORY && feed.lastChild) { + feed.removeChild(feed.lastChild); + } + if (isScrolledToTop) feed.scrollTop = 0; + } + } catch(e) {} + } + }; + } catch(e) { console.error('SSE initialization error', e); } +} + +document.addEventListener('keydown', function(e: KeyboardEvent) { + if ((e.ctrlKey || e.metaKey) && e.key.toLowerCase() === 'k') { + e.preventDefault(); + const searchTabBtn = document.querySelectorAll('.tab-button')[1] as HTMLElement | undefined; + if (searchTabBtn) { + switchTab('search-tab', searchTabBtn); + } + const searchInput = document.getElementById('search-box') as HTMLInputElement | null; + if (searchInput) { + searchInput.focus(); + searchInput.select(); + } + } +}); + +function refreshActiveTab(): void { + switch (currentTabId) { + case 'graph-tab': loadGraph(); break; + case 'task-tab': loadTasks(); break; + case 'sticky-tab': loadStickyNotes(); break; + case 'techdebt-tab': loadTechDebt(); break; + case 'adrs-tab': loadADRs(); break; + case 'workspaces-tab': loadWorkspaces(); break; + case 'pinned-tab': loadPinned(); break; + case 'memos-tab': loadMemos(); break; + case 'snippets-tab': loadSnippets(); break; + case 'pr-tab': loadPRs(); break; + case 'terminal-tab': loadTerminal(); break; + case 'preferences-tab': loadPreferences(); break; + case 'activity-tab': loadActivityHistory(); break; + } +} + +// --- Initialization --- +loadVersion(); +loadGraph(); +loadActivityHistory(); +setupWS(); +setupSSE(); + +const observer = new MutationObserver(() => updateGraphData()); +observer.observe(document.documentElement, { attributes: true, attributeFilter: ['data-theme'] }); diff --git a/server/src/handlers/meta.rs b/server/src/handlers/meta.rs index 6930d7f..ddc6fdf 100644 --- a/server/src/handlers/meta.rs +++ b/server/src/handlers/meta.rs @@ -2028,10 +2028,7 @@ mod tests { #[tokio::test] async fn test_all_meta_handlers_comprehensive() { - use crate::handlers::tasks::{ - AddMilestoneHandler, ListMilestonesHandler, SetAcceptanceCriteriaHandler, - UpdateMilestoneHandler, VerifyAcceptanceCriteriaHandler, - }; + use crate::handlers::tasks::{MilestonesHandler, TasksHandler}; use crate::handlers::graph::SweepGraphHealthHandler; use crate::handlers::git::QueryGitDiffsHandler; @@ -2116,10 +2113,11 @@ mod tests { assert!(q_hyp_res.contains("Caching improves response speed")); // AddMilestone & UpdateMilestone & ListMilestones - let add_ms = AddMilestoneHandler; - let ms_res = add_ms + let handler_ms = MilestonesHandler; + let ms_res = handler_ms .execute( serde_json::json!({ + "action": "add", "title": "v1.0 Release", "description": "First major release" }), @@ -2130,10 +2128,10 @@ mod tests { assert!(ms_res.contains("Milestone added")); let ms_id = state.project.milestones.read_with(|ms| ms[0].id.clone()); - let upd_ms = UpdateMilestoneHandler; - let upd_ms_res = upd_ms + let upd_ms_res = handler_ms .execute( serde_json::json!({ + "action": "update", "id": ms_id, "status": "completed" }), @@ -2142,10 +2140,7 @@ mod tests { .await; assert!(upd_ms_res.is_ok()); - - - let list_ms = ListMilestonesHandler; - let list_ms_res = list_ms.execute(serde_json::json!({}), state.clone()).await.unwrap(); + let list_ms_res = handler_ms.execute(serde_json::json!({"action": "list"}), state.clone()).await.unwrap(); assert!(list_ms_res.contains("v1.0 Release")); // SetAcceptanceCriteria & VerifyAcceptanceCriteria @@ -2164,11 +2159,12 @@ mod tests { }; state.project.tasks.modify(|t| t.push(task)); - let set_ac = SetAcceptanceCriteriaHandler; - let set_ac_res = set_ac + let handler_tasks = TasksHandler; + let set_ac_res = handler_tasks .execute( serde_json::json!({ - "task_title": "AC Task", + "action": "set_criteria", + "id": "t_ac_1", "criteria": ["Code compiles cleanly", "Tests pass"] }), state.clone(), @@ -2177,12 +2173,11 @@ mod tests { .unwrap(); assert!(set_ac_res.contains("Acceptance criteria set successfully.")); - let ver_ac = VerifyAcceptanceCriteriaHandler; - let ver_ac_res = ver_ac + let ver_ac_res = handler_tasks .execute( serde_json::json!({ - "task_id": "t_ac_1", - "criteria": "Code compiles cleanly", + "action": "verify", + "id": "t_ac_1", "proof": "cargo test passed" }), state.clone(), diff --git a/server/src/handlers/notes.rs b/server/src/handlers/notes.rs index 0b73692..cfabf88 100644 --- a/server/src/handlers/notes.rs +++ b/server/src/handlers/notes.rs @@ -91,130 +91,7 @@ impl McpTool for StickyNotesHandler { } } -pub struct AddStickyNoteHandler; -#[async_trait] -impl McpTool for AddStickyNoteHandler { - fn name(&self) -> &'static str { - "add_sticky_note" - } - - fn schema(&self) -> Value { - crate::mcp::tool_def::("add_sticky_note", "Execute add_sticky_note") - } - - async fn execute(&self, args: Value, state: Arc) -> crate::error::Result { - let req: AddStickyNoteTool = serde_json::from_value(args).map_err(|e| e.to_string())?; - let now = crate::handlers::utils::now_secs(); - let expires_at = if let Some(ttl) = req.ttl_seconds { - Some(now + ttl) - } else if req.session_only.unwrap_or(false) { - Some(now + 14400) // Default 4-hour session TTL - } else { - None - }; - - state.code.sticky.modify(|notes| { - notes.push(StickyNote { - timestamp: now, - content: req.content, - expires_at, - }); - }); - Ok("Sticky note added.".to_string()) - } -} - -pub struct ReadStickyNotesHandler; - -#[async_trait] -impl McpTool for ReadStickyNotesHandler { - fn name(&self) -> &'static str { - "read_sticky_notes" - } - - fn schema(&self) -> Value { - crate::mcp::tool_def::( - "read_sticky_notes", - "Execute read_sticky_notes", - ) - } - - async fn execute(&self, _args: Value, state: Arc) -> crate::error::Result { - let now = crate::handlers::utils::now_secs(); - let mut active_notes = Vec::new(); - - state.code.sticky.modify(|notes| { - notes.retain(|n| { - if let Some(exp) = n.expires_at { - exp > now - } else { - true - } - }); - active_notes = notes.clone(); - }); - - Ok(serde_json::to_string(&active_notes)?) - } -} - -pub struct DeleteStickyNoteHandler; - -#[async_trait] -impl McpTool for DeleteStickyNoteHandler { - fn name(&self) -> &'static str { - "delete_sticky_note" - } - - fn schema(&self) -> Value { - crate::mcp::tool_def::( - "delete_sticky_note", - "Execute delete_sticky_note", - ) - } - - async fn execute(&self, args: Value, state: Arc) -> crate::error::Result { - let req: DeleteStickyNoteTool = serde_json::from_value(args).map_err(|e| e.to_string())?; - let mut success = false; - state.code.sticky.modify(|notes| { - if req.index > 0 && req.index <= notes.len() { - notes.remove(req.index - 1); - success = true; - } - }); - if success { - Ok("Sticky note deleted.".to_string()) - } else { - Err(crate::error::AppError::Internal( - "Invalid sticky note index.".to_string(), - )) - } - } -} - -pub struct ClearStickyNotesHandler; - -#[async_trait] -impl McpTool for ClearStickyNotesHandler { - fn name(&self) -> &'static str { - "clear_sticky_notes" - } - - fn schema(&self) -> Value { - crate::mcp::tool_def::( - "clear_sticky_notes", - "Execute clear_sticky_notes", - ) - } - - async fn execute(&self, _args: Value, state: Arc) -> crate::error::Result { - state.code.sticky.modify(|notes| { - notes.clear(); - }); - Ok("All sticky notes cleared.".to_string()) - } -} pub struct HandoffMemosHandler; @@ -414,37 +291,36 @@ mod tests { let dir = tempdir().unwrap(); let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap())); - let add_handler = AddStickyNoteHandler; + let handler = StickyNotesHandler; let args = json!({ + "action": "add", "content": "Buy milk", }); - let res = add_handler + let res = handler .execute(args, state.clone()) .await .map_err(|e| crate::error::AppError::Internal(e.to_string())) .unwrap(); assert!(res.contains("Sticky note added")); - let read_handler = ReadStickyNotesHandler; - let res2 = read_handler - .execute(json!({}), state.clone()) + let res2 = handler + .execute(json!({"action": "read"}), state.clone()) .await .map_err(|e| crate::error::AppError::Internal(e.to_string())) .unwrap(); assert!(res2.contains("Buy milk")); - let delete_handler = DeleteStickyNoteHandler; - let args2 = json!({"index": 1}); - let res3 = delete_handler + let args2 = json!({"action": "delete", "index": 1}); + let res3 = handler .execute(args2, state.clone()) .await .map_err(|e| crate::error::AppError::Internal(e.to_string())) .unwrap(); assert_eq!(res3, "Sticky note deleted."); - let res4 = read_handler - .execute(json!({}), state.clone()) + let res4 = handler + .execute(json!({"action": "read"}), state.clone()) .await .map_err(|e| crate::error::AppError::Internal(e.to_string())) .unwrap(); diff --git a/server/src/handlers/tasks.rs b/server/src/handlers/tasks.rs index 1488443..023a8ff 100644 --- a/server/src/handlers/tasks.rs +++ b/server/src/handlers/tasks.rs @@ -6,512 +6,6 @@ use async_trait::async_trait; use serde_json::Value; use std::sync::Arc; -pub struct AddTaskHandler; - -#[async_trait] -impl McpTool for AddTaskHandler { - fn name(&self) -> &'static str { - "add_task" - } - - fn schema(&self) -> Value { - crate::mcp::tool_def::("add_task", "Execute add_task") - } - - async fn execute(&self, args: Value, state: Arc) -> crate::error::Result { - let req: AddTaskTool = serde_json::from_value(args).map_err(|e| e.to_string())?; - let now = crate::handlers::utils::now_secs(); - let task_id = uuid::Uuid::new_v4().to_string(); - - let deps = req.dependencies.unwrap_or_default(); - - let task = Task { - id: task_id.clone(), - title: req.title, - status: "pending".to_string(), - description: req.description, - created_at: now, - updated_at: now, - git_branch: req.git_branch, - parent_id: req.parent_id, - dependencies: deps, - acceptance_criteria: vec![], - expires_at: None, - }; - let idx = state.get_search_index(); - drop(idx.index_task(&task)); - state.project.tasks.modify(|tasks| { - tasks.push(task.clone()); - }); - state.record_activity("task_create", &format!("Created task: {}", task.title), Some(&task.description)); - Ok(format!("Task added with ID: {}", task_id).to_string()) - } -} - -pub struct DeleteTaskHandler; - -#[async_trait] -impl McpTool for DeleteTaskHandler { - fn name(&self) -> &'static str { - "delete_task" - } - - fn schema(&self) -> Value { - crate::mcp::tool_def::("delete_task", "Execute delete_task") - } - - async fn execute(&self, args: Value, state: Arc) -> crate::error::Result { - let req: DeleteTaskTool = serde_json::from_value(args).map_err(|e| e.to_string())?; - let mut deleted_count = 0; - let mut actually_deleted = Vec::new(); - state.project.tasks.modify(|tasks| { - let initial_len = tasks.len(); - - // Build index-based children map - let mut children_map: std::collections::HashMap> = - std::collections::HashMap::new(); - let mut id_to_index = std::collections::HashMap::new(); - for (idx, t) in tasks.iter().enumerate() { - id_to_index.insert(t.id.as_str(), idx); - } - - for (idx, t) in tasks.iter().enumerate() { - if let Some(pid) = &t.parent_id - && let Some(&parent_idx) = id_to_index.get(pid.as_str()) - { - children_map.entry(parent_idx).or_default().push(idx); - } - } - - let mut to_delete_idx = std::collections::HashSet::new(); - if let Some(&start_idx) = id_to_index.get(req.id.as_str()) { - let mut queue = std::collections::VecDeque::new(); - queue.push_back(start_idx); - - while let Some(curr) = queue.pop_front() { - if to_delete_idx.insert(curr) - && let Some(children) = children_map.get(&curr) - { - queue.extend(children.iter().copied()); - } - } - } - - for &idx in &to_delete_idx { - actually_deleted.push(tasks[idx].id.clone()); - } - - tasks.retain(|t| !actually_deleted.contains(&t.id)); - deleted_count = initial_len - tasks.len(); - }); - - if deleted_count > 0 { - let idx = state.get_search_index(); - for id in actually_deleted { - drop(idx.delete_document(&id)); - } - Ok(vec![ - format!("Deleted task and its children ({} total).", deleted_count).to_string(), - ][0] - .clone()) - } else { - Err(crate::error::AppError::Internal("Task not found. Please use the list_active_tasks tool to verify the correct task ID.".to_string())) - } - } -} - -pub struct UpdateTaskStatusHandler; - -#[async_trait] -impl McpTool for UpdateTaskStatusHandler { - fn name(&self) -> &'static str { - "update_task_status" - } - - fn schema(&self) -> Value { - crate::mcp::tool_def::( - "update_task_status", - "Execute update_task_status", - ) - } - - async fn execute(&self, args: Value, state: Arc) -> crate::error::Result { - let req: UpdateTaskStatusTool = serde_json::from_value(args).map_err(|e| e.to_string())?; - let mut found = false; - let mut blocked = false; - let mut blocker_details = String::new(); - let target_status = req.status.to_lowercase(); - - state.project.tasks.modify(|tasks| { - // Find target task - let target_idx = tasks - .iter() - .position(|t| t.id == req.id || t.title == req.id); - let target_idx = match target_idx { - Some(idx) => idx, - None => return, - }; - - found = true; - - if target_status == "done" || target_status == "completed" { - // 1. Check Acceptance Criteria - if tasks[target_idx] - .acceptance_criteria - .iter() - .any(|c| !c.is_met) - { - blocked = true; - blocker_details = "Unmet acceptance criteria exist.".to_string(); - } - - // 2. Check dependencies - if !blocked { - let mut uncompleted_deps = Vec::new(); - for dep_id in &tasks[target_idx].dependencies { - if let Some(dep_task) = tasks.iter().find(|dt| dt.id == *dep_id) - && dep_task.status != "completed" - && dep_task.status != "done" - { - uncompleted_deps.push(dep_task.title.as_str()); - } - } - if !uncompleted_deps.is_empty() { - blocked = true; - blocker_details = - format!("Blocked by dependencies: {}", uncompleted_deps.join(", ")); - } - } - - // 3. Check child tasks - if !blocked { - let target_id_ref = tasks[target_idx].id.as_str(); - let mut uncompleted_children = Vec::new(); - for child in tasks - .iter() - .filter(|t| t.parent_id.as_deref() == Some(target_id_ref)) - { - if child.status != "completed" && child.status != "done" { - uncompleted_children.push(child.title.as_str()); - } - } - if !uncompleted_children.is_empty() { - blocked = true; - blocker_details = format!( - "Blocked by child tasks: {}", - uncompleted_children.join(", ") - ); - } - } - } - - if !blocked { - // Apply update - tasks[target_idx].status = target_status.clone(); - tasks[target_idx].updated_at = crate::handlers::utils::now_secs(); - - // Cascade cancellation to children - if target_status == "cancelled" || target_status == "abandoned" { - let mut children_map: std::collections::HashMap> = - std::collections::HashMap::new(); - - // First pass: map string ID to index to build the adjacency list by index - let mut id_to_idx = std::collections::HashMap::new(); - for (idx, t) in tasks.iter().enumerate() { - id_to_idx.insert(t.id.as_str(), idx); - } - - for (idx, t) in tasks.iter().enumerate() { - if let Some(pid) = &t.parent_id - && let Some(&p_idx) = id_to_idx.get(pid.as_str()) - { - children_map.entry(p_idx).or_default().push(idx); - } - } - - if let Some(&start_idx) = id_to_idx.get(tasks[target_idx].id.as_str()) { - let mut queue = std::collections::VecDeque::new(); - queue.push_back(start_idx); - - while let Some(curr) = queue.pop_front() { - if let Some(child_indices) = children_map.get(&curr) { - for &idx in child_indices { - if tasks[idx].status != "completed" - && tasks[idx].status != target_status - { - tasks[idx].status = target_status.clone(); - queue.push_back(idx); - } - } - } - } - } - } - } - }); - - if blocked { - Err(crate::error::AppError::Internal(format!( - "Error: Cannot transition task. {}", - blocker_details - ))) - } else if found { - state.record_activity("task_update", &format!("Task {} status -> {}", req.id, req.status), None); - Ok("Task status updated.".to_string()) - } else { - Err(crate::error::AppError::Internal("Task not found. Please use the list_active_tasks tool to verify the correct task ID.".to_string())) - } - } -} - -pub struct ListActiveTasksHandler; - -#[async_trait] -impl McpTool for ListActiveTasksHandler { - fn name(&self) -> &'static str { - "list_active_tasks" - } - - fn schema(&self) -> Value { - crate::mcp::tool_def::( - "list_active_tasks", - "Execute list_active_tasks", - ) - } - - async fn execute(&self, args: Value, state: Arc) -> crate::error::Result { - let req: ListActiveTasksTool = serde_json::from_value(args).map_err(|e| e.to_string())?; - let level = req.summary_level.as_deref().unwrap_or("detailed"); - let data = state.project.tasks.read_with(|tasks| { - let filtered: Vec<_> = tasks - .iter() - .filter(|t| { - let status_match = t.status != "done" && t.status != "completed"; - let branch_match = match &req.git_branch { - Some(branch) => { - t.git_branch.is_none() - || t.git_branch.as_deref() == Some(branch.as_str()) - } - None => true, - }; - status_match && branch_match - }) - .map(|t| match level { - "compact" => serde_json::json!({ - "id": t.id, - "title": t.title, - "status": t.status, - }), - "full" => serde_json::to_value(t).unwrap_or_default(), - _ => serde_json::json!({ - "id": t.id, - "title": t.title, - "status": t.status, - "description": t.description, - "git_branch": t.git_branch, - }), - }) - .collect(); - let mut json_str = serde_json::to_string(&filtered)?; - if let Some(max_t) = req.max_tokens { - let char_limit = max_t * 4; - if json_str.len() > char_limit { - json_str.truncate(char_limit); - json_str.push_str(" ...[truncated due to max_tokens]"); - } - } - Ok::(json_str) - })?; - Ok(data) - } -} - -pub struct SetAcceptanceCriteriaHandler; - -#[async_trait] -impl McpTool for SetAcceptanceCriteriaHandler { - fn name(&self) -> &'static str { - "set_acceptance_criteria" - } - - fn schema(&self) -> Value { - crate::mcp::tool_def::( - "set_acceptance_criteria", - "Execute set_acceptance_criteria", - ) - } - - async fn execute(&self, args: Value, state: Arc) -> crate::error::Result { - let req: SetAcceptanceCriteriaTool = - serde_json::from_value(args).map_err(|e| e.to_string())?; - let mut success = false; - state.project.tasks.modify(|tasks| { - if let Some(task) = tasks.iter_mut().rev().find(|t| t.title == req.task_title) { - task.acceptance_criteria = req - .criteria - .into_iter() - .map(|desc| crate::models::AcceptanceCriteria { - id: uuid::Uuid::new_v4().to_string(), - description: desc, - is_met: false, - }) - .collect(); - task.updated_at = crate::handlers::utils::now_secs(); - success = true; - } - }); - if success { - Ok("Acceptance criteria set successfully.".to_string()) - } else { - Err(crate::error::AppError::Internal("Task not found. Please use the list_active_tasks tool to verify the correct task ID.".to_string())) - } - } -} - -pub struct VerifyAcceptanceCriteriaHandler; - -#[async_trait] -impl McpTool for VerifyAcceptanceCriteriaHandler { - fn name(&self) -> &'static str { - "verify_acceptance_criteria" - } - - fn schema(&self) -> Value { - crate::mcp::tool_def::( - "verify_acceptance_criteria", - "Execute verify_acceptance_criteria", - ) - } - - async fn execute(&self, args: Value, state: Arc) -> crate::error::Result { - let req: VerifyAcceptanceCriteriaTool = - serde_json::from_value(args).map_err(|e| e.to_string())?; - let mut success = false; - let mut already_met = false; - state.project.tasks.modify(|tasks| { - if let Some(task) = tasks.iter_mut().find(|t| t.id == req.task_id) - && let Some(ac) = task - .acceptance_criteria - .iter_mut() - .find(|c| c.id == req.criteria || c.description == req.criteria) - { - if ac.is_met { - already_met = true; - } else { - ac.is_met = true; - success = true; - task.updated_at = crate::handlers::utils::now_secs(); - } - } - }); - if success { - Ok(format!( - "Acceptance criteria verified with proof: {}", - req.proof - )) - } else if already_met { - Ok("Acceptance criteria was already met.".to_string()) - } else { - Err(crate::error::AppError::Internal( - "Acceptance criteria or task not found.".to_string(), - )) - } - } -} - -pub struct AddMilestoneHandler; - -#[async_trait] -impl McpTool for AddMilestoneHandler { - fn name(&self) -> &'static str { - "add_milestone" - } - - fn schema(&self) -> Value { - crate::mcp::tool_def::("add_milestone", "Execute add_milestone") - } - - async fn execute(&self, args: Value, state: Arc) -> crate::error::Result { - let req: AddMilestoneTool = serde_json::from_value(args).map_err(|e| e.to_string())?; - state.project.milestones.modify(|ms| { - ms.push(crate::models::Milestone { - id: uuid::Uuid::new_v4().to_string(), - title: req.title, - status: "pending".to_string(), - namespace: req.namespace, - target_date: None, - }) - }); - Ok("Milestone added".to_string()) - } -} - -pub struct UpdateMilestoneHandler; - -#[async_trait] -impl McpTool for UpdateMilestoneHandler { - fn name(&self) -> &'static str { - "update_milestone" - } - - fn schema(&self) -> Value { - crate::mcp::tool_def::("update_milestone", "Execute update_milestone") - } - - async fn execute(&self, args: Value, state: Arc) -> crate::error::Result { - let req: UpdateMilestoneTool = serde_json::from_value(args).map_err(|e| e.to_string())?; - let mut found = false; - state.project.milestones.modify(|ms| { - for m in ms.iter_mut() { - if m.id == req.id { - m.status = req.status.clone(); - found = true; - break; - } - } - }); - if found { - Ok("Milestone updated".to_string()) - } else { - Err(crate::error::AppError::Internal( - "Milestone not found. Please verify the milestone ID using list_milestones." - .to_string(), - )) - } - } -} - -pub struct ListMilestonesHandler; - -#[async_trait] -impl McpTool for ListMilestonesHandler { - fn name(&self) -> &'static str { - "list_milestones" - } - - fn schema(&self) -> Value { - crate::mcp::tool_def::("list_milestones", "Execute list_milestones") - } - - async fn execute(&self, args: Value, state: Arc) -> crate::error::Result { - let req: ListMilestonesTool = serde_json::from_value(args).map_err(|e| e.to_string())?; - let data = state.project.milestones.read_with(|items| { - let filtered: Vec<_> = items - .iter() - .filter(|i| { - if let Some(ns) = &req.namespace { - &i.namespace == ns - } else { - true - } - }) - .collect(); - Ok::(serde_json::to_string(&filtered)?) - })?; - Ok(data) - } -} - pub struct TasksHandler; #[async_trait] @@ -535,14 +29,39 @@ impl McpTool for TasksHandler { crate::error::AppError::Internal("Missing required parameter 'title' for action 'add'. Next step: Provide non-empty 'title' string in request and retry.".to_string()) })?; let description = req.description.unwrap_or_default(); - let add_args = serde_json::json!({ - "title": title, - "description": description, - "git_branch": req.git_branch, - "parent_id": req.parent_id, - "dependencies": req.dependencies, + let now = crate::handlers::utils::now_secs(); + let task_id = uuid::Uuid::new_v4().to_string(); + let deps = req.dependencies.unwrap_or_default(); + + let task = Task { + id: task_id.clone(), + title, + status: "pending".to_string(), + description, + created_at: now, + updated_at: now, + git_branch: req.git_branch, + parent_id: req.parent_id, + dependencies: deps, + acceptance_criteria: vec![], + expires_at: None, + }; + let idx = state.get_search_index(); + drop(idx.index_task(&task)); + state.project.tasks.modify(|tasks| { + tasks.push(task.clone()); }); - AddTaskHandler.execute(add_args, state).await + state.record_activity("task_create", &format!("Created task: {}", task.title), Some(&task.description)); + state.broadcast_task_event(TaskEvent { + task_id: task_id.clone(), + status: "created".to_string(), + action: Some("add".to_string()), + result: Some(serde_json::json!({ "title": task.title, "git_branch": task.git_branch })), + error: None, + timestamp: now, + session_id: None, + }); + Ok(format!("Task added with ID: {}", task_id)) } TaskAction::Update => { let id = req.id.ok_or_else(|| { @@ -551,50 +70,253 @@ impl McpTool for TasksHandler { let status = req.status.ok_or_else(|| { crate::error::AppError::Internal("Missing required parameter 'status' for action 'update'. Next step: Provide valid 'status' ('pending', 'completed', or 'cancelled') in request and retry.".to_string()) })?; - let update_args = serde_json::json!({ - "id": id, - "status": status, + let target_status = status.to_lowercase(); + let mut found = false; + let mut blocked = false; + let mut blocker_details = String::new(); + + state.project.tasks.modify(|tasks| { + let target_idx = tasks.iter().position(|t| t.id == id || t.title == id); + let target_idx = match target_idx { + Some(i) => i, + None => return, + }; + found = true; + + if target_status == "done" || target_status == "completed" { + if tasks[target_idx].acceptance_criteria.iter().any(|c| !c.is_met) { + blocked = true; + blocker_details = "Unmet acceptance criteria exist.".to_string(); + } + if !blocked { + let mut uncompleted_deps = Vec::new(); + for dep_id in &tasks[target_idx].dependencies { + if let Some(dep_task) = tasks.iter().find(|dt| dt.id == *dep_id) + && dep_task.status != "completed" + && dep_task.status != "done" + { + uncompleted_deps.push(dep_task.title.as_str()); + } + } + if !uncompleted_deps.is_empty() { + blocked = true; + blocker_details = format!("Blocked by dependencies: {}", uncompleted_deps.join(", ")); + } + } + if !blocked { + let target_id_ref = tasks[target_idx].id.as_str(); + let mut uncompleted_children = Vec::new(); + for child in tasks.iter().filter(|t| t.parent_id.as_deref() == Some(target_id_ref)) { + if child.status != "completed" && child.status != "done" { + uncompleted_children.push(child.title.as_str()); + } + } + if !uncompleted_children.is_empty() { + blocked = true; + blocker_details = format!("Blocked by child tasks: {}", uncompleted_children.join(", ")); + } + } + } + + if !blocked { + tasks[target_idx].status = target_status.clone(); + tasks[target_idx].updated_at = crate::handlers::utils::now_secs(); + } }); - UpdateTaskStatusHandler.execute(update_args, state).await + + if blocked { + state.broadcast_task_event(TaskEvent { + task_id: id.clone(), + status: "blocked".to_string(), + action: Some("update".to_string()), + result: None, + error: Some(blocker_details.clone()), + timestamp: crate::handlers::utils::now_secs(), + session_id: None, + }); + Err(crate::error::AppError::Internal(format!("Error: Cannot transition task. {}", blocker_details))) + } else if found { + state.record_activity("task_update", &format!("Task {} status -> {}", id, status), None); + state.broadcast_task_event(TaskEvent { + task_id: id.clone(), + status: target_status.clone(), + action: Some("update".to_string()), + result: Some(serde_json::json!({ "status": target_status })), + error: None, + timestamp: crate::handlers::utils::now_secs(), + session_id: None, + }); + Ok("Task status updated.".to_string()) + } else { + Err(crate::error::AppError::Internal("Task not found. Please verify the task ID.".to_string())) + } } TaskAction::Delete => { let id = req.id.ok_or_else(|| { crate::error::AppError::Internal("Missing required parameter 'id' for action 'delete'. Next step: Provide task 'id' string in request and retry.".to_string()) })?; - let del_args = serde_json::json!({ - "id": id, + let mut deleted_count = 0; + let mut actually_deleted = Vec::new(); + state.project.tasks.modify(|tasks| { + let initial_len = tasks.len(); + let mut id_to_index = std::collections::HashMap::new(); + for (idx, t) in tasks.iter().enumerate() { + id_to_index.insert(t.id.as_str(), idx); + } + let mut children_map: std::collections::HashMap> = std::collections::HashMap::new(); + for (idx, t) in tasks.iter().enumerate() { + if let Some(pid) = &t.parent_id && let Some(&p_idx) = id_to_index.get(pid.as_str()) { + children_map.entry(p_idx).or_default().push(idx); + } + } + let mut to_delete_idx = std::collections::HashSet::new(); + if let Some(&start_idx) = id_to_index.get(id.as_str()) { + let mut queue = std::collections::VecDeque::new(); + queue.push_back(start_idx); + while let Some(curr) = queue.pop_front() { + if to_delete_idx.insert(curr) && let Some(children) = children_map.get(&curr) { + queue.extend(children.iter().copied()); + } + } + } + for &idx in &to_delete_idx { + actually_deleted.push(tasks[idx].id.clone()); + } + tasks.retain(|t| !actually_deleted.contains(&t.id)); + deleted_count = initial_len - tasks.len(); }); - DeleteTaskHandler.execute(del_args, state).await + + if deleted_count > 0 { + let idx = state.get_search_index(); + for deleted_id in actually_deleted { + drop(idx.delete_document(&deleted_id)); + } + state.broadcast_task_event(TaskEvent { + task_id: id.clone(), + status: "deleted".to_string(), + action: Some("delete".to_string()), + result: Some(serde_json::json!({ "deleted_count": deleted_count })), + error: None, + timestamp: crate::handlers::utils::now_secs(), + session_id: None, + }); + Ok(format!("Deleted task and its children ({} total).", deleted_count)) + } else { + Err(crate::error::AppError::Internal("Task not found. Please verify the task ID.".to_string())) + } } TaskAction::List => { - let list_args = serde_json::json!({ - "git_branch": req.git_branch, - "summary_level": req.summary_level, - "max_tokens": req.max_tokens, - }); - ListActiveTasksHandler.execute(list_args, state).await + let level = req.summary_level.as_deref().unwrap_or("detailed"); + let data = state.project.tasks.read_with(|tasks| { + let filtered: Vec<_> = tasks + .iter() + .filter(|t| { + let status_match = t.status != "done" && t.status != "completed"; + let branch_match = match &req.git_branch { + Some(branch) => t.git_branch.is_none() || t.git_branch.as_deref() == Some(branch.as_str()), + None => true, + }; + status_match && branch_match + }) + .map(|t| match level { + "compact" => serde_json::json!({ "id": t.id, "title": t.title, "status": t.status }), + "full" => serde_json::to_value(t).unwrap_or_default(), + _ => serde_json::json!({ "id": t.id, "title": t.title, "status": t.status, "description": t.description, "git_branch": t.git_branch }), + }) + .collect(); + let mut json_str = serde_json::to_string(&filtered)?; + if let Some(max_t) = req.max_tokens { + let char_limit = max_t * 4; + if json_str.len() > char_limit { + json_str.truncate(char_limit); + json_str.push_str(" ...[truncated due to max_tokens]"); + } + } + Ok::(json_str) + })?; + Ok(data) } TaskAction::SetCriteria => { - let id = req.id.ok_or_else(|| { - crate::error::AppError::Internal("Missing required parameter 'id' for action 'set_criteria'. Next step: Provide task 'id' string in request and retry.".to_string()) + let id = req.id.or(req.title.clone()).ok_or_else(|| { + crate::error::AppError::Internal("Missing required parameter 'id' or 'title' for action 'set_criteria'. Next step: Provide task 'id' string in request and retry.".to_string()) })?; - let criteria = req.criteria.ok_or_else(|| { - crate::error::AppError::Internal("Missing required parameter 'criteria' for action 'set_criteria'. Next step: Provide array of acceptance criteria descriptions in request and retry.".to_string()) + let criteria_list = req.criteria.ok_or_else(|| { + crate::error::AppError::Internal("Missing required parameter 'criteria' for action 'set_criteria'. Next step: Provide array of acceptance criteria strings in request and retry.".to_string()) })?; - let set_args = serde_json::json!({ - "id": id, - "acceptance_criteria": criteria, + let mut success = false; + state.project.tasks.modify(|tasks| { + if let Some(task) = tasks.iter_mut().find(|t| t.id == id || t.title == id) { + task.acceptance_criteria = criteria_list + .into_iter() + .map(|desc| crate::models::AcceptanceCriteria { + id: uuid::Uuid::new_v4().to_string(), + description: desc, + is_met: false, + }) + .collect(); + task.updated_at = crate::handlers::utils::now_secs(); + success = true; + } }); - SetAcceptanceCriteriaHandler.execute(set_args, state).await + if success { + state.broadcast_task_event(TaskEvent { + task_id: id.clone(), + status: "criteria_set".to_string(), + action: Some("set_criteria".to_string()), + result: Some(serde_json::json!({ "id": id })), + error: None, + timestamp: crate::handlers::utils::now_secs(), + session_id: None, + }); + Ok("Acceptance criteria set successfully.".to_string()) + } else { + Err(crate::error::AppError::Internal("Task not found. Please verify the task ID.".to_string())) + } } TaskAction::Verify => { let id = req.id.ok_or_else(|| { crate::error::AppError::Internal("Missing required parameter 'id' for action 'verify'. Next step: Provide task 'id' string in request and retry.".to_string()) })?; - let verify_args = serde_json::json!({ - "id": id, + let proof_str = req.proof.unwrap_or_else(|| "Verified".to_string()); + let mut success = false; + let mut already_met = false; + state.project.tasks.modify(|tasks| { + if let Some(task) = tasks.iter_mut().find(|t| t.id == id || t.title == id) { + if let Some(ac) = task.acceptance_criteria.iter_mut().next() { + if ac.is_met { + already_met = true; + } else { + ac.is_met = true; + success = true; + task.updated_at = crate::handlers::utils::now_secs(); + } + } else { + task.acceptance_criteria.push(crate::models::AcceptanceCriteria { + id: uuid::Uuid::new_v4().to_string(), + description: proof_str.clone(), + is_met: true, + }); + task.updated_at = crate::handlers::utils::now_secs(); + success = true; + } + } }); - VerifyAcceptanceCriteriaHandler.execute(verify_args, state).await + if success { + state.broadcast_task_event(TaskEvent { + task_id: id.clone(), + status: "verified".to_string(), + action: Some("verify".to_string()), + result: Some(serde_json::json!({ "proof": proof_str })), + error: None, + timestamp: crate::handlers::utils::now_secs(), + session_id: None, + }); + Ok(format!("Acceptance criteria verified with proof: {}", proof_str)) + } else if already_met { + Ok("Acceptance criteria was already met.".to_string()) + } else { + Err(crate::error::AppError::Internal("Task not found. Please verify the task ID.".to_string())) + } } } } @@ -622,11 +344,17 @@ impl McpTool for MilestonesHandler { let title = req.title.ok_or_else(|| { crate::error::AppError::Internal("Missing required parameter 'title' for action 'add'. Next step: Provide non-empty 'title' string in request and retry.".to_string()) })?; - let add_args = serde_json::json!({ - "title": title, - "namespace": req.namespace.unwrap_or_else(|| crate::models::default_namespace()), + let ns = req.namespace.unwrap_or_else(|| crate::models::default_namespace()); + state.project.milestones.modify(|ms| { + ms.push(crate::models::Milestone { + id: uuid::Uuid::new_v4().to_string(), + title, + status: "pending".to_string(), + namespace: ns, + target_date: None, + }) }); - AddMilestoneHandler.execute(add_args, state).await + Ok("Milestone added".to_string()) } MilestoneAction::Update => { let id = req.id.ok_or_else(|| { @@ -635,17 +363,37 @@ impl McpTool for MilestonesHandler { let status = req.status.ok_or_else(|| { crate::error::AppError::Internal("Missing required parameter 'status' for action 'update'. Next step: Provide milestone 'status' in request and retry.".to_string()) })?; - let update_args = serde_json::json!({ - "id": id, - "status": status, + let mut found = false; + state.project.milestones.modify(|ms| { + for m in ms.iter_mut() { + if m.id == id { + m.status = status.clone(); + found = true; + break; + } + } }); - UpdateMilestoneHandler.execute(update_args, state).await + if found { + Ok("Milestone updated".to_string()) + } else { + Err(crate::error::AppError::Internal("Milestone not found. Please verify the milestone ID.".to_string())) + } } MilestoneAction::List => { - let list_args = serde_json::json!({ - "namespace": req.namespace, - }); - ListMilestonesHandler.execute(list_args, state).await + let data = state.project.milestones.read_with(|items| { + let filtered: Vec<_> = items + .iter() + .filter(|i| { + if let Some(ns) = &req.namespace { + &i.namespace == ns + } else { + true + } + }) + .collect(); + Ok::(serde_json::to_string(&filtered)?) + })?; + Ok(data) } } } @@ -662,23 +410,22 @@ mod tests { let dir = tempdir().unwrap(); let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap())); - let add_handler = AddTaskHandler; + let handler = TasksHandler; let args = json!({ + "action": "add", "title": "Fix the hyperdrive", "description": "It's making a strange noise", - "acceptance_criteria": ["Stop the noise", "Reach lightspeed"], }); - let res = add_handler + let res = handler .execute(args, state.clone()) .await .map_err(|e| crate::error::AppError::Internal(e.to_string())) .unwrap(); assert!(res.contains("Task added with ID:")); - let list_handler = ListActiveTasksHandler; - let res2 = list_handler - .execute(json!({}), state.clone()) + let res2 = handler + .execute(json!({"action": "list"}), state.clone()) .await .map_err(|e| crate::error::AppError::Internal(e.to_string())) .unwrap(); @@ -690,10 +437,10 @@ mod tests { let dir = tempdir().unwrap(); let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap())); - let add_handler = AddTaskHandler; - let res = add_handler + let handler = TasksHandler; + let res = handler .execute( - json!({"title": "Test", "description": "test"}), + json!({"action": "add", "title": "Test", "description": "test"}), state.clone(), ) .await @@ -703,21 +450,20 @@ mod tests { let id_start = res.find("ID: ").unwrap() + 4; let task_id = res[id_start..].trim(); - let update_handler = UpdateTaskStatusHandler; let args = json!({ + "action": "update", "id": task_id, "status": "done" }); - let res3 = update_handler + let res3 = handler .execute(args, state.clone()) .await .map_err(|e| crate::error::AppError::Internal(e.to_string())) .unwrap(); assert_eq!(res3, "Task status updated."); - let list_handler = ListActiveTasksHandler; - let res4 = list_handler - .execute(json!({}), state.clone()) + let res4 = handler + .execute(json!({"action": "list"}), state.clone()) .await .map_err(|e| crate::error::AppError::Internal(e.to_string())) .unwrap(); @@ -729,54 +475,45 @@ mod tests { let dir = tempdir().unwrap(); let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap())); - // Add Milestone - let add_milestone = AddMilestoneHandler; + let handler_ms = MilestonesHandler; let args_ms = json!({ - "name": "v1.0", + "action": "add", "title": "Release 1.0", - "description": "First release", - "target_date": 1700000000, - "end_date": 1700000000, "namespace": "global" }); - let res1 = add_milestone + let res1 = handler_ms .execute(args_ms, state.clone()) .await .map_err(|e| crate::error::AppError::Internal(e.to_string())) .unwrap(); assert!(res1.contains("Milestone added")); - // Fetch milestone ID from state directly to update let ms_id = state.project.milestones.read_with(|ms| ms[0].id.clone()); - // Update Milestone - let update_ms = UpdateMilestoneHandler; let args_ums = json!({ + "action": "update", "id": ms_id, "status": "completed" }); - let res2 = update_ms + let res2 = handler_ms .execute(args_ums, state.clone()) .await .map_err(|e| crate::error::AppError::Internal(e.to_string())) .unwrap(); assert_eq!(res2, "Milestone updated"); - // List Milestones - let list_ms = ListMilestonesHandler; - let res3 = list_ms - .execute(json!({"namespace": "global"}), state.clone()) + let res3 = handler_ms + .execute(json!({"action": "list", "namespace": "global"}), state.clone()) .await .map_err(|e| crate::error::AppError::Internal(e.to_string())) .unwrap(); assert!(res3.contains("completed")); assert!(res3.contains("Release 1.0")); - // Task Acceptance Criteria - let add_task = AddTaskHandler; - let res_task = add_task + let handler_t = TasksHandler; + let res_task = handler_t .execute( - json!({"title": "Test", "description": "desc"}), + json!({"action": "add", "title": "Test", "description": "desc"}), state.clone(), ) .await @@ -784,26 +521,24 @@ mod tests { .unwrap(); let task_id = res_task[res_task.find("ID: ").unwrap() + 4..].trim(); - let set_ac = SetAcceptanceCriteriaHandler; let args_ac = json!({ - "task_id": task_id, - "task_title": "Test", + "action": "set_criteria", + "id": task_id, "criteria": ["Do X", "Do Y"] }); - let res4 = set_ac + let res4 = handler_t .execute(args_ac, state.clone()) .await .map_err(|e| crate::error::AppError::Internal(e.to_string())) .unwrap(); assert_eq!(res4, "Acceptance criteria set successfully."); - let verify_ac = VerifyAcceptanceCriteriaHandler; let args_vac = json!({ - "task_id": task_id, - "criteria": "Do X", + "action": "verify", + "id": task_id, "proof": "I did X" }); - let res5 = verify_ac + let res5 = handler_t .execute(args_vac, state.clone()) .await .map_err(|e| crate::error::AppError::Internal(e.to_string())) @@ -816,10 +551,10 @@ mod tests { let dir = tempfile::tempdir().unwrap(); let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap())); - let add_task = AddTaskHandler; - let parent = add_task + let handler = TasksHandler; + let parent = handler .execute( - json!({"title": "Parent", "description": "p"}), + json!({"action": "add", "title": "Parent", "description": "p"}), state.clone(), ) .await @@ -829,9 +564,9 @@ mod tests { .trim() .to_string(); - let child = add_task + let child = handler .execute( - json!({"title": "Child", "description": "c", "parent_id": parent_id}), + json!({"action": "add", "title": "Child", "description": "c", "parent_id": parent_id}), state.clone(), ) .await @@ -839,9 +574,8 @@ mod tests { .unwrap(); let _child_id = child[child.find("ID: ").unwrap() + 4..].trim().to_string(); - let del_task = DeleteTaskHandler; - let res_del = del_task - .execute(json!({"id": parent_id}), state.clone()) + let res_del = handler + .execute(json!({"action": "delete", "id": parent_id}), state.clone()) .await .map_err(|e| crate::error::AppError::Internal(e.to_string())) .unwrap(); @@ -849,34 +583,28 @@ mod tests { } #[tokio::test] - async fn test_list_milestones_with_namespace() { let dir = tempfile::tempdir().unwrap(); let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap())); - let add_milestone = AddMilestoneHandler; + let handler = MilestonesHandler; let args_ms = serde_json::json!({ - "name": "v1.0", + "action": "add", "title": "Release 1.0", - "description": "First release", - "target_date": 1700000000, - "end_date": 1700000000, "namespace": "global" }); - let res1 = add_milestone + let res1 = handler .execute(args_ms, state.clone()) .await .map_err(|e| crate::error::AppError::Internal(e.to_string())) .unwrap(); assert!(res1.contains("Milestone added")); - let list_ms = ListMilestonesHandler; - let res2 = list_ms - .execute(serde_json::json!({"namespace": "global"}), state.clone()) + let res2 = handler + .execute(serde_json::json!({"action": "list", "namespace": "global"}), state.clone()) .await .map_err(|e| crate::error::AppError::Internal(e.to_string())) .unwrap(); assert!(res2.contains("Release 1.0")); } } - diff --git a/server/src/handlers/vision.rs b/server/src/handlers/vision.rs index 5242f7f..8979787 100644 --- a/server/src/handlers/vision.rs +++ b/server/src/handlers/vision.rs @@ -206,6 +206,8 @@ impl McpTool for ToggleClipboardWatchModeHandler { let mut watch_mode = state.clipboard_watch_mode.write().await; *watch_mode = tool_args.enable; + drop(watch_mode); + state.clipboard_notify.notify_waiters(); let status_msg = if tool_args.enable { "Clipboard watch mode enabled. Changes will be ingested as StickyNotes." diff --git a/server/src/indexer.rs b/server/src/indexer.rs index 898ba9a..cc16b37 100644 --- a/server/src/indexer.rs +++ b/server/src/indexer.rs @@ -13,32 +13,38 @@ pub async fn start_background_indexer(state: Arc) { tokio::spawn(async move { tracing::info!("Starting background indexer in {:?}", workspace_root); - let walker = WalkBuilder::new(&workspace_root) - .hidden(true) - .git_ignore(true) - .build(); + let root_clone = workspace_root.clone(); + let files_to_process = tokio::task::spawn_blocking(move || { + let walker = WalkBuilder::new(&root_clone) + .hidden(true) + .git_ignore(true) + .build(); - let mut files_to_process = Vec::new(); - for result in walker { - match result { - Ok(entry) => { - if entry.file_type().is_some_and(|ft| ft.is_file()) { - let path = entry.path().to_path_buf(); - let ext = path.extension().and_then(|e| e.to_str()).unwrap_or(""); - if [ - "rs", "ts", "js", "jsx", "tsx", "py", "java", "c", "cpp", "go", - ] - .contains(&ext) - { - files_to_process.push(path); + let mut files = Vec::new(); + for result in walker { + match result { + Ok(entry) => { + if entry.file_type().is_some_and(|ft| ft.is_file()) { + let path = entry.path().to_path_buf(); + let ext = path.extension().and_then(|e| e.to_str()).unwrap_or(""); + if [ + "rs", "ts", "js", "jsx", "tsx", "py", "java", "c", "cpp", "go", + ] + .contains(&ext) + { + files.push(path); + } } } - } - Err(e) => { - tracing::warn!("Error walking directory: {}", e); + Err(e) => { + tracing::warn!("Error walking directory: {}", e); + } } } - } + files + }) + .await + .unwrap_or_default(); let idx = state.get_search_index(); diff --git a/server/src/lib.rs b/server/src/lib.rs index edf812f..6a1c6d1 100644 --- a/server/src/lib.rs +++ b/server/src/lib.rs @@ -98,7 +98,7 @@ pub struct AppState { pub async fn ttl_sweeper_worker(state: Arc) { loop { - tokio::time::sleep(Duration::from_secs(3600)).await; + state.ttl_notify.notified().await; let now = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .unwrap_or_default() @@ -121,7 +121,7 @@ pub async fn ttl_sweeper_worker(state: Arc) { pub async fn index_committer_worker(state: Arc) { loop { - tokio::time::sleep(Duration::from_secs(5)).await; + state.index_commit_notify.notified().await; let idx_opt = state.search_index.read().ok().map(|idx| idx.clone()); if let Some(idx) = idx_opt { let _ = idx.commit().await; @@ -131,7 +131,7 @@ pub async fn index_committer_worker(state: Arc) { pub async fn condense_graph_worker(state: Arc) { loop { - tokio::time::sleep(Duration::from_secs(3600)).await; + state.condense_notify.notified().await; let threshold: usize = std::env::var("MCP_MEMORY_CONDENSE_THRESHOLD") .unwrap_or_else(|_| "100".to_string()) @@ -659,6 +659,10 @@ mod tests { let handle2 = tokio::spawn(index_committer_worker(state.clone())); let handle3 = tokio::spawn(condense_graph_worker(state.clone())); + state.ttl_notify.notify_one(); + state.index_commit_notify.notify_one(); + state.condense_notify.notify_one(); + tokio::time::sleep(Duration::from_millis(50)).await; handle1.abort(); diff --git a/server/src/models.rs b/server/src/models.rs index 51fdcaf..c93a310 100644 --- a/server/src/models.rs +++ b/server/src/models.rs @@ -309,6 +309,17 @@ pub struct AgentSignal { pub ttl_seconds: Option, } +#[derive(Debug, Clone, Serialize, Deserialize, Default)] +pub struct TaskEvent { + pub task_id: String, + pub status: String, + pub action: Option, + pub result: Option, + pub error: Option, + pub timestamp: u64, + pub session_id: Option, +} + #[cfg(test)] mod tests { use axum::http::StatusCode; diff --git a/server/src/router.rs b/server/src/router.rs index 4380323..6ff058e 100644 --- a/server/src/router.rs +++ b/server/src/router.rs @@ -495,94 +495,46 @@ impl MemoryHandler { register!(tasks::TasksHandler); register!(tasks::MilestonesHandler); - register!(tasks::AddTaskHandler); - register!(tasks::DeleteTaskHandler); - register!(tasks::UpdateTaskStatusHandler); - register!(tasks::ListActiveTasksHandler); - register!(tasks::SetAcceptanceCriteriaHandler); - register!(tasks::VerifyAcceptanceCriteriaHandler); - register!(tasks::AddMilestoneHandler); - register!(tasks::UpdateMilestoneHandler); - register!(tasks::ListMilestonesHandler); register!(notes::StickyNotesHandler); register!(notes::HandoffMemosHandler); - register!(notes::AddStickyNoteHandler); - register!(notes::ReadStickyNotesHandler); - register!(notes::DeleteStickyNoteHandler); - register!(notes::ClearStickyNotesHandler); register!(notes::AddSessionSummaryHandler); register!(notes::GenerateStandupReportHandler); register!(notes::PromoteToEntityHandler); register!(meta::DecisionsHandler); register!(meta::TechDebtHandler); - register!(meta::LogDecisionHandler); - register!(meta::QueryDecisionsHandler); - register!(meta::DeleteDecisionHandler); register!(meta::LogErrorFixHandler); register!(meta::SearchErrorFixesHandler); register!(meta::LogCodeChangeHandler); register!(meta::QueryRecentChangesHandler); register!(meta::LearnPreferenceHandler); register!(meta::ReadPreferencesHandler); - register!(meta::LogTechDebtHandler); - register!(meta::ResolveTechDebtHandler); - register!(meta::ListTechDebtHandler); register!(meta::OmniSearchHandler); register!(meta::GetProjectHealthHandler); register!(env::EnvironmentHandler); - register!(env::UpdateEnvFingerprintHandler); - register!(env::ReadEnvFingerprintHandler); - register!(env::LogEnvRequirementHandler); - register!(env::RegisterEnvironmentHandler); - register!(env::GetEnvironmentDetailsHandler); register!(workspaces::PinnedFilesHandler); register!(workspaces::ContextWorkspacesHandler); register!(workspaces::PrChecklistHandler); register!(workspaces::SnippetsHandler); - register!(workspaces::PinFileHandler); - register!(workspaces::UnpinFileHandler); - register!(workspaces::ListPinnedFilesHandler); - register!(workspaces::StoreSnippetHandler); - register!(workspaces::SearchSnippetsHandler); - register!(workspaces::DeleteSnippetHandler); - register!(workspaces::SaveContextWorkspaceHandler); - register!(workspaces::LoadContextWorkspaceHandler); - register!(workspaces::ListContextWorkspacesHandler); - register!(workspaces::DeleteContextWorkspaceHandler); - register!(workspaces::AddPrChecklistItemHandler); - register!(workspaces::GetPrChecklistHandler); - register!(workspaces::ClearPrChecklistHandler); register!(vision::ClipboardHandler); - register!(vision::ReadClipboardHandler); - register!(vision::WriteClipboardHandler); + register!(git::GetActiveWorktreeContextHandler); register!(git::QueryGitDiffsHandler); register!(logs::WatchProcessLogsHandler); register!(logs::GetRecentLogsHandler); register!(ast::ReadFileSkeletonHandler); - register!(vision::ToggleClipboardWatchModeHandler); register!(ast::ReplaceAstNodeHandler); register!(ast::FindSymbolReferencesHandler); register!(ast::GetCallersHandler); register!(ast::AnalyzeImpactHandler); register!(workspaces::ReadDirectoryArchitectureHandler); register!(workspaces::SemanticCodeSearchHandler); - register!(workspaces::CreateSnapshotHandler); - register!(workspaces::RestoreSnapshotHandler); - register!(workspaces::CreateSubagentNamespaceHandler); register!(workspaces::ManageSubagentNamespaceHandler); register!(graph::GetSubgraphHandler); - register!(meta::SuggestErrorFixHandler); - register!(meta::CheckpointStateHandler); register!(meta::ManageCheckpointHandler); - register!(meta::RestoreStateHandler); - register!(workspaces::TagSnippetHandler); - register!(workspaces::PurgeSubagentNamespaceHandler); - register!(workspaces::CondenseSubagentNamespaceHandler); register!(graph::SweepGraphHealthHandler); register!(meta::QueryLineageHandler); register!(meta::GetNextActionableTasksHandler); @@ -595,7 +547,6 @@ impl MemoryHandler { register!(meta::BroadcastAgentSignalHandler); register!(meta::QueryAgentSignalsHandler); register!(meta::AutoSessionCheckpointHandler); - register!(meta::SearchSnippetsHybridHandler); Self { state, @@ -763,19 +714,21 @@ impl MemoryHandler { .unwrap_or(serde_json::Value::Object(Default::default())); let category = match name { - "read_clipboard" | "write_clipboard" | "toggle_clipboard_watch_mode" => "CLIPBOARD", + "clipboard" => "CLIPBOARD", "create_entities" | "create_relations" | "read_graph" | "get_subgraph" | "search_graph" | "get_schema" => "GRAPH", - "log_decision" => "DECISION", + "decisions" => "DECISION", "log_code_change" => "CODE", "log_error_fix" => "ERROR_FIX", - "log_tech_debt" => "TECH_DEBT", - "add_task" | "update_task_status" | "delete_task" | "add_milestone" => "TASK", - "manage_sticky_notes" | "read_notes" => "STICKY_NOTE", + "tech_debt" => "TECH_DEBT", + "tasks" | "milestones" => "TASK", + "sticky_notes" | "handoff_memos" => "STICKY_NOTE", "manage_checkpoint" => "CHECKPOINT", "manage_subagent_namespace" => "SUBAGENT", - "search_snippets" => "SNIPPET", + "snippets" => "SNIPPET", "search_web" => "WEB_SEARCH", "omni_search" => "OMNI_SEARCH", + "environment" => "ENVIRONMENT", + "pinned_files" | "context_workspaces" | "pr_checklist" => "WORKSPACE", _ => "TOOL", }; @@ -829,10 +782,38 @@ impl MemoryHandler { pub fn format_tool_activity_description(name: &str, args: &serde_json::Value) -> String { let (action, detail) = match name { + "tasks" => { + let act = args.get("action").and_then(|v| v.as_str()).unwrap_or("manage"); + let title = args.get("title").or_else(|| args.get("id")).and_then(|v| v.as_str()).unwrap_or(""); + ("Tasks", format!("{}: {}", act, title).trim_end_matches(": ").to_string()) + } + "decisions" => { + let act = args.get("action").and_then(|v| v.as_str()).unwrap_or("log"); + let title = args.get("title").or_else(|| args.get("query")).and_then(|v| v.as_str()).unwrap_or(""); + ("Decisions", format!("{}: {}", act, title).trim_end_matches(": ").to_string()) + } + "tech_debt" => { + let act = args.get("action").and_then(|v| v.as_str()).unwrap_or("log"); + let desc = args.get("description").or_else(|| args.get("id")).and_then(|v| v.as_str()).unwrap_or(""); + ("Tech Debt", format!("{}: {}", act, desc).trim_end_matches(": ").to_string()) + } + "sticky_notes" => { + let act = args.get("action").and_then(|v| v.as_str()).unwrap_or("add"); + let preview = args.get("content").and_then(|v| v.as_str()).map(|c| c.chars().take(40).collect::()).unwrap_or_default(); + ("Sticky Notes", format!("{}: {}", act, preview).trim_end_matches(": ").to_string()) + } + "clipboard" => { + let act = args.get("action").and_then(|v| v.as_str()).unwrap_or("read"); + ("Clipboard", act.to_string()) + } + "snippets" => { + let act = args.get("action").and_then(|v| v.as_str()).unwrap_or("search"); + let q = args.get("query").and_then(|v| v.as_str()).unwrap_or(""); + ("Snippets", format!("{}: {}", act, q).trim_end_matches(": ").to_string()) + } "log_code_change" => { let file = args.get("file_path") .or_else(|| args.get("file")) - .or_else(|| args.get("path")) .or_else(|| args.get("target_file")) .and_then(|v| v.as_str()); let summary = args.get("summary") @@ -847,15 +828,6 @@ pub fn format_tool_activity_description(name: &str, args: &serde_json::Value) -> }; ("Log Code Change", d) } - "log_decision" => { - let d = args.get("title") - .or_else(|| args.get("decision")) - .or_else(|| args.get("summary")) - .and_then(|v| v.as_str()) - .unwrap_or("") - .to_string(); - ("Log Decision", d) - } "log_error_fix" => { let d = args.get("error") .or_else(|| args.get("summary")) @@ -865,16 +837,6 @@ pub fn format_tool_activity_description(name: &str, args: &serde_json::Value) -> .to_string(); ("Log Error Fix", d) } - "log_tech_debt" => { - let d = if let Some(summary) = args.get("summary").or_else(|| args.get("description")).and_then(|v| v.as_str()) { - summary.to_string() - } else if let Some(file) = args.get("file_path").or_else(|| args.get("file")).and_then(|v| v.as_str()) { - file.to_string() - } else { - String::new() - }; - ("Log Tech Debt", d) - } "create_entities" => { let d = if let Some(entities) = args.get("entities").and_then(|v| v.as_array()) { let names: Vec<&str> = entities @@ -917,54 +879,17 @@ pub fn format_tool_activity_description(name: &str, args: &serde_json::Value) -> }; ("Create Relations", d) } - "add_task" => { - let d = args.get("title") - .or_else(|| args.get("name")) - .and_then(|v| v.as_str()) - .unwrap_or("") - .to_string(); - ("Add Task", d) - } - "update_task_status" => { - let d = if let (Some(id), Some(status)) = ( - args.get("task_id").or_else(|| args.get("id")).and_then(|v| v.as_str()), - args.get("status").and_then(|v| v.as_str()), - ) { - format!("Task {} -> {}", id, status) - } else { - String::new() - }; - ("Update Task Status", d) - } - "omni_search" | "search_graph" | "search_snippets" | "search_web" => { + "omni_search" | "search_graph" | "search_web" => { let d = args.get("query") .and_then(|v| v.as_str()) .map(|q| format!("\"{}\"", q)) .unwrap_or_default(); ("Search", d) } - "manage_sticky_notes" | "add_sticky_note" => { - let action = args.get("action").and_then(|v| v.as_str()).unwrap_or("add"); - let d = if let Some(content) = args.get("content").and_then(|v| v.as_str()) { - let preview: String = content.chars().take(40).collect(); - format!("{} \"{}\"", action, preview) - } else { - action.to_string() - }; - ("Sticky Note", d) - } - "write_clipboard" => { - let d = if let Some(text) = args.get("text").or_else(|| args.get("content")).and_then(|v| v.as_str()) { - let preview: String = text.chars().take(40).collect(); - format!("\"{}\"", preview) - } else { - String::new() - }; - ("Write Clipboard", d) - } _ => { - let d = if let Some(title) = args + let d = args .get("title") + .or_else(|| args.get("action")) .or_else(|| args.get("summary")) .or_else(|| args.get("description")) .or_else(|| args.get("name")) @@ -973,12 +898,8 @@ pub fn format_tool_activity_description(name: &str, args: &serde_json::Value) -> .or_else(|| args.get("file")) .or_else(|| args.get("path")) .and_then(|v| v.as_str()) - { - title.to_string() - } else { - String::new() - }; - (name, d) + .unwrap_or(""); + (name, d.to_string()) } }; @@ -1021,7 +942,7 @@ mod tests { // Assert some known tools are registered assert!(handler.tools.contains_key("create_entities")); - assert!(handler.tools.contains_key("add_task")); + assert!(handler.tools.contains_key("tasks")); // Ensure we can fetch list of tools let list_tools_req = json!({ @@ -1145,10 +1066,11 @@ mod tests { "id": 3, "method": "tools/call", "params": { - "name": "update_task_status", + "name": "tasks", "arguments": { + "action": "update", "id": "nonexistent_task_123", - "status": "in_progress" + "status": "completed" } } }); diff --git a/server/src/state.rs b/server/src/state.rs index 8f832c8..34fb3b8 100644 --- a/server/src/state.rs +++ b/server/src/state.rs @@ -51,6 +51,10 @@ pub struct TelemetryStores { pub struct MemoryState { pub base_dir: PathBuf, pub clipboard_watch_mode: tokio::sync::RwLock, + pub clipboard_notify: Arc, + pub index_commit_notify: Arc, + pub ttl_notify: Arc, + pub condense_notify: Arc, pub graph: Store, pub search_index: RwLock, pub vector_db: tokio::sync::RwLock>, @@ -75,6 +79,10 @@ impl MemoryState { let state = Self { ollama: Arc::new(crate::ollama::OllamaClient::new_from_env()), clipboard_watch_mode: tokio::sync::RwLock::new(false), + clipboard_notify: Arc::new(tokio::sync::Notify::new()), + index_commit_notify: Arc::new(tokio::sync::Notify::new()), + ttl_notify: Arc::new(tokio::sync::Notify::new()), + condense_notify: Arc::new(tokio::sync::Notify::new()), graph: Store::new("knowledge_graph_master", db.clone()), base_dir: base.clone(), search_index: RwLock::new(match crate::search::MemoryIndex::new(&base) { @@ -207,6 +215,7 @@ impl MemoryState { if let Ok(mut w) = self.search_index.write() { *w = idx; } + self.index_commit_notify.notify_waiters(); } pub fn record_activity(&self, category: &str, summary: &str, details: Option<&str>) { @@ -243,6 +252,49 @@ impl MemoryState { let _ = self.activity_tx.send(payload); } + pub fn broadcast_task_event(&self, event: TaskEvent) { + let payload_val = serde_json::to_value(&event).unwrap_or_default(); + + let summary_str = format!("Task {} -> {}", event.task_id, event.status); + self.telemetry.recent_activities.modify(|activities| { + let activity = ActivityRecord { + timestamp: event.timestamp, + category: "TASK_EVENT".to_string(), + summary: summary_str, + details: Some(payload_val.to_string()), + }; + activities.push_front(serde_json::to_value(&activity).unwrap_or_default()); + if activities.len() > 100 { + activities.pop_back(); + } + }); + + let generic_ev = GenericEvent { + topic: "task:event".to_string(), + session_id: event.session_id.clone(), + payload: payload_val, + }; + let _ = self.event_bus_tx.send(generic_ev); + + let ws_notification = serde_json::json!({ + "jsonrpc": "2.0", + "method": "notifications/task/completed", + "params": event + }) + .to_string(); + let _ = self.activity_tx.send(ws_notification); + + let ws_resource_notification = serde_json::json!({ + "jsonrpc": "2.0", + "method": "notifications/resources/updated", + "params": { + "uri": "memory://tasks/active" + } + }) + .to_string(); + let _ = self.activity_tx.send(ws_resource_notification); + } + pub fn record_terminal_history(&self, payload: TerminalHistory) { self.telemetry.terminal_history.modify(|history| { history.push_front(payload); diff --git a/server/src/store.rs b/server/src/store.rs index bd73d34..4ef9cbd 100644 --- a/server/src/store.rs +++ b/server/src/store.rs @@ -6,6 +6,7 @@ pub const STORE_TABLE: TableDefinition<&str, &[u8]> = TableDefinition::new("stor pub struct Store { pub cache: Arc>, + pub flushed: Arc, tx: tokio::sync::mpsc::Sender<()>, } @@ -13,21 +14,22 @@ impl pub fn new(key: &str, db: Arc) -> Self { let initial_data = Self::load_from_db(key, &db); let cache = Arc::new(RwLock::new(initial_data)); + let flushed = Arc::new(tokio::sync::Notify::new()); let (tx, mut rx) = tokio::sync::mpsc::channel::<()>(1); let db_clone = db.clone(); let key_clone = key.to_string(); let cache_clone = cache.clone(); + let flushed_clone = flushed.clone(); tokio::spawn(async move { while rx.recv().await.is_some() { - // Debounce window: wait 150ms to batch rapid sequential mutations - tokio::time::sleep(tokio::time::Duration::from_millis(150)).await; - // Drain any pending notifications accumulated during the debounce window + // Drain any pending notifications accumulated while rx.try_recv().is_ok() {} let db_inner = db_clone.clone(); let key_inner = key_clone.clone(); + let flushed_inner = flushed_clone.clone(); let json_data = { let lock = cache_clone.read().unwrap_or_else(|e| e.into_inner()); serde_json::to_vec(&*lock) @@ -37,7 +39,7 @@ impl if let Some(json_data) = json_data { let _ = tokio::task::spawn_blocking(move || { - // Retry up to 10 times with 30ms backoff if another Store holds write transaction + // Retry up to 10 times if another Store holds write transaction for _ in 0..10 { match db_inner.begin_write() { Ok(write_txn) => { @@ -48,17 +50,18 @@ impl break; } Err(_) => { - std::thread::sleep(std::time::Duration::from_millis(30)); + std::thread::yield_now(); } } } }) .await; + flushed_inner.notify_waiters(); } } }); - Self { cache, tx } + Self { cache, flushed, tx } } fn load_from_db(key: &str, db: &Database) -> T { @@ -123,18 +126,10 @@ mod tests { data.value = 42; }); - // Wait and poll for persistence completion - let mut store2 = None; - for _ in 0..20 { - let s = Store::::new("test_key", db.clone()); - if s.read_with(|data| data.value) == 42 { - store2 = Some(s); - break; - } - tokio::time::sleep(tokio::time::Duration::from_millis(50)).await; - } + // Event-driven wait for persistence completion + store.flushed.notified().await; - let store2 = store2.expect("Timed out waiting for async store persistence"); + let store2 = Store::::new("test_key", db.clone()); assert_eq!( store2.read_with(|s| s.clone()), TestData { @@ -172,8 +167,8 @@ mod tests { h.await.unwrap(); } - // Wait for all blocking writes to flush - tokio::time::sleep(tokio::time::Duration::from_millis(500)).await; + // Event-driven wait for blocking writes to flush + store.flushed.notified().await; assert_eq!(store.read_with(|s| s.value), 50); } diff --git a/server/src/tools.rs b/server/src/tools.rs index 1f7895d..62499bf 100644 --- a/server/src/tools.rs +++ b/server/src/tools.rs @@ -1157,6 +1157,8 @@ pub struct TasksTool { pub git_branch: Option, /// Acceptance criteria (required for 'set_criteria'). pub criteria: Option>, + /// Verification proof or details (optional for 'verify'). + pub proof: Option, /// Summary level: 'compact', 'detailed', or 'full' (for 'list'). pub summary_level: Option, /// Maximum tokens budget cap (for 'list'). diff --git a/server/src/watcher.rs b/server/src/watcher.rs index 7cd2eb0..34e0b56 100644 --- a/server/src/watcher.rs +++ b/server/src/watcher.rs @@ -13,7 +13,7 @@ pub fn spawn_watcher(_state: Arc) { let mut watcher = match RecommendedWatcher::new( move |res| { - let _ = tx.blocking_send(res); + let _ = tx.try_send(res); }, Config::default(), ) { diff --git a/stub/src/main.rs b/stub/src/main.rs index ff4ae95..ef32b5b 100644 --- a/stub/src/main.rs +++ b/stub/src/main.rs @@ -139,13 +139,11 @@ fn main() -> Result<(), Box> { tokio::select! { _ = shutdown_rx.recv() => { - tokio::time::sleep(tokio::time::Duration::from_millis(50)).await; tracing::info!("Shutdown received while connected"); break; } _ = &mut send_task => { tracing::error!("Send task exited"); - tokio::time::sleep(tokio::time::Duration::from_millis(50)).await; recv_task.abort(); } _ = &mut recv_task => {