use crate::Role;
use crate::message_router::{AgentJob, JobKind};
use crate::turso::{self, Connection, Row, TxGuard, Value, params};
use anyhow::{Context, Result};
use std::time::Duration;
use tracing::{debug, info, warn};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RowStatus {
Launched,
Done,
Failed,
}
impl RowStatus {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Launched => "launched",
Self::Done => "done",
Self::Failed => "failed",
}
}
}
impl std::str::FromStr for RowStatus {
type Err = anyhow::Error;
fn from_str(s: &str) -> Result<Self, Self::Err> {
match s {
"launched" => Ok(Self::Launched),
"done" => Ok(Self::Done),
"failed" => Ok(Self::Failed),
_ => Err(anyhow::anyhow!(
"Invalid row status '{s}'. Valid statuses: launched, done, failed"
)),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AgentKind {
Analyst,
Verifier,
Engineer,
Sanitation,
}
impl AgentKind {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Analyst => "analyst",
Self::Verifier => "verifier",
Self::Engineer => "engineer",
Self::Sanitation => "sanitation",
}
}
}
pub const MAX_BOOT_REDISPATCH: i64 = 3;
pub const DRAIN_CAP_SECS: u64 = 10 * 60;
pub const PURGE_CUTOFF_HOURS: i64 = 8;
#[derive(Debug, Clone)]
pub(crate) struct JobRow {
pub id: String,
pub kind: String,
pub workspace_name: String,
pub retry_count: i64,
}
#[derive(Debug, Clone)]
pub(crate) struct AgentRow {
pub agent_id: String,
pub idx: Option<i64>,
pub status: String,
pub outcome: Option<String>,
pub task: String,
}
#[derive(Debug, Clone)]
pub(crate) struct PendingJobRow {
pub id: String,
pub target_agent_id: String,
pub envelope: String,
pub created_at: String,
}
fn job_row_from(row: &Row) -> anyhow::Result<JobRow> {
Ok(JobRow {
id: row.get(0)?,
kind: row.get(1)?,
workspace_name: row.get(2)?,
retry_count: row.get(3)?,
})
}
fn agent_row_from(row: &Row) -> anyhow::Result<AgentRow> {
Ok(AgentRow {
agent_id: row.get(0)?,
idx: row.get(1)?,
status: row.get(2)?,
outcome: row.get(3)?,
task: row.get(4)?,
})
}
fn pending_row_from(row: &Row) -> anyhow::Result<PendingJobRow> {
Ok(PendingJobRow {
id: row.get(0)?,
target_agent_id: row.get(1)?,
envelope: row.get(2)?,
created_at: row.get(3)?,
})
}
pub(crate) const PENDING_JOB_INSERT_SQL: &str = "INSERT INTO pending_jobs \
(id, target_agent_id, envelope, created_at) \
VALUES (?1, ?2, ?3, ?4)";
pub(crate) fn pending_job_params(id: &str, envelope: &AgentJob, now: &str) -> Result<[Value; 4]> {
let envelope_json = serde_json::to_string(envelope)?;
Ok([
id.into(),
envelope_target(envelope).into(),
envelope_json.into(),
now.into(),
])
}
pub(crate) const AGENT_INSERT_SQL: &str = "INSERT INTO agents \
(job_id, agent_id, kind, idx, status, task) \
VALUES (?1, ?2, ?3, ?4, 'launched', ?5)";
pub(crate) fn agent_params(
job_id: &str,
agent_id: &str,
kind: AgentKind,
idx: Option<i64>,
task: &str,
) -> [Value; 5] {
[
job_id.into(),
agent_id.into(),
kind.as_str().into(),
idx.into(),
task.into(),
]
}
#[derive(Debug, Clone)]
pub(crate) enum SpawnChild {
Analyze,
Research,
ResearchCleanup,
TempCleanup,
TicketStage {
ticket_id: String,
stage: String,
phase: String,
round: i64,
},
}
impl SpawnChild {
#[must_use]
pub const fn kind_str(&self) -> &'static str {
match self {
Self::Analyze => "analyze",
Self::Research => "research",
Self::ResearchCleanup => "research_cleanup",
Self::TempCleanup => "temp_cleanup",
Self::TicketStage { .. } => "ticket_stage",
}
}
}
#[expect(clippy::too_many_arguments)]
pub(crate) async fn spawn_job(
conn: &Connection,
id: &str,
task: &str,
workspace_name: &str,
user_name: &str,
channel: &str,
role: Role,
agents: &[NewAgent],
child: &SpawnChild,
) -> Result<()> {
let kind = child.kind_str();
let now = turso::now();
let tx = conn.begin_tx().await?;
tx.execute(
"INSERT INTO jobs (id, kind, status, task, workspace_name, user_name, channel, role, \
retry_count, created_at, updated_at) \
VALUES (?1, ?2, 'launched', ?3, ?4, ?5, ?6, ?7, 0, ?8, ?8)",
params![
id,
kind,
task,
workspace_name,
user_name,
channel,
role.as_str(),
now.clone(),
],
)
.await
.with_context(|| format!("failed to spawn job {id}"))?;
for a in agents {
tx.execute(
AGENT_INSERT_SQL,
agent_params(id, &a.agent_id, a.kind, a.idx, &a.task),
)
.await
.with_context(|| format!("failed to insert agent roster for job {id}"))?;
}
match child {
SpawnChild::Analyze | SpawnChild::ResearchCleanup | SpawnChild::TempCleanup => {}
SpawnChild::Research => {
tx.execute(
"INSERT INTO research_jobs (id, state) VALUES (?1, '{}')",
params![id],
)
.await
.with_context(|| format!("failed to insert research_jobs row for job {id}"))?;
}
SpawnChild::TicketStage {
ticket_id,
stage,
phase,
round,
} => {
tx.execute(
"INSERT INTO ticket_stage_jobs (id, ticket_id, stage, phase, round) \
VALUES (?1, ?2, ?3, ?4, ?5)",
params![id, ticket_id.clone(), stage.clone(), phase.clone(), round],
)
.await
.with_context(|| format!("failed to insert ticket_stage_jobs row for job {id}"))?;
}
}
tx.commit().await?;
Ok(())
}
pub(crate) struct NewAgent {
pub agent_id: String,
pub kind: AgentKind,
pub idx: Option<i64>,
pub task: String,
}
pub(crate) async fn checkpoint_job(conn: &Connection, id: &str, retry_count: i64) -> Result<()> {
let status = RowStatus::Launched;
let now = turso::now();
conn.execute(
"UPDATE jobs SET status = ?1, retry_count = ?2, updated_at = ?3 WHERE id = ?4",
params![status.as_str(), retry_count, now, id],
)
.await
.with_context(|| format!("failed to checkpoint job {id}"))?;
Ok(())
}
pub(crate) async fn write_agent_outcome(
conn: &Connection,
job_id: &str,
agent_id: &str,
status: RowStatus,
outcome: Option<&str>,
) -> Result<()> {
conn.execute(
"UPDATE agents SET status = ?1, outcome = ?2 WHERE job_id = ?3 AND agent_id = ?4",
params![status.as_str(), outcome, job_id, agent_id],
)
.await
.with_context(|| format!("failed to write agent outcome for {agent_id}"))?;
Ok(())
}
pub(crate) async fn complete_job_with_envelope(
conn: &Connection,
job_id: &str,
envelope: &AgentJob,
) -> Result<()> {
let now = turso::now();
let tx = conn.begin_tx().await?;
tx.execute(
PENDING_JOB_INSERT_SQL,
pending_job_params(job_id, envelope, &now).context("serialize envelope")?,
)
.await
.with_context(|| format!("failed to persist envelope for job {job_id}"))?;
delete_job_tx(&tx, job_id).await?;
if crate::research_cancel::is_cancelled(job_id) {
tracing::info!(job = %job_id, "Job completion rolled back — run manually cancelled");
return Ok(());
}
tx.commit().await?;
Ok(())
}
pub(crate) async fn complete_durable_job(
job_id: &str,
content: String,
kind: JobKind,
caller_role: Role,
user_name: &str,
channel: &str,
workspace_name: &str,
) -> AgentJob {
let mut envelope = AgentJob {
content,
workspace_name: workspace_name.to_string(),
user_name: user_name.to_string(),
channel: channel.to_string(),
kind,
role: caller_role,
reply_target: None,
pending_job_id: Some(job_id.to_string()),
};
if complete_job_with_envelope(&crate::session::store().conn, job_id, &envelope)
.await
.is_err()
{
envelope.pending_job_id = None;
}
envelope
}
pub(crate) async fn complete_ticket_stage_job(conn: &Connection, job_id: &str) -> Result<()> {
let tx = conn.begin_tx().await?;
tx.execute(
"UPDATE jobs SET status = 'done', updated_at = ?1 WHERE id = ?2",
params![turso::now(), job_id],
)
.await?;
tx.execute("DELETE FROM agents WHERE job_id = ?1", params![job_id])
.await?;
tx.commit().await?;
Ok(())
}
pub(crate) async fn terminalize_job(conn: &Connection, job_id: &str) -> Result<()> {
conn.execute("DELETE FROM jobs WHERE id = ?1", params![job_id])
.await
.context("terminalize job")?;
Ok(())
}
async fn delete_job_tx(tx: &TxGuard<'_>, job_id: &str) -> Result<()> {
tx.execute("DELETE FROM jobs WHERE id = ?1", params![job_id])
.await
.with_context(|| format!("failed to delete job {job_id}"))?;
Ok(())
}
pub(crate) fn envelope_target(job: &AgentJob) -> String {
crate::session::resolve_agent_id(&job.user_name, job.role.as_str(), &job.workspace_name)
}
pub(crate) struct JobCaller {
pub task: String,
pub role: String,
pub user_name: String,
pub channel: String,
}
pub(crate) async fn job_caller(conn: &Connection, job_id: &str) -> Result<Option<JobCaller>> {
conn.query_optional(
"SELECT task, role, user_name, channel FROM jobs WHERE id = ?1",
params![job_id],
|r| {
Ok::<_, anyhow::Error>(JobCaller {
task: r.get::<String>(0)?,
role: r.get::<String>(1)?,
user_name: r.get::<String>(2)?,
channel: r.get::<String>(3)?,
})
},
)
.await
.context("load job caller")
}
pub(crate) async fn resume_job_preamble(
conn: &Connection,
job_id: &str,
abort_site: &str,
missing_site: &str,
) -> Option<(JobCaller, Role)> {
if crate::shutdown::aborting() {
tracing::info!(job = %job_id, "{abort_site} aborted — drain/shutdown in progress");
return None;
}
let Ok(Some(caller)) = job_caller(conn, job_id).await else {
tracing::warn!(job = %job_id, "{missing_site}: missing job row — terminalizing job");
let _ = terminalize_job(conn, job_id).await;
return None;
};
let caller_role = std::str::FromStr::from_str(&caller.role).unwrap_or(Role::Manager);
Some((caller, caller_role))
}
pub(crate) async fn job_retry_count(conn: &Connection, job_id: &str) -> i64 {
match conn
.query_optional(
"SELECT retry_count FROM jobs WHERE id = ?1",
params![job_id],
|r| r.get::<i64>(0),
)
.await
{
Ok(Some(n)) => n,
Ok(None) => 0,
Err(e) => {
warn!(job = %job_id, error = %e, "Failed to read job retry_count — assuming 0");
0
}
}
}
pub(crate) async fn list_active_jobs(conn: &Connection) -> Result<Vec<JobRow>> {
let rows = conn
.query(
"SELECT id, kind, workspace_name, retry_count FROM jobs \
WHERE status != 'done' ORDER BY created_at",
(),
)
.await
.context("list active jobs")?;
rows.iter().map(job_row_from).collect()
}
pub(crate) async fn list_agents_for_job(conn: &Connection, job_id: &str) -> Result<Vec<AgentRow>> {
let rows = conn
.query(
"SELECT agent_id, idx, status, outcome, task \
FROM agents WHERE job_id = ?1 ORDER BY idx",
params![job_id],
)
.await
.context("list agents for job")?;
rows.iter().map(agent_row_from).collect()
}
pub async fn run_drain_watch() {
use crate::shutdown::{aborting, force_cancel, shutdown, shutdown_token};
use std::time::Instant;
loop {
if aborting() {
break;
}
tokio::time::sleep(Duration::from_millis(500)).await;
}
if shutdown_token().is_cancelled() {
return;
}
let start = Instant::now();
loop {
if crate::registry::AGENT_REGISTRY.list().is_empty()
&& crate::call_registry::NON_AGENT_CALLS.list().is_empty()
{
info!("Drain complete — no in-flight agents or orchestrator calls; exiting");
shutdown();
return;
}
if start.elapsed() >= Duration::from_secs(DRAIN_CAP_SECS) {
warn!("Drain cap ({DRAIN_CAP_SECS}s) reached — force-cancelling in-flight work");
crate::registry::AGENT_REGISTRY.shutdown_all();
force_cancel();
return;
}
if shutdown_token().is_cancelled() {
return;
}
tokio::time::sleep(Duration::from_secs(2)).await;
}
}
pub(crate) async fn list_pending_jobs(conn: &Connection) -> Result<Vec<PendingJobRow>> {
let rows = conn
.query(
"SELECT id, target_agent_id, envelope, created_at \
FROM pending_jobs ORDER BY created_at",
(),
)
.await
.context("list pending jobs")?;
rows.iter().map(pending_row_from).collect()
}
pub(crate) async fn delete_pending_job(conn: &Connection, id: &str) -> Result<()> {
conn.execute("DELETE FROM pending_jobs WHERE id = ?1", params![id])
.await
.context("delete pending job")?;
Ok(())
}
pub(crate) enum ResumableStage {
TicketStage {
job_id: String,
ticket_id: String,
stage: String,
workspace_name: String,
},
Research {
job_id: String,
workspace_name: String,
capped: bool,
},
Analyze {
job_id: String,
workspace_name: String,
capped: bool,
},
ResearchCleanup {
job_id: String,
workspace_name: String,
},
}
struct StageCandidate {
job_id: String,
ticket_id: String,
stage: String,
workspace_name: String,
round: i64,
}
fn strip_timestamp_wrapper(content: &str) -> &str {
if !content.starts_with("<timestamp>")
&& let Some(body) = content.strip_suffix("</timestamp>")
&& let Some(start) = body.rfind("\n\n<timestamp>")
&& is_timestamp_body(&body[start + "\n\n<timestamp>".len()..])
{
return &body[..start];
}
if let Some(body) = content.strip_prefix("<timestamp>")
&& let Some(end) = body.find("</timestamp>")
&& is_timestamp_body(&body[..end])
&& let Some(after) = body[end..].strip_prefix("</timestamp>\n\n")
{
return after;
}
content
}
fn is_timestamp_body(s: &str) -> bool {
let Some((datetime, tz)) = s.strip_suffix(')').and_then(|rest| rest.rsplit_once(" (")) else {
return false;
};
!tz.is_empty() && chrono::NaiveDateTime::parse_from_str(datetime, "%Y-%m-%d %H:%M:%S").is_ok()
}
async fn pending_already_appended(conn: &Connection, row: &PendingJobRow) -> bool {
let Ok(Some((last_content, last_created))) = conn
.query_optional(
"SELECT content, created_at FROM sessions \
WHERE agent_id = ?1 AND role = 'user' ORDER BY id DESC LIMIT 1",
params![row.target_agent_id.clone()],
|r| Ok::<_, anyhow::Error>((r.get::<String>(0)?, r.get::<String>(1)?)),
)
.await
else {
return false;
};
let Ok(envelope) = serde_json::from_str::<AgentJob>(&row.envelope) else {
return false;
};
let appended_before = match (
crate::turso::parse_utc_timestamp(&row.created_at).ok(),
crate::turso::parse_utc_timestamp(&last_created).ok(),
) {
(Some(row_ts), Some(last_ts)) => row_ts <= last_ts,
_ => row.created_at <= last_created,
};
strip_timestamp_wrapper(&last_content).ends_with(&envelope.content) && appended_before
}
async fn replay_pending_jobs(conn: &Connection) -> Result<usize> {
let rows = list_pending_jobs(conn).await?;
let mut replayed = 0usize;
for row in &rows {
let Ok(mut job) = serde_json::from_str::<AgentJob>(&row.envelope) else {
warn!(pending_job = %row.id, "Pending job envelope unreadable — skipping");
continue;
};
job.pending_job_id = Some(row.id.clone());
if pending_already_appended(conn, row).await {
debug!(pending_job = %row.id, "Pending job already appended — skipping");
if let Err(e) = delete_pending_job(conn, &row.id).await {
warn!(pending_job = %row.id, error = %e, "Failed to delete deduped pending job");
}
continue;
}
if job.kind == JobKind::ResearchResult
&& !crate::research_cleanup::research_cleanup_row_exists(conn, &row.id)
.await
.unwrap_or(true)
{
crate::research_cleanup::dispatch_cleanup_for_pending_envelope(&row.id, &job).await;
}
crate::message_router::route(&row.target_agent_id, job);
replayed += 1;
}
Ok(replayed)
}
pub async fn purge_stale_jobs(cutoff: &str) -> Result<u64> {
let conn = &crate::session::store().conn;
let rows = conn
.query(
"SELECT j.id, j.kind, ts.ticket_id, ts.phase, ts.round, ts.stage FROM jobs j \
LEFT JOIN ticket_stage_jobs ts ON ts.id = j.id \
WHERE j.updated_at < ?1 \
AND NOT EXISTS ( \
SELECT 1 FROM agents a \
JOIN session_metadata sm ON sm.agent_id = a.agent_id \
WHERE a.job_id = j.id AND sm.last_activity >= ?1)",
params![cutoff],
)
.await
.context("select stale jobs for purge")?;
if rows.is_empty() {
return Ok(0);
}
let mut purge_ids: Vec<String> = Vec::with_capacity(rows.len());
let mut ticket_stage_ids: std::collections::HashSet<String> =
std::collections::HashSet::with_capacity(rows.len());
let mut rollbacks: Vec<(String, String, bool)> = Vec::new();
for row in &rows {
let id: String = row.get(0)?;
let kind: String = row.get(1)?;
purge_ids.push(id.clone());
if kind == "ticket_stage" {
ticket_stage_ids.insert(id);
let ticket_id: Option<String> = row.get(2).ok();
let phase: Option<String> = row.get(3).ok();
let round: Option<i64> = row.get(4).ok();
let stage: Option<String> = row.get(5).ok();
if let (Some(t), Some(p), Some(r), Some(s)) = (ticket_id, phase, round, stage) {
let latest: Option<i64> = conn
.query_row(
"SELECT MAX(round) FROM ticket_stage_jobs \
WHERE ticket_id = ?1 AND stage = ?2",
params![t.clone(), s.clone()],
|row| row.get::<i64>(0),
)
.await
.ok();
rollbacks.push((t, p, latest == Some(r)));
}
}
}
let mut rollback_ok = true;
if !rollbacks.is_empty() {
if crate::board::BOARD.get().is_some() {
rollback_ok = rollback_stranded_tickets(&rollbacks).await;
} else {
rollback_ok = false;
}
}
let tx = conn.begin_tx().await?;
let mut deleted = 0usize;
for id in &purge_ids {
if !rollback_ok && ticket_stage_ids.contains(id) {
continue;
}
delete_job_tx(&tx, id).await?;
deleted += 1;
}
tx.commit().await?;
if deleted > 0 {
tracing::debug!(deleted, "Purged stale jobs");
}
Ok(deleted as u64)
}
fn rollback_transition(phase: &str) -> Option<(String, bool)> {
let from = phase.parse::<crate::board::TicketPhase>().ok()?;
if from == crate::board::TicketPhase::InDiagnostics {
return None;
}
crate::board::BoardStore::reset_transition(from)
.map(|(to, pipeline_reservation)| (to.to_string(), pipeline_reservation))
}
async fn rollback_stranded_tickets(rollbacks: &[(String, String, bool)]) -> bool {
let board = crate::board::BOARD.get().expect("BOARD initialized");
let tx = match board.conn.begin_tx().await {
Ok(tx) => tx,
Err(e) => {
warn!(error = %e, "Purge rollback: failed to begin board tx");
return false;
}
};
let now = crate::turso::now();
let mut rolled_back = 0usize;
for (ticket_id, phase, is_latest) in rollbacks {
let Some((to, pipeline_reservation)) = rollback_transition(phase) else {
continue;
};
if !is_latest {
debug!(ticket = %ticket_id, "Purge rollback skipped (not latest round)");
continue;
}
let sql = format!(
"UPDATE tickets SET {} WHERE id = ?3 AND phase = ?5",
crate::board::BoardStore::RESET_TICKET_SET_CLAUSE
);
let updated = tx
.execute(
&sql,
crate::turso::params![
to,
now.clone(),
ticket_id.clone(),
i64::from(pipeline_reservation),
phase.clone(),
],
)
.await;
match updated {
Ok(n) if n > 0 => rolled_back += 1,
Ok(_) => {
debug!(ticket = %ticket_id, "Purge rollback CAS no-op (phase moved)");
}
Err(e) => {
warn!(ticket = %ticket_id, error = %e, "Purge rollback failed");
}
}
}
match tx.commit().await {
Ok(()) => {
if rolled_back > 0 {
info!(rolled_back, "Purge: rolled back stranded tickets in place");
}
true
}
Err(e) => {
warn!(error = %e, "Purge rollback commit failed");
false
}
}
}
#[expect(clippy::cast_possible_truncation, clippy::too_many_lines)]
pub(crate) async fn recover_from_restart() -> Result<Vec<ResumableStage>> {
let start = std::time::Instant::now();
let conn = &crate::session::store().conn;
let replayed = replay_pending_jobs(conn).await?;
let jobs = list_active_jobs(conn).await?;
let mut exclusion: Vec<String> = Vec::new();
let mut to_complete: Vec<String> = Vec::new();
let mut resumable: Vec<ResumableStage> = Vec::new();
let mut candidates: Vec<StageCandidate> = Vec::new();
for job in &jobs {
if job.kind != "ticket_stage" {
continue;
}
let stage_rows = conn
.query(
"SELECT ticket_id, stage, phase, round FROM ticket_stage_jobs WHERE id = ?1",
params![job.id.clone()],
)
.await
.unwrap_or_default();
let Some(row) = stage_rows.first() else {
continue;
};
let ticket_id: String = row.get(0).expect("ticket_stage_jobs.ticket_id is NOT NULL");
let stage: String = row.get(1).expect("ticket_stage_jobs.stage is NOT NULL");
let phase: String = row.get(2).expect("ticket_stage_jobs.phase is NOT NULL");
let round: i64 = row.get(3).expect("ticket_stage_jobs.round is NOT NULL");
if job.retry_count >= MAX_BOOT_REDISPATCH {
warn!(
job = %job.id,
ticket = %ticket_id,
retry_count = job.retry_count,
"Ticket-stage job exceeded boot re-dispatch cap — marking done",
);
let _ = complete_ticket_stage_job(conn, &job.id).await;
continue;
}
let _ = checkpoint_job(conn, &job.id, job.retry_count + 1).await;
let in_phase = crate::board::store()
.get_ticket_phase(&ticket_id)
.await
.is_ok_and(|p| p.is_some_and(|ph| ph.as_ref() == phase.as_str()));
if in_phase {
candidates.push(StageCandidate {
job_id: job.id.clone(),
ticket_id,
stage,
workspace_name: job.workspace_name.clone(),
round,
});
} else {
to_complete.push(job.id.clone());
}
}
let mut best: std::collections::HashMap<(String, String), (i64, usize)> =
std::collections::HashMap::new();
for (idx, cand) in candidates.iter().enumerate() {
let key = (cand.ticket_id.clone(), cand.stage.clone());
match best.get(&key) {
Some(&(best_round, _)) if best_round >= cand.round => {
to_complete.push(cand.job_id.clone());
}
_ => {
if let Some(&(_, prev_idx)) = best.get(&key) {
to_complete.push(candidates[prev_idx].job_id.clone());
}
best.insert(key, (cand.round, idx));
}
}
}
for (_round, idx) in best.values() {
let cand = &candidates[*idx];
exclusion.push(cand.ticket_id.clone());
resumable.push(ResumableStage::TicketStage {
job_id: cand.job_id.clone(),
ticket_id: cand.ticket_id.clone(),
stage: cand.stage.clone(),
workspace_name: cand.workspace_name.clone(),
});
}
if let Some(board) = crate::board::BOARD.get()
&& let Err(e) = board.reset_inflight_tickets(&exclusion).await
{
warn!(error = %e, "Failed to reset in-flight tickets");
}
for id in &to_complete {
if let Err(e) = complete_ticket_stage_job(conn, id).await {
warn!(job = %id, error = %e, "Failed to mark stale ticket_stage job done");
}
}
let mut resumed_other = 0usize;
for job in &jobs {
match job.kind.as_str() {
"research" | "analyze" => {
let kind = job.kind.as_str();
let capped = job.retry_count >= MAX_BOOT_REDISPATCH;
if capped {
let msg = if kind == "research" {
"Research job exceeded boot re-dispatch cap — delivering partial report"
} else {
"Analyze job exceeded boot re-dispatch cap — delivering failure envelope"
};
warn!(job = %job.id, kind = %kind, "{msg}");
}
let _ = checkpoint_job(
conn,
&job.id,
if capped {
job.retry_count
} else {
job.retry_count + 1
},
)
.await;
resumable.push(if kind == "research" {
ResumableStage::Research {
job_id: job.id.clone(),
workspace_name: job.workspace_name.clone(),
capped,
}
} else {
ResumableStage::Analyze {
job_id: job.id.clone(),
workspace_name: job.workspace_name.clone(),
capped,
}
});
resumed_other += 1;
}
"research_cleanup" => {
if job.retry_count >= MAX_BOOT_REDISPATCH {
warn!(
job = %job.id,
"Research cleanup job exceeded boot re-dispatch cap — run folder released in-process, row left for the purge",
);
crate::research_cleanup::release_run_folder(&job.id).await;
continue;
}
let _ = checkpoint_job(conn, &job.id, job.retry_count + 1).await;
resumable.push(ResumableStage::ResearchCleanup {
job_id: job.id.clone(),
workspace_name: job.workspace_name.clone(),
});
resumed_other += 1;
}
"ticket_stage" => {
}
"temp_cleanup" => {
info!(
job = %job.id,
"Temp-dir cleanup row left over from a previous lifetime — terminalizing (fire-and-forget, no resume)",
);
let _ = terminalize_job(conn, &job.id).await;
}
_ => {
warn!(job = %job.id, kind = %job.kind, "Unknown job kind — skipping");
}
}
}
let elapsed = start.elapsed();
info!(
duration_ms = elapsed.as_millis() as u64,
resumed_tickets = resumable.len() - resumed_other,
resumed_other,
replayed_pending = replayed,
"Boot recovery scan complete",
);
let _ = purge_terminal_engineer_anchors().await;
Ok(resumable)
}
#[must_use]
pub fn engineer_anchor_id(ticket_id: &str) -> String {
crate::session::ticket_agent_id(ticket_id, crate::Role::Engineer.as_str())
}
pub(crate) async fn upsert_engineer_anchor(
conn: &Connection,
ticket_id: &str,
task: &str,
status: RowStatus,
) -> Result<()> {
let anchor_id = engineer_anchor_id(ticket_id);
conn.execute(
"INSERT INTO agents (job_id, agent_id, kind, idx, status, outcome, task) \
VALUES (NULL, ?1, 'engineer', NULL, ?2, NULL, ?3) \
ON CONFLICT(agent_id) WHERE job_id IS NULL \
DO UPDATE SET status = ?2, task = ?3, outcome = NULL",
params![anchor_id, status.as_str(), task],
)
.await
.with_context(|| format!("failed to upsert engineer anchor for ticket {ticket_id}"))?;
Ok(())
}
#[must_use]
pub fn engineer_anchor_ticket_id(agent_id: &str) -> Option<String> {
agent_id
.strip_prefix("ticket_")
.and_then(|rest| rest.strip_suffix("_engineer"))
.filter(|id| !id.is_empty())
.map(str::to_string)
}
pub(crate) async fn purge_terminal_engineer_anchors() -> usize {
let conn = &crate::session::store().conn;
let Ok(rows) = conn
.query(
"SELECT agent_id FROM agents WHERE job_id IS NULL AND kind = 'engineer'",
(),
)
.await
else {
return 0;
};
let mut deleted = 0usize;
let board_ready = crate::board::BOARD.get().is_some();
for row in &rows {
let Ok(agent_id) = row.get::<String>(0) else {
continue;
};
let Some(ticket_id) = engineer_anchor_ticket_id(&agent_id) else {
continue;
};
let terminal = board_ready
&& crate::board::store()
.get_ticket_phase(&ticket_id)
.await
.is_ok_and(|p| p.is_some_and(|ph| ph.is_terminal()));
if !terminal {
continue;
}
match conn
.execute(
"DELETE FROM agents WHERE agent_id = ?1 AND job_id IS NULL",
params![agent_id.clone()],
)
.await
{
Ok(_) => deleted += 1,
Err(e) => {
warn!(agent = %agent_id, error = %e, "Failed to delete terminal engineer anchor");
}
}
}
if deleted > 0 {
info!(deleted, "Removed engineer anchors for terminal tickets");
}
deleted
}
#[cfg(test)]
mod tests {
use super::*;
use crate::turso::params;
#[tokio::test]
async fn anchor_upsert_partial_index_semantics() {
let (store, _tmp) = crate::open_test_store!(crate::session::SessionStore, "session");
let conn = &store.conn;
let anchor_id = engineer_anchor_id("t-1400");
upsert_engineer_anchor(conn, "t-1400", "task-1", RowStatus::Launched)
.await
.expect("anchor upsert (first)");
let rows = conn
.query(
"SELECT agent_id, job_id, status, task FROM agents WHERE agent_id = ?1",
params![anchor_id.clone()],
)
.await
.unwrap();
assert_eq!(rows.len(), 1, "exactly one anchor row after first upsert");
assert!(rows[0].get::<Option<String>>(1).unwrap().is_none());
assert_eq!(rows[0].get::<String>(2).unwrap(), "launched");
assert_eq!(rows[0].get::<String>(3).unwrap(), "task-1");
upsert_engineer_anchor(conn, "t-1400", "task-2", RowStatus::Launched)
.await
.expect("anchor upsert (second)");
let rows = conn
.query(
"SELECT agent_id, job_id, status, task FROM agents WHERE agent_id = ?1",
params![anchor_id.clone()],
)
.await
.unwrap();
assert_eq!(rows.len(), 1, "anchor must stay a single row per agent_id");
assert!(rows[0].get::<Option<String>>(1).unwrap().is_none());
assert_eq!(rows[0].get::<String>(3).unwrap(), "task-2");
crate::util::test::JobRowBuilder::new(conn, "j1", "ticket_stage", "engineer", "ws")
.timestamps(turso::now())
.insert()
.await
.expect("insert job for roster row");
conn.execute(
"INSERT INTO agents (job_id, agent_id, kind, idx, status, task) \
VALUES ('j1', ?1, 'engineer', NULL, 'launched', 'roster-task')",
params![anchor_id.clone()],
)
.await
.expect("roster row with same agent_id must coexist");
let rows = conn
.query(
"SELECT job_id, task FROM agents WHERE agent_id = ?1 ORDER BY job_id IS NULL",
params![anchor_id],
)
.await
.unwrap();
assert_eq!(rows.len(), 2, "anchor + roster row coexist");
let err = conn
.execute(
"INSERT INTO agents (job_id, agent_id, kind, idx, status, task) \
VALUES (NULL, ?1, 'engineer', NULL, 'launched', 'dup')",
params![engineer_anchor_id("t-1400")],
)
.await
.expect_err("duplicate NULL-seat anchor must violate the partial unique index");
debug!("anchor duplicate conflict: {err:?}");
conn.execute("DELETE FROM jobs WHERE id = 'j1'", ())
.await
.unwrap();
let rows = conn
.query(
"SELECT job_id FROM agents WHERE agent_id = ?1",
params![engineer_anchor_id("t-1400")],
)
.await
.unwrap();
assert_eq!(
rows.len(),
1,
"anchor survives job deletion (NULL FK child)"
);
assert!(rows[0].get::<Option<String>>(0).unwrap().is_none());
}
#[test]
fn engineer_anchor_id_roundtrip() {
assert_eq!(
engineer_anchor_ticket_id(&engineer_anchor_id("mahbot-1400")).as_deref(),
Some("mahbot-1400")
);
assert_eq!(
engineer_anchor_ticket_id("ticket_my_ws-42_engineer").as_deref(),
Some("my_ws-42")
);
assert!(engineer_anchor_ticket_id("ticket_1400_reviewer").is_none());
assert!(engineer_anchor_ticket_id("engineer_1400").is_none());
}
#[tokio::test]
async fn complete_job_with_envelope_atomicity() {
let (store, _tmp) = crate::open_test_store!(crate::session::SessionStore, "session");
let conn = &store.conn;
spawn_job(
conn,
"j1",
"q",
"ws",
"",
"",
crate::Role::Manager,
&[NewAgent {
agent_id: "analyze_a1".to_string(),
kind: AgentKind::Analyst,
idx: Some(0),
task: "t1".to_string(),
}],
&SpawnChild::Analyze,
)
.await
.unwrap();
let envelope = AgentJob {
content: "<analyze-tool-result>\n\nok</analyze-tool-result>".to_string(),
workspace_name: "ws".to_string(),
user_name: String::new(),
channel: String::new(),
kind: JobKind::AnalyzeToolResult,
role: crate::Role::Manager,
reply_target: None,
pending_job_id: Some("j1".to_string()),
};
complete_job_with_envelope(conn, "j1", &envelope)
.await
.unwrap();
let pending = list_pending_jobs(conn).await.unwrap();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].id, "j1");
let jobs = conn.query("SELECT COUNT(*) FROM jobs", ()).await.unwrap();
assert_eq!(jobs[0].get::<i64>(0).unwrap(), 0);
let agents = conn.query("SELECT COUNT(*) FROM agents", ()).await.unwrap();
assert_eq!(agents[0].get::<i64>(0).unwrap(), 0);
}
#[tokio::test]
#[expect(clippy::await_holding_lock)] async fn resume_job_preamble_abort_and_missing_quiet_return() {
let _lock = crate::util::test::retry_tests_lock();
let (store, _tmp) = crate::open_test_store!(crate::session::SessionStore, "session");
let conn = &store.conn;
spawn_job(
conn,
"j-preamble",
"question?",
"ws",
"caller-user",
"telegram",
Role::Manager,
&[],
&SpawnChild::Analyze,
)
.await
.unwrap();
crate::shutdown::drain_begin();
let r = resume_job_preamble(conn, "j-preamble", "Analyze resume", "Analyze resume").await;
crate::shutdown::drain_clear();
assert!(r.is_none(), "drain-abort must quiet-return");
let rows = conn
.query("SELECT COUNT(*) FROM jobs WHERE id = 'j-preamble'", ())
.await
.unwrap();
assert_eq!(
rows[0].get::<i64>(0).unwrap(),
1,
"drain-abort must not terminalize the job"
);
conn.execute("DELETE FROM jobs WHERE id = 'j-preamble'", ())
.await
.unwrap();
let r = resume_job_preamble(conn, "j-preamble", "Analyze resume", "Analyze resume").await;
assert!(r.is_none(), "missing job row must quiet-return");
let rows = conn
.query("SELECT COUNT(*) FROM jobs WHERE id = 'j-preamble'", ())
.await
.unwrap();
assert_eq!(
rows[0].get::<i64>(0).unwrap(),
0,
"terminalize must not recreate the row"
);
}
#[tokio::test]
async fn pending_replay_dedup_skips_appended() {
let (store, _tmp) = crate::open_test_store!(crate::session::SessionStore, "session");
let conn = &store.conn;
let agent_id = "manager_ws";
let content = "hello manager";
conn.execute(
"INSERT INTO sessions (agent_id, role, content, created_at) \
VALUES (?1, 'user', ?2, ?3)",
params![
agent_id,
format!("{content}\n\n<timestamp>2026-01-01 00:00:00 (UTC)</timestamp>"),
turso::now()
],
)
.await
.unwrap();
conn.execute(
"INSERT INTO pending_jobs (id, target_agent_id, envelope, created_at) \
VALUES ('p1', ?1, ?2, ?3)",
params![
agent_id,
serde_json::to_string(&AgentJob {
content: content.to_string(),
workspace_name: "ws".to_string(),
user_name: "u".to_string(),
channel: "telegram".to_string(),
kind: JobKind::UserMessage,
role: crate::Role::Manager,
reply_target: None,
pending_job_id: None,
})
.unwrap(),
"2025-12-31T00:00:00Z",
],
)
.await
.unwrap();
let rows = list_pending_jobs(conn).await.unwrap();
assert!(pending_already_appended(conn, &rows[0]).await);
let replayed = replay_pending_jobs(conn).await.unwrap();
assert_eq!(replayed, 0, "deduped row must not be re-routed");
assert_eq!(
list_pending_jobs(conn).await.unwrap().len(),
0,
"deduped row reclaimed"
);
}
#[test]
fn strip_timestamp_wrapper_both_formats() {
let body = "task text";
let suffix = format!("{body}\n\n<timestamp>2026-01-01 00:00:00 (UTC)</timestamp>");
assert_eq!(strip_timestamp_wrapper(&suffix), body);
let legacy = format!("<timestamp>2026-01-01 00:00:00 (UTC)</timestamp>\n\n{body}");
assert_eq!(strip_timestamp_wrapper(&legacy), body);
assert_eq!(
strip_timestamp_wrapper(body),
body,
"no wrapper — unchanged"
);
let drained =
format!("drained\n{body}\n\n<timestamp>2026-01-01 00:00:00 (UTC)</timestamp>");
assert!(
strip_timestamp_wrapper(&drained).ends_with(body),
"drain prefix survives stripping"
);
}
#[test]
fn strip_timestamp_wrapper_ignores_lookalikes() {
let body = "report\n\n<timestamp>2026-01-01 00:00:00 (UTC)</timestamp>";
let legacy = format!("<timestamp>2026-01-01 00:00:00 (UTC)</timestamp>\n\n{body}");
assert_eq!(strip_timestamp_wrapper(&legacy), body);
let raw = "<timestamp>config options</timestamp>\n\nrest of the message";
assert_eq!(strip_timestamp_wrapper(raw), raw);
let leading = "<timestamp>2026-01-01 00:00:00 (UTC)</timestamp> more\n\n\
<timestamp>2026-01-01 00:00:00 (UTC)</timestamp>";
assert_eq!(
strip_timestamp_wrapper(leading),
leading,
"suffix-format message whose content starts with a timestamp block stays intact"
);
}
#[tokio::test]
async fn pending_replay_dedup_recognizes_legacy_prefix_format() {
let (store, _tmp) = crate::open_test_store!(crate::session::SessionStore, "session");
let conn = &store.conn;
let agent_id = "manager_legacy";
let content = "legacy hello";
conn.execute(
"INSERT INTO sessions (agent_id, role, content, created_at) \
VALUES (?1, 'user', ?2, ?3)",
params![
agent_id,
format!("<timestamp>2026-01-01 00:00:00 (UTC)</timestamp>\n\n{content}"),
turso::now()
],
)
.await
.unwrap();
conn.execute(
"INSERT INTO pending_jobs (id, target_agent_id, envelope, created_at) \
VALUES ('p-legacy', ?1, ?2, ?3)",
params![
agent_id,
serde_json::to_string(&AgentJob {
content: content.to_string(),
workspace_name: "ws".to_string(),
user_name: "u".to_string(),
channel: "telegram".to_string(),
kind: JobKind::UserMessage,
role: crate::Role::Manager,
reply_target: None,
pending_job_id: None,
})
.unwrap(),
"2025-12-31T00:00:00Z",
],
)
.await
.unwrap();
let rows = list_pending_jobs(conn).await.unwrap();
assert!(
pending_already_appended(conn, &rows[0]).await,
"legacy prefix-format message must still dedup"
);
}
#[tokio::test]
async fn pending_replay_not_in_session_routes_full_content() {
let _ = crate::message_router::init_global();
let (store, _tmp) = crate::open_test_store!(crate::session::SessionStore, "session");
let conn = &store.conn;
let agent_id = "f1_replay_target";
let mut rx = crate::message_router::register_agent(agent_id);
let content = "FULL CONTENT";
conn.execute(
"INSERT INTO pending_jobs (id, target_agent_id, envelope, created_at) \
VALUES ('p2', ?1, ?2, ?3)",
params![
agent_id,
serde_json::to_string(&AgentJob {
content: content.to_string(),
workspace_name: "ws".to_string(),
user_name: String::new(),
channel: String::new(),
kind: JobKind::AnalyzeToolResult,
role: crate::Role::Manager,
reply_target: None,
pending_job_id: None,
})
.unwrap(),
turso::now(),
],
)
.await
.unwrap();
let replayed = replay_pending_jobs(conn).await.unwrap();
assert_eq!(replayed, 1, "row not in session must be re-routed");
let routed = rx
.try_recv()
.expect("routed job must carry the full content");
assert_eq!(
routed.content, content,
"not-in-session routes the FULL content, never empty"
);
crate::message_router::unregister_agent(agent_id);
}
#[tokio::test]
async fn ttl_guard_protects_job_tracked_sessions() {
crate::util::test::init_test_stores().await;
let store = crate::session::store();
let stale = (chrono::Utc::now() - chrono::Duration::hours(20)).to_rfc3339();
store
.batch_append_with_context(
"ticket_t1_0_suf_engineer",
&[crate::ChatMessage::user("u")],
"gui",
"u",
"ws",
"engineer",
)
.await
.unwrap();
store
.conn
.execute(
"UPDATE session_metadata SET last_activity = ?1 WHERE agent_id = 'ticket_t1_0_suf_engineer'",
params![stale.clone()],
)
.await
.unwrap();
let conn = &store.conn;
crate::util::test::JobRowBuilder::new(conn, "jttl", "ticket_stage", "engineer", "ws")
.timestamps(stale.clone())
.insert()
.await
.unwrap();
conn.execute(
"INSERT INTO agents (job_id, agent_id, kind, status, task) \
VALUES ('jttl', 'ticket_t1_0_suf_engineer', 'engineer', 'launched', '')",
(),
)
.await
.unwrap();
let cutoff =
(chrono::Utc::now() - chrono::Duration::hours(PURGE_CUTOFF_HOURS)).to_rfc3339();
crate::session::cleanup_old_transient_sessions(&cutoff)
.await
.unwrap();
let msgs = store.load("ticket_t1_0_suf_engineer").await;
assert_eq!(
msgs.len(),
1,
"job-tracked session must survive the TTL guard"
);
conn.execute("DELETE FROM jobs WHERE id = 'jttl'", ())
.await
.unwrap();
crate::session::cleanup_old_transient_sessions(&cutoff)
.await
.unwrap();
let msgs = store.load("ticket_t1_0_suf_engineer").await;
assert_eq!(
msgs.len(),
0,
"unprotected session is cleaned up after cascade"
);
}
#[tokio::test]
#[serial_test::serial(reset_inflight)]
async fn purge_rolls_back_stranded_ticket() {
crate::util::test::init_test_stores().await;
let store = crate::session::store();
let conn = &store.conn;
let board = crate::board::store();
let ws = crate::workspace::test_ws_named("/tmp/purge_ws", "purge_ws");
let ticket_id = crate::util::test::make_ticket(
board,
&ws,
"Purge me",
crate::board::TicketPhase::InDevelopment,
)
.await;
let stale = (chrono::Utc::now() - chrono::Duration::hours(20)).to_rfc3339();
crate::util::test::JobRowBuilder::new(
conn,
"jstale",
"ticket_stage",
"engineer",
"purge_ws",
)
.timestamps(stale.clone())
.insert()
.await
.unwrap();
conn.execute(
"INSERT INTO ticket_stage_jobs (id, ticket_id, stage, phase, round) \
VALUES ('jstale', ?1, 'engineer', 'in_development', 1)",
params![ticket_id.clone()],
)
.await
.unwrap();
conn.execute(
"INSERT INTO agents (job_id, agent_id, kind, status, task) \
VALUES ('jstale', 'ticket_purge-t1_engineer', 'engineer', 'launched', '')",
(),
)
.await
.unwrap();
conn.execute(
"INSERT INTO session_metadata (agent_id, last_activity) \
VALUES ('ticket_purge-t1_engineer', ?1)",
params![stale],
)
.await
.unwrap();
let cutoff =
(chrono::Utc::now() - chrono::Duration::hours(PURGE_CUTOFF_HOURS)).to_rfc3339();
let purged = purge_stale_jobs(&cutoff).await.unwrap();
assert!(purged >= 1, "the stale job must be purged");
let t = board.get_ticket(&ticket_id).await.unwrap().unwrap();
assert_eq!(t.phase, crate::board::TicketPhase::ReadyForDevelopment);
assert!(t.assigned_to.is_none());
let n = conn
.query("SELECT COUNT(*) FROM jobs WHERE id = 'jstale'", ())
.await
.unwrap();
assert_eq!(n[0].get::<i64>(0).unwrap(), 0);
let _ = &ws;
}
#[test]
fn rollback_transition_covers_reset_table() {
let cases: &[(&str, &str, bool)] = &[
("analysis", "backlog", false),
("in_development", "ready_for_development", true),
("in_review", "diagnostics_done", false),
("in_qa", "reviewed", false),
("in_sanitation", "qa_passed", true),
];
for (phase, to, reservation) in cases {
assert_eq!(
rollback_transition(phase),
Some((to.to_string(), *reservation))
);
}
for phase in ["in_diagnostics", "ready_for_development", "bogus"] {
assert_eq!(rollback_transition(phase), None);
}
}
#[tokio::test]
#[serial_test::serial(reset_inflight)]
async fn purge_round_guard_skips_older_round() {
crate::util::test::init_test_stores().await;
let conn = &crate::session::store().conn;
let board = crate::board::store();
let ws = crate::workspace::test_ws_named("/tmp/purge_ws2", "purge_ws2");
let ticket_id = crate::util::test::make_ticket(
board,
&ws,
"Round guard",
crate::board::TicketPhase::InReview,
)
.await;
let stale = (chrono::Utc::now() - chrono::Duration::hours(20)).to_rfc3339();
for (id, round) in [("jr1", 1i64), ("jr2", 2i64)] {
let ts = if round == 1 {
stale.clone()
} else {
turso::now()
};
crate::util::test::JobRowBuilder::new(
conn,
id,
"ticket_stage",
"reviewer",
"purge_ws2",
)
.timestamps(ts.clone())
.insert()
.await
.unwrap();
conn.execute(
"INSERT INTO ticket_stage_jobs (id, ticket_id, stage, phase, round) \
VALUES (?1, ?2, 'review', 'in_review', ?3)",
params![id, ticket_id.clone(), round],
)
.await
.unwrap();
}
let cutoff =
(chrono::Utc::now() - chrono::Duration::hours(PURGE_CUTOFF_HOURS)).to_rfc3339();
purge_stale_jobs(&cutoff).await.unwrap();
let t = board.get_ticket(&ticket_id).await.unwrap().unwrap();
assert_eq!(
t.phase,
crate::board::TicketPhase::InReview,
"older round's purge must not roll back the ticket (round guard)"
);
}
#[tokio::test]
#[serial_test::serial(reset_inflight)]
async fn purge_round_guard_scopes_per_stage() {
crate::util::test::init_test_stores().await;
let conn = &crate::session::store().conn;
let board = crate::board::store();
let ws = crate::workspace::test_ws_named("/tmp/purge_ws4", "purge_ws4");
let ticket_id = crate::util::test::make_ticket(
board,
&ws,
"Per-stage round guard",
crate::board::TicketPhase::InReview,
)
.await;
let stale = (chrono::Utc::now() - chrono::Duration::hours(20)).to_rfc3339();
for (id, stage, phase, role, round) in [
("jsa2", "analysis", "analysis", "analyst", 2i64),
("jsr1", "review", "in_review", "reviewer", 1i64),
] {
crate::util::test::JobRowBuilder::new(conn, id, "ticket_stage", role, "purge_ws4")
.timestamps(stale.clone())
.insert()
.await
.unwrap();
conn.execute(
"INSERT INTO ticket_stage_jobs (id, ticket_id, stage, phase, round) \
VALUES (?1, ?2, ?3, ?4, ?5)",
params![id, ticket_id.clone(), stage, phase, round],
)
.await
.unwrap();
}
let cutoff =
(chrono::Utc::now() - chrono::Duration::hours(PURGE_CUTOFF_HOURS)).to_rfc3339();
purge_stale_jobs(&cutoff).await.unwrap();
let t = board.get_ticket(&ticket_id).await.unwrap().unwrap();
assert_eq!(
t.phase,
crate::board::TicketPhase::DiagnosticsDone,
"stale review round must roll the ticket back out of in_review"
);
for id in ["jsa2", "jsr1"] {
let n = conn
.query("SELECT COUNT(*) FROM jobs WHERE id = ?1", params![id])
.await
.unwrap();
assert_eq!(n[0].get::<i64>(0).unwrap(), 0, "{id} must be purged");
}
}
#[tokio::test]
#[serial_test::serial(reset_inflight)]
async fn purge_keeps_ticket_stage_rows_when_rollback_fails() {
crate::util::test::init_test_stores().await;
let conn = &crate::session::store().conn;
let stale = (chrono::Utc::now() - chrono::Duration::hours(20)).to_rfc3339();
crate::util::test::JobRowBuilder::new(
conn,
"jfail_ts",
"ticket_stage",
"engineer",
"purge_ws3",
)
.timestamps(stale.clone())
.insert()
.await
.unwrap();
conn.execute(
"INSERT INTO ticket_stage_jobs (id, ticket_id, stage, phase, round) \
VALUES ('jfail_ts', 'ticket-missing', 'engineer', 'in_development', 1)",
(),
)
.await
.unwrap();
crate::util::test::JobRowBuilder::new(
conn,
"jfail_analyze",
"analyze",
"assistant",
"purge_ws3",
)
.timestamps(stale.clone())
.insert()
.await
.unwrap();
crate::board::store()
.conn
.execute("BEGIN", ())
.await
.expect("raw board BEGIN");
let cutoff =
(chrono::Utc::now() - chrono::Duration::hours(PURGE_CUTOFF_HOURS)).to_rfc3339();
purge_stale_jobs(&cutoff).await.unwrap();
crate::board::store()
.conn
.execute("ROLLBACK", ())
.await
.expect("restore raw board tx");
let ts_left = conn
.query("SELECT COUNT(*) FROM jobs WHERE id = 'jfail_ts'", ())
.await
.unwrap();
assert_eq!(
ts_left[0].get::<i64>(0).unwrap(),
1,
"ticket_stage row must survive a failed rollback (next tick retries)"
);
let analyze_left = conn
.query("SELECT COUNT(*) FROM jobs WHERE id = 'jfail_analyze'", ())
.await
.unwrap();
assert_eq!(
analyze_left[0].get::<i64>(0).unwrap(),
0,
"non-ticket_stage stale rows are purged regardless of the rollback"
);
}
#[tokio::test]
#[serial_test::serial(reset_inflight)]
async fn reset_inflight_tickets_exclusion() {
crate::util::test::init_test_stores().await;
let board = crate::board::store();
let ws = crate::workspace::test_ws_named("/tmp/ws_x", "ws_x");
let resumed = crate::util::test::make_ticket(
board,
&ws,
"resume me",
crate::board::TicketPhase::InDevelopment,
)
.await;
let reset = crate::util::test::make_ticket(
board,
&ws,
"reset me",
crate::board::TicketPhase::InDevelopment,
)
.await;
board
.reset_inflight_tickets(std::slice::from_ref(&resumed))
.await
.unwrap();
let t_resumed = board.get_ticket(&resumed).await.unwrap().unwrap();
let t_reset = board.get_ticket(&reset).await.unwrap().unwrap();
assert_eq!(t_resumed.phase, crate::board::TicketPhase::InDevelopment);
assert_eq!(
t_reset.phase,
crate::board::TicketPhase::ReadyForDevelopment
);
}
#[tokio::test]
#[serial_test::serial(reset_inflight)]
async fn recover_from_restart_excludes_resumed_tickets() {
crate::util::test::init_management_test_stores().await;
let conn = &crate::session::store().conn;
let board = crate::board::store();
let ws = crate::workspace::test_ws_named("/tmp/ws_scan", "ws_scan");
let resumed_ticket = crate::util::test::make_ticket(
board,
&ws,
"resumed",
crate::board::TicketPhase::InDevelopment,
)
.await;
let reset_ticket = crate::util::test::make_ticket(
board,
&ws,
"reset",
crate::board::TicketPhase::InDevelopment,
)
.await;
let now = crate::turso::now();
crate::util::test::JobRowBuilder::new(conn, "jscan", "ticket_stage", "engineer", "ws_scan")
.timestamps(now.clone())
.insert()
.await
.unwrap();
conn.execute(
"INSERT INTO ticket_stage_jobs (id, ticket_id, stage, phase, round) \
VALUES ('jscan', ?1, 'engineer', 'in_development', 1)",
params![resumed_ticket.clone()],
)
.await
.unwrap();
let resumable = recover_from_restart().await.unwrap();
assert!(
resumable.iter().any(|r| matches!(
r,
ResumableStage::TicketStage { job_id, .. } if job_id.as_str() == "jscan"
)),
"the in-phase ticket_stage job must be selected for resume"
);
assert!(resumable.iter().any(|r| matches!(
r,
ResumableStage::TicketStage { ticket_id, .. } if ticket_id == &resumed_ticket
)));
let t1 = board.get_ticket(&resumed_ticket).await.unwrap().unwrap();
assert_eq!(
t1.phase,
crate::board::TicketPhase::InDevelopment,
"excluded ticket must NOT reset at boot"
);
let t2 = board.get_ticket(&reset_ticket).await.unwrap().unwrap();
assert_eq!(
t2.phase,
crate::board::TicketPhase::ReadyForDevelopment,
"unexcluded ticket resets normally"
);
let rc = conn
.query_row("SELECT retry_count FROM jobs WHERE id = 'jscan'", (), |r| {
r.get::<i64>(0)
})
.await
.unwrap();
assert_eq!(rc, 1, "retry_count must be bumped per boot resume");
let bumped = conn
.query_row("SELECT updated_at FROM jobs WHERE id = 'jscan'", (), |r| {
r.get::<String>(0)
})
.await
.unwrap();
let before = crate::turso::parse_utc_timestamp(&now).unwrap();
let after = crate::turso::parse_utc_timestamp(&bumped).unwrap();
assert!(after > before, "boot bump must refresh updated_at");
}
#[tokio::test]
#[serial_test::serial(reset_inflight)] async fn recover_from_restart_caps_research_cleanup_releases_folder() {
crate::util::test::init_management_test_stores().await;
let conn = &crate::session::store().conn;
let now = crate::turso::now();
crate::util::test::JobRowBuilder::new(
conn,
"jclean",
"research_cleanup",
"sanitation",
"ws_scan",
)
.retry_count(MAX_BOOT_REDISPATCH)
.timestamps(now.clone())
.insert()
.await
.unwrap();
let run_folder = crate::research_cleanup::run_root_path("jclean");
tokio::fs::create_dir_all(&run_folder).await.unwrap();
let resumable = recover_from_restart().await.unwrap();
assert!(
!resumable.iter().any(|r| matches!(
r,
ResumableStage::ResearchCleanup { job_id, .. } if job_id.as_str() == "jclean"
)),
"capped cleanup must NOT be selected for resume (left for the purge)"
);
assert!(
!run_folder.exists(),
"the cap branch releases the run folder in-process (no sweep exists)"
);
let updated: String = conn
.query_row("SELECT updated_at FROM jobs WHERE id = 'jclean'", (), |r| {
r.get::<String>(0)
})
.await
.unwrap();
assert_eq!(
updated, now,
"capped cleanup must not be checkpointed — updated_at unchanged so the 8h purge can reclaim the row"
);
}
#[tokio::test]
#[serial_test::serial(reset_inflight)] async fn recover_from_restart_terminalizes_temp_cleanup_rows() {
crate::util::test::init_management_test_stores().await;
let conn = &crate::session::store().conn;
let now = crate::turso::now();
crate::util::test::JobRowBuilder::new(
conn,
"jtmpclean",
"temp_cleanup",
"sanitation",
"tmp",
)
.timestamps(now.clone())
.insert()
.await
.unwrap();
let resumable = recover_from_restart().await.unwrap();
let remaining = conn
.query(
"SELECT COUNT(*) FROM jobs WHERE id = 'jtmpclean' AND kind = 'temp_cleanup'",
(),
)
.await
.unwrap();
assert_eq!(
remaining[0].get::<i64>(0).unwrap(),
0,
"temp_cleanup row must be terminalized at boot (fire-and-forget, no resume)"
);
assert!(
!resumable.iter().any(|r| matches!(
r,
ResumableStage::ResearchCleanup { job_id, .. } if job_id.as_str() == "jtmpclean"
)),
"no resumable entry may refer to the temp_cleanup row"
);
}
#[tokio::test]
#[serial_test::serial(reset_inflight)] async fn replay_skips_cleanup_dispatch_when_cleanup_row_exists() {
crate::util::test::init_management_test_stores().await;
let conn = &crate::session::store().conn;
crate::util::test::create_test_workspace("/tmp/test_ws_replay_cleanup", "ws_replay").await;
let now = crate::turso::now();
crate::util::test::JobRowBuilder::new(
conn,
"rcln",
"research_cleanup",
"sanitation",
"ws_replay",
)
.timestamps(now.clone())
.insert()
.await
.unwrap();
conn.execute(
"INSERT INTO pending_jobs (id, target_agent_id, envelope, created_at) \
VALUES ('rcln', 'manager_ws_replay', ?1, ?2)",
params![
serde_json::to_string(&AgentJob {
content: "<research-result>done</research-result>".to_string(),
workspace_name: "ws_replay".to_string(),
user_name: "u".to_string(),
channel: "telegram".to_string(),
kind: JobKind::ResearchResult,
role: crate::Role::Manager,
reply_target: None,
pending_job_id: None,
})
.unwrap(),
now
],
)
.await
.unwrap();
let replayed = replay_pending_jobs(conn).await.unwrap();
assert_eq!(replayed, 1, "the envelope itself is replayed");
let rows = conn
.query(
"SELECT COUNT(*) FROM jobs WHERE id = 'rcln' AND kind = 'research_cleanup'",
(),
)
.await
.unwrap();
assert_eq!(
rows[0].get::<i64>(0).unwrap(),
1,
"no duplicate cleanup dispatch — the existing row is the dedup marker"
);
let sessions = conn
.query(
"SELECT COUNT(*) FROM session_metadata WHERE agent_id = 'cleanup_rcln'",
(),
)
.await
.unwrap();
assert_eq!(
sessions[0].get::<i64>(0).unwrap(),
0,
"the scan must not spawn a second cleanup agent"
);
conn.execute("DELETE FROM pending_jobs WHERE id = 'rcln'", ())
.await
.unwrap();
conn.execute("DELETE FROM jobs WHERE id = 'rcln'", ())
.await
.unwrap();
}
#[tokio::test]
#[serial_test::serial(reset_inflight)] async fn replay_creates_cleanup_row_then_scan_resumes() {
crate::util::test::init_management_test_stores().await;
let conn = &crate::session::store().conn;
crate::util::test::create_test_workspace("/tmp/test_ws_replay_resume", "ws_replay2").await;
let now = crate::turso::now();
conn.execute(
"INSERT INTO pending_jobs (id, target_agent_id, envelope, created_at) \
VALUES ('rcln2', 'manager_ws_replay2', ?1, ?2)",
params![
serde_json::to_string(&AgentJob {
content: "<research-result>done</research-result>".to_string(),
workspace_name: "ws_replay2".to_string(),
user_name: "u".to_string(),
channel: "telegram".to_string(),
kind: JobKind::ResearchResult,
role: crate::Role::Manager,
reply_target: None,
pending_job_id: None,
})
.unwrap(),
now
],
)
.await
.unwrap();
let run_folder = crate::research_cleanup::run_root_path("rcln2");
tokio::fs::create_dir_all(&run_folder).await.unwrap();
let resumable = recover_from_restart().await.unwrap();
assert!(
crate::research_cleanup::research_cleanup_row_exists(conn, "rcln2")
.await
.unwrap(),
"the crash-window replay must create the cleanup row"
);
assert!(
resumable.iter().any(|r| matches!(
r,
ResumableStage::ResearchCleanup { job_id, .. } if job_id.as_str() == "rcln2"
)),
"the boot-scan arm must resume the replay-created cleanup row"
);
let sessions = conn
.query(
"SELECT COUNT(*) FROM session_metadata WHERE agent_id = 'cleanup_rcln2'",
(),
)
.await
.unwrap();
assert_eq!(
sessions[0].get::<i64>(0).unwrap(),
0,
"the replay must not spawn the cleanup agent — the boot-scan arm resumes it"
);
conn.execute("DELETE FROM pending_jobs WHERE id = 'rcln2'", ())
.await
.unwrap();
conn.execute("DELETE FROM jobs WHERE id = 'rcln2'", ())
.await
.unwrap();
let _ = tokio::fs::remove_dir_all(&run_folder).await;
}
#[tokio::test]
#[serial_test::serial(reset_inflight)] async fn replay_skips_cleanup_dispatch_when_run_folder_gone() {
crate::util::test::init_management_test_stores().await;
let conn = &crate::session::store().conn;
crate::util::test::create_test_workspace("/tmp/test_ws_replay_archived", "ws_replay3")
.await;
let now = crate::turso::now();
let run_folder = crate::research_cleanup::run_root_path("rcln3");
tokio::fs::create_dir_all(&run_folder).await.unwrap();
crate::research_cleanup::release_run_folder("rcln3").await;
assert!(!run_folder.exists(), "fixture: run folder released");
conn.execute(
"INSERT INTO pending_jobs (id, target_agent_id, envelope, created_at) \
VALUES ('rcln3', 'manager_ws_replay3', ?1, ?2)",
params![
serde_json::to_string(&AgentJob {
content: "<research-result>done</research-result>".to_string(),
workspace_name: "ws_replay3".to_string(),
user_name: "u".to_string(),
channel: "telegram".to_string(),
kind: JobKind::ResearchResult,
role: crate::Role::Manager,
reply_target: None,
pending_job_id: None,
})
.unwrap(),
now
],
)
.await
.unwrap();
let replayed = replay_pending_jobs(conn).await.unwrap();
assert_eq!(replayed, 1, "the envelope itself is replayed");
assert!(
!crate::research_cleanup::research_cleanup_row_exists(conn, "rcln3")
.await
.unwrap(),
"a released run folder means the cleanup completed — no re-dispatch"
);
let sessions = conn
.query(
"SELECT COUNT(*) FROM session_metadata WHERE agent_id = 'cleanup_rcln3'",
(),
)
.await
.unwrap();
assert_eq!(
sessions[0].get::<i64>(0).unwrap(),
0,
"no cleanup agent for a completed cleanup"
);
conn.execute("DELETE FROM pending_jobs WHERE id = 'rcln3'", ())
.await
.unwrap();
}
#[tokio::test]
#[serial_test::serial(reset_inflight)]
async fn recover_from_restart_classifies_done_phase_left_dedupes_rounds() {
crate::util::test::init_management_test_stores().await;
let conn = &crate::session::store().conn;
let board = crate::board::store();
let ws = crate::workspace::test_ws_named("/tmp/ws_cls", "ws_cls");
let now = crate::turso::now();
let dup_ticket = crate::util::test::make_ticket(
board,
&ws,
"dup",
crate::board::TicketPhase::InDevelopment,
)
.await;
for (id, round) in [("jdup1", 1i64), ("jdup2", 2i64)] {
crate::util::test::JobRowBuilder::new(conn, id, "ticket_stage", "engineer", "ws_cls")
.timestamps(now.clone())
.insert()
.await
.unwrap();
conn.execute(
"INSERT INTO ticket_stage_jobs (id, ticket_id, stage, phase, round) \
VALUES (?1, ?2, 'engineer', 'in_development', ?3)",
params![id, dup_ticket.clone(), round],
)
.await
.unwrap();
}
let phase_left = crate::util::test::make_ticket(
board,
&ws,
"phase_left",
crate::board::TicketPhase::InQa,
)
.await;
crate::util::test::JobRowBuilder::new(conn, "jcls2", "ticket_stage", "reviewer", "ws_cls")
.timestamps(now.clone())
.insert()
.await
.unwrap();
conn.execute(
"INSERT INTO ticket_stage_jobs (id, ticket_id, stage, phase, round) \
VALUES ('jcls2', ?1, 'review', 'in_review', 1)",
params![phase_left.clone()],
)
.await
.unwrap();
let done_ticket = crate::util::test::make_ticket(
board,
&ws,
"done",
crate::board::TicketPhase::InDevelopment,
)
.await;
crate::util::test::JobRowBuilder::new(conn, "jcls3", "ticket_stage", "engineer", "ws_cls")
.status("done")
.timestamps(now.clone())
.insert()
.await
.unwrap();
conn.execute(
"INSERT INTO ticket_stage_jobs (id, ticket_id, stage, phase, round) \
VALUES ('jcls3', ?1, 'engineer', 'in_development', 1)",
params![done_ticket.clone()],
)
.await
.unwrap();
let resumable = recover_from_restart().await.unwrap();
assert!(resumable.iter().any(|r| matches!(
r,
ResumableStage::TicketStage { job_id, .. } if job_id.as_str() == "jdup2"
)));
assert!(!resumable.iter().any(|r| matches!(
r,
ResumableStage::TicketStage { job_id, .. } if job_id.as_str() == "jdup1"
)));
let status1 = conn
.query_row("SELECT status FROM jobs WHERE id = 'jdup1'", (), |r| {
r.get::<String>(0)
})
.await
.unwrap();
assert_eq!(
status1, "done",
"superseded older round must be marked done"
);
assert!(!resumable.iter().any(|r| matches!(
r,
ResumableStage::TicketStage { job_id, .. }
if job_id.as_str() == "jcls2" || job_id.as_str() == "jcls3"
)));
let status2 = conn
.query_row("SELECT status FROM jobs WHERE id = 'jcls2'", (), |r| {
r.get::<String>(0)
})
.await
.unwrap();
assert_eq!(status2, "done");
let t_dup = board.get_ticket(&dup_ticket).await.unwrap().unwrap();
assert_eq!(t_dup.phase, crate::board::TicketPhase::InDevelopment);
let t2 = board.get_ticket(&phase_left).await.unwrap().unwrap();
assert_eq!(t2.phase, crate::board::TicketPhase::Reviewed);
let t3 = board.get_ticket(&done_ticket).await.unwrap().unwrap();
assert_eq!(t3.phase, crate::board::TicketPhase::ReadyForDevelopment);
}
}