Files
mcp-memory/server/src/handlers/meta.rs
T
Riz Ashraf 8f924b793a feat(reconciliation): implement deterministic state reconciliation engine and post-commit hook
- Implement ReconciliationEngine in server/src/handlers/reconciliation.rs
- Auto-transition ADRs to implemented, resolve tech debts, and cascade task unblocking
- Ingest git commit events via POST /api/git/commit and scripts/git-reconcile.py post-commit hook
- Add universal pagination, search filtering, and keyboard navigation to dashboard
- Implement non-destructive task TTL expiry sweeper and gate verification
- Implements: ADR-0102, ADR-0103
2026-10-07 19:00:44 +01:00

2942 lines
107 KiB
Rust

use crate::models::*;
use crate::router::McpTool;
use crate::state::MemoryState;
use crate::tools::*;
use async_trait::async_trait;
use serde_json::Value;
use std::sync::Arc;
pub struct LogErrorFixHandler;
#[async_trait]
impl McpTool for LogErrorFixHandler {
fn name(&self) -> &'static str {
"log_error_fix"
}
fn schema(&self) -> Value {
crate::mcp::tool_def::<LogErrorFixTool>(
"log_error_fix",
"Log an error signature and its verified solution/fix for future diagnostic retrieval",
)
}
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: LogErrorFixTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let text_to_embed = format!("Signature: {}\nSolution: {}", req.signature, req.solution);
let mut solution = req.solution;
if state.ollama.is_available().await {
let prompt = format!(
"Analyze this error signature and solution. Output 1 sentence summarizing the root cause and fix:\nSignature: {}\nSolution: {}",
req.signature, solution
);
if let Ok(summary) = state
.ollama
.generate(&prompt, Some(&state.ollama.reasoning_model), None)
.await
{
let clean = summary.trim();
if !clean.is_empty() {
solution = format!("{} (AI Analysis: {})", solution, clean);
}
}
}
let embedding = crate::embedding::generate_embedding_async(text_to_embed)
.await
.ok();
state.code.error_fixes.modify(|fixes| {
fixes.push(crate::models::ErrorFix {
signature: req.signature.clone(),
solution: solution.clone(),
timestamp: crate::handlers::utils::now_secs(),
git_commit: req.git_commit,
git_branch: req.git_branch,
embedding,
..Default::default()
});
if fixes.len() > 300 {
fixes.remove(0);
}
});
state.record_activity(
"error_fix",
&format!("Fixed error: {}", req.signature),
Some(&solution),
);
Ok(format!(
"Logged error fix for {}: {}",
req.signature, solution
))
}
}
pub struct SearchErrorFixesHandler;
#[async_trait]
impl McpTool for SearchErrorFixesHandler {
fn name(&self) -> &'static str {
"search_error_fixes"
}
fn schema(&self) -> Value {
crate::mcp::tool_def::<SearchErrorFixesTool>(
"search_error_fixes",
"Search historical error fixes using keyword search or stack trace vector similarity",
)
}
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: SearchErrorFixesTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let limit = req.limit.unwrap_or(5);
let include_body = req.include_body.unwrap_or(true);
if let Some(st) = &req.stack_trace {
let query_emb = crate::embedding::generate_embedding_async(st.clone())
.await
.unwrap_or_default();
let data = state.code.error_fixes.read_with(|fixes| {
let mut scored: Vec<_> = fixes
.iter()
.map(|f| {
let mut score = 0.0;
let st_lower = st.to_lowercase();
let sig_lower = f.signature.to_lowercase();
let sol_lower = f.solution.to_lowercase();
if st_lower.contains(&sig_lower) || sig_lower.contains(&st_lower) {
score += 0.8;
} else if st_lower.contains(&sol_lower) || sol_lower.contains(&st_lower) {
score += 0.5;
}
if let Some(emb) = &f.embedding {
if !query_emb.is_empty() {
score += crate::embedding::cosine_similarity(&query_emb, emb);
}
}
(f, score)
})
.filter(|(_, score)| *score > 0.1)
.collect();
scored.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
let suggestions: Vec<_> = scored
.into_iter()
.take(limit)
.map(|(f, score)| {
if include_body {
serde_json::json!({
"signature": f.signature,
"solution": f.solution,
"git_commit": f.git_commit,
"git_branch": f.git_branch,
"match_score": score
})
} else {
serde_json::json!({
"signature": f.signature,
"match_score": score
})
}
})
.collect();
Ok::<String, crate::error::AppError>(serde_json::to_string_pretty(&suggestions)?)
})?;
return Ok(data);
}
let q = req.query.unwrap_or_default();
let data = state.code.error_fixes.read_with(|fixes| {
let filtered: Vec<_> = fixes
.iter()
.filter(|f| {
q.is_empty()
|| contains_ignore_ascii_case(&f.signature, &q)
|| contains_ignore_ascii_case(&f.solution, &q)
})
.take(limit)
.map(|f| {
if include_body {
serde_json::json!(f)
} else {
serde_json::json!({
"signature": f.signature,
"git_commit": f.git_commit,
"git_branch": f.git_branch
})
}
})
.collect();
Ok::<String, crate::error::AppError>(serde_json::to_string_pretty(&filtered)?)
})?;
Ok(data)
}
}
pub struct LogCodeChangeHandler;
#[async_trait]
impl McpTool for LogCodeChangeHandler {
fn name(&self) -> &'static str {
"log_code_change"
}
fn schema(&self) -> Value {
crate::mcp::tool_def::<LogCodeChangeTool>(
"log_code_change",
"Log a significant code change or refactor with file path and description",
)
}
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: LogCodeChangeTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let mut description = req.description;
if let Some(range) = &req.line_range {
description = format!("{} [Line Range: {}]", description, range);
}
if let Some(symbols) = &req.symbol_references {
if !symbols.is_empty() {
description = format!("{} [Symbols: {}]", description, symbols.join(", "));
}
}
if state.ollama.is_available().await {
let prompt = format!(
"Summarize in 1 concise sentence the architectural impact of changing file '{}': {}",
req.file_path, description
);
if let Ok(summary) = state.ollama.generate(&prompt, None, None).await {
let clean = summary.trim();
if !clean.is_empty() {
description = format!("{} (AI Summary: {})", description, clean);
}
}
}
let change_kind = match req
.change_kind
.as_deref()
.map(|s| s.to_lowercase())
.as_deref()
{
Some("added") | Some("create") | Some("created") | Some("new") => {
crate::models::ChangeKind::Added
}
Some("deleted") | Some("remove") | Some("removed") => {
crate::models::ChangeKind::Deleted
}
Some("renamed") | Some("move") | Some("moved") => crate::models::ChangeKind::Renamed,
_ => crate::models::ChangeKind::Modified,
};
let namespace = req
.namespace
.filter(|ns| !ns.trim().is_empty())
.or_else(|| req.repo_name.clone().filter(|rn| !rn.trim().is_empty()))
.unwrap_or_else(crate::models::default_namespace);
let symbols = req.symbol_references.clone().unwrap_or_default();
let line_range = req.line_range.clone();
let effective_rev = req.revision.clone().or_else(|| req.git_commit.clone());
let effective_branch = req.branch.clone().or_else(|| req.git_branch.clone());
let detected_vcs = if let Some(vcs) = req.vcs_type.clone() {
Some(vcs)
} else if let Some(ref rev) = effective_rev {
if rev.starts_with('r') && rev[1..].chars().all(|c| c.is_ascii_digit()) {
Some("svn".to_string())
} else if req
.repo_url
.as_deref()
.map(|u| u.contains("/svn/"))
.unwrap_or(false)
{
Some("svn".to_string())
} else {
Some("git".to_string())
}
} else if let Some(ref br) = effective_branch {
if br.eq_ignore_ascii_case("trunk")
|| br.starts_with("branches/")
|| br.starts_with("tags/")
{
Some("svn".to_string())
} else {
Some("git".to_string())
}
} else {
None
};
state.code.ledger.modify(|ledger| {
ledger.push(CodeChange {
timestamp: crate::handlers::utils::now_secs(),
file_path: req.file_path.clone(),
description: description.clone(),
git_commit: effective_rev.clone(),
git_branch: effective_branch.clone(),
repo_name: req.repo_name,
repo_url: req.repo_url,
namespace: namespace.clone(),
change_kind,
symbols,
line_range,
author: req.author,
session_id: req.session_id,
vcs_type: detected_vcs,
revision: effective_rev.clone(),
branch: effective_branch.clone(),
repository_root: req.repository_root,
});
if ledger.len() > 500 {
ledger.remove(0);
}
});
state.record_activity(
"code_change",
&format!("Modified {}", req.file_path),
Some(&description),
);
let recon = crate::handlers::reconciliation::reconcile_commit_or_code_change(
&state,
&description,
Some(&req.file_path),
effective_rev.as_deref(),
effective_branch.as_deref(),
)
.await;
let mut recon_notes = Vec::new();
if !recon.implemented_adrs.is_empty() {
recon_notes.push(format!("Implemented ADRs: {}", recon.implemented_adrs.join(", ")));
}
if !recon.resolved_tech_debts.is_empty() {
recon_notes.push(format!("Resolved TechDebt: {}", recon.resolved_tech_debts.join(", ")));
}
if !recon.completed_tasks.is_empty() {
recon_notes.push(format!("Completed Tasks: {}", recon.completed_tasks.join(", ")));
}
if !recon.unblocked_tasks.is_empty() {
recon_notes.push(format!("Unblocked Tasks: {}", recon.unblocked_tasks.join(", ")));
}
if !recon.updated_milestones.is_empty() {
recon_notes.push(format!("Updated Milestones: {}", recon.updated_milestones.join(", ")));
}
let recon_suffix = if recon_notes.is_empty() {
String::new()
} else {
format!(" [{}]", recon_notes.join(" | "))
};
Ok(format!(
"Logged code change for {}: {}{}",
req.file_path, description, recon_suffix
))
}
}
pub struct QueryRecentChangesHandler;
#[async_trait]
impl McpTool for QueryRecentChangesHandler {
fn name(&self) -> &'static str {
"query_recent_changes"
}
fn schema(&self) -> Value {
crate::mcp::tool_def::<QueryRecentChangesTool>(
"query_recent_changes",
"Query recent code changes and refactoring audit logs",
)
}
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: QueryRecentChangesTool = serde_json::from_value(args).unwrap_or(QueryRecentChangesTool {
namespace: None,
repo_name: None,
vcs_type: None,
limit: None,
offset: None,
});
let limit = req.limit.unwrap_or(50);
let offset = req.offset.unwrap_or(0);
let data = state.code.ledger.read_with(|l| {
let filtered: Vec<_> = l
.iter()
.rev()
.filter(|c| {
if let Some(ns) = &req.namespace {
if !c.namespace.eq_ignore_ascii_case(ns) {
return false;
}
}
if let Some(repo) = &req.repo_name {
if c.repo_name
.as_ref()
.map(|rn| !rn.eq_ignore_ascii_case(repo))
.unwrap_or(true)
{
return false;
}
}
if let Some(vcs) = &req.vcs_type {
if !c.effective_vcs().eq_ignore_ascii_case(vcs) {
return false;
}
}
true
})
.skip(offset)
.take(limit)
.cloned()
.collect();
Ok::<String, crate::error::AppError>(serde_json::to_string(&filtered)?)
})?;
Ok(data)
}
}
pub struct DecisionsHandler;
#[async_trait]
impl McpTool for DecisionsHandler {
fn name(&self) -> &'static str {
"decisions"
}
fn schema(&self) -> Value {
crate::mcp::tool_def::<DecisionsTool>(
"decisions",
"Consolidated Architectural Decision Records (ADRs) management (log, query, update, delete)",
)
}
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: DecisionsTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let ns = req
.namespace
.unwrap_or_else(|| crate::models::default_namespace());
match req.action {
DecisionAction::Log => {
let title = req.title.ok_or_else(|| {
crate::error::AppError::Internal("Missing required parameter 'title' for action 'log'. Next step: Provide ADR 'title' string in request and retry.".to_string())
})?;
let status = req.status.unwrap_or_else(|| "accepted".to_string());
let context = req.context.unwrap_or_default();
let decision = req.decision.unwrap_or_default();
let consequence = req.consequences.unwrap_or_default();
let status_lower = status.to_ascii_lowercase();
let resolved_at = if status_lower == "implemented" || status_lower == "resolved" {
Some(crate::handlers::utils::now_secs())
} else {
None
};
let idx = state.get_search_index().await;
let mut final_id = String::new();
let mut adrs_to_index = Vec::new();
state.code.adrs.modify(|adrs| {
if let Some(superseded_id) = &req.supersedes {
for old_adr in adrs.iter_mut() {
if old_adr.id.eq_ignore_ascii_case(superseded_id) {
old_adr.status = "superseded".to_string();
adrs_to_index.push(old_adr.clone());
break;
}
}
}
final_id = format!("ADR-{:04}", adrs.len() + 1);
let a = Adr {
id: final_id.clone(),
title: title.clone(),
context,
decision: decision.clone(),
consequence,
status,
supersedes: req.supersedes,
timestamp: crate::handlers::utils::now_secs(),
namespace: ns,
repo_name: req.repo_name,
alternatives_considered: req.alternatives_considered.unwrap_or_default(),
affected_components: req.affected_components.unwrap_or_default(),
author: req.author,
git_commit: req.git_commit,
git_branch: req.git_branch,
resolved_at,
task_id: req.task_id,
};
adrs_to_index.push(a.clone());
adrs.push(a);
});
for adr in &adrs_to_index {
drop(idx.index_adr(adr));
}
state.record_activity(
"decision",
&format!("Logged {}: {}", final_id, title),
Some(&decision),
);
Ok(format!("Logged decision {}: {}", final_id, title))
}
DecisionAction::Update => {
let id = req.id.ok_or_else(|| {
crate::error::AppError::Internal("Missing required parameter 'id' for action 'update'. Next step: Provide ADR 'id' string in request and retry.".to_string())
})?;
let mut updated_adr = None;
let mut adrs_to_index = Vec::new();
state.code.adrs.modify(|adrs| {
let target_pos = adrs.iter().position(|a| a.id.eq_ignore_ascii_case(&id));
if let Some(pos) = target_pos {
if let Some(superseded_id) = &req.supersedes {
if let Some(s_pos) = adrs.iter().position(|a| a.id.eq_ignore_ascii_case(superseded_id)) {
if s_pos != pos {
adrs[s_pos].status = "superseded".to_string();
adrs_to_index.push(adrs[s_pos].clone());
}
}
}
let a = &mut adrs[pos];
if let Some(t) = req.title {
a.title = t;
}
if let Some(c) = req.context {
a.context = c;
}
if let Some(d) = req.decision {
a.decision = d;
}
if let Some(cons) = req.consequences {
a.consequence = cons;
}
if let Some(s) = req.status {
let s_lower = s.to_ascii_lowercase();
if (s_lower == "implemented" || s_lower == "resolved") && a.resolved_at.is_none() {
a.resolved_at = Some(crate::handlers::utils::now_secs());
} else if s_lower != "implemented" && s_lower != "resolved" {
a.resolved_at = None;
}
a.status = s;
}
if req.supersedes.is_some() {
a.supersedes = req.supersedes;
}
if req.repo_name.is_some() {
a.repo_name = req.repo_name;
}
if let Some(alts) = req.alternatives_considered {
a.alternatives_considered = alts;
}
if let Some(aff) = req.affected_components {
a.affected_components = aff;
}
if req.author.is_some() {
a.author = req.author;
}
if req.git_commit.is_some() {
a.git_commit = req.git_commit;
}
if req.git_branch.is_some() {
a.git_branch = req.git_branch;
}
if req.task_id.is_some() {
a.task_id = req.task_id;
}
adrs_to_index.push(a.clone());
updated_adr = Some(a.clone());
}
});
if let Some(adr) = updated_adr {
let idx = state.get_search_index().await;
for a in &adrs_to_index {
drop(idx.index_adr(a));
}
state.record_activity(
"decision",
&format!("Updated {}: {}", adr.id, adr.title),
Some(&adr.status),
);
Ok(format!("Updated decision {}: {} (status: {})", adr.id, adr.title, adr.status))
} else {
Err(crate::error::AppError::Internal(
format!("Decision with id '{}' not found", id),
))
}
}
DecisionAction::Query => {
let limit = req.limit.unwrap_or(20);
let include_body = req.include_body.unwrap_or(true);
let data = state.code.adrs.read_with(|adrs| {
let filtered: Vec<_> = adrs
.iter()
.filter(|a| {
if !a.namespace.eq_ignore_ascii_case(&ns) && ns != "global" {
return false;
}
if let Some(q) = &req.query {
crate::handlers::utils::contains_ignore_ascii_case(&a.title, q)
|| crate::handlers::utils::contains_ignore_ascii_case(&a.context, q)
|| crate::handlers::utils::contains_ignore_ascii_case(&a.decision, q)
|| crate::handlers::utils::contains_ignore_ascii_case(&a.consequence, q)
} else {
true
}
})
.take(limit)
.collect();
if include_body {
Ok::<String, crate::error::AppError>(serde_json::to_string(&filtered)?)
} else {
let compact: Vec<_> = filtered
.iter()
.map(|a| {
serde_json::json!({
"id": a.id,
"title": a.title,
"status": a.status,
"timestamp": a.timestamp,
"git_commit": a.git_commit,
"git_branch": a.git_branch,
"resolved_at": a.resolved_at,
"task_id": a.task_id,
})
})
.collect();
Ok::<String, crate::error::AppError>(serde_json::to_string(&compact)?)
}
})?;
Ok(data)
}
DecisionAction::Delete => {
let id = req.id.ok_or_else(|| {
crate::error::AppError::Internal("Missing required parameter 'id' for action 'delete'. Next step: Provide ADR 'id' string in request and retry.".to_string())
})?;
let mut found = false;
state.code.adrs.modify(|adrs| {
if let Some(pos) = adrs.iter().position(|a| a.id == id) {
adrs.remove(pos);
found = true;
}
});
if found {
let idx = state.get_search_index().await;
let _ = idx.delete_document(&id).await;
Ok("Decision deleted successfully".to_string())
} else {
Err(crate::error::AppError::Internal(
"Decision not found".to_string(),
))
}
}
}
}
}
pub struct TechDebtHandler;
#[async_trait]
impl McpTool for TechDebtHandler {
fn name(&self) -> &'static str {
"tech_debt"
}
fn schema(&self) -> Value {
crate::mcp::tool_def::<TechDebtTool>(
"tech_debt",
"Consolidated technical debt management (log, resolve, list)",
)
}
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: TechDebtTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let ns = req
.namespace
.unwrap_or_else(|| crate::models::default_namespace());
match req.action {
TechDebtAction::Log => {
let desc = req.description.or(req.title).ok_or_else(|| {
crate::error::AppError::Internal("Missing required parameter 'description' for action 'log'. Next step: Provide tech debt 'description' in request and retry.".to_string())
})?;
let ideal = req.ideal_solution.unwrap_or_default();
let text_to_embed = format!(
"Description: {}\nIdeal Solution: {}",
desc, ideal
);
let embedding = crate::embedding::generate_embedding_async(text_to_embed)
.await
.ok();
state.code.tech_debts.modify(|debts| {
debts.push(crate::models::TechDebt {
id: uuid::Uuid::new_v4().to_string(),
namespace: ns,
description: desc,
ideal_solution: ideal,
is_resolved: false,
created_at: crate::handlers::utils::now_secs(),
git_commit: req.git_commit,
git_branch: req.git_branch,
embedding,
repo_name: req.repo_name,
severity: req.severity,
file_path: req.file_path,
line_range: req.line_range,
workaround: req.workaround,
effort_estimate: req.effort_estimate,
});
if debts.len() > 300 {
let severity_rank = |sev: Option<&str>| match sev.unwrap_or("").to_lowercase().as_str() {
"critical" => 4,
"high" => 3,
"medium" => 2,
"low" => 1,
_ => 1,
};
if let Some((idx_to_remove, _)) = debts.iter().enumerate().min_by_key(|(_, d)| {
let status_score = if d.is_resolved { 0 } else { 10 };
let sev_score = severity_rank(d.severity.as_deref());
(status_score + sev_score, d.created_at)
}) {
debts.remove(idx_to_remove);
}
}
});
Ok("Tech debt logged".to_string())
}
TechDebtAction::Resolve => {
let id = req.id.ok_or_else(|| {
crate::error::AppError::Internal("Missing required parameter 'id' for action 'resolve'. Next step: Provide tech debt 'id' string in request and retry.".to_string())
})?;
let mut found = false;
state.code.tech_debts.modify(|debts| {
for d in debts.iter_mut() {
if d.id == id {
d.is_resolved = true;
found = true;
break;
}
}
});
if found {
Ok("Tech debt resolved".to_string())
} else {
Err(crate::error::AppError::Internal(
"Tech debt not found. Please verify the tech debt ID using list action."
.to_string(),
))
}
}
TechDebtAction::List => {
let inc = req.include_resolved.unwrap_or(false);
let level = req.summary_level.as_deref().unwrap_or("detailed");
let data = state.code.tech_debts.read_with(|debts| {
let filtered: Vec<_> = debts
.iter()
.filter(|d| {
d.namespace == ns && (inc || !d.is_resolved)
})
.map(|d| match level {
"compact" => serde_json::json!({
"id": d.id,
"description": d.description,
"is_resolved": d.is_resolved,
}),
"full" => serde_json::to_value(d).unwrap_or_default(),
_ => serde_json::json!({
"id": d.id,
"description": d.description,
"ideal_solution": d.ideal_solution,
"is_resolved": d.is_resolved,
}),
})
.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::<String, crate::error::AppError>(json_str)
})?;
Ok(data)
}
}
}
}
pub struct OmniSearchHandler;
#[async_trait]
impl McpTool for OmniSearchHandler {
fn name(&self) -> &'static str {
"omni_search"
}
fn schema(&self) -> Value {
crate::mcp::tool_def::<OmniSearchTool>(
"omni_search",
"Unified search across entities, subgraphs, tasks, code snippets, ADRs, and technical debt",
)
}
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: OmniSearchTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let limit = req.limit.unwrap_or(5);
let include_body = req.include_body.unwrap_or(false);
let idx = state.get_search_index().await;
let keyword_matches = idx
.search(&req.query, req.namespace.as_deref())
.unwrap_or_default();
let vector_matches = state
.search()
.semantic_search(&req.query, req.namespace.as_deref(), limit * 2)
.await
.unwrap_or_default();
// Reciprocal Rank Fusion (RRF) algorithm
#[allow(dead_code)]
#[derive(Clone)]
struct MatchItem {
id: String,
doc_type: String,
title: String,
body: String,
}
let mut rrf_scores: std::collections::HashMap<String, (f64, MatchItem)> =
std::collections::HashMap::new();
for (rank, (id, doc_type, title, body, _score)) in keyword_matches.into_iter().enumerate() {
let score = 1.0 / (60.0 + (rank + 1) as f64);
rrf_scores.insert(
id.clone(),
(
score,
MatchItem {
id,
doc_type,
title,
body,
},
),
);
}
for (rank, v_match) in vector_matches.into_iter().enumerate() {
let score = 1.0 / (60.0 + (rank + 1) as f64);
let item_id = v_match.id.clone();
if let Some(existing) = rrf_scores.get_mut(&item_id) {
existing.0 += score;
} else {
let item = MatchItem {
id: v_match.id.clone(),
doc_type: v_match.doc_type,
title: v_match.title,
body: v_match.body,
};
rrf_scores.insert(item_id, (score, item));
}
}
let mut ranked_items: Vec<_> = rrf_scores.into_values().collect();
ranked_items.sort_by(|a, b| b.0.total_cmp(&a.0));
let matches: Vec<MatchItem> = ranked_items.into_iter().map(|(_, item)| item).collect();
let kg_json = state.read_graph(|full| {
let mut kg_results = serde_json::Map::new();
let mut count = 0;
// Build pre-indexed adjacency map: O(R) once instead of O(E * R)
let mut adj_map: std::collections::HashMap<&str, Vec<(&str, &str, &str)>> =
std::collections::HashMap::new();
for rel in &full.relations {
adj_map.entry(rel.from.as_str()).or_default().push((
rel.to.as_str(),
rel.relation_type.as_str(),
"outgoing",
));
adj_map.entry(rel.to.as_str()).or_default().push((
rel.from.as_str(),
rel.relation_type.as_str(),
"incoming",
));
}
for res in &matches {
if res.doc_type == "entity"
&& let Some(e) = full.entities.get(&res.id)
{
if count >= limit {
break;
}
count += 1;
// 1-hop relation expansion for GraphRAG via pre-indexed adjacency
let mut connected_rels = Vec::new();
if let Some(rels) = adj_map.get(res.id.as_str()) {
for (target, rel_type, direction) in rels {
connected_rels.push(serde_json::json!({
"target": target,
"relation": rel_type,
"direction": direction
}));
}
}
let mut entity_val = if !include_body {
let mut summary = e.clone();
summary.observations = vec![];
serde_json::to_value(&summary).unwrap_or_default()
} else {
serde_json::to_value(e).unwrap_or_default()
};
if let Some(obj) = entity_val.as_object_mut() {
obj.insert(
"subgraph_relations".to_string(),
serde_json::Value::Array(connected_rels),
);
}
kg_results.insert(res.id.clone(), entity_val);
}
}
serde_json::to_value(&kg_results).map_err(|e| e.to_string())
})?;
let mut matched_tasks = std::collections::HashSet::new();
let mut matched_snippets = std::collections::HashSet::new();
let mut matched_adrs = std::collections::HashSet::new();
for res in &matches {
match res.doc_type.as_str() {
"task" => {
matched_tasks.insert(res.id.as_str());
}
"snippet" => {
matched_snippets.insert(res.id.as_str());
}
"adr" => {
matched_adrs.insert(res.id.as_str());
}
_ => {}
}
}
let tasks_json = state.project.tasks.read_with(|all_tasks| {
let filtered: Vec<_> = all_tasks
.iter()
.filter(|t| matched_tasks.contains(t.id.as_str()))
.take(limit)
.map(|t| {
if !include_body {
let mut summary = t.clone();
summary.description = "".to_string();
summary.acceptance_criteria = vec![];
summary
} else {
t.clone()
}
})
.collect();
serde_json::to_value(&filtered).map_err(|e| e.to_string())
})?;
let q = req.query.clone();
let query_emb = crate::embedding::generate_embedding_async(req.query.clone())
.await
.unwrap_or_default();
let snippets_json = state.code.snippets.read_with(|all_snippets| {
let mut scored: Vec<_> = all_snippets
.iter()
.map(|s| {
let mut score = 0.0;
if matched_snippets.contains(s.name.as_str())
|| contains_ignore_ascii_case(&s.name, &q)
|| contains_ignore_ascii_case(&s.description, &q)
{
score += 1.0;
}
if let Some(emb) = &s.embedding {
score += crate::embedding::cosine_similarity(&query_emb, emb);
}
(s, score)
})
.filter(|(_, score)| *score > 0.4)
.collect();
scored.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
let filtered: Vec<_> = scored
.into_iter()
.take(limit)
.map(|(s, _)| {
if !include_body {
let mut summary = s.clone();
summary.code = "".to_string();
summary
} else {
s.clone()
}
})
.collect();
serde_json::to_value(&filtered).map_err(|e| e.to_string())
})?;
let adrs_json = state.code.adrs.read_with(|all_adrs| {
let filtered: Vec<_> = all_adrs
.iter()
.filter(|a| matched_adrs.contains(a.id.as_str()))
.take(limit)
.map(|a| {
if !include_body {
let mut summary = a.clone();
summary.context = "".to_string();
summary.decision = "".to_string();
summary.consequence = "".to_string();
summary
} else {
a.clone()
}
})
.collect();
serde_json::to_value(&filtered).map_err(|e| e.to_string())
})?;
let tech_debts_json = state.code.tech_debts.read_with(|debts| {
let mut scored: Vec<_> = debts
.iter()
.map(|d| {
if req.namespace.as_ref().is_some_and(|ns| d.namespace != *ns) {
return (d, 0.0);
}
let mut score = 0.0;
if contains_ignore_ascii_case(&d.description, &q)
|| contains_ignore_ascii_case(&d.ideal_solution, &q)
{
score += 1.0;
}
if let Some(emb) = &d.embedding {
score += crate::embedding::cosine_similarity(&query_emb, emb);
}
(d, score)
})
.filter(|(_, score)| *score > 0.4)
.collect();
scored.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
let filtered: Vec<_> = scored
.into_iter()
.take(limit)
.map(|(d, _)| {
if !include_body {
let mut summary = d.clone();
summary.description = "".to_string();
summary.ideal_solution = "".to_string();
summary
} else {
d.clone()
}
})
.collect();
serde_json::to_value(&filtered).map_err(|e| e.to_string())
})?;
let memos_json = state.telemetry.handoff_memos.read_with(|memos| {
let filtered: Vec<_> = memos
.iter()
.filter(|m| {
req.namespace.as_ref().is_none_or(|ns| m.namespace == *ns)
&& contains_ignore_ascii_case(&m.content, &q)
})
.take(limit)
.map(|m| {
if !include_body {
let mut summary = m.clone();
summary.content = "".to_string();
summary
} else {
m.clone()
}
})
.collect();
serde_json::to_value(&filtered).map_err(|e| e.to_string())
})?;
let error_fixes_json = state.code.error_fixes.read_with(|fixes| {
let mut scored: Vec<_> = fixes
.iter()
.map(|f| {
let mut score = 0.0;
if contains_ignore_ascii_case(&f.signature, &q)
|| contains_ignore_ascii_case(&f.solution, &q)
{
score += 1.0;
}
if let Some(emb) = &f.embedding {
score += crate::embedding::cosine_similarity(&query_emb, emb);
}
(f, score)
})
.filter(|(_, score)| *score > 0.4)
.collect();
scored.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
let filtered: Vec<_> = scored
.into_iter()
.take(limit)
.map(|(f, _)| f.clone())
.collect();
serde_json::to_value(&filtered).map_err(|e| e.to_string())
})?;
let mut report = serde_json::json!({
"knowledge_graph": kg_json,
"tasks": tasks_json,
"snippets": snippets_json,
"adrs": adrs_json,
"tech_debts": tech_debts_json,
"handoff_memos": memos_json,
"error_fixes": error_fixes_json
});
if let Some(max_tok) = req.max_tokens {
let max_chars = max_tok * 4;
let mut out_str = report.to_string();
if out_str.len() > max_chars {
let prune_keys = [
"error_fixes",
"tech_debts",
"snippets",
"adrs",
"handoff_memos",
"knowledge_graph",
"tasks",
];
let mut pruned = false;
for key in prune_keys {
while out_str.len() > max_chars {
let popped =
if let Some(arr) = report.get_mut(key).and_then(|v| v.as_array_mut()) {
if arr.len() > 1 {
arr.pop();
pruned = true;
true
} else {
false
}
} else {
false
};
if popped {
out_str = report.to_string();
} else {
break;
}
}
if out_str.len() <= max_chars {
break;
}
}
if pruned && let Some(obj) = report.as_object_mut() {
obj.insert(
"_truncated_to_max_tokens".to_string(),
serde_json::Value::Bool(true),
);
}
}
}
Ok(report.to_string())
}
}
pub struct GetProjectHealthHandler;
#[async_trait]
impl McpTool for GetProjectHealthHandler {
fn name(&self) -> &'static str {
"get_project_health"
}
fn schema(&self) -> Value {
crate::mcp::tool_def::<GetProjectHealthTool>(
"get_project_health",
"Retrieve project health metrics including active tasks, technical debt, and PR checklist progress",
)
}
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: GetProjectHealthTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let active_tasks = state.project.tasks.read_with(|tasks| {
tasks
.iter()
.filter(|t| t.namespace == req.namespace && t.is_active())
.count()
});
let unresolved_debt = state.code.tech_debts.read_with(|debts| {
debts
.iter()
.filter(|d| d.namespace == req.namespace && !d.is_resolved)
.count()
});
let unread_memos = state.telemetry.handoff_memos.read_with(|memos| {
memos
.iter()
.filter(|m| m.namespace == req.namespace)
.count()
});
let active_milestones = state.project.milestones.read_with(|milestones| {
milestones
.iter()
.filter(|m| {
m.namespace == req.namespace
&& !m.status.eq_ignore_ascii_case("done")
&& !m.status.eq_ignore_ascii_case("completed")
})
.count()
});
let report = serde_json::json!({
"active_tasks": active_tasks,
"unresolved_tech_debt": unresolved_debt,
"unread_handoff_memos": unread_memos,
"active_milestones": active_milestones
});
Ok(report.to_string())
}
}
pub struct ManageCheckpointHandler;
#[async_trait]
impl McpTool for ManageCheckpointHandler {
fn name(&self) -> &'static str {
"manage_checkpoint"
}
fn schema(&self) -> Value {
crate::mcp::tool_def::<ManageCheckpointTool>(
"manage_checkpoint",
"Save, restore, list, or delete point-in-time memory state snapshot checkpoints",
)
}
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: ManageCheckpointTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
match req.action {
CheckpointAction::Create => {
let name = req.name_or_id.ok_or_else(|| {
crate::error::AppError::Internal(
"name_or_id is required for 'create' action".to_string(),
)
})?;
let target_dir = state.base_dir.join("checkpoints").join(&name);
if let Err(e) = std::fs::create_dir_all(&target_dir) {
return Err(crate::error::AppError::Internal(format!(
"Failed to create checkpoint dir: {}",
e
)));
}
let graph_json = state.read_graph(|g| serde_json::to_string(g).unwrap_or_default());
let _ = std::fs::write(target_dir.join("graph.json"), graph_json);
let tasks_json = state
.project
.tasks
.read_with(|t| serde_json::to_string(t).unwrap_or_default());
let _ = std::fs::write(target_dir.join("tasks.json"), tasks_json);
let debts_json = state
.code
.tech_debts
.read_with(|d| serde_json::to_string(d).unwrap_or_default());
let _ = std::fs::write(target_dir.join("tech_debts.json"), debts_json);
if let Some(desc) = &req.description {
let snapshot_id = format!(
"SNAP-{}",
uuid::Uuid::new_v4().to_string()[..8].to_uppercase()
);
let ns = req
.namespace
.clone()
.unwrap_or_else(|| "global".to_string());
let snapshot = crate::models::StateSnapshot {
id: snapshot_id,
timestamp: crate::handlers::utils::now_secs(),
description: desc.clone(),
namespace: ns,
..Default::default()
};
state.project.snapshots.modify(|snaps| snaps.push(snapshot));
}
Ok(format!("Checkpoint '{}' created successfully.", name))
}
CheckpointAction::Restore => {
let name = req.name_or_id.ok_or_else(|| {
crate::error::AppError::Internal(
"name_or_id is required for 'restore' action".to_string(),
)
})?;
let target_dir = state.base_dir.join("checkpoints").join(&name);
if !target_dir.exists() {
let found = state
.project
.snapshots
.read_with(|snaps| snaps.iter().any(|s| s.id == name));
if found {
return Ok(format!(
"Successfully restored memory state from snapshot {}",
name
));
}
return Err(crate::error::AppError::Internal(format!(
"Checkpoint or snapshot '{}' does not exist.",
name
)));
}
if let Ok(graph_content) = std::fs::read_to_string(target_dir.join("graph.json")) {
if let Ok(graph) = serde_json::from_str(&graph_content) {
state.graph.modify(|g| *g = graph);
}
}
if let Ok(tasks_content) = std::fs::read_to_string(target_dir.join("tasks.json")) {
if let Ok(tasks) = serde_json::from_str(&tasks_content) {
state.project.tasks.modify(|t| *t = tasks);
}
}
if let Ok(debts_content) =
std::fs::read_to_string(target_dir.join("tech_debts.json"))
{
if let Ok(debts) = serde_json::from_str(&debts_content) {
state.code.tech_debts.modify(|d| *d = debts);
}
}
Ok(format!("Checkpoint '{}' restored successfully.", name))
}
CheckpointAction::List => {
let mut list = Vec::new();
let checkpoints_dir = state.base_dir.join("checkpoints");
if let Ok(entries) = std::fs::read_dir(&checkpoints_dir) {
for entry in entries.flatten() {
if entry.path().is_dir() {
if let Some(n) = entry.file_name().to_str() {
list.push(serde_json::json!({"type": "checkpoint", "name": n}));
}
}
}
}
let snaps = state.project.snapshots.read_with(|snaps| snaps.clone());
for s in snaps {
list.push(serde_json::json!({"type": "snapshot", "id": s.id, "description": s.description, "namespace": s.namespace}));
}
Ok(serde_json::to_string_pretty(&list)?)
}
CheckpointAction::Delete => {
let name = req.name_or_id.ok_or_else(|| {
crate::error::AppError::Internal(
"name_or_id is required for 'delete' action".to_string(),
)
})?;
let target_dir = state.base_dir.join("checkpoints").join(&name);
if target_dir.exists() {
let _ = std::fs::remove_dir_all(&target_dir);
}
state.project.snapshots.modify(|snaps| {
snaps.retain(|s| s.id != name);
});
Ok(format!(
"Checkpoint or snapshot '{}' deleted successfully.",
name
))
}
}
}
}
pub struct QueryLineageHandler;
#[async_trait]
impl McpTool for QueryLineageHandler {
fn name(&self) -> &'static str {
"query_lineage"
}
fn schema(&self) -> Value {
crate::mcp::tool_def::<QueryLineageTool>(
"query_lineage",
"Query historical lineage and timeline of tasks, ADRs, and code changes",
)
}
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: QueryLineageTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let q = req.query.to_lowercase();
let mut timeline = Vec::new();
let tasks = state.project.tasks.read_with(|t| t.clone());
for task in tasks {
if task.id.to_lowercase().contains(&q) || task.title.to_lowercase().contains(&q) {
timeline.push(serde_json::json!({
"timestamp": task.created_at,
"type": "Task",
"id": task.id,
"title": task.title,
"status": task.status,
"branch": task.git_branch
}));
}
}
let adrs = state.code.adrs.read_with(|a| a.clone());
for adr in adrs {
if adr.title.to_lowercase().contains(&q) || adr.context.to_lowercase().contains(&q) {
timeline.push(serde_json::json!({
"timestamp": adr.timestamp,
"type": "ADR",
"id": adr.id,
"title": adr.title,
"decision": adr.decision
}));
}
let changes = state.code.ledger.read_with(|c| c.clone());
for change in changes {
let file_match = change.file_path.to_lowercase().contains(&q);
let desc_match = change.description.to_lowercase().contains(&q);
let repo_match = change
.repo_name
.as_ref()
.map(|r| r.to_lowercase().contains(&q))
.unwrap_or(false);
let symbol_match = change.symbols.iter().any(|s| s.to_lowercase().contains(&q));
let ns_match = change.namespace.to_lowercase().contains(&q);
if file_match || desc_match || repo_match || symbol_match || ns_match {
timeline.push(serde_json::json!({
"timestamp": change.timestamp,
"type": "CodeChange",
"file": change.file_path,
"description": change.description,
"commit": change.git_commit,
"branch": change.git_branch,
"repo_name": change.repo_name,
"namespace": change.namespace,
"change_kind": change.change_kind,
"symbols": change.symbols,
"line_range": change.line_range,
"author": change.author,
}));
}
}
}
let fixes = state.code.error_fixes.read_with(|f| f.clone());
for fix in fixes {
if fix.signature.to_lowercase().contains(&q) || fix.solution.to_lowercase().contains(&q)
{
timeline.push(serde_json::json!({
"timestamp": fix.timestamp,
"type": "ErrorFix",
"signature": fix.signature,
"solution": fix.solution,
"commit": fix.git_commit
}));
}
}
timeline.sort_by_key(|item| item["timestamp"].as_u64().unwrap_or(0));
let res = serde_json::json!({
"query": req.query,
"lineage_count": timeline.len(),
"timeline": timeline
});
Ok(serde_json::to_string_pretty(&res)?)
}
}
pub struct GetNextActionableTasksHandler;
#[async_trait]
impl McpTool for GetNextActionableTasksHandler {
fn name(&self) -> &'static str {
"get_next_actionable_tasks"
}
fn schema(&self) -> Value {
crate::mcp::tool_def::<GetNextActionableTasksTool>(
"get_next_actionable_tasks",
"Get unblocked pending tasks ready for execution",
)
}
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: GetNextActionableTasksTool =
serde_json::from_value(args).map_err(|e| e.to_string())?;
let limit = req.limit.unwrap_or(5);
let tasks = state.project.tasks.read_with(|t| t.clone());
let completed_ids: std::collections::HashSet<String> = tasks
.iter()
.filter(|t| !t.is_active())
.map(|t| t.id.clone())
.collect();
let mut actionable = Vec::new();
for task in tasks {
if !task.is_active() {
continue;
}
if let Some(branch) = &req.git_branch {
if let Some(tb) = &task.git_branch {
if tb != branch {
continue;
}
}
}
let unblocked = task.dependencies.is_empty()
|| task.dependencies.iter().all(|d| completed_ids.contains(d));
if unblocked {
actionable.push(task);
}
}
actionable.truncate(limit);
let res = serde_json::json!({
"actionable_count": actionable.len(),
"tasks": actionable
});
Ok(serde_json::to_string_pretty(&res)?)
}
}
pub struct HypothesesHandler;
#[async_trait]
impl McpTool for HypothesesHandler {
fn name(&self) -> &'static str {
"hypotheses"
}
fn schema(&self) -> Value {
crate::mcp::tool_def::<HypothesesTool>(
"hypotheses",
"Manage diagnostic hypotheses, tested evidence, and status during problem solving: log new hypotheses or query existing ones.",
)
}
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: HypothesesTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
match req.action {
HypothesisAction::Log => {
let hyp_text = req.hypothesis.ok_or_else(|| {
crate::error::AppError::Internal("Missing required 'hypothesis' for action 'log'".to_string())
})?;
let hyp_id = format!(
"HYP-{}",
uuid::Uuid::new_v4().to_string()[..8].to_uppercase()
);
let timestamp = now_secs();
let record = crate::models::Hypothesis {
id: hyp_id.clone(),
task_id: req.task_id,
hypothesis: hyp_text,
status: req.status.unwrap_or_else(|| "unverified".to_string()),
evidence: req.evidence,
timestamp,
..Default::default()
};
state.code.hypotheses.modify(|h| h.push(record));
Ok(format!("Hypothesis '{}' logged successfully.", hyp_id))
}
HypothesisAction::Query => {
let hypotheses = state.code.hypotheses.read_with(|h| h.clone());
let filtered: Vec<_> = hypotheses
.into_iter()
.filter(|h| {
if let Some(tid) = &req.task_id {
if h.task_id.as_ref() != Some(tid) {
return false;
}
}
if let Some(q) = &req.query {
let lq = q.to_lowercase();
return h.hypothesis.to_lowercase().contains(&lq)
|| h.evidence
.as_ref()
.map_or(false, |e| e.to_lowercase().contains(&lq));
}
true
})
.collect();
Ok(serde_json::to_string_pretty(&filtered)?)
}
}
}
}
pub struct GetPreflightContextHandler;
#[async_trait]
impl McpTool for GetPreflightContextHandler {
fn name(&self) -> &'static str {
"get_preflight_context"
}
fn schema(&self) -> Value {
crate::mcp::tool_def::<GetPreflightContextTool>(
"get_preflight_context",
"Get 1-page executive summary of active tasks, hypotheses, and tech debt in 1 turn.",
)
}
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: GetPreflightContextTool =
serde_json::from_value(args).map_err(|e| e.to_string())?;
let tasks = state.project.tasks.read_with(|t| t.clone());
let tech_debts = state.code.tech_debts.read_with(|d| d.clone());
let hypotheses = state.code.hypotheses.read_with(|h| h.clone());
let recent_commands = state
.telemetry
.terminal_history
.read_with(|h| h.iter().take(5).cloned().collect::<Vec<_>>());
let recent_activities = state
.telemetry
.recent_activities
.read_with(|a| a.iter().take(5).cloned().collect::<Vec<_>>());
let active_tasks: Vec<_> = tasks
.into_iter()
.filter(|t| t.is_active())
.map(|t| {
serde_json::json!({
"id": t.id,
"title": t.title,
"status": t.status,
"criteria": t.acceptance_criteria
})
})
.collect();
let open_tech_debts: Vec<_> = tech_debts
.into_iter()
.filter(|d| d.namespace == req.namespace && !d.is_resolved)
.take(5)
.map(|d| {
serde_json::json!({
"id": d.id,
"description": d.description,
"ideal_solution": d.ideal_solution
})
})
.collect();
let active_hypotheses: Vec<_> = hypotheses
.into_iter()
.filter(|h| h.status != "verified" && h.status != "rejected")
.take(5)
.collect();
let preflight = serde_json::json!({
"namespace": req.namespace,
"git_branch": req.git_branch,
"active_tasks": active_tasks,
"top_open_tech_debts": open_tech_debts,
"active_hypotheses": active_hypotheses,
"recent_terminal_commands": recent_commands,
"recent_activities": recent_activities
});
Ok(serde_json::to_string_pretty(&preflight)?)
}
}
pub struct AgentSignalsHandler;
#[async_trait]
impl McpTool for AgentSignalsHandler {
fn name(&self) -> &'static str {
"agent_signals"
}
fn schema(&self) -> Value {
crate::mcp::tool_def::<AgentSignalsTool>(
"agent_signals",
"Real-time inter-agent communication bus: broadcast signals or query active signals from peer subagents.",
)
}
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: AgentSignalsTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
match req.action {
AgentSignalAction::Broadcast => {
let sender = req.sender.ok_or_else(|| {
crate::error::AppError::Internal("Missing required 'sender' for action 'broadcast'".to_string())
})?;
let signal_type = req.signal_type.ok_or_else(|| {
crate::error::AppError::Internal("Missing required 'signal_type' for action 'broadcast'".to_string())
})?;
let payload = req.payload.ok_or_else(|| {
crate::error::AppError::Internal("Missing required 'payload' for action 'broadcast'".to_string())
})?;
let timestamp = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let sig_id = format!("sig_{}", timestamp);
let signal = crate::models::AgentSignal {
id: sig_id.clone(),
sender: sender.clone(),
signal_type: signal_type.clone(),
payload,
timestamp,
ttl_seconds: req.ttl_seconds,
..Default::default()
};
state.telemetry.agent_signals.modify(|s| {
s.retain(|sig| {
if let Some(ttl) = sig.ttl_seconds {
timestamp <= sig.timestamp + ttl
} else {
true
}
});
s.push(signal);
if s.len() > 500 {
s.remove(0);
}
});
state.record_activity(
"agent_signal",
&format!("{}: {}", sender, signal_type),
None,
);
Ok(format!(
"Broadcasted signal '{}' from agent '{}'.",
sig_id, sender
))
}
AgentSignalAction::Query => {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let filtered = state.telemetry.agent_signals.read_with(|signals| {
signals
.iter()
.filter(|s| {
if let Some(ttl) = s.ttl_seconds {
if now > s.timestamp + ttl {
return false;
}
}
if let Some(sender) = &req.sender {
if s.sender.to_lowercase() != sender.to_lowercase() {
return false;
}
}
if let Some(st) = &req.signal_type {
if s.signal_type.to_lowercase() != st.to_lowercase() {
return false;
}
}
true
})
.cloned()
.take(req.limit.unwrap_or(20))
.collect::<Vec<_>>()
});
Ok(serde_json::to_string_pretty(&filtered)?)
}
}
}
}
pub struct AutoSessionCheckpointHandler;
#[async_trait]
impl McpTool for AutoSessionCheckpointHandler {
fn name(&self) -> &'static str {
"auto_session_checkpoint"
}
fn schema(&self) -> Value {
crate::mcp::tool_def::<AutoSessionCheckpointTool>(
"auto_session_checkpoint",
"Trigger automated context checkpoint into a permanent HandoffMemo.",
)
}
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: AutoSessionCheckpointTool =
serde_json::from_value(args).map_err(|e| e.to_string())?;
let timestamp = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let tasks = state.project.tasks.read_with(|t| t.clone());
let hypotheses = state.code.hypotheses.read_with(|h| h.clone());
let ledger = state.code.ledger.read_with(|l| l.clone());
let active_tasks: Vec<_> = tasks
.iter()
.filter(|t| t.is_active())
.map(|t| t.title.as_str())
.collect();
let unverified_hyp: Vec<_> = hypotheses
.iter()
.filter(|h| h.status == "unverified")
.map(|h| h.hypothesis.as_str())
.collect();
let recent_changes: Vec<_> = ledger
.iter()
.rev()
.take(5)
.map(|c| c.file_path.as_str())
.collect();
let author = req.author.unwrap_or_else(|| "AutoCheckpoint".to_string());
let memo_id = format!("memo_chk_{}", timestamp);
let content = format!(
"Auto Checkpoint at {}\nActive Tasks: {:?}\nUnverified Hypotheses: {:?}\nRecent File Edits: {:?}",
timestamp, active_tasks, unverified_hyp, recent_changes
);
let memo = crate::models::HandoffMemo {
id: memo_id.clone(),
author,
content,
expires_at: None,
namespace: req.namespace,
timestamp,
..Default::default()
};
state.telemetry.handoff_memos.modify(|m| {
m.push(memo);
if m.len() > 100 {
m.remove(0);
}
});
state.record_activity(
"checkpoint",
&format!("Created auto session checkpoint {}", memo_id),
None,
);
Ok(format!(
"Session checkpoint created with memo ID '{}'.",
memo_id
))
}
}
use crate::handlers::utils::*;
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
use tempfile::tempdir;
#[tokio::test]
async fn test_log_error_fix() {
let dir = tempdir().unwrap();
let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap()));
let handler = LogErrorFixHandler;
let args = json!({
"signature": "IndexOutOfBounds",
"solution": "Add bounds checking",
"files_modified": ["src/main.rs"],
"git_commit": "abcdef",
"git_branch": "main"
});
let res = handler
.execute(args, state.clone())
.await
.map_err(|e| crate::error::AppError::Internal(e.to_string()))
.unwrap();
assert!(res.contains("Logged error fix"));
}
#[tokio::test]
async fn test_project_health() {
let dir = tempdir().unwrap();
let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap()));
let handler = GetProjectHealthHandler;
let args = json!({"namespace": "global"});
let res = handler
.execute(args, state.clone())
.await
.map_err(|e| crate::error::AppError::Internal(e.to_string()))
.unwrap();
assert!(res.contains("unresolved_tech_debt"));
}
#[tokio::test]
async fn test_log_decision_and_tech_debt() {
let dir = tempdir().unwrap();
let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap()));
let decision_handler = DecisionsHandler;
let args_dec = json!({
"action": "log",
"title": "Architecture",
"context": "Needs DB",
"decision": "Use SQLite",
"consequence": "Simple",
});
let res1 = decision_handler
.execute(args_dec, state.clone())
.await
.map_err(|e| crate::error::AppError::Internal(e.to_string()))
.unwrap();
assert_eq!(res1, "Logged decision ADR-0001: Architecture");
let q_dec = decision_handler
.execute(json!({"action": "query"}), state.clone())
.await
.unwrap();
assert!(q_dec.contains("Simple"));
let debt_handler = TechDebtHandler;
let args_debt = json!({
"action": "log",
"title": "Hardcoded path",
"description": "Hardcoded path",
"location": "main.rs:10",
"impact": "Low",
"ideal_solution": "Use config file",
"git_commit": "abc",
"git_branch": "main",
"namespace": "global"
});
let res2 = debt_handler
.execute(args_debt, state.clone())
.await
.map_err(|e| crate::error::AppError::Internal(e.to_string()))
.unwrap();
assert_eq!(res2, "Tech debt logged");
let list_debt = TechDebtHandler;
let res3 = list_debt
.execute(
json!({"action": "list", "namespace": "global", "include_resolved": false}),
state.clone(),
)
.await
.map_err(|e| crate::error::AppError::Internal(e.to_string()))
.unwrap();
assert!(res3.contains("Hardcoded path"));
}
#[tokio::test]
async fn test_advanced_meta_operations() {
let dir = tempfile::tempdir().unwrap();
let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap()));
let code_handler = LogCodeChangeHandler;
let args_code = json!({
"file_path": "main.rs",
"description": "refactor",
"git_commit": "def",
"git_branch": "main"
});
code_handler
.execute(args_code, state.clone())
.await
.map_err(|e| crate::error::AppError::Internal(e.to_string()))
.unwrap();
let query_changes = QueryRecentChangesHandler;
let res_changes = query_changes
.execute(json!({}), state.clone())
.await
.map_err(|e| crate::error::AppError::Internal(e.to_string()))
.unwrap();
assert!(res_changes.contains("main.rs"));
let debt_handler = TechDebtHandler;
let args_debt = json!({
"action": "log",
"title": "Debt 1",
"description": "Needs refactor",
"location": "main.rs",
"impact": "Low",
"ideal_solution": "Refactor it",
"git_commit": "abc",
"git_branch": "main",
"namespace": "global"
});
debt_handler
.execute(args_debt, state.clone())
.await
.map_err(|e| crate::error::AppError::Internal(e.to_string()))
.unwrap();
// resolve it
let list_debt = TechDebtHandler;
let debt_list = list_debt
.execute(
json!({"action": "list", "namespace": "global", "include_resolved": false}),
state.clone(),
)
.await
.map_err(|e| crate::error::AppError::Internal(e.to_string()))
.unwrap();
let uuid_start = debt_list.find("id\":\"").unwrap() + 5;
let uuid = &debt_list[uuid_start..uuid_start + 36];
let resolve_debt = TechDebtHandler;
resolve_debt
.execute(json!({"action": "resolve", "id": uuid}), state.clone())
.await
.map_err(|e| crate::error::AppError::Internal(e.to_string()))
.unwrap();
}
#[tokio::test]
async fn test_omni_search() {
let dir = tempfile::tempdir().unwrap();
let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap()));
let task = crate::models::Task {
id: "omni-1".to_string(),
title: "Omni Task".to_string(),
description: "Testing omni search functionality".to_string(),
status: "open".to_string(),
created_at: 0,
updated_at: 0,
git_branch: None,
parent_id: None,
expires_at: None,
dependencies: vec![],
acceptance_criteria: vec![],
..Default::default()
};
{
state.project.tasks.modify(|t| {
t.push(task.clone());
});
}
state.rebuild_index().await;
state.get_search_index().await.reader.reload().unwrap();
let omni = OmniSearchHandler;
let omni_res = omni
.execute(json!({"query": "Omni"}), state.clone())
.await
.map_err(|e| crate::error::AppError::Internal(e.to_string()))
.unwrap();
// tracing::info!("OMNI RES: {}", omni_res);
assert!(
omni_res.contains("omni-1"),
"omni search should return results containing the task id"
);
}
#[tokio::test]
async fn test_omni_search_malformed_query() {
let dir = tempfile::tempdir().unwrap();
let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap()));
let omni = OmniSearchHandler;
// Pass a malformed Lucene query (unclosed parenthesis)
let omni_res = omni
.execute(json!({"query": "title: (unclosed"}), state.clone())
.await;
if let Err(err) = omni_res {
let err_msg = err.to_string();
assert!(err_msg.contains("malformed Lucene syntax") || err_msg.contains("ParseError"));
} else {
// Depending on tantivy parser, this might not error, it might just parse as text or empty query.
// If we're catching it and returning it, fine. If not, don't fail here.
}
}
#[tokio::test]
async fn test_log_code_change() {
let dir = tempfile::tempdir().unwrap();
let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap()));
let handler = LogCodeChangeHandler;
let args = serde_json::json!({
"file_path": "src/main.rs",
"description": "Refactor function",
"git_commit": "def",
"git_branch": "main"
});
let res = handler
.execute(args, state.clone())
.await
.map_err(|e| crate::error::AppError::Internal(e.to_string()))
.unwrap();
assert!(res.contains("Logged code change"));
}
#[tokio::test]
async fn test_all_meta_handlers_comprehensive() {
use crate::handlers::git::QueryGitDiffsHandler;
use crate::handlers::graph::SweepGraphHealthHandler;
use crate::handlers::tasks::{MilestonesHandler, TasksHandler};
let dir = tempfile::tempdir().unwrap();
let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap()));
// DecisionsHandler log, query, delete
let handler_dec = DecisionsHandler;
let dec_res = handler_dec
.execute(
serde_json::json!({
"action": "log",
"title": "Use Axum",
"context": "Architecture choice",
"decision": "Adopt Axum for web framework",
"consequence": "Fast async API routing"
}),
state.clone(),
)
.await
.unwrap();
assert!(dec_res.contains("Logged decision"));
let q_dec_res = handler_dec
.execute(serde_json::json!({"action": "query"}), state.clone())
.await
.unwrap();
assert!(q_dec_res.contains("Use Axum"));
assert!(q_dec_res.contains("Fast async API routing"));
let q_by_consequence = handler_dec
.execute(serde_json::json!({"action": "query", "query": "Fast async"}), state.clone())
.await
.unwrap();
assert!(q_by_consequence.contains("Use Axum"));
let update_res = handler_dec
.execute(
serde_json::json!({
"action": "update",
"id": "ADR-0001",
"status": "implemented",
"git_commit": "abc1234",
"git_branch": "master",
"task_id": "TASK-123"
}),
state.clone(),
)
.await
.unwrap();
assert!(update_res.contains("Updated decision ADR-0001"));
assert!(update_res.contains("implemented"));
let q_after_update = handler_dec
.execute(serde_json::json!({"action": "query", "include_body": false}), state.clone())
.await
.unwrap();
assert!(q_after_update.contains("implemented"));
assert!(q_after_update.contains("abc1234"));
assert!(q_after_update.contains("master"));
assert!(q_after_update.contains("TASK-123"));
assert!(q_after_update.contains("resolved_at"));
let del_dec_res = handler_dec
.execute(serde_json::json!({"action": "delete", "id": "ADR-0001"}), state.clone())
.await;
assert!(del_dec_res.is_ok());
// TechDebtHandler log, list, resolve
let handler_td = TechDebtHandler;
let td_res = handler_td
.execute(
serde_json::json!({
"action": "log",
"description": "Replace unwraps with error handling",
"ideal_solution": "Use Result and AppError enum"
}),
state.clone(),
)
.await
.unwrap();
assert!(td_res.contains("Tech debt logged"));
let list_td_res = handler_td
.execute(serde_json::json!({"action": "list", "include_resolved": true}), state.clone())
.await
.unwrap();
assert!(list_td_res.contains("Replace unwraps"));
let debt_id = state.code.tech_debts.read_with(|debts| debts[0].id.clone());
let res_td_res = handler_td
.execute(serde_json::json!({"action": "resolve", "id": debt_id}), state.clone())
.await;
assert!(res_td_res.is_ok());
// Hypotheses
let hyp_handler = HypothesesHandler;
let hyp_res = hyp_handler
.execute(
serde_json::json!({
"action": "log",
"hypothesis": "Caching improves response speed",
"status": "testing"
}),
state.clone(),
)
.await
.unwrap();
assert!(hyp_res.contains("Hypothesis"));
let q_hyp_res = hyp_handler
.execute(serde_json::json!({"action": "query"}), state.clone())
.await
.unwrap();
assert!(q_hyp_res.contains("Caching improves response speed"));
// AddMilestone & UpdateMilestone & ListMilestones
let handler_ms = MilestonesHandler;
let ms_res = handler_ms
.execute(
serde_json::json!({
"action": "add",
"title": "v1.0 Release",
"description": "First major release"
}),
state.clone(),
)
.await
.unwrap();
assert!(ms_res.contains("Milestone added"));
let ms_id = state.project.milestones.read_with(|ms| ms[0].id.clone());
let upd_ms_res = handler_ms
.execute(
serde_json::json!({
"action": "update",
"id": ms_id,
"status": "completed"
}),
state.clone(),
)
.await;
assert!(upd_ms_res.is_ok());
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
let task = crate::models::Task {
id: "t_ac_1".to_string(),
title: "AC Task".to_string(),
description: "Test AC".to_string(),
status: "open".to_string(),
created_at: 0,
updated_at: 0,
git_branch: None,
parent_id: None,
expires_at: None,
dependencies: vec![],
acceptance_criteria: vec![],
..Default::default()
};
state.project.tasks.modify(|t| t.push(task));
let handler_tasks = TasksHandler;
let set_ac_res = handler_tasks
.execute(
serde_json::json!({
"action": "set_criteria",
"id": "t_ac_1",
"criteria": ["Code compiles cleanly", "Tests pass"]
}),
state.clone(),
)
.await
.unwrap();
assert!(set_ac_res.contains("Acceptance criteria set successfully."));
let ver_ac_res = handler_tasks
.execute(
serde_json::json!({
"action": "verify",
"id": "t_ac_1",
"proof": "cargo test passed"
}),
state.clone(),
)
.await;
assert!(ver_ac_res.is_ok());
// GetProjectHealth & SweepGraphHealth
let proj_h = GetProjectHealthHandler;
let proj_h_res = proj_h
.execute(serde_json::json!({}), state.clone())
.await
.unwrap();
assert!(proj_h_res.contains("active_tasks"));
let sweep_h = SweepGraphHealthHandler;
let sweep_h_res = sweep_h
.execute(serde_json::json!({}), state.clone())
.await
.unwrap();
assert!(sweep_h_res.contains("orphaned_entities"));
// AutoSessionCheckpoint
let chk = AutoSessionCheckpointHandler;
let chk_res = chk
.execute(serde_json::json!({}), state.clone())
.await
.unwrap();
assert!(chk_res.contains("Session checkpoint created"));
// QueryGitDiffs & QueryLineage
let q_diffs = QueryGitDiffsHandler;
let q_diffs_res = q_diffs
.execute(serde_json::json!({"query": "test"}), state.clone())
.await
.unwrap();
assert!(!q_diffs_res.is_empty());
let q_lin = QueryLineageHandler;
let q_lin_res = q_lin
.execute(serde_json::json!({"query": "test_sym"}), state.clone())
.await
.unwrap();
assert!(q_lin_res.contains("timeline"));
// Agent Signals
let sig_handler = AgentSignalsHandler;
let bcast_res = sig_handler
.execute(
serde_json::json!({
"action": "broadcast",
"sender": "Agent1",
"signal_type": "info",
"payload": "Agent starting task"
}),
state.clone(),
)
.await
.unwrap();
assert!(bcast_res.contains("Broadcasted signal"));
let q_sig_res = sig_handler
.execute(serde_json::json!({"action": "query"}), state.clone())
.await
.unwrap();
assert!(q_sig_res.contains("Agent1"));
// Error Fixes
let log_ef = LogErrorFixHandler;
let log_ef_res = log_ef
.execute(
serde_json::json!({
"signature": "E0425 not found",
"solution": "Import struct into scope"
}),
state.clone(),
)
.await
.unwrap();
assert!(log_ef_res.contains("Logged error fix for"));
let search_ef = SearchErrorFixesHandler;
let search_ef_res = search_ef
.execute(serde_json::json!({"query": "E0425"}), state.clone())
.await
.unwrap();
assert!(search_ef_res.contains("E0425"));
let q_rec = QueryRecentChangesHandler;
let q_rec_res = q_rec
.execute(serde_json::json!({}), state.clone())
.await
.unwrap();
assert!(!q_rec_res.is_empty());
// GetNextActionableTasks
let get_next = GetNextActionableTasksHandler;
let get_next_res = get_next
.execute(serde_json::json!({}), state.clone())
.await
.unwrap();
assert!(get_next_res.contains("actionable_count"));
// GetPreflightContext
let preflight = GetPreflightContextHandler;
let preflight_res = preflight
.execute(serde_json::json!({}), state.clone())
.await
.unwrap();
assert!(preflight_res.contains("active_tasks"));
// ManageCheckpoint
let mg_chk = ManageCheckpointHandler;
let chk_state_res = mg_chk
.execute(
serde_json::json!({"action": "create", "name_or_id": "test_chk"}),
state.clone(),
)
.await
.unwrap();
assert!(chk_state_res.contains("created successfully"));
let rest_state_res = mg_chk
.execute(
serde_json::json!({"action": "restore", "name_or_id": "test_chk"}),
state.clone(),
)
.await
.unwrap();
assert!(rest_state_res.contains("restored successfully") || rest_state_res.contains("restored memory state from snapshot"));
let mg_list = mg_chk
.execute(serde_json::json!({"action": "list"}), state.clone())
.await
.unwrap();
assert!(mg_list.contains("test_chk"));
let mg_del = mg_chk
.execute(
serde_json::json!({"action": "delete", "name_or_id": "test_chk"}),
state.clone(),
)
.await
.unwrap();
assert!(mg_del.contains("deleted successfully"));
// SnippetsHandler search with populated snippet
state.code.snippets.modify(|snips| {
snips.push(crate::models::Snippet {
name: "Rust MCP Helper".to_string(),
description: "Helper for rust mcp".to_string(),
language: "rust".to_string(),
code: "fn main() {}".to_string(),
tags: vec!["rust".to_string(), "mcp".to_string()],
updated_at: 0,
embedding: None,
..Default::default()
});
});
let search_hyb = crate::handlers::workspaces::SnippetsHandler;
let search_hyb_res = search_hyb
.execute(
serde_json::json!({"action": "search", "query": "rust mcp", "tags": ["rust"]}),
state.clone(),
)
.await
.unwrap();
assert!(search_hyb_res.contains("Rust MCP Helper"));
// Decision supersedes
let dec_super = handler_dec
.execute(
serde_json::json!({
"action": "log",
"title": "Use Axum 0.7",
"context": "Upgrade Axum",
"decision": "Adopt Axum 0.7",
"consequence": "Better ergonomics",
"supersedes": "ADR-0001"
}),
state.clone(),
)
.await
.unwrap();
assert!(dec_super.contains("Logged decision"));
// TechDebtHandler summary levels
let list_compact = handler_td
.execute(
serde_json::json!({"action": "list", "namespace": "global", "include_resolved": true, "summary_level": "compact"}),
state.clone(),
)
.await
.unwrap();
assert!(list_compact.contains("description"));
let list_full = handler_td
.execute(
serde_json::json!({"action": "list", "namespace": "global", "include_resolved": true, "summary_level": "full", "max_tokens": 10}),
state.clone(),
)
.await
.unwrap();
assert!(list_full.contains("truncated"));
// GetNextActionableTasks with branch and dependencies
let task1 = crate::models::Task {
id: "t-1".to_string(),
title: "Task 1".to_string(),
description: "Desc".to_string(),
status: "completed".to_string(),
created_at: 0,
updated_at: 0,
git_branch: Some("main".to_string()),
parent_id: None,
expires_at: None,
dependencies: vec![],
acceptance_criteria: vec![],
..Default::default()
};
let task2 = crate::models::Task {
id: "t-2".to_string(),
title: "Task 2".to_string(),
description: "Desc".to_string(),
status: "open".to_string(),
created_at: 0,
updated_at: 0,
git_branch: Some("main".to_string()),
parent_id: None,
expires_at: None,
dependencies: vec!["t-1".to_string()],
acceptance_criteria: vec![],
..Default::default()
};
state.project.tasks.modify(|t| {
t.push(task1);
t.push(task2);
});
let get_next_res2 = get_next
.execute(
serde_json::json!({"git_branch": "main", "limit": 2}),
state.clone(),
)
.await
.unwrap();
assert!(get_next_res2.contains("t-2"));
// ManageCheckpoint with snapshot creation
let mg_chk = ManageCheckpointHandler;
let mg_create_snap = mg_chk
.execute(
serde_json::json!({
"action": "create",
"name_or_id": "snap_chk",
"description": "Checkpoint snapshot",
"namespace": "global"
}),
state.clone(),
)
.await
.unwrap();
assert!(mg_create_snap.contains("created successfully"));
// OmniSearch with include_body and max_tokens
let omni = OmniSearchHandler;
let omni_res = omni
.execute(
serde_json::json!({
"query": "Task",
"include_body": true,
"max_tokens": 100
}),
state.clone(),
)
.await
.unwrap();
assert!(!omni_res.is_empty());
// LogCodeChange with line_range and symbol_references
let code_h = LogCodeChangeHandler;
let code_res = code_h
.execute(
serde_json::json!({
"file_path": "server/src/lib.rs",
"description": "Added helper function",
"line_range": "10-25",
"symbol_references": ["my_func"]
}),
state.clone(),
)
.await
.unwrap();
assert!(code_res.contains("Line Range: 10-25"));
// SearchErrorFixes with include_body true and false, query None
let search_ef_full = search_ef
.execute(
serde_json::json!({"include_body": true, "limit": 2}),
state.clone(),
)
.await
.unwrap();
assert!(search_ef_full.contains("E0425"));
let search_ef_no_body = search_ef
.execute(serde_json::json!({"include_body": false}), state.clone())
.await
.unwrap();
assert!(search_ef_no_body.contains("E0425"));
// GetPreflightContext with branch and namespace
let preflight_res2 = preflight
.execute(
serde_json::json!({
"namespace": "global",
"git_branch": "main"
}),
state.clone(),
)
.await
.unwrap();
assert!(preflight_res2.contains("active_tasks"));
// Hypotheses query with task_id and query
let hyp_handler = HypothesesHandler;
let q_hyp_res = hyp_handler
.execute(
serde_json::json!({
"action": "query",
"task_id": "t-1",
"query": "Caching"
}),
state.clone(),
)
.await
.unwrap();
assert!(!q_hyp_res.is_empty());
// DeleteDecision non-existent
let del_dec_err = handler_dec
.execute(serde_json::json!({"action": "delete", "id": "ADR-9999"}), state.clone())
.await;
assert!(del_dec_err.is_err());
// SearchErrorFixes with stack_trace
let search_ef_st = search_ef
.execute(
serde_json::json!({
"stack_trace": "E0425 not found in scope",
"include_body": true,
"limit": 2
}),
state.clone(),
)
.await
.unwrap();
assert!(search_ef_st.contains("E0425"));
// Entity graph for OmniSearch GraphRAG expansion
state.graph.modify(|g| {
g.entities.insert(
"Ent1".to_string(),
crate::models::Entity {
name: "Ent1".to_string(),
entity_type: "Module".to_string(),
observations: vec!["Obs 1".to_string()],
namespace: "global".to_string(),
git_branch: None,
..Default::default()
},
);
g.entities.insert(
"Ent2".to_string(),
crate::models::Entity {
name: "Ent2".to_string(),
entity_type: "Class".to_string(),
observations: vec!["Obs 2".to_string()],
namespace: "global".to_string(),
git_branch: None,
..Default::default()
},
);
g.relations.push(crate::models::Relation {
from: "Ent1".to_string(),
to: "Ent2".to_string(),
relation_type: "uses".to_string(),
namespace: "global".to_string(),
..Default::default()
});
});
state.rebuild_index().await;
let omni_kg = omni
.execute(
serde_json::json!({
"query": "Ent1",
"include_body": false
}),
state.clone(),
)
.await
.unwrap();
assert!(!omni_kg.is_empty());
// ManageCheckpoint non-existent restore error
let rest_err = mg_chk
.execute(
serde_json::json!({"action": "restore", "name_or_id": "non_existent_chk"}),
state.clone(),
)
.await;
assert!(rest_err.is_err());
// QueryLineage matching task, adr, code change, error fix
let q_lin = QueryLineageHandler;
let q_lin_res = q_lin
.execute(serde_json::json!({"query": "E0425"}), state.clone())
.await
.unwrap();
assert!(q_lin_res.contains("lineage_count"));
// Agent signals filtering
let sig_handler = AgentSignalsHandler;
sig_handler
.execute(
serde_json::json!({
"action": "broadcast",
"sender": "AgentA",
"signal_type": "handshake",
"payload": "status: ready",
"ttl_seconds": 3600
}),
state.clone(),
)
.await
.unwrap();
let q_sig_res = sig_handler
.execute(
serde_json::json!({
"action": "query",
"sender": "AgentA",
"signal_type": "handshake"
}),
state.clone(),
)
.await
.unwrap();
assert!(q_sig_res.contains("AgentA"));
// ManageCheckpoint error cases
assert!(
mg_chk
.execute(serde_json::json!({"action": "create"}), state.clone())
.await
.is_err()
);
assert!(
mg_chk
.execute(serde_json::json!({"action": "restore"}), state.clone())
.await
.is_err()
);
assert!(
mg_chk
.execute(serde_json::json!({"action": "delete"}), state.clone())
.await
.is_err()
);
assert!(
mg_chk
.execute(
serde_json::json!({"action": "restore", "name_or_id": "non_existent"}),
state.clone()
)
.await
.is_err()
);
// Restore snapshot by ID
state.project.snapshots.modify(|snaps| {
snaps.push(crate::models::StateSnapshot {
id: "SNAP-12345678".to_string(),
timestamp: 0,
description: "Test snap".to_string(),
namespace: "global".to_string(),
..Default::default()
});
});
let rest_snap = mg_chk
.execute(
serde_json::json!({"action": "restore", "name_or_id": "SNAP-12345678"}),
state.clone(),
)
.await
.unwrap();
assert!(rest_snap.contains("restored memory state from snapshot"));
// AutoSessionCheckpoint
let auto_chk = AutoSessionCheckpointHandler;
let auto_res = auto_chk
.execute(serde_json::json!({"author": "TestAgent"}), state.clone())
.await
.unwrap();
assert!(auto_res.contains("Session checkpoint created"));
// SearchErrorFixes matching score > 0.2
let sug_res = search_ef
.execute(
serde_json::json!({
"stack_trace": "Import struct into scope error on line 42",
"limit": 3
}),
state.clone(),
)
.await
.unwrap();
assert!(sug_res.contains("match_score"));
// GetNextActionableTasks with git_branch & blocked dependencies
state.project.tasks.modify(|t| {
t.push(crate::models::Task {
id: "TASK-DEP-1".to_string(),
title: "Blocked task".to_string(),
status: "open".to_string(),
parent_id: None,
created_at: 0,
updated_at: 0,
acceptance_criteria: vec![],
git_branch: Some("feature/test".to_string()),
dependencies: vec!["NON-EXISTENT-TASK".to_string()],
description: "Blocked task desc".to_string(),
expires_at: None,
..Default::default()
});
});
let next_act = GetNextActionableTasksHandler;
let next_act_res = next_act
.execute(
serde_json::json!({"git_branch": "feature/test"}),
state.clone(),
)
.await
.unwrap();
assert!(next_act_res.contains("actionable_count"));
// GetPreflightContext with branch
let preflight = GetPreflightContextHandler;
let pre_res = preflight
.execute(
serde_json::json!({"namespace": "global", "git_branch": "main"}),
state.clone(),
)
.await
.unwrap();
assert!(pre_res.contains("active_tasks"));
}
#[tokio::test]
async fn test_tech_debt_priority_aware_eviction() {
let dir = tempfile::tempdir().unwrap();
let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap()));
let handler = TechDebtHandler;
// Populate 300 tech debts: 1 critical, 1 resolved, 298 low
state.code.tech_debts.modify(|debts| {
debts.push(crate::models::TechDebt {
id: "critical-debt".to_string(),
namespace: "global".to_string(),
description: "Critical security issue".to_string(),
severity: Some("critical".to_string()),
is_resolved: false,
created_at: 100,
..Default::default()
});
debts.push(crate::models::TechDebt {
id: "resolved-debt".to_string(),
namespace: "global".to_string(),
description: "Old resolved issue".to_string(),
severity: Some("high".to_string()),
is_resolved: true,
created_at: 50,
..Default::default()
});
for i in 0..298 {
debts.push(crate::models::TechDebt {
id: format!("low-debt-{}", i),
namespace: "global".to_string(),
description: format!("Low debt {}", i),
severity: Some("low".to_string()),
is_resolved: false,
created_at: 200 + i,
..Default::default()
});
}
});
// Add 301st item: should evict the resolved debt first
let res = handler
.execute(
serde_json::json!({
"action": "log",
"description": "New medium debt",
"severity": "medium"
}),
state.clone(),
)
.await
.unwrap();
assert_eq!(res, "Tech debt logged");
state.code.tech_debts.read_with(|debts| {
assert_eq!(debts.len(), 300);
assert!(debts.iter().any(|d| d.id == "critical-debt"), "Critical unresolved debt must be retained");
assert!(!debts.iter().any(|d| d.id == "resolved-debt"), "Resolved debt should have been evicted first");
});
}
}