Refactor: Dismantle MemoryState into semantic domain sub-structs

This commit is contained in:
Riz Ashraf committed 2026-09-30 21:13:03 +01:00
1 parent 0e866f2465
commit a34554b7ff
13 files changed
+176 -148

No files matched your search

+2 -2
View File
@@ -38,7 +38,7 @@ pub async fn gate_verify_handler(
) -> Result<impl IntoResponse, AppError> { ) -> Result<impl IntoResponse, AppError> {
let mut found = None; let mut found = None;
let mut to_remove = None; let mut to_remove = None;
app_state.handler.state.gates.modify(|gates| { app_state.handler.state.env.gates.modify(|gates| {
if let Some(idx) = gates.iter().position(|g| { if let Some(idx) = gates.iter().position(|g| {
g.action == q.action g.action == q.action
&& g.target == q.target && g.target == q.target
@@ -96,7 +96,7 @@ pub async fn gate_set_handler(
reason: body.reason.clone(), reason: body.reason.clone(),
timestamp: crate::handlers::utils::now_secs(), timestamp: crate::handlers::utils::now_secs(),
}; };
app_state.handler.state.gates.modify(|gates| { app_state.handler.state.env.gates.modify(|gates| {
gates.retain(|g| !(g.action == record.action && g.target == record.target)); gates.retain(|g| !(g.action == record.action && g.target == record.target));
gates.push(record); gates.push(record);
}); });
+29 -29
View File
@@ -80,7 +80,7 @@ pub fn create_router(app_state: Arc<AppState>) -> Router {
post({ post({
let state_clone = app_state.handler.state.clone(); let state_clone = app_state.handler.state.clone();
move |axum::extract::Path(id): axum::extract::Path<String>| async move { move |axum::extract::Path(id): axum::extract::Path<String>| async move {
state_clone.tasks.modify(|tasks| { state_clone.project.tasks.modify(|tasks| {
for t in tasks.iter_mut() { for t in tasks.iter_mut() {
if t.id == id { if t.id == id {
t.status = "completed".to_string(); t.status = "completed".to_string();
@@ -97,7 +97,7 @@ pub fn create_router(app_state: Arc<AppState>) -> Router {
get({ get({
let state_clone = app_state.handler.state.clone(); let state_clone = app_state.handler.state.clone();
move || async move { move || async move {
let tasks_json = state_clone.tasks.read_with(|t| serde_json::to_string(t).unwrap_or_else(|_| "[]".to_string())); let tasks_json = state_clone.project.tasks.read_with(|t| serde_json::to_string(t).unwrap_or_else(|_| "[]".to_string()));
([(axum::http::header::CONTENT_TYPE, "application/json")], tasks_json) ([(axum::http::header::CONTENT_TYPE, "application/json")], tasks_json)
} }
}), }),
@@ -107,7 +107,7 @@ pub fn create_router(app_state: Arc<AppState>) -> Router {
get({ get({
let state_clone = app_state.handler.state.clone(); let state_clone = app_state.handler.state.clone();
move || async move { move || async move {
let sticky_json = state_clone.sticky.read_with(|s| serde_json::to_string(s).unwrap_or_else(|_| "[]".to_string())); let sticky_json = state_clone.code.sticky.read_with(|s| serde_json::to_string(s).unwrap_or_else(|_| "[]".to_string()));
([(axum::http::header::CONTENT_TYPE, "application/json")], sticky_json) ([(axum::http::header::CONTENT_TYPE, "application/json")], sticky_json)
} }
}), }),
@@ -145,7 +145,7 @@ pub fn create_router(app_state: Arc<AppState>) -> Router {
get({ get({
let state_clone = app_state.handler.state.clone(); let state_clone = app_state.handler.state.clone();
move || async move { move || async move {
let activities_json = state_clone.recent_activities.read_with(|a| serde_json::to_string(a).unwrap_or_else(|_| "[]".to_string())); let activities_json = state_clone.telemetry.recent_activities.read_with(|a| serde_json::to_string(a).unwrap_or_else(|_| "[]".to_string()));
([(axum::http::header::CONTENT_TYPE, "application/json")], activities_json) ([(axum::http::header::CONTENT_TYPE, "application/json")], activities_json)
} }
}), }),
@@ -155,7 +155,7 @@ pub fn create_router(app_state: Arc<AppState>) -> Router {
get({ get({
let state_clone = app_state.handler.state.clone(); let state_clone = app_state.handler.state.clone();
move || async move { move || async move {
let json = state_clone.tech_debts.read_with(|items| serde_json::to_string(items).unwrap_or_else(|_| "[]".to_string())); let json = state_clone.code.tech_debts.read_with(|items| serde_json::to_string(items).unwrap_or_else(|_| "[]".to_string()));
([(axum::http::header::CONTENT_TYPE, "application/json")], json) ([(axum::http::header::CONTENT_TYPE, "application/json")], json)
} }
}), }),
@@ -165,7 +165,7 @@ pub fn create_router(app_state: Arc<AppState>) -> Router {
get({ get({
let state_clone = app_state.handler.state.clone(); let state_clone = app_state.handler.state.clone();
move || async move { move || async move {
let json = state_clone.adrs.read_with(|items| serde_json::to_string(items).unwrap_or_else(|_| "[]".to_string())); let json = state_clone.code.adrs.read_with(|items| serde_json::to_string(items).unwrap_or_else(|_| "[]".to_string()));
([(axum::http::header::CONTENT_TYPE, "application/json")], json) ([(axum::http::header::CONTENT_TYPE, "application/json")], json)
} }
}), }),
@@ -175,7 +175,7 @@ pub fn create_router(app_state: Arc<AppState>) -> Router {
get({ get({
let state_clone = app_state.handler.state.clone(); let state_clone = app_state.handler.state.clone();
move || async move { move || async move {
let json = state_clone.context_workspaces.read_with(|items| serde_json::to_string(items).unwrap_or_else(|_| "[]".to_string())); let json = state_clone.project.context_workspaces.read_with(|items| serde_json::to_string(items).unwrap_or_else(|_| "[]".to_string()));
([(axum::http::header::CONTENT_TYPE, "application/json")], json) ([(axum::http::header::CONTENT_TYPE, "application/json")], json)
} }
}), }),
@@ -185,7 +185,7 @@ pub fn create_router(app_state: Arc<AppState>) -> Router {
get({ get({
let state_clone = app_state.handler.state.clone(); let state_clone = app_state.handler.state.clone();
move || async move { move || async move {
let json = state_clone.handoff_memos.read_with(|items| serde_json::to_string(items).unwrap_or_else(|_| "[]".to_string())); let json = state_clone.telemetry.handoff_memos.read_with(|items| serde_json::to_string(items).unwrap_or_else(|_| "[]".to_string()));
([(axum::http::header::CONTENT_TYPE, "application/json")], json) ([(axum::http::header::CONTENT_TYPE, "application/json")], json)
} }
}), }),
@@ -195,7 +195,7 @@ pub fn create_router(app_state: Arc<AppState>) -> Router {
get({ get({
let state_clone = app_state.handler.state.clone(); let state_clone = app_state.handler.state.clone();
move || async move { move || async move {
let json = state_clone.milestones.read_with(|items| serde_json::to_string(items).unwrap_or_else(|_| "[]".to_string())); let json = state_clone.project.milestones.read_with(|items| serde_json::to_string(items).unwrap_or_else(|_| "[]".to_string()));
([(axum::http::header::CONTENT_TYPE, "application/json")], json) ([(axum::http::header::CONTENT_TYPE, "application/json")], json)
} }
}), }),
@@ -205,7 +205,7 @@ pub fn create_router(app_state: Arc<AppState>) -> Router {
get({ get({
let state_clone = app_state.handler.state.clone(); let state_clone = app_state.handler.state.clone();
move || async move { move || async move {
let json = state_clone.snippets.read_with(|items| serde_json::to_string(items).unwrap_or_else(|_| "[]".to_string())); let json = state_clone.code.snippets.read_with(|items| serde_json::to_string(items).unwrap_or_else(|_| "[]".to_string()));
([(axum::http::header::CONTENT_TYPE, "application/json")], json) ([(axum::http::header::CONTENT_TYPE, "application/json")], json)
} }
}), }),
@@ -215,7 +215,7 @@ pub fn create_router(app_state: Arc<AppState>) -> Router {
get({ get({
let state_clone = app_state.handler.state.clone(); let state_clone = app_state.handler.state.clone();
move || async move { move || async move {
let json = state_clone.pr_checklists.read_with(|items| serde_json::to_string(items).unwrap_or_else(|_| "[]".to_string())); let json = state_clone.project.pr_checklists.read_with(|items| serde_json::to_string(items).unwrap_or_else(|_| "[]".to_string()));
([(axum::http::header::CONTENT_TYPE, "application/json")], json) ([(axum::http::header::CONTENT_TYPE, "application/json")], json)
} }
}), }),
@@ -225,7 +225,7 @@ pub fn create_router(app_state: Arc<AppState>) -> Router {
get({ get({
let state_clone = app_state.handler.state.clone(); let state_clone = app_state.handler.state.clone();
move || async move { move || async move {
let json = state_clone.error_fixes.read_with(|items| serde_json::to_string(items).unwrap_or_else(|_| "[]".to_string())); let json = state_clone.code.error_fixes.read_with(|items| serde_json::to_string(items).unwrap_or_else(|_| "[]".to_string()));
([(axum::http::header::CONTENT_TYPE, "application/json")], json) ([(axum::http::header::CONTENT_TYPE, "application/json")], json)
} }
}), }),
@@ -236,24 +236,24 @@ pub fn create_router(app_state: Arc<AppState>) -> Router {
let state_clone = app_state.handler.state.clone(); let state_clone = app_state.handler.state.clone();
move || async move { move || async move {
let (entities, relations) = state_clone.read_graph(|g| (g.entities.len(), g.relations.len())); let (entities, relations) = state_clone.read_graph(|g| (g.entities.len(), g.relations.len()));
let tasks = state_clone.tasks.read_with(|items| items.len()); let tasks = state_clone.project.tasks.read_with(|items| items.len());
let snippets = state_clone.snippets.read_with(|items| items.len()); let snippets = state_clone.code.snippets.read_with(|items| items.len());
let tech_debts = state_clone.tech_debts.read_with(|items| items.len()); let tech_debts = state_clone.code.tech_debts.read_with(|items| items.len());
let adrs = state_clone.adrs.read_with(|items| items.len()); let adrs = state_clone.code.adrs.read_with(|items| items.len());
let ledger = state_clone.ledger.read_with(|items| items.len()); let ledger = state_clone.code.ledger.read_with(|items| items.len());
let sticky = state_clone.sticky.read_with(|items| items.len()); let sticky = state_clone.code.sticky.read_with(|items| items.len());
let error_fixes = state_clone.error_fixes.read_with(|items| items.len()); let error_fixes = state_clone.code.error_fixes.read_with(|items| items.len());
let pinned_files = state_clone.pinned_files.read_with(|items| items.len()); let pinned_files = state_clone.project.pinned_files.read_with(|items| items.len());
let session_summaries = state_clone.session_summaries.read_with(|items| items.len()); let session_summaries = state_clone.telemetry.session_summaries.read_with(|items| items.len());
let handoff_memos = state_clone.handoff_memos.read_with(|items| items.len()); let handoff_memos = state_clone.telemetry.handoff_memos.read_with(|items| items.len());
let env_fingerprints = state_clone.env_fingerprints.read_with(|items| items.len()); let env_fingerprints = state_clone.env.env_fingerprints.read_with(|items| items.len());
let env_requirements = state_clone.env_requirements.read_with(|items| items.len()); let env_requirements = state_clone.env.env_requirements.read_with(|items| items.len());
let milestones = state_clone.milestones.read_with(|items| items.len()); let milestones = state_clone.project.milestones.read_with(|items| items.len());
let environments = state_clone.environments.read_with(|items| items.len()); let environments = state_clone.env.environments.read_with(|items| items.len());
let pr_checklists = state_clone.pr_checklists.read_with(|items| items.len()); let pr_checklists = state_clone.project.pr_checklists.read_with(|items| items.len());
let gates = state_clone.gates.read_with(|items| items.len()); let gates = state_clone.env.gates.read_with(|items| items.len());
let context_workspaces = state_clone.context_workspaces.read_with(|items| items.len()); let context_workspaces = state_clone.project.context_workspaces.read_with(|items| items.len());
axum::Json(serde_json::json!({ axum::Json(serde_json::json!({
"entities": entities, "entities": entities,
+2 -2
View File
@@ -79,7 +79,7 @@ use crate::models::TerminalHistory;
pub async fn get_terminal_history_handler( pub async fn get_terminal_history_handler(
State(state): State<Arc<AppState>>, State(state): State<Arc<AppState>>,
) -> impl axum::response::IntoResponse { ) -> impl axum::response::IntoResponse {
let history_json = state.handler.state.terminal_history.read_with(|h| serde_json::to_string(h).unwrap_or_else(|_| "[]".to_string())); let history_json = state.handler.state.telemetry.terminal_history.read_with(|h| serde_json::to_string(h).unwrap_or_else(|_| "[]".to_string()));
([(axum::http::header::CONTENT_TYPE, "application/json")], history_json) ([(axum::http::header::CONTENT_TYPE, "application/json")], history_json)
} }
@@ -87,7 +87,7 @@ pub async fn terminal_telemetry_handler(
State(state): State<Arc<AppState>>, State(state): State<Arc<AppState>>,
axum::Json(payload): axum::Json<TerminalHistory>, axum::Json(payload): axum::Json<TerminalHistory>,
) -> impl axum::response::IntoResponse { ) -> impl axum::response::IntoResponse {
state.handler.state.terminal_history.modify(|history| { state.handler.state.telemetry.terminal_history.modify(|history| {
history.push_front(payload.clone()); history.push_front(payload.clone());
if history.len() > 100 { if history.len() > 100 {
history.pop_back(); history.pop_back();
+1 -1
View File
@@ -33,7 +33,7 @@ pub fn spawn_watcher(state: Arc<MemoryState>) {
expires_at: None, expires_at: None,
}; };
state.sticky.modify(|notes| { state.code.sticky.modify(|notes| {
notes.push(note.clone()); notes.push(note.clone());
}); });
+6 -6
View File
@@ -23,7 +23,7 @@ impl McpTool for UpdateEnvFingerprintHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: UpdateEnvFingerprintTool = let req: UpdateEnvFingerprintTool =
serde_json::from_value(args).map_err(|e| e.to_string())?; serde_json::from_value(args).map_err(|e| e.to_string())?;
state.env_fingerprints.modify(|fps| { state.env.env_fingerprints.modify(|fps| {
fps.insert( fps.insert(
req.namespace.clone(), req.namespace.clone(),
crate::models::EnvFingerprint { crate::models::EnvFingerprint {
@@ -58,7 +58,7 @@ impl McpTool for ReadEnvFingerprintHandler {
let req: ReadEnvFingerprintTool = let req: ReadEnvFingerprintTool =
serde_json::from_value(args).map_err(|e| e.to_string())?; serde_json::from_value(args).map_err(|e| e.to_string())?;
let data = state let data = state
.env_fingerprints .env.env_fingerprints
.read_with(|fps| fps.get(&req.namespace).cloned()); .read_with(|fps| fps.get(&req.namespace).cloned());
if let Some(fp) = data { if let Some(fp) = data {
let data = Ok::<String, crate::error::AppError>(serde_json::to_string(&fp)?)?; let data = Ok::<String, crate::error::AppError>(serde_json::to_string(&fp)?)?;
@@ -86,7 +86,7 @@ impl McpTool for LogEnvRequirementHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: LogEnvRequirementTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: LogEnvRequirementTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
state.env_requirements.modify(|reqs| { state.env.env_requirements.modify(|reqs| {
reqs.retain(|r| !(r.namespace == req.namespace && r.key == req.key)); reqs.retain(|r| !(r.namespace == req.namespace && r.key == req.key));
reqs.push(crate::models::EnvRequirement { reqs.push(crate::models::EnvRequirement {
namespace: req.namespace, namespace: req.namespace,
@@ -117,7 +117,7 @@ impl McpTool for RegisterEnvironmentHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: RegisterEnvironmentTool = let req: RegisterEnvironmentTool =
serde_json::from_value(args).map_err(|e| e.to_string())?; serde_json::from_value(args).map_err(|e| e.to_string())?;
state.environments.modify(|envs| { state.env.environments.modify(|envs| {
envs.retain(|e| !(e.namespace == req.namespace && e.name == req.name)); envs.retain(|e| !(e.namespace == req.namespace && e.name == req.name));
envs.push(crate::models::EnvironmentDetail { envs.push(crate::models::EnvironmentDetail {
namespace: req.namespace, namespace: req.namespace,
@@ -150,7 +150,7 @@ impl McpTool for GetEnvironmentDetailsHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: GetEnvironmentDetailsTool = let req: GetEnvironmentDetailsTool =
serde_json::from_value(args).map_err(|e| e.to_string())?; serde_json::from_value(args).map_err(|e| e.to_string())?;
let data = state.environments.read_with(|envs| { let data = state.env.environments.read_with(|envs| {
let filtered: Vec<_> = envs let filtered: Vec<_> = envs
.iter() .iter()
.filter(|e| e.namespace == req.namespace) .filter(|e| e.namespace == req.namespace)
@@ -197,7 +197,7 @@ mod tests {
let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap())); let state = Arc::new(MemoryState::new(dir.path().to_str().unwrap()));
// Ensure namespace is present in test setup // Ensure namespace is present in test setup
state.environments.modify(|e| { state.env.environments.modify(|e| {
e.push(crate::models::EnvironmentDetail { e.push(crate::models::EnvironmentDetail {
namespace: "global".to_string(), namespace: "global".to_string(),
name: "test".to_string(), name: "test".to_string(),
+24 -24
View File
@@ -24,7 +24,7 @@ impl McpTool for LogDecisionHandler {
let idx = state.get_search_index(); let idx = state.get_search_index();
let mut final_id = String::new(); let mut final_id = String::new();
state.adrs.modify(|adrs| { state.code.adrs.modify(|adrs| {
if let Some(superseded_id) = &req.supersedes { if let Some(superseded_id) = &req.supersedes {
for old_adr in adrs.iter_mut() { for old_adr in adrs.iter_mut() {
if old_adr.id == *superseded_id { if old_adr.id == *superseded_id {
@@ -70,7 +70,7 @@ impl McpTool for QueryDecisionsHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: QueryDecisionsTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: QueryDecisionsTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let data = state.adrs.read_with(|adrs| { let data = state.code.adrs.read_with(|adrs| {
let filtered: Vec<_> = adrs let filtered: Vec<_> = adrs
.iter() .iter()
.filter(|a| { .filter(|a| {
@@ -108,7 +108,7 @@ impl McpTool for DeleteDecisionHandler {
let req: crate::tools::DeleteDecisionTool = let req: crate::tools::DeleteDecisionTool =
serde_json::from_value(args).map_err(|e| e.to_string())?; serde_json::from_value(args).map_err(|e| e.to_string())?;
let mut found = false; let mut found = false;
state.adrs.modify(|adrs| { state.code.adrs.modify(|adrs| {
if let Some(pos) = adrs.iter().position(|a| a.id == req.id) { if let Some(pos) = adrs.iter().position(|a| a.id == req.id) {
adrs.remove(pos); adrs.remove(pos);
found = true; found = true;
@@ -140,7 +140,7 @@ impl McpTool for LogErrorFixHandler {
let req: LogErrorFixTool = serde_json::from_value(args).map_err(|e| e.to_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 text_to_embed = format!("Signature: {}\nSolution: {}", req.signature, req.solution);
let embedding = crate::embedding::generate_embedding_async(text_to_embed).await.ok(); let embedding = crate::embedding::generate_embedding_async(text_to_embed).await.ok();
state.error_fixes.modify(|fixes| { state.code.error_fixes.modify(|fixes| {
fixes.push(crate::models::ErrorFix { fixes.push(crate::models::ErrorFix {
signature: req.signature, signature: req.signature,
solution: req.solution, solution: req.solution,
@@ -172,7 +172,7 @@ impl McpTool for SearchErrorFixesHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { 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 req: SearchErrorFixesTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let q = req.query; let q = req.query;
let data = state.error_fixes.read_with(|fixes| { let data = state.code.error_fixes.read_with(|fixes| {
let filtered: Vec<_> = fixes let filtered: Vec<_> = fixes
.iter() .iter()
.filter(|f| { .filter(|f| {
@@ -200,7 +200,7 @@ impl McpTool for LogCodeChangeHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { 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 req: LogCodeChangeTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
state.ledger.modify(|ledger| { state.code.ledger.modify(|ledger| {
ledger.push(CodeChange { ledger.push(CodeChange {
timestamp: crate::handlers::utils::now_secs(), timestamp: crate::handlers::utils::now_secs(),
file_path: req.file_path, file_path: req.file_path,
@@ -230,7 +230,7 @@ impl McpTool for QueryRecentChangesHandler {
async fn execute(&self, _args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, _args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let data = state let data = state
.ledger .code.ledger
.read_with(|l| Ok::<String, crate::error::AppError>(serde_json::to_string(l)?))?; .read_with(|l| Ok::<String, crate::error::AppError>(serde_json::to_string(l)?))?;
Ok(data) Ok(data)
} }
@@ -250,7 +250,7 @@ impl McpTool for LearnPreferenceHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: LearnPreferenceTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: LearnPreferenceTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
state.prefs.modify(|prefs| { state.env.prefs.modify(|prefs| {
prefs.insert( prefs.insert(
req.key.clone(), req.key.clone(),
crate::models::Preference { crate::models::Preference {
@@ -278,7 +278,7 @@ impl McpTool for ReadPreferencesHandler {
async fn execute(&self, _args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, _args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
state state
.prefs .env.prefs
.read_with(|prefs| Ok::<String, crate::error::AppError>(serde_json::to_string(prefs)?)) .read_with(|prefs| Ok::<String, crate::error::AppError>(serde_json::to_string(prefs)?))
} }
} }
@@ -299,7 +299,7 @@ impl McpTool for LogTechDebtHandler {
let req: LogTechDebtTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: LogTechDebtTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let text_to_embed = format!("Description: {}\nIdeal Solution: {}", req.description, req.ideal_solution); let text_to_embed = format!("Description: {}\nIdeal Solution: {}", req.description, req.ideal_solution);
let embedding = crate::embedding::generate_embedding_async(text_to_embed).await.ok(); let embedding = crate::embedding::generate_embedding_async(text_to_embed).await.ok();
state.tech_debts.modify(|debts| { state.code.tech_debts.modify(|debts| {
debts.push(crate::models::TechDebt { debts.push(crate::models::TechDebt {
id: uuid::Uuid::new_v4().to_string(), id: uuid::Uuid::new_v4().to_string(),
namespace: req.namespace, namespace: req.namespace,
@@ -334,7 +334,7 @@ impl McpTool for ResolveTechDebtHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: ResolveTechDebtTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: ResolveTechDebtTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let mut found = false; let mut found = false;
state.tech_debts.modify(|debts| { state.code.tech_debts.modify(|debts| {
for d in debts.iter_mut() { for d in debts.iter_mut() {
if d.id == req.id { if d.id == req.id {
d.is_resolved = true; d.is_resolved = true;
@@ -365,7 +365,7 @@ impl McpTool for ListTechDebtHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: ListTechDebtTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: ListTechDebtTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let data = state.tech_debts.read_with(|debts| { let data = state.code.tech_debts.read_with(|debts| {
let filtered: Vec<_> = debts let filtered: Vec<_> = debts
.iter() .iter()
.filter(|d| { .filter(|d| {
@@ -451,7 +451,7 @@ impl McpTool for OmniSearchHandler {
} }
} }
let tasks_json = state.tasks.read_with(|all_tasks| { let tasks_json = state.project.tasks.read_with(|all_tasks| {
let filtered: Vec<_> = all_tasks let filtered: Vec<_> = all_tasks
.iter() .iter()
.filter(|t| matched_tasks.contains(t.id.as_str())) .filter(|t| matched_tasks.contains(t.id.as_str()))
@@ -470,7 +470,7 @@ impl McpTool for OmniSearchHandler {
serde_json::to_value(&filtered).map_err(|e| e.to_string()) serde_json::to_value(&filtered).map_err(|e| e.to_string())
})?; })?;
let snippets_json = state.snippets.read_with(|all_snippets| { let snippets_json = state.code.snippets.read_with(|all_snippets| {
let mut scored: Vec<_> = all_snippets.iter().map(|s| { let mut scored: Vec<_> = all_snippets.iter().map(|s| {
let mut score = 0.0; 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) { if matched_snippets.contains(s.name.as_str()) || contains_ignore_ascii_case(&s.name, &q) || contains_ignore_ascii_case(&s.description, &q) {
@@ -495,7 +495,7 @@ impl McpTool for OmniSearchHandler {
serde_json::to_value(&filtered).map_err(|e| e.to_string()) serde_json::to_value(&filtered).map_err(|e| e.to_string())
})?; })?;
let adrs_json = state.adrs.read_with(|all_adrs| { let adrs_json = state.code.adrs.read_with(|all_adrs| {
let filtered: Vec<_> = all_adrs let filtered: Vec<_> = all_adrs
.iter() .iter()
.filter(|a| matched_adrs.contains(a.id.as_str())) .filter(|a| matched_adrs.contains(a.id.as_str()))
@@ -516,7 +516,7 @@ impl McpTool for OmniSearchHandler {
})?; })?;
let q = req.query; let q = req.query;
let tech_debts_json = state.tech_debts.read_with(|debts| { let tech_debts_json = state.code.tech_debts.read_with(|debts| {
let mut scored: Vec<_> = debts.iter().map(|d| { let mut scored: Vec<_> = debts.iter().map(|d| {
if req.namespace.as_ref().is_some_and(|ns| d.namespace != *ns) { if req.namespace.as_ref().is_some_and(|ns| d.namespace != *ns) {
return (d, 0.0); return (d, 0.0);
@@ -545,7 +545,7 @@ impl McpTool for OmniSearchHandler {
serde_json::to_value(&filtered).map_err(|e| e.to_string()) serde_json::to_value(&filtered).map_err(|e| e.to_string())
})?; })?;
let memos_json = state.handoff_memos.read_with(|memos| { let memos_json = state.telemetry.handoff_memos.read_with(|memos| {
let filtered: Vec<_> = memos let filtered: Vec<_> = memos
.iter() .iter()
.filter(|m| { .filter(|m| {
@@ -566,7 +566,7 @@ impl McpTool for OmniSearchHandler {
serde_json::to_value(&filtered).map_err(|e| e.to_string()) serde_json::to_value(&filtered).map_err(|e| e.to_string())
})?; })?;
let error_fixes_json = state.error_fixes.read_with(|fixes| { let error_fixes_json = state.code.error_fixes.read_with(|fixes| {
let mut scored: Vec<_> = fixes.iter().map(|f| { let mut scored: Vec<_> = fixes.iter().map(|f| {
let mut score = 0.0; let mut score = 0.0;
if contains_ignore_ascii_case(&f.signature, &q) || contains_ignore_ascii_case(&f.solution, &q) { if contains_ignore_ascii_case(&f.signature, &q) || contains_ignore_ascii_case(&f.solution, &q) {
@@ -614,27 +614,27 @@ impl McpTool for GetProjectHealthHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { 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 req: GetProjectHealthTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let active_tasks = state let active_tasks = state
.tasks .project.tasks
.read_with(|tasks| tasks.iter().filter(|t| t.status != "done").count()); .read_with(|tasks| tasks.iter().filter(|t| t.status != "done").count());
let unresolved_debt = state.tech_debts.read_with(|debts| { let unresolved_debt = state.code.tech_debts.read_with(|debts| {
debts debts
.iter() .iter()
.filter(|d| d.namespace == req.namespace && !d.is_resolved) .filter(|d| d.namespace == req.namespace && !d.is_resolved)
.count() .count()
}); });
let unread_memos = state.handoff_memos.read_with(|memos| { let unread_memos = state.telemetry.handoff_memos.read_with(|memos| {
memos memos
.iter() .iter()
.filter(|m| m.namespace == req.namespace) .filter(|m| m.namespace == req.namespace)
.count() .count()
}); });
let active_milestones = state.milestones.read_with(|milestones| { let active_milestones = state.project.milestones.read_with(|milestones| {
milestones milestones
.iter() .iter()
.filter(|m| m.namespace == req.namespace && m.status != "done") .filter(|m| m.namespace == req.namespace && m.status != "done")
.count() .count()
}); });
let remaining_checklists = state.pr_checklists.read_with(|checklists| { let remaining_checklists = state.project.pr_checklists.read_with(|checklists| {
checklists checklists
.iter() .iter()
.filter(|c| c.namespace == req.namespace) .filter(|c| c.namespace == req.namespace)
@@ -821,7 +821,7 @@ mod tests {
}; };
{ {
state.tasks.modify(|t| { state.project.tasks.modify(|t| {
t.push(task.clone()); t.push(task.clone());
}); });
} }
+11 -11
View File
@@ -21,7 +21,7 @@ impl McpTool for AddStickyNoteHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: AddStickyNoteTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: AddStickyNoteTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
state.sticky.modify(|notes| { state.code.sticky.modify(|notes| {
notes.push(StickyNote { notes.push(StickyNote {
timestamp: crate::handlers::utils::now_secs(), timestamp: crate::handlers::utils::now_secs(),
content: req.content, content: req.content,
@@ -49,7 +49,7 @@ impl McpTool for ReadStickyNotesHandler {
async fn execute(&self, _args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, _args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let data = state let data = state
.sticky .code.sticky
.read_with(|s| Ok::<String, crate::error::AppError>(serde_json::to_string(s)?))?; .read_with(|s| Ok::<String, crate::error::AppError>(serde_json::to_string(s)?))?;
Ok(data) Ok(data)
} }
@@ -73,7 +73,7 @@ impl McpTool for DeleteStickyNoteHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: DeleteStickyNoteTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: DeleteStickyNoteTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let mut success = false; let mut success = false;
state.sticky.modify(|notes| { state.code.sticky.modify(|notes| {
if req.index > 0 && req.index <= notes.len() { if req.index > 0 && req.index <= notes.len() {
notes.remove(req.index - 1); notes.remove(req.index - 1);
success = true; success = true;
@@ -103,7 +103,7 @@ impl McpTool for ClearStickyNotesHandler {
} }
async fn execute(&self, _args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, _args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
state.sticky.modify(|notes| { state.code.sticky.modify(|notes| {
notes.clear(); notes.clear();
}); });
Ok("All sticky notes cleared.".to_string()) Ok("All sticky notes cleared.".to_string())
@@ -127,7 +127,7 @@ impl McpTool for LeaveHandoffMemoHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: LeaveHandoffMemoTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: LeaveHandoffMemoTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
state.handoff_memos.modify(|memos| { state.telemetry.handoff_memos.modify(|memos| {
memos.push(crate::models::HandoffMemo { memos.push(crate::models::HandoffMemo {
id: uuid::Uuid::new_v4().to_string(), id: uuid::Uuid::new_v4().to_string(),
author: "agy".to_string(), author: "agy".to_string(),
@@ -158,7 +158,7 @@ impl McpTool for ReadHandoffMemosHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: ReadHandoffMemosTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: ReadHandoffMemosTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let data = state.handoff_memos.read_with(|items| { let data = state.telemetry.handoff_memos.read_with(|items| {
let filtered: Vec<_> = items let filtered: Vec<_> = items
.iter() .iter()
.filter(|i| { .filter(|i| {
@@ -194,7 +194,7 @@ impl McpTool for ClearHandoffMemosHandler {
let req: ClearHandoffMemosTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: ClearHandoffMemosTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let ids: HashSet<_> = req.ids.into_iter().collect(); let ids: HashSet<_> = req.ids.into_iter().collect();
state state
.handoff_memos .telemetry.handoff_memos
.modify(|memos| memos.retain(|m| !ids.contains(&m.id))); .modify(|memos| memos.retain(|m| !ids.contains(&m.id)));
Ok("Handoff memos cleared".to_string()) Ok("Handoff memos cleared".to_string())
} }
@@ -217,7 +217,7 @@ impl McpTool for AddSessionSummaryHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: AddSessionSummaryTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: AddSessionSummaryTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
state.session_summaries.modify(|summaries| { state.telemetry.session_summaries.modify(|summaries| {
summaries.push(crate::models::SessionSummary { summaries.push(crate::models::SessionSummary {
summary: req.summary, summary: req.summary,
namespace: req.namespace, namespace: req.namespace,
@@ -249,9 +249,9 @@ impl McpTool for GenerateStandupReportHandler {
serde_json::from_value(args).map_err(|e| e.to_string())?; serde_json::from_value(args).map_err(|e| e.to_string())?;
let cutoff = crate::handlers::utils::now_secs().saturating_sub(req.hours_lookback * 3600); let cutoff = crate::handlers::utils::now_secs().saturating_sub(req.hours_lookback * 3600);
let report_str = state.tasks.read_with(|items| { let report_str = state.project.tasks.read_with(|items| {
state.ledger.read_with(|changes| { state.code.ledger.read_with(|changes| {
state.session_summaries.read_with(|summaries| { state.telemetry.session_summaries.read_with(|summaries| {
let filtered_tasks: Vec<_> = items.iter().filter(|t| t.updated_at >= cutoff).collect(); let filtered_tasks: Vec<_> = items.iter().filter(|t| t.updated_at >= cutoff).collect();
let filtered_changes: Vec<_> = changes.iter().filter(|c| c.timestamp >= cutoff).collect(); let filtered_changes: Vec<_> = changes.iter().filter(|c| c.timestamp >= cutoff).collect();
let filtered_summaries: Vec<_> = summaries.iter().filter(|s| s.namespace == req.namespace && s.timestamp >= cutoff).collect(); let filtered_summaries: Vec<_> = summaries.iter().filter(|s| s.namespace == req.namespace && s.timestamp >= cutoff).collect();
+10 -10
View File
@@ -40,7 +40,7 @@ impl McpTool for AddTaskHandler {
}; };
let idx = state.get_search_index(); let idx = state.get_search_index();
drop(idx.index_task(&task)); drop(idx.index_task(&task));
state.tasks.modify(|tasks| { state.project.tasks.modify(|tasks| {
tasks.push(task); tasks.push(task);
}); });
Ok(format!("Task added with ID: {}", task_id).to_string()) Ok(format!("Task added with ID: {}", task_id).to_string())
@@ -63,7 +63,7 @@ impl McpTool for DeleteTaskHandler {
let req: DeleteTaskTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: DeleteTaskTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let mut deleted_count = 0; let mut deleted_count = 0;
let mut actually_deleted = Vec::new(); let mut actually_deleted = Vec::new();
state.tasks.modify(|tasks| { state.project.tasks.modify(|tasks| {
let initial_len = tasks.len(); let initial_len = tasks.len();
// Build index-based children map // Build index-based children map
@@ -141,7 +141,7 @@ impl McpTool for UpdateTaskStatusHandler {
let mut blocker_details = String::new(); let mut blocker_details = String::new();
let target_status = req.status.to_lowercase(); let target_status = req.status.to_lowercase();
state.tasks.modify(|tasks| { state.project.tasks.modify(|tasks| {
// Find target task // Find target task
let target_idx = tasks let target_idx = tasks
.iter() .iter()
@@ -279,7 +279,7 @@ impl McpTool for ListActiveTasksHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: ListActiveTasksTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: ListActiveTasksTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let data = state.tasks.read_with(|tasks| { let data = state.project.tasks.read_with(|tasks| {
let filtered: Vec<_> = tasks let filtered: Vec<_> = tasks
.iter() .iter()
.filter(|t| { .filter(|t| {
@@ -319,7 +319,7 @@ impl McpTool for SetAcceptanceCriteriaHandler {
let req: SetAcceptanceCriteriaTool = let req: SetAcceptanceCriteriaTool =
serde_json::from_value(args).map_err(|e| e.to_string())?; serde_json::from_value(args).map_err(|e| e.to_string())?;
let mut success = false; let mut success = false;
state.tasks.modify(|tasks| { state.project.tasks.modify(|tasks| {
if let Some(task) = tasks.iter_mut().rev().find(|t| t.title == req.task_title) { if let Some(task) = tasks.iter_mut().rev().find(|t| t.title == req.task_title) {
task.acceptance_criteria = req task.acceptance_criteria = req
.criteria .criteria
@@ -362,7 +362,7 @@ impl McpTool for VerifyAcceptanceCriteriaHandler {
serde_json::from_value(args).map_err(|e| e.to_string())?; serde_json::from_value(args).map_err(|e| e.to_string())?;
let mut success = false; let mut success = false;
let mut already_met = false; let mut already_met = false;
state.tasks.modify(|tasks| { state.project.tasks.modify(|tasks| {
if let Some(task) = tasks.iter_mut().find(|t| t.id == req.task_id) if let Some(task) = tasks.iter_mut().find(|t| t.id == req.task_id)
&& let Some(ac) = task && let Some(ac) = task
.acceptance_criteria .acceptance_criteria
@@ -405,7 +405,7 @@ impl McpTool for AddMilestoneHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: AddMilestoneTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: AddMilestoneTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
state.milestones.modify(|ms| { state.project.milestones.modify(|ms| {
ms.push(crate::models::Milestone { ms.push(crate::models::Milestone {
id: uuid::Uuid::new_v4().to_string(), id: uuid::Uuid::new_v4().to_string(),
title: req.title, title: req.title,
@@ -433,7 +433,7 @@ impl McpTool for UpdateMilestoneHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: UpdateMilestoneTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: UpdateMilestoneTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let mut found = false; let mut found = false;
state.milestones.modify(|ms| { state.project.milestones.modify(|ms| {
for m in ms.iter_mut() { for m in ms.iter_mut() {
if m.id == req.id { if m.id == req.id {
m.status = req.status.clone(); m.status = req.status.clone();
@@ -464,7 +464,7 @@ impl McpTool for ListMilestonesHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: ListMilestonesTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: ListMilestonesTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let data = state.milestones.read_with(|items| { let data = state.project.milestones.read_with(|items| {
let filtered: Vec<_> = items let filtered: Vec<_> = items
.iter() .iter()
.filter(|i| { .filter(|i| {
@@ -559,7 +559,7 @@ mod tests {
assert!(res1.contains("Milestone added")); assert!(res1.contains("Milestone added"));
// Fetch milestone ID from state directly to update // Fetch milestone ID from state directly to update
let ms_id = state.milestones.read_with(|ms| ms[0].id.clone()); let ms_id = state.project.milestones.read_with(|ms| ms[0].id.clone());
// Update Milestone // Update Milestone
let update_ms = UpdateMilestoneHandler; let update_ms = UpdateMilestoneHandler;
+15 -15
View File
@@ -20,7 +20,7 @@ impl McpTool for PinFileHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: PinFileTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: PinFileTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
state.pinned_files.modify(|pinned| { state.project.pinned_files.modify(|pinned| {
pinned.retain(|p| p.namespace != req.namespace || p.file_path != req.file_path); pinned.retain(|p| p.namespace != req.namespace || p.file_path != req.file_path);
pinned.push(crate::models::PinnedFile { pinned.push(crate::models::PinnedFile {
namespace: req.namespace, namespace: req.namespace,
@@ -47,7 +47,7 @@ impl McpTool for UnpinFileHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: UnpinFileTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: UnpinFileTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
state.pinned_files.modify(|pinned| { state.project.pinned_files.modify(|pinned| {
pinned.retain(|p| p.namespace != req.namespace || p.file_path != req.file_path) pinned.retain(|p| p.namespace != req.namespace || p.file_path != req.file_path)
}); });
Ok("File unpinned".to_string()) Ok("File unpinned".to_string())
@@ -71,7 +71,7 @@ impl McpTool for ListPinnedFilesHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: ListPinnedFilesTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: ListPinnedFilesTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let data = state.pinned_files.read_with(|pinned| { let data = state.project.pinned_files.read_with(|pinned| {
let filtered: Vec<_> = pinned let filtered: Vec<_> = pinned
.iter() .iter()
.filter(|p| { .filter(|p| {
@@ -124,7 +124,7 @@ impl McpTool for StoreSnippetHandler {
let idx = state.get_search_index(); let idx = state.get_search_index();
drop(idx.index_snippet(&snippet)); drop(idx.index_snippet(&snippet));
state.snippets.modify(|snippets| { state.code.snippets.modify(|snippets| {
snippets.retain(|s| s.name != req_name); snippets.retain(|s| s.name != req_name);
snippets.push(snippet); snippets.push(snippet);
}); });
@@ -148,7 +148,7 @@ impl McpTool for SearchSnippetsHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: SearchSnippetsTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: SearchSnippetsTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let query = req.query; let query = req.query;
let data = state.snippets.read_with(|snippets| { let data = state.code.snippets.read_with(|snippets| {
let results: Vec<_> = snippets let results: Vec<_> = snippets
.iter() .iter()
.filter(|s| { .filter(|s| {
@@ -178,7 +178,7 @@ impl McpTool for DeleteSnippetHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: DeleteSnippetTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: DeleteSnippetTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let mut deleted = false; let mut deleted = false;
state.snippets.modify(|snippets| { state.code.snippets.modify(|snippets| {
let orig = snippets.len(); let orig = snippets.len();
snippets.retain(|s| s.name != req.name); snippets.retain(|s| s.name != req.name);
deleted = snippets.len() < orig; deleted = snippets.len() < orig;
@@ -211,7 +211,7 @@ impl McpTool for SaveContextWorkspaceHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: SaveContextWorkspaceTool = let req: SaveContextWorkspaceTool =
serde_json::from_value(args).map_err(|e| e.to_string())?; serde_json::from_value(args).map_err(|e| e.to_string())?;
state.context_workspaces.modify(|ws| { state.project.context_workspaces.modify(|ws| {
ws.retain(|w| w.namespace != req.namespace || w.name != req.name); ws.retain(|w| w.namespace != req.namespace || w.name != req.name);
ws.push(crate::models::ContextWorkspace { ws.push(crate::models::ContextWorkspace {
namespace: req.namespace, namespace: req.namespace,
@@ -243,7 +243,7 @@ impl McpTool for LoadContextWorkspaceHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: LoadContextWorkspaceTool = let req: LoadContextWorkspaceTool =
serde_json::from_value(args).map_err(|e| e.to_string())?; serde_json::from_value(args).map_err(|e| e.to_string())?;
let data = state.context_workspaces.read_with(|ws| { let data = state.project.context_workspaces.read_with(|ws| {
let filtered: Vec<_> = ws let filtered: Vec<_> = ws
.iter() .iter()
.filter(|w| w.namespace == req.namespace && w.name == req.name) .filter(|w| w.namespace == req.namespace && w.name == req.name)
@@ -272,7 +272,7 @@ impl McpTool for ListContextWorkspacesHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: ListContextWorkspacesTool = let req: ListContextWorkspacesTool =
serde_json::from_value(args).map_err(|e| e.to_string())?; serde_json::from_value(args).map_err(|e| e.to_string())?;
let data = state.context_workspaces.read_with(|ws| { let data = state.project.context_workspaces.read_with(|ws| {
let filtered: Vec<_> = ws let filtered: Vec<_> = ws
.iter() .iter()
.filter(|w| req.namespace.as_ref().is_none_or(|ns| &w.namespace == ns)) .filter(|w| req.namespace.as_ref().is_none_or(|ns| &w.namespace == ns))
@@ -303,7 +303,7 @@ impl McpTool for DeleteContextWorkspaceHandler {
serde_json::from_value(args).map_err(|e| e.to_string())?; serde_json::from_value(args).map_err(|e| e.to_string())?;
let mut found = false; let mut found = false;
state.context_workspaces.modify(|ws| { state.project.context_workspaces.modify(|ws| {
if let Some(pos) = ws if let Some(pos) = ws
.iter() .iter()
.position(|w| w.namespace == req.namespace && w.name == req.name) .position(|w| w.namespace == req.namespace && w.name == req.name)
@@ -339,7 +339,7 @@ impl McpTool for AddPrChecklistItemHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: AddPrChecklistItemTool = let req: AddPrChecklistItemTool =
serde_json::from_value(args).map_err(|e| e.to_string())?; serde_json::from_value(args).map_err(|e| e.to_string())?;
state.pr_checklists.modify(|items| { state.project.pr_checklists.modify(|items| {
items.push(crate::models::PrChecklistItem { items.push(crate::models::PrChecklistItem {
namespace: req.namespace, namespace: req.namespace,
id: uuid::Uuid::new_v4().to_string(), id: uuid::Uuid::new_v4().to_string(),
@@ -364,7 +364,7 @@ impl McpTool for GetPrChecklistHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: GetPrChecklistTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: GetPrChecklistTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
let data = state.pr_checklists.read_with(|items| { let data = state.project.pr_checklists.read_with(|items| {
let filtered: Vec<_> = items let filtered: Vec<_> = items
.iter() .iter()
.filter(|i| i.namespace == req.namespace) .filter(|i| i.namespace == req.namespace)
@@ -393,7 +393,7 @@ impl McpTool for ClearPrChecklistHandler {
async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn execute(&self, args: Value, state: Arc<MemoryState>) -> crate::error::Result<String> {
let req: ClearPrChecklistTool = serde_json::from_value(args).map_err(|e| e.to_string())?; let req: ClearPrChecklistTool = serde_json::from_value(args).map_err(|e| e.to_string())?;
state state
.pr_checklists .project.pr_checklists
.modify(|items| items.retain(|i| i.namespace != req.namespace)); .modify(|items| items.retain(|i| i.namespace != req.namespace));
Ok("PR checklist cleared".to_string()) Ok("PR checklist cleared".to_string())
} }
@@ -625,14 +625,14 @@ impl McpTool for SemanticCodeSearchHandler {
let mut texts_to_embed = Vec::new(); let mut texts_to_embed = Vec::new();
let mut metadata = Vec::new(); let mut metadata = Vec::new();
let snippets = state.snippets.read_with(|snips| snips.clone()); let snippets = state.code.snippets.read_with(|snips| snips.clone());
for snippet in snippets { for snippet in snippets {
let combined = format!("{} {} {}", snippet.name, snippet.description, snippet.code); let combined = format!("{} {} {}", snippet.name, snippet.description, snippet.code);
texts_to_embed.push(combined); texts_to_embed.push(combined);
metadata.push((snippet.name, snippet.description)); metadata.push((snippet.name, snippet.description));
} }
let sticky = state.sticky.read_with(|s| s.clone()); let sticky = state.code.sticky.read_with(|s| s.clone());
for note in sticky { for note in sticky {
texts_to_embed.push(note.content.clone()); texts_to_embed.push(note.content.clone());
metadata.push(("StickyNote".to_string(), note.content.chars().take(200).collect::<String>())); metadata.push(("StickyNote".to_string(), note.content.chars().take(200).collect::<String>()));
+1 -1
View File
@@ -81,7 +81,7 @@ pub async fn start_background_indexer(state: Arc<MemoryState>) {
embedding, embedding,
}; };
state.snippets.modify(|snippets| { state.code.snippets.modify(|snippets| {
// Prevent duplicates if already indexed // Prevent duplicates if already indexed
if !snippets.iter().any(|s| s.name == snippet.name) { if !snippets.iter().any(|s| s.name == snippet.name) {
snippets.push(snippet.clone()); snippets.push(snippet.clone());
+7 -7
View File
@@ -103,16 +103,16 @@ async fn ttl_sweeper_worker(state: Arc<MemoryState>) {
.unwrap_or_default() .unwrap_or_default()
.as_secs(); .as_secs();
state.tasks.modify(|tasks| { state.project.tasks.modify(|tasks| {
tasks.retain(|t| t.expires_at.is_none_or(|exp| exp > now)); tasks.retain(|t| t.expires_at.is_none_or(|exp| exp > now));
}); });
state.sticky.modify(|notes| { state.code.sticky.modify(|notes| {
notes.retain(|n| n.expires_at.is_none_or(|exp| exp > now)); notes.retain(|n| n.expires_at.is_none_or(|exp| exp > now));
}); });
state.handoff_memos.modify(|memos| { state.telemetry.handoff_memos.modify(|memos| {
memos.retain(|m| m.expires_at.is_none_or(|exp| exp > now)); memos.retain(|m| m.expires_at.is_none_or(|exp| exp > now));
}); });
state.session_summaries.modify(|summaries| { state.telemetry.session_summaries.modify(|summaries| {
summaries.retain(|s| s.expires_at.is_none_or(|exp| exp > now)); summaries.retain(|s| s.expires_at.is_none_or(|exp| exp > now));
}); });
} }
@@ -144,7 +144,7 @@ async fn condense_graph_worker(state: Arc<MemoryState>) {
// Condense sticky notes // Condense sticky notes
let mut condensed_sticky_content = String::new(); let mut condensed_sticky_content = String::new();
state.sticky.modify(|notes| { state.code.sticky.modify(|notes| {
if notes.len() > threshold { if notes.len() > threshold {
notes.sort_by_key(|n| n.timestamp); notes.sort_by_key(|n| n.timestamp);
let to_remove = notes.len() - (threshold / 2); let to_remove = notes.len() - (threshold / 2);
@@ -174,7 +174,7 @@ async fn condense_graph_worker(state: Arc<MemoryState>) {
// Condense snippets // Condense snippets
let mut condensed_snippet_content = String::new(); let mut condensed_snippet_content = String::new();
state.snippets.modify(|snippets| { state.code.snippets.modify(|snippets| {
if snippets.len() > threshold { if snippets.len() > threshold {
snippets.sort_by_key(|s| s.updated_at); snippets.sort_by_key(|s| s.updated_at);
let to_remove = snippets.len() - (threshold / 2); let to_remove = snippets.len() - (threshold / 2);
@@ -260,7 +260,7 @@ async fn run_server(state: Arc<MemoryState>) -> Result<(), Box<dyn std::error::E
if let Ok((len, _addr)) = socket.recv_from(&mut buf).await if let Ok((len, _addr)) = socket.recv_from(&mut buf).await
&& let Ok(payload) = serde_json::from_slice::<crate::models::TerminalHistory>(&buf[..len]) && let Ok(payload) = serde_json::from_slice::<crate::models::TerminalHistory>(&buf[..len])
{ {
udp_state.handler.state.terminal_history.modify(|history| { udp_state.handler.state.telemetry.terminal_history.modify(|history| {
history.push_front(payload.clone()); history.push_front(payload.clone());
if history.len() > 100 { if history.len() > 100 {
history.pop_back(); history.pop_back();
+4 -4
View File
@@ -101,7 +101,7 @@ impl McpResource for TasksActiveResource {
async fn read(&self, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn read(&self, state: Arc<MemoryState>) -> crate::error::Result<String> {
let state_clone = Arc::clone(&state); let state_clone = Arc::clone(&state);
tokio::task::spawn_blocking(move || -> crate::error::Result<String> { tokio::task::spawn_blocking(move || -> crate::error::Result<String> {
let tasks = state_clone.tasks.cache.read().unwrap(); let tasks = state_clone.project.tasks.cache.read().unwrap();
let data: Vec<_> = tasks let data: Vec<_> = tasks
.iter() .iter()
.filter(|t| t.status != "completed" && t.status != "done") .filter(|t| t.status != "completed" && t.status != "done")
@@ -190,7 +190,7 @@ impl MemoryHandler {
async fn read(&self, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn read(&self, state: Arc<MemoryState>) -> crate::error::Result<String> {
let state_clone = Arc::clone(&state); let state_clone = Arc::clone(&state);
tokio::task::spawn_blocking(move || -> crate::error::Result<String> { tokio::task::spawn_blocking(move || -> crate::error::Result<String> {
let items = state_clone.terminal_history.cache.read().unwrap(); let items = state_clone.telemetry.terminal_history.cache.read().unwrap();
Ok(serde_json::to_string_pretty(&*items)?) Ok(serde_json::to_string_pretty(&*items)?)
}) })
.await.map_err(|e| crate::error::AppError::Internal(e.to_string()))? .await.map_err(|e| crate::error::AppError::Internal(e.to_string()))?
@@ -211,7 +211,7 @@ impl MemoryHandler {
async fn read(&self, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn read(&self, state: Arc<MemoryState>) -> crate::error::Result<String> {
let state_clone = Arc::clone(&state); let state_clone = Arc::clone(&state);
tokio::task::spawn_blocking(move || -> crate::error::Result<String> { tokio::task::spawn_blocking(move || -> crate::error::Result<String> {
let items = state_clone.pinned_files.cache.read().unwrap(); let items = state_clone.project.pinned_files.cache.read().unwrap();
Ok(serde_json::to_string_pretty(&*items)?) Ok(serde_json::to_string_pretty(&*items)?)
}) })
.await.map_err(|e| crate::error::AppError::Internal(e.to_string()))? .await.map_err(|e| crate::error::AppError::Internal(e.to_string()))?
@@ -233,7 +233,7 @@ impl MemoryHandler {
async fn read(&self, state: Arc<MemoryState>) -> crate::error::Result<String> { async fn read(&self, state: Arc<MemoryState>) -> crate::error::Result<String> {
let state_clone = Arc::clone(&state); let state_clone = Arc::clone(&state);
tokio::task::spawn_blocking(move || -> crate::error::Result<String> { tokio::task::spawn_blocking(move || -> crate::error::Result<String> {
let items = state_clone.milestones.cache.read().unwrap(); let items = state_clone.project.milestones.cache.read().unwrap();
Ok(serde_json::to_string_pretty(&*items)?) Ok(serde_json::to_string_pretty(&*items)?)
}) })
.await.map_err(|e| crate::error::AppError::Internal(e.to_string()))? .await.map_err(|e| crate::error::AppError::Internal(e.to_string()))?
+64 -36
View File
@@ -13,32 +13,50 @@ pub struct GenericEvent {
pub payload: serde_json::Value, pub payload: serde_json::Value,
} }
pub struct ProjectStores {
pub tasks: Store<Vec<Task>>,
pub milestones: Store<Vec<Milestone>>,
pub pr_checklists: Store<Vec<PrChecklistItem>>,
pub context_workspaces: Store<Vec<ContextWorkspace>>,
pub pinned_files: Store<Vec<PinnedFile>>,
}
pub struct CodeStores {
pub ledger: Store<Vec<CodeChange>>,
pub snippets: Store<Vec<Snippet>>,
pub adrs: Store<Vec<Adr>>,
pub error_fixes: Store<Vec<ErrorFix>>,
pub tech_debts: Store<Vec<TechDebt>>,
pub sticky: Store<Vec<StickyNote>>,
}
pub struct EnvironmentStores {
pub env_fingerprints: Store<HashMap<String, EnvFingerprint>>,
pub env_requirements: Store<Vec<EnvRequirement>>,
pub environments: Store<Vec<EnvironmentDetail>>,
pub gates: Store<Vec<GateRecord>>,
pub prefs: Store<HashMap<String, Preference>>,
}
pub struct TelemetryStores {
pub session_summaries: Store<Vec<SessionSummary>>,
pub handoff_memos: Store<Vec<HandoffMemo>>,
pub recent_activities: Store<std::collections::VecDeque<serde_json::Value>>,
pub terminal_history: Store<std::collections::VecDeque<TerminalHistory>>,
}
pub struct MemoryState { pub struct MemoryState {
pub base_dir: PathBuf, pub base_dir: PathBuf,
pub clipboard_watch_mode: tokio::sync::RwLock<bool>, pub clipboard_watch_mode: tokio::sync::RwLock<bool>,
pub graph: Store<KnowledgeGraph>, pub graph: Store<KnowledgeGraph>,
pub search_index: RwLock<MemoryIndex>, pub search_index: RwLock<MemoryIndex>,
pub vector_db: tokio::sync::RwLock<Option<VectorDB>>, pub vector_db: tokio::sync::RwLock<Option<VectorDB>>,
pub ledger: Store<Vec<CodeChange>>,
pub sticky: Store<Vec<StickyNote>>, pub project: ProjectStores,
pub tasks: Store<Vec<Task>>, pub code: CodeStores,
pub snippets: Store<Vec<Snippet>>, pub env: EnvironmentStores,
pub adrs: Store<Vec<Adr>>, pub telemetry: TelemetryStores,
pub prefs: Store<HashMap<String, Preference>>,
pub error_fixes: Store<Vec<ErrorFix>>,
pub pinned_files: Store<Vec<PinnedFile>>,
pub session_summaries: Store<Vec<SessionSummary>>,
pub handoff_memos: Store<Vec<HandoffMemo>>,
pub env_fingerprints: Store<HashMap<String, EnvFingerprint>>,
pub env_requirements: Store<Vec<EnvRequirement>>,
pub milestones: Store<Vec<Milestone>>,
pub environments: Store<Vec<EnvironmentDetail>>,
pub pr_checklists: Store<Vec<PrChecklistItem>>,
pub tech_debts: Store<Vec<TechDebt>>,
pub gates: Store<Vec<GateRecord>>,
pub context_workspaces: Store<Vec<ContextWorkspace>>,
pub recent_activities: Store<std::collections::VecDeque<serde_json::Value>>,
pub terminal_history: Store<std::collections::VecDeque<TerminalHistory>>,
pub activity_tx: tokio::sync::broadcast::Sender<String>, pub activity_tx: tokio::sync::broadcast::Sender<String>,
pub event_bus_tx: tokio::sync::broadcast::Sender<GenericEvent>, pub event_bus_tx: tokio::sync::broadcast::Sender<GenericEvent>,
} }
@@ -66,26 +84,36 @@ impl MemoryState {
} }
}), }),
vector_db: tokio::sync::RwLock::new(None), vector_db: tokio::sync::RwLock::new(None),
ledger: Store::new("audit_ledger", db.clone()),
sticky: Store::new("sticky_notes", db.clone()), project: ProjectStores {
tasks: Store::new("tasks", db.clone()), tasks: Store::new("tasks", db.clone()),
milestones: Store::new("milestones", db.clone()),
pr_checklists: Store::new("pr_checklists", db.clone()),
context_workspaces: Store::new("context_workspaces", db.clone()),
pinned_files: Store::new("pinned_files", db.clone()),
},
code: CodeStores {
ledger: Store::new("audit_ledger", db.clone()),
snippets: Store::new("snippets", db.clone()), snippets: Store::new("snippets", db.clone()),
adrs: Store::new("adrs", db.clone()), adrs: Store::new("adrs", db.clone()),
prefs: Store::new("preferences", db.clone()),
error_fixes: Store::new("error_fixes", db.clone()), error_fixes: Store::new("error_fixes", db.clone()),
pinned_files: Store::new("pinned_files", db.clone()), tech_debts: Store::new("tech_debts", db.clone()),
session_summaries: Store::new("session_summaries", db.clone()), sticky: Store::new("sticky_notes", db.clone()),
handoff_memos: Store::new("handoff_memos", db.clone()), },
env: EnvironmentStores {
env_fingerprints: Store::new("env_fingerprints", db.clone()), env_fingerprints: Store::new("env_fingerprints", db.clone()),
env_requirements: Store::new("env_requirements", db.clone()), env_requirements: Store::new("env_requirements", db.clone()),
milestones: Store::new("milestones", db.clone()),
environments: Store::new("environments", db.clone()), environments: Store::new("environments", db.clone()),
pr_checklists: Store::new("pr_checklists", db.clone()),
tech_debts: Store::new("tech_debts", db.clone()),
gates: Store::new("gates", db.clone()), gates: Store::new("gates", db.clone()),
context_workspaces: Store::new("context_workspaces", db.clone()), prefs: Store::new("preferences", db.clone()),
},
telemetry: TelemetryStores {
session_summaries: Store::new("session_summaries", db.clone()),
handoff_memos: Store::new("handoff_memos", db.clone()),
recent_activities: Store::new("recent_activities", db.clone()), recent_activities: Store::new("recent_activities", db.clone()),
terminal_history: Store::new("terminal_history", db.clone()), terminal_history: Store::new("terminal_history", db.clone()),
},
activity_tx: tokio::sync::broadcast::channel(100).0, activity_tx: tokio::sync::broadcast::channel(100).0,
event_bus_tx: tokio::sync::broadcast::channel(1000).0, event_bus_tx: tokio::sync::broadcast::channel(1000).0,
} }
@@ -102,7 +130,7 @@ impl MemoryState {
"message": message "message": message
}); });
self.recent_activities.modify(|activities| { self.telemetry.recent_activities.modify(|activities| {
activities.push_back(item.clone()); activities.push_back(item.clone());
if activities.len() > 100 { if activities.len() > 100 {
activities.pop_front(); activities.pop_front();
@@ -143,9 +171,9 @@ impl MemoryState {
let entities: Vec<_> = self let entities: Vec<_> = self
.graph .graph
.read_with(|g| g.entities.values().cloned().collect()); .read_with(|g| g.entities.values().cloned().collect());
let tasks = self.tasks.read_with(|t| t.clone()); let tasks = self.project.tasks.read_with(|t| t.clone());
let snippets = self.snippets.read_with(|s| s.clone()); let snippets = self.code.snippets.read_with(|s| s.clone());
let adrs = self.adrs.read_with(|a| a.clone()); let adrs = self.code.adrs.read_with(|a| a.clone());
tracing::info!( tracing::info!(
"rebuild_index: found {} entities, {} tasks", "rebuild_index: found {} entities, {} tasks",
@@ -195,7 +223,7 @@ mod tests {
assert_eq!(state.base_dir, dir.path()); assert_eq!(state.base_dir, dir.path());
// Write a test value // Write a test value
state.tasks.modify(|tasks| { state.project.tasks.modify(|tasks| {
tasks.push(Task { tasks.push(Task {
id: "123".to_string(), id: "123".to_string(),
title: "Test Task".to_string(), title: "Test Task".to_string(),
@@ -212,7 +240,7 @@ mod tests {
}); });
// Ensure it is saved // Ensure it is saved
state.tasks.read_with(|tasks| { state.project.tasks.read_with(|tasks| {
assert_eq!(tasks.len(), 1); assert_eq!(tasks.len(), 1);
assert_eq!(tasks[0].id, "123"); assert_eq!(tasks[0].id, "123");
}); });