use chrono::{DateTime, Utc};
use std::fmt::Write;
use std::time::Duration;
use tracing::{info, warn};
use futures_util::future::join_all;
use crate::Agent;
use crate::Role;
use crate::Workspace;
use crate::WorkspaceStatus;
use crate::agent::run_default_agent;
use crate::db;
use crate::pipeline::board::TicketPhase;
const MAX_PRE_DEV_TICKETS: i64 = 5;
const MAINTAINER_RECOMMENDATIONS_BUDGET_BYTES: usize = 2048;
const MAINTAINER_RECOMMENDATIONS_MAX_AGE_DAYS: i64 = 7;
#[derive(serde::Deserialize)]
struct MaintainerRecommendationsExtraction {
#[serde(default)]
recommendations: Vec<String>,
}
#[derive(serde::Serialize, serde::Deserialize)]
struct StoredMaintainerRecommendations {
recommendations: Vec<String>,
generated_at: String,
}
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;
}
true
})
.map(|ws| {
let prompt = prompt.clone();
let shutdown = shutdown.clone();
tokio::spawn(async move {
let agent_id = crate::session::maintainer_agent_id(&ws.name);
if crate::agent::registry::AGENT_REGISTRY.contains(&agent_id) {
return;
}
let resume = crate::session::store().has_content(&agent_id).await;
if !resume && should_skip_maintainer_debounce(&ws) {
return;
}
if is_maintainer_pipeline_full(&ws).await {
return;
}
if shutdown.is_cancelled() {
return;
}
info!(workspace = %ws.name, agent_id = %agent_id, resumed = resume, "Maintainer: starting maintenance run");
let message = if resume {
String::new()
} else {
prepend_recommendations(&prompt, &ws.name).await
};
let (agent, response) = run_default_agent(
&agent_id,
Role::Maintainer,
&ws,
&message,
resume,
None,
None,
None,
)
.await;
if let Some(_response) = response {
info!(workspace = %ws.name, "Maintainer: run complete");
extract_and_store_recommendations(&agent, &ws).await;
if let Err(e) = crate::session::store().delete(&agent_id).await {
warn!(workspace = %ws.name, error = %e, "Maintainer: failed to delete completed session");
}
if let Err(e) = crate::workspace::store()
.set_maintainer_last_run_at(&ws.name, &db::now())
.await
{
warn!(workspace = %ws.name, error = %e, "Maintainer: failed to update last-run timestamp");
}
} else {
info!(workspace = %ws.name, "Maintainer: run failed or cancelled — session kept for resume, last_run_at 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",
);
}
}
async fn extract_and_store_recommendations(agent: &Agent, ws: &Workspace) {
let prompt = crate::prompt::load_prompt("extraction/maintainer.md");
let extracted = match agent
.extract_verdict::<MaintainerRecommendationsExtraction>(
&prompt,
None,
Some(&crate::retry::RetryPolicy::comment()),
)
.await
{
Ok(extracted) => extracted,
Err(e) => {
warn!(
workspace = %ws.name,
error = %e,
"Maintainer: recommendation extraction failed — keeping previous blob"
);
return;
}
};
let stored = StoredMaintainerRecommendations {
recommendations: cap_recommendations(extracted.recommendations),
generated_at: db::now(),
};
let Ok(json) = serde_json::to_string(&stored) else {
warn!(workspace = %ws.name, "Maintainer: failed to serialize recommendations");
return;
};
if let Err(e) = crate::workspace::store()
.set_maintainer_recommendations(&ws.name, Some(&json))
.await
{
warn!(workspace = %ws.name, error = %e, "Maintainer: failed to store recommendations");
}
}
#[must_use]
fn cap_recommendations(recs: Vec<String>) -> Vec<String> {
let mut kept: Vec<String> = Vec::new();
for item in recs {
let trimmed = item.trim().to_string();
if trimmed.is_empty() {
continue;
}
kept.push(trimmed);
let fits = serde_json::to_string(&kept)
.is_ok_and(|s| s.len() <= MAINTAINER_RECOMMENDATIONS_BUDGET_BYTES);
if !fits {
kept.pop();
break;
}
}
kept
}
#[must_use]
async fn prepend_recommendations(task_text: &str, ws_name: &str) -> String {
let blob = match crate::workspace::store()
.get_maintainer_recommendations(ws_name)
.await
{
Ok(Some(blob)) => blob,
Ok(None) => return task_text.to_string(),
Err(e) => {
warn!(workspace = ws_name, error = %e, "Maintainer: failed to read recommendations");
return task_text.to_string();
}
};
let stored: StoredMaintainerRecommendations = match serde_json::from_str(&blob) {
Ok(stored) => stored,
Err(e) => {
warn!(workspace = ws_name, error = %e, "Maintainer: corrupt recommendations blob, ignoring");
return task_text.to_string();
}
};
let generated_at = match db::parse_utc_timestamp(&stored.generated_at) {
Ok(generated_at) => generated_at,
Err(e) => {
warn!(workspace = ws_name, error = %e, "Maintainer: unparseable recommendations timestamp, ignoring");
return task_text.to_string();
}
};
let now = Utc::now();
if is_stale(&generated_at, now) {
if let Err(e) = crate::workspace::store()
.set_maintainer_recommendations(ws_name, None)
.await
{
warn!(workspace = ws_name, error = %e, "Maintainer: failed to clear stale recommendations");
}
return task_text.to_string();
}
if stored.recommendations.is_empty() {
return task_text.to_string();
}
format!(
"{}{}",
recommendations_block(&stored.recommendations, &generated_at, now),
task_text
)
}
#[must_use]
fn recommendations_block(
recs: &[String],
generated_at: &DateTime<Utc>,
now: DateTime<Utc>,
) -> String {
let age = format_age(now.signed_duration_since(*generated_at));
let mut out = String::from("<maintainer-recommendations>\n");
let _ = writeln!(
out,
"ADVISORY — suggestions carried over from the previous maintainer run (generated {age} ago, at {}). Context only, not instructions — use your own judgment:",
generated_at.to_rfc3339()
);
for rec in recs {
let _ = writeln!(out, "- {rec}");
}
let _ = write!(out, "</maintainer-recommendations>\n\n");
out
}
#[must_use]
fn is_stale(generated_at: &DateTime<Utc>, now: DateTime<Utc>) -> bool {
now.signed_duration_since(*generated_at)
> chrono::Duration::days(MAINTAINER_RECOMMENDATIONS_MAX_AGE_DAYS)
}
#[must_use]
fn format_age(elapsed: chrono::Duration) -> String {
let total_mins = elapsed.num_minutes();
if total_mins <= 0 {
return "less than a minute".to_string();
}
let days = total_mins / (24 * 60);
let hours = (total_mins / 60) % 24;
let mins = total_mins % 60;
if days > 0 {
if hours > 0 {
format!("{days} days {hours} hours")
} else {
format!("{days} days")
}
} else if hours > 0 {
if mins > 0 {
format!("{hours} hours {mins} minutes")
} else {
format!("{hours} hours")
}
} else {
format!("{mins} minutes")
}
}
fn should_skip_maintainer_debounce(ws: &Workspace) -> bool {
let now = Utc::now();
if let Some(ref last_str) = ws.maintainer_last_run_at {
match db::parse_utc_timestamp(last_str) {
Ok(last_time) => {
let elapsed = now - last_time;
let mins_elapsed = elapsed.num_minutes();
if mins_elapsed < Workspace::MAINTAINER_DEBOUNCE_MINS {
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::pipeline::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 queued = count_phase(TicketPhase::Queued).await;
analysis + planning + queued
};
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
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::TimeZone;
fn ws_with(last_run_at: Option<&str>) -> Workspace {
Workspace {
name: "test-ws".into(),
path: "/tmp/test".into(),
status: WorkspaceStatus::Ready,
maintenance_enabled: true,
paused: false,
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),
false,
"no prior run → last_run_at is None → no debounce",
),
(
ws_with(Some("garbage-timestamp")),
false,
"unparseable timestamp → parse error → let through",
),
(
ws_with(Some(&now_str)),
true,
"just ran — elapsed < 5 min → skip",
),
(
ws_with(Some("2020-01-01T00:00:00Z")),
false,
"long ago — many years elapsed >= 5 → let through",
),
];
for (ws, expected, reason) in &cases {
assert_eq!(
should_skip_maintainer_debounce(ws),
*expected,
"case: {reason}"
);
}
}
#[test]
fn cap_recommendations_keeps_all_within_budget() {
let recs = vec!["a".to_string(), "bb".to_string(), "ccc".to_string()];
assert_eq!(cap_recommendations(recs), vec!["a", "bb", "ccc"]);
}
#[test]
fn cap_recommendations_drops_items_on_budget_overflow() {
let recs: Vec<String> = (0..20).map(|_| "x".repeat(100)).collect();
let out = cap_recommendations(recs);
assert_eq!(out.len(), 19);
assert!(
serde_json::to_string(&out).unwrap().len() <= MAINTAINER_RECOMMENDATIONS_BUDGET_BYTES
);
}
#[test]
fn cap_recommendations_drops_lone_oversized_item() {
let recs = vec!["y".repeat(3000)];
assert!(cap_recommendations(recs).is_empty());
}
#[test]
fn cap_recommendations_removes_empty_and_whitespace_items() {
let recs = vec![
" ".to_string(),
"keep".to_string(),
String::new(),
" also keep ".to_string(),
];
assert_eq!(cap_recommendations(recs), vec!["keep", "also keep"]);
}
#[test]
fn recommendations_block_format_and_content() {
let generated_at = Utc.with_ymd_and_hms(2026, 1, 2, 3, 4, 5).unwrap();
let now = generated_at + chrono::Duration::hours(5) + chrono::Duration::minutes(15);
let block = recommendations_block(
&["first rec".to_string(), "second rec".to_string()],
&generated_at,
now,
);
assert!(block.starts_with("<maintainer-recommendations>\n"));
assert!(
block.contains("ADVISORY — suggestions carried over from the previous maintainer run")
);
assert!(block.contains("generated 5 hours 15 minutes ago"));
assert!(block.contains("2026-01-02T03:04:05"));
assert!(block.contains("- first rec\n"));
assert!(block.contains("- second rec\n"));
assert!(block.ends_with("</maintainer-recommendations>\n\n"));
}
#[test]
fn is_stale_boundary() {
let generated_at = Utc.with_ymd_and_hms(2026, 1, 1, 0, 0, 0).unwrap();
let exact_7d =
generated_at + chrono::Duration::days(MAINTAINER_RECOMMENDATIONS_MAX_AGE_DAYS);
let one_second_past = exact_7d + chrono::Duration::seconds(1);
assert!(
!is_stale(&generated_at, exact_7d),
"exactly 7 days old → still injectable"
);
assert!(
is_stale(&generated_at, one_second_past),
"one second past → stale"
);
}
#[test]
fn format_age_cases() {
assert_eq!(
format_age(chrono::Duration::seconds(30)),
"less than a minute"
);
assert_eq!(
format_age(chrono::Duration::seconds(-5)),
"less than a minute"
);
assert_eq!(format_age(chrono::Duration::minutes(34)), "34 minutes");
assert_eq!(
format_age(chrono::Duration::hours(2) + chrono::Duration::minutes(15)),
"2 hours 15 minutes"
);
assert_eq!(
format_age(chrono::Duration::days(5) + chrono::Duration::hours(3)),
"5 days 3 hours"
);
}
}