use chrono::Utc;
use std::time::Duration;
use tracing::{info, warn};
use futures_util::future::join_all;
use crate::Role;
use crate::Workspace;
use crate::WorkspaceStatus;
use crate::agent::run_default_agent;
use crate::board::TicketPhase;
use crate::turso;
const MAX_PRE_DEV_TICKETS: i64 = 5;
pub async fn run_maintainer_loop() {
let interval = Duration::from_mins(1);
let shutdown = crate::shutdown::shutdown_token();
loop {
if !crate::shutdown::sleep_or_shutdown_or_drain(interval).await {
break;
}
let workspaces = match crate::workspace::get_workspaces().await {
Ok(list) => list,
Err(e) => {
warn!(error = %e, "Maintainer: failed to list workspaces");
continue;
}
};
if workspaces.is_empty() {
info!("Maintainer: no workspaces configured, skipping cycle");
continue;
}
let prompt = crate::prompt::load_prompt("maintain.md");
let tasks: Vec<_> = workspaces
.into_iter()
.filter(|ws| {
if !ws.maintenance_enabled {
return false;
}
if ws.status != WorkspaceStatus::Ready {
info!(workspace = %ws.name, status = %ws.status, "Maintainer: skipping — workspace not ready");
return false;
}
if should_skip_maintainer_debounce(ws) {
return false;
}
true
})
.map(|ws| {
let prompt = prompt.clone();
let shutdown = shutdown.clone();
tokio::spawn(async move {
if is_maintainer_pipeline_full(&ws).await {
return;
}
if shutdown.is_cancelled() {
return;
}
let agent_id = crate::session::maintainer_agent_id(&ws.name);
info!(workspace = %ws.name, agent_id = %agent_id, "Maintainer: starting maintenance run");
let (agent, response) =
run_default_agent(&agent_id, Role::Maintainer, &ws, &prompt, None, None, None)
.await;
if let Some(_response) = response {
info!(workspace = %ws.name, "Maintainer: run complete");
let now_str = turso::now();
let new_debounce = compute_debounce(
&agent.agent_id,
ws.maintainer_debounce_mins,
ws.name.as_str(),
)
.await;
if let Err(e) = crate::workspace::store()
.set_maintenance_debounce(&ws.name, new_debounce, &now_str)
.await
{
warn!(workspace = %ws.name, error = %e, "Maintainer: failed to update debounce state");
}
} else {
info!(workspace = %ws.name, "Maintainer: run failed or cancelled — debounce unchanged");
}
})
})
.collect();
let results = join_all(tasks).await;
crate::util::log_join_failures(
results,
"Panic in maintainer task — maintainer loop continues",
"Maintainer task was cancelled — maintainer loop continues",
);
}
}
fn should_skip_maintainer_debounce(ws: &Workspace) -> bool {
let now = Utc::now();
let debounce = ws
.maintainer_debounce_mins
.clamp(0, Workspace::MAX_MAINTAINER_DEBOUNCE_MINS);
if let Some(ref last_str) = ws.maintainer_last_run_at {
match turso::parse_utc_timestamp(last_str) {
Ok(last_time) => {
let elapsed = now - last_time;
let mins_elapsed = elapsed.num_minutes();
if mins_elapsed < debounce {
return true;
}
}
Err(e) => {
warn!(
maintainer_last_run_at = %last_str,
error = %e,
"Failed to parse maintainer_last_run_at, letting through"
);
}
}
}
false
}
async fn is_maintainer_pipeline_full(ws: &Workspace) -> bool {
let Some(board) = crate::board::BOARD.get() else {
return false;
};
let count_phase = |phase: TicketPhase| async move {
match board.count_by_phase(phase, Some(&ws.name)).await {
Ok(c) => c,
Err(e) => {
warn!(workspace = %ws.name, %phase, error = %e, "Maintainer: failed to count tickets");
0
}
}
};
let pre_dev_count = {
let analysis = count_phase(TicketPhase::Analysis).await;
let planning = count_phase(TicketPhase::Planning).await;
let ready = count_phase(TicketPhase::ReadyForDevelopment).await;
analysis + planning + ready
};
if pre_dev_count >= MAX_PRE_DEV_TICKETS {
info!(
workspace = %ws.name,
pre_dev = pre_dev_count,
"Maintainer: skipping — pre-development pipeline has >= {} tickets",
MAX_PRE_DEV_TICKETS,
);
return true;
}
false
}
async fn compute_debounce(agent_id: &str, current: i64, ws_name: &str) -> i64 {
let Some(store) = crate::logs::LOG_STORE.get() else {
return advance_debounce(current);
};
match store.query_tool_usage(agent_id, "create_ticket").await {
Ok(call_count) if call_count > 0 => {
info!(workspace = %ws_name, "Maintainer: produced tickets — reset debounce to 1");
1
}
Ok(_) => {
let new_val = advance_debounce(current);
if new_val >= Workspace::MAX_MAINTAINER_DEBOUNCE_MINS
&& current < Workspace::MAX_MAINTAINER_DEBOUNCE_MINS
{
info!(workspace = %ws_name, "Maintainer: no tickets created — debounce capped at {}", Workspace::MAX_MAINTAINER_DEBOUNCE_MINS);
} else {
info!(workspace = %ws_name, "Maintainer: no tickets created — debounce advanced to {new_val}");
}
new_val
}
Err(e) => {
warn!(workspace = %ws_name, error = %e, "Maintainer: stats query failed, advancing debounce");
advance_debounce(current)
}
}
}
fn advance_debounce(mins: i64) -> i64 {
(mins.clamp(5, Workspace::MAX_MAINTAINER_DEBOUNCE_MINS) * 2)
.min(Workspace::MAX_MAINTAINER_DEBOUNCE_MINS)
}
#[cfg(test)]
mod tests {
use super::*;
fn ws_with(last_run_at: Option<&str>, debounce_mins: i64) -> Workspace {
Workspace {
name: "test-ws".into(),
path: "/tmp/test".into(),
status: WorkspaceStatus::Ready,
maintenance_enabled: true,
paused: false,
maintainer_debounce_mins: debounce_mins,
maintainer_last_run_at: last_run_at.map(String::from),
diagnostics: None,
notes: String::new(),
last_analyzed_commit: None,
ephemeral: false,
}
}
#[test]
fn should_skip_maintainer_debounce_cases() {
let now_str = Utc::now().to_rfc3339();
let cases = [
(
ws_with(None, 5),
false,
"no prior run → last_run_at is None → no debounce",
),
(
ws_with(Some("garbage-timestamp"), 5),
false,
"unparseable timestamp → parse error → let through",
),
(
ws_with(Some(&now_str), Workspace::MAX_MAINTAINER_DEBOUNCE_MINS),
true,
"just ran — elapsed ~0s < 240 → skip",
),
(
ws_with(Some("2020-01-01T00:00:00Z"), 5),
false,
"long ago — many years elapsed >= 5 → let through",
),
(
ws_with(Some("2020-01-01T00:00:00Z"), -5),
false,
"debounce clamped from -5 to 0 → mins_elapsed < 0 never true",
),
(
ws_with(Some(&now_str), 500),
true,
"debounce clamped from 500 to 240 → elapsed ~0s < 240 → skip",
),
];
for (ws, expected, reason) in &cases {
assert_eq!(
should_skip_maintainer_debounce(ws),
*expected,
"case: {reason}"
);
}
}
#[test]
fn advance_debounce_edges() {
assert_eq!(advance_debounce(0), 10);
assert_eq!(advance_debounce(4), 10);
assert_eq!(advance_debounce(5), 10);
assert_eq!(advance_debounce(6), 12);
assert_eq!(advance_debounce(60), 120);
assert_eq!(advance_debounce(119), 238);
assert_eq!(
advance_debounce(120),
Workspace::MAX_MAINTAINER_DEBOUNCE_MINS
);
assert_eq!(
advance_debounce(121),
Workspace::MAX_MAINTAINER_DEBOUNCE_MINS
);
assert_eq!(
advance_debounce(240),
Workspace::MAX_MAINTAINER_DEBOUNCE_MINS
);
assert_eq!(
advance_debounce(300),
Workspace::MAX_MAINTAINER_DEBOUNCE_MINS
);
}
}