Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
44 changes: 41 additions & 3 deletions m1nd-mcp/src/daemon_handlers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -158,6 +158,8 @@ pub fn handle_daemon_start(
state.daemon_state.last_tick_changed_files = 0;
state.daemon_state.last_tick_deleted_files = 0;
state.daemon_state.last_tick_alerts_emitted = 0;
state.daemon_state.idle_streak = 0;
state.daemon_state.max_backoff_multiplier = 8;
state.persist_daemon_state()?;
Ok(json!({
"status": "started",
Expand Down Expand Up @@ -188,19 +190,42 @@ pub fn handle_daemon_status(
state: &mut SessionState,
_input: layers::DaemonStatusInput,
) -> M1ndResult<serde_json::Value> {
let now = now_ms();
let next_tick_due_ms = if state.daemon_state.active && state.daemon_state.poll_interval_ms > 0 {
state
.daemon_state
.last_tick_ms
.map(|last| last.saturating_add(state.daemon_state.poll_interval_ms))
} else {
None
};
let overdue_ms = next_tick_due_ms.map(|due| now.saturating_sub(due));
let effective_poll_interval_ms = state.daemon_state.poll_interval_ms.saturating_mul(
2u64.pow(
state
.daemon_state
.idle_streak
.min(state.daemon_state.max_backoff_multiplier.saturating_sub(1)),
),
);
Ok(json!({
"active": state.daemon_state.active,
"started_at_ms": state.daemon_state.started_at_ms,
"last_tick_ms": state.daemon_state.last_tick_ms,
"next_tick_due_ms": next_tick_due_ms,
"overdue_ms": overdue_ms,
"watch_paths": state.daemon_state.watch_paths,
"poll_interval_ms": state.daemon_state.poll_interval_ms,
"effective_poll_interval_ms": effective_poll_interval_ms,
"alert_count": state.daemon_alerts.len(),
"tracked_files": state.daemon_state.tracked_files.len(),
"tick_count": state.daemon_state.tick_count,
"last_tick_duration_ms": state.daemon_state.last_tick_duration_ms,
"last_tick_changed_files": state.daemon_state.last_tick_changed_files,
"last_tick_deleted_files": state.daemon_state.last_tick_deleted_files,
"last_tick_alerts_emitted": state.daemon_state.last_tick_alerts_emitted,
"idle_streak": state.daemon_state.idle_streak,
"max_backoff_multiplier": state.daemon_state.max_backoff_multiplier,
"runtime_root": state.runtime_root,
"graph_generation": state.graph_generation,
"cache_generation": state.cache_generation,
Expand Down Expand Up @@ -334,13 +359,18 @@ pub fn handle_daemon_tick(
}

let tick_ms = now_ms();
let emitted_alerts_total = emitted_alert_ids.len() + heuristic_alerts_emitted;
state.daemon_state.last_tick_ms = Some(tick_ms);
state.daemon_state.tick_count = state.daemon_state.tick_count.saturating_add(1);
state.daemon_state.last_tick_duration_ms = Some(start.elapsed().as_secs_f64() * 1000.0);
state.daemon_state.last_tick_changed_files = changed_entries.len();
state.daemon_state.last_tick_deleted_files = deleted_entries.len();
state.daemon_state.last_tick_alerts_emitted =
emitted_alert_ids.len() + heuristic_alerts_emitted;
state.daemon_state.last_tick_alerts_emitted = emitted_alerts_total;
if changed_entries.is_empty() && deleted_entries.is_empty() && emitted_alerts_total == 0 {
state.daemon_state.idle_streak = state.daemon_state.idle_streak.saturating_add(1);
} else {
state.daemon_state.idle_streak = 0;
}
state.persist_daemon_state()?;
state.persist_daemon_alerts()?;

Expand All @@ -356,7 +386,7 @@ pub fn handle_daemon_tick(
"file_path": entry.file_path,
"external_id": entry.external_id,
})).collect::<Vec<_>>(),
"alerts_emitted": emitted_alert_ids.len() + heuristic_alerts_emitted,
"alerts_emitted": emitted_alerts_total,
"alert_ids": emitted_alert_ids,
}))
}
Expand Down Expand Up @@ -552,6 +582,9 @@ mod tests {
assert_eq!(status["active"], true);
assert_eq!(status["alert_count"], 1);
assert_eq!(status["tick_count"], 0);
assert!(status["next_tick_due_ms"].as_u64().is_some());
assert_eq!(status["overdue_ms"], 0);
assert_eq!(status["idle_streak"], 0);

let stopped = handle_daemon_stop(
&mut state,
Expand Down Expand Up @@ -633,6 +666,8 @@ mod tests {
assert_eq!(status["tick_count"], 2);
assert_eq!(status["last_tick_changed_files"], 1);
assert_eq!(status["last_tick_deleted_files"], 0);
assert!(status["next_tick_due_ms"].as_u64().is_some());
assert_eq!(status["idle_streak"], 0);
}

#[test]
Expand Down Expand Up @@ -739,6 +774,7 @@ mod tests {
status["last_tick_alerts_emitted"].as_u64().unwrap_or(0) >= 1,
"risky daemon tick should emit at least one alert"
);
assert_eq!(status["idle_streak"], 0);
}

#[test]
Expand Down Expand Up @@ -800,5 +836,7 @@ mod tests {
assert_eq!(status["last_tick_deleted_files"], 1);
assert_eq!(status["last_tick_alerts_emitted"], 1);
assert!(status["last_tick_duration_ms"].as_f64().is_some());
assert!(status["next_tick_due_ms"].as_u64().is_some());
assert_eq!(status["idle_streak"], 0);
}
}
96 changes: 80 additions & 16 deletions m1nd-mcp/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2162,6 +2162,39 @@ fn background_tick_if_due(state: &mut SessionState) {
);
}

fn daemon_wait_duration_ms(state: &SessionState) -> u64 {
if !state.daemon_state.active {
return 1000;
}
if state.daemon_state.poll_interval_ms == 0 {
return 1000;
}

let exponent = state
.daemon_state
.idle_streak
.min(state.daemon_state.max_backoff_multiplier.saturating_sub(1));
let effective_poll_interval_ms = state
.daemon_state
.poll_interval_ms
.saturating_mul(2u64.pow(exponent))
.clamp(25, 10_000);

match state.daemon_state.last_tick_ms {
Some(last_tick_ms) => {
let elapsed = now_ms().saturating_sub(last_tick_ms);
if elapsed >= effective_poll_interval_ms {
25
} else {
effective_poll_interval_ms
.saturating_sub(elapsed)
.clamp(25, 1000)
}
}
None => 25,
}
}

/// Dispatch perspective tools (12 tools).
fn dispatch_perspective_tool(
state: &mut SessionState,
Expand Down Expand Up @@ -2479,22 +2512,17 @@ impl McpServer {
});

loop {
let daemon_wait = if self.state.daemon_state.active {
self.state.daemon_state.poll_interval_ms.clamp(25, 1000)
} else {
1000
let (payload, transport_mode) = match rx
.recv_timeout(Duration::from_millis(daemon_wait_duration_ms(&self.state)))
{
Ok(Some(value)) => value,
Ok(None) => break,
Err(mpsc::RecvTimeoutError::Timeout) => {
background_tick_if_due(&mut self.state);
continue;
}
Err(mpsc::RecvTimeoutError::Disconnected) => break,
};

let (payload, transport_mode) =
match rx.recv_timeout(Duration::from_millis(daemon_wait)) {
Ok(Some(value)) => value,
Ok(None) => break,
Err(mpsc::RecvTimeoutError::Timeout) => {
background_tick_if_due(&mut self.state);
continue;
}
Err(mpsc::RecvTimeoutError::Disconnected) => break,
};
let trimmed = payload.trim();
if trimmed.is_empty() {
continue;
Expand Down Expand Up @@ -2686,7 +2714,9 @@ impl McpServer {

#[cfg(test)]
mod tests {
use super::{background_tick_if_due, should_autotick_daemon, tool_schemas};
use super::{
background_tick_if_due, daemon_wait_duration_ms, should_autotick_daemon, tool_schemas,
};
use crate::server::McpConfig;
use crate::session::SessionState;
use m1nd_core::domain::DomainConfig;
Expand Down Expand Up @@ -2829,4 +2859,38 @@ mod tests {
"background tick should refresh the graph before the next explicit tool call"
);
}

#[test]
fn daemon_wait_duration_uses_remaining_time_until_next_tick() {
let (_temp, mut state) = build_state();
state.daemon_state.active = true;
state.daemon_state.poll_interval_ms = 500;
state.daemon_state.last_tick_ms = Some(super::now_ms().saturating_sub(125));

let wait_ms = daemon_wait_duration_ms(&state);
assert!(
(300..=400).contains(&wait_ms),
"remaining wait should be close to the poll interval remainder"
);

state.daemon_state.last_tick_ms = Some(0);
let overdue_wait_ms = daemon_wait_duration_ms(&state);
assert_eq!(overdue_wait_ms, 25);
}

#[test]
fn daemon_wait_duration_expands_with_idle_backoff() {
let (_temp, mut state) = build_state();
state.daemon_state.active = true;
state.daemon_state.poll_interval_ms = 200;
state.daemon_state.last_tick_ms = Some(super::now_ms());
state.daemon_state.idle_streak = 2;
state.daemon_state.max_backoff_multiplier = 8;

let wait_ms = daemon_wait_duration_ms(&state);
assert!(
(700..=800).contains(&wait_ms),
"idle streak should expand effective wait close to 4x the base interval"
);
}
}
2 changes: 2 additions & 0 deletions m1nd-mcp/src/session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,8 @@ pub struct DaemonRuntimeState {
pub last_tick_changed_files: usize,
pub last_tick_deleted_files: usize,
pub last_tick_alerts_emitted: usize,
pub idle_streak: u32,
pub max_backoff_multiplier: u32,
}

#[derive(Clone, Debug, Serialize, Deserialize)]
Expand Down