use crate::Role;
use crate::agent::message_router::{AgentJob, MessageKind};
use crate::db::{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,
Coder,
Sanitation,
Diagnostics,
}
impl AgentKind {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Analyst => "analyst",
Self::Verifier => "verifier",
Self::Engineer => "engineer",
Self::Coder => "coder",
Self::Sanitation => "sanitation",
Self::Diagnostics => "diagnostics",
}
}
}
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,
pub ticket_id: Option<String>,
}
#[derive(Debug, Clone)]
pub(crate) struct AgentRow {
pub idx: Option<i64>,
pub agent_id: String,
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)?,
ticket_id: row.get(4)?,
})
}
fn agent_row_from(row: &Row) -> anyhow::Result<AgentRow> {
Ok(AgentRow {
idx: row.get(0)?,
agent_id: 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,
Implement,
Research,
ResearchCleanup,
TempCleanup,
Phase {
phase: crate::pipeline::board::TicketPhase,
ticket_id: String,
},
}
impl SpawnChild {
#[must_use]
pub fn kind_str(&self) -> &'static str {
match self {
Self::Analyze => "analyze",
Self::Implement => "implement",
Self::Research => "research",
Self::ResearchCleanup => "research_cleanup",
Self::TempCleanup => "temp_cleanup",
Self::Phase { phase, .. } => match phase {
crate::pipeline::board::TicketPhase::Analysis => "analysis",
crate::pipeline::board::TicketPhase::InDevelopment => "in_development",
crate::pipeline::board::TicketPhase::InDiagnostics => "in_diagnostics",
crate::pipeline::board::TicketPhase::InReview => "in_review",
crate::pipeline::board::TicketPhase::InQa => "in_qa",
crate::pipeline::board::TicketPhase::InSanitation => "in_sanitation",
_ => unreachable!("non-working phase as a job kind"),
},
}
}
}
#[must_use]
pub(crate) fn child_ticket_id(child: &SpawnChild) -> Option<&str> {
match child {
SpawnChild::Phase { ticket_id, .. } => Some(ticket_id.as_str()),
_ => None,
}
}
#[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 ticket_id = child_ticket_id(child);
let now = db::now();
let tx = conn.begin_tx().await?;
tx.execute(
"INSERT INTO jobs (id, kind, status, task, workspace_name, user_name, channel, role, \
ticket_id, retry_count, created_at, updated_at) \
VALUES (?1, ?2, 'launched', ?3, ?4, ?5, ?6, ?7, ?8, 0, ?9, ?9)",
params![
id,
kind,
task,
workspace_name,
user_name,
channel,
role.as_str(),
ticket_id,
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::Implement
| SpawnChild::ResearchCleanup
| SpawnChild::TempCleanup
| SpawnChild::Phase { .. } => {}
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}"))?;
}
}
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 = db::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 = db::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: MessageKind,
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 terminalize_job(conn: &Connection, job_id: &str) -> Result<()> {
conn.execute("DELETE FROM jobs WHERE id = ?1", params![job_id])
.await
.context("terminalize job")?;
Ok(())
}
#[derive(Debug, Clone)]
pub(crate) struct TicketJobRow {
pub id: String,
}
fn ticket_job_row_from(row: &Row) -> anyhow::Result<TicketJobRow> {
Ok(TicketJobRow { id: row.get(0)? })
}
pub(crate) async fn find_phase_job(
conn: &Connection,
ticket_id: &str,
phase: crate::pipeline::board::TicketPhase,
) -> Result<Option<TicketJobRow>> {
conn.query_optional_cached(
"SELECT id FROM jobs \
WHERE ticket_id = ?1 AND kind = ?2 AND status = 'launched' LIMIT 1",
params![ticket_id, phase.as_ref()],
ticket_job_row_from,
)
.await
.context("find ticket phase job")
}
pub(crate) async fn complete_ticket_phase_jobs(conn: &Connection, ticket_id: &str) -> Result<()> {
let rows = conn
.query(
"SELECT id FROM jobs WHERE ticket_id = ?1 AND status = 'launched'",
params![ticket_id],
)
.await
.context("find ticket phase jobs for completion")?;
for row in rows {
let id: String = row.get(0)?;
terminalize_job(conn, &id).await?;
}
Ok(())
}
pub(crate) async fn update_phase_job_task(
conn: &Connection,
job_id: &str,
task: &str,
) -> Result<()> {
conn.execute(
"UPDATE jobs SET task = ?2, updated_at = ?3 WHERE id = ?1",
params![job_id, task, db::now()],
)
.await
.context("update ticket phase job task")?;
Ok(())
}
pub(crate) async fn upsert_job_agent(
conn: &Connection,
job_id: &str,
agent_id: &str,
kind: AgentKind,
status: RowStatus,
) -> Result<()> {
conn.execute(
"INSERT INTO agents (job_id, agent_id, kind, idx, status, task) \
VALUES (?1, ?2, ?3, NULL, ?4, '') \
ON CONFLICT(job_id, agent_id) DO UPDATE SET status = ?4",
params![job_id, agent_id, kind.as_str(), status.as_str()],
)
.await
.with_context(|| format!("failed to upsert job agent {agent_id} for job {job_id}"))?;
Ok(())
}
pub(crate) async fn list_running_agents_for_ticket(
conn: &Connection,
ticket_id: &str,
) -> Result<Vec<String>> {
let rows = conn
.query(
"SELECT a.agent_id FROM agents a \
JOIN jobs j ON j.id = a.job_id \
WHERE j.ticket_id = ?1 AND a.status = 'launched'",
params![ticket_id],
)
.await
.context("list running agents for ticket")?;
let mut out = Vec::with_capacity(rows.len());
for row in rows {
out.push(row.get::<String>(0)?);
}
Ok(out)
}
pub(crate) async fn job_has_launched_agents(conn: &Connection, job_id: &str) -> Result<bool> {
Ok(conn
.query_optional_cached(
"SELECT 1 FROM agents WHERE job_id = ?1 AND status = 'launched' LIMIT 1",
params![job_id],
|_| Ok::<_, anyhow::Error>(1_i64),
)
.await?
.is_some())
}
#[derive(Debug)]
pub(crate) struct SlotResume<'a> {
pub done: Vec<&'a AgentRow>,
pub not_done: Vec<&'a AgentRow>,
}
#[must_use]
pub(crate) fn split_slot_resume(roster: &[AgentRow]) -> SlotResume<'_> {
let mut done = Vec::new();
let mut not_done = Vec::new();
for row in roster {
if row.status == RowStatus::Done.as_str() {
done.push(row);
} else {
not_done.push(row);
}
}
SlotResume { done, not_done }
}
pub(crate) async fn interrupt_phase_job_roster(conn: &Connection, job_id: &str) -> Result<()> {
conn.execute(
"UPDATE agents SET status = 'failed' WHERE job_id = ?1 AND status = 'launched'",
params![job_id],
)
.await
.context("mark interrupted phase roster rows failed")?;
Ok(())
}
pub(crate) async fn rearm_roster_launched(
conn: &Connection,
job_id: &str,
agent_ids: &[String],
) -> Result<()> {
if agent_ids.is_empty() {
return Ok(());
}
let placeholders = db::sql_in_placeholders(agent_ids.len());
let mut params: Vec<Value> = vec![job_id.into()];
params.extend(agent_ids.iter().map(|id| Value::from(id.as_str())));
conn.execute(
&format!(
"UPDATE agents SET status = 'launched' WHERE job_id = ?1 AND agent_id IN ({placeholders})"
),
params,
)
.await
.context("re-arm resumed roster slots as launched")?;
conn.execute(
"UPDATE jobs SET updated_at = ?2 WHERE id = ?1",
params![job_id, db::now()],
)
.await
.context("touch job updated_at on slot resume")?;
Ok(())
}
pub(crate) async fn clear_launched_agents_for_job(conn: &Connection, job_id: &str) -> Result<()> {
conn.execute(
"DELETE FROM agents WHERE job_id = ?1 AND status = 'launched'",
params![job_id],
)
.await
.context("clear launched agents for 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, ticket_id 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 idx, agent_id, 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::{force_cancel, shutdown, shutdown_token};
use std::time::Instant;
let token = shutdown_token();
tokio::select! {
() = crate::shutdown::drain_wait() => {}
() = token.cancelled() => return,
}
if shutdown_token().is_cancelled() {
return;
}
let start = Instant::now();
loop {
if crate::agent::registry::AGENT_REGISTRY.list().is_empty()
&& crate::agent::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::agent::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 ResumableJob {
Research {
job_id: String,
workspace_name: String,
},
Analyze {
job_id: String,
workspace_name: String,
},
Implement {
job_id: String,
workspace_name: String,
},
ResearchCleanup {
job_id: String,
workspace_name: String,
},
}
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::db::parse_utc_timestamp(&row.created_at).ok(),
crate::db::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 == MessageKind::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::agent::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, j.workspace_name, j.ticket_id \
FROM jobs j \
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 paused_ws: std::collections::HashSet<String> = match crate::workspace::store().list().await
{
Ok(list) => list
.into_iter()
.filter(|w| w.paused)
.map(|w| w.name)
.collect(),
Err(_) => std::collections::HashSet::new(),
};
let mut purge_ids: Vec<String> = Vec::with_capacity(rows.len());
for row in &rows {
let id: String = row.get(0)?;
let kind: String = row.get(1)?;
let workspace_name: String = row.get(2)?;
if is_ticket_phase_kind(&kind) && paused_ws.contains(&workspace_name) {
continue;
}
purge_ids.push(id);
}
if purge_ids.is_empty() {
return Ok(0);
}
let tx = conn.begin_tx().await?;
let mut deleted = 0usize;
for id in &purge_ids {
delete_job_tx(&tx, id).await?;
deleted += 1;
}
tx.commit().await?;
if deleted > 0 {
tracing::debug!(deleted, "Purged stale jobs");
}
Ok(deleted as u64)
}
pub(crate) fn is_ticket_phase_kind(kind: &str) -> bool {
matches!(
kind,
"analysis" | "in_development" | "in_diagnostics" | "in_review" | "in_qa" | "in_sanitation"
)
}
#[expect(clippy::cast_possible_truncation)]
pub(crate) async fn recover_from_restart() -> Result<Vec<ResumableJob>> {
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 resumable: Vec<ResumableJob> = Vec::new();
let mut resumed_other = 0usize;
for job in &jobs {
if is_ticket_phase_kind(&job.kind) {
if let Err(e) = interrupt_phase_job_roster(conn, &job.id).await {
warn!(
job = %job.id,
ticket = %job.ticket_id.as_deref().unwrap_or("?"),
error = %e,
"Failed to mark stale launched agents on ticket phase boot scan",
);
}
} else if job.kind == "research" || job.kind == "analyze" {
let kind = job.kind.as_str();
let _ = checkpoint_job(conn, &job.id, job.retry_count + 1).await;
resumable.push(if kind == "research" {
ResumableJob::Research {
job_id: job.id.clone(),
workspace_name: job.workspace_name.clone(),
}
} else {
ResumableJob::Analyze {
job_id: job.id.clone(),
workspace_name: job.workspace_name.clone(),
}
});
resumed_other += 1;
} else if job.kind == "implement" {
let _ = checkpoint_job(conn, &job.id, job.retry_count + 1).await;
resumable.push(ResumableJob::Implement {
job_id: job.id.clone(),
workspace_name: job.workspace_name.clone(),
});
resumed_other += 1;
} else if job.kind == "research_cleanup" {
let _ = checkpoint_job(conn, &job.id, job.retry_count + 1).await;
resumable.push(ResumableJob::ResearchCleanup {
job_id: job.id.clone(),
workspace_name: job.workspace_name.clone(),
});
resumed_other += 1;
} else if job.kind == "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;
} else {
warn!(job = %job.id, kind = %job.kind, "Unknown job kind — skipping");
}
}
let elapsed = start.elapsed();
info!(
duration_ms = elapsed.as_millis() as u64,
resumed_other,
replayed_pending = replayed,
"Boot recovery scan complete",
);
let _ = purge_terminal_session_pins().await;
Ok(resumable)
}
#[must_use]
pub(crate) fn session_pin_id(ticket_id: &str, role: Role) -> String {
crate::session::ticket_agent_id(ticket_id, role.as_str())
}
pub(crate) async fn upsert_session_pin(
conn: &Connection,
ticket_id: &str,
task: &str,
status: RowStatus,
role: Role,
) -> Result<()> {
let pin_id = session_pin_id(ticket_id, role);
conn.execute(
"INSERT INTO agents (job_id, agent_id, kind, idx, status, outcome, task) \
VALUES (NULL, ?1, ?4, NULL, ?2, NULL, ?3) \
ON CONFLICT(agent_id) WHERE job_id IS NULL \
DO UPDATE SET status = ?2, task = ?3, outcome = NULL",
params![pin_id, status.as_str(), task, role.as_str()],
)
.await
.with_context(|| format!("failed to upsert {role} session pin for ticket {ticket_id}"))?;
Ok(())
}
#[must_use]
pub(crate) fn session_pin_ticket_id(agent_id: &str, role: Role) -> Option<String> {
agent_id
.strip_prefix("ticket_")
.and_then(|rest| {
rest.strip_suffix(role.as_str())
.and_then(|r| r.strip_suffix('_'))
})
.filter(|id| !id.is_empty())
.map(str::to_string)
}
pub(crate) async fn purge_terminal_session_pins() -> usize {
let conn = &crate::session::store().conn;
let Ok(rows) = conn
.query(
"SELECT agent_id, kind FROM agents WHERE job_id IS NULL \
AND kind IN ('engineer', 'sanitation')",
(),
)
.await
else {
return 0;
};
let mut deleted = 0usize;
let board_ready = crate::pipeline::board::BOARD.get().is_some();
for row in &rows {
let Ok(agent_id) = row.get::<String>(0) else {
continue;
};
let Ok(kind) = row.get::<String>(1) else {
continue;
};
let role = match kind.as_str() {
"engineer" => Role::Engineer,
_ => Role::Sanitation,
};
let ticket_id = session_pin_ticket_id(&agent_id, role);
let Some(ticket_id) = ticket_id else {
continue;
};
let terminal = board_ready
&& crate::pipeline::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 session pin");
}
}
}
if deleted > 0 {
info!(deleted, "Removed session pins for terminal tickets");
}
deleted
}