use crate::Role;
use crate::agent::message_router::{AgentJob, MessageKind};
use crate::db::{self, Connection, Row, TxGuard, Value, params};
use crate::pipeline::board::{TicketPhase, WORKING_PHASES};
use anyhow::{Context, Result};
use std::time::Duration;
use tracing::{debug, error, info, warn};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) 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)]
enum JobMode {
Sync,
Async,
}
impl JobMode {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Sync => "sync",
Self::Async => "async",
}
}
}
impl std::str::FromStr for JobMode {
type Err = anyhow::Error;
fn from_str(s: &str) -> Result<Self, Self::Err> {
match s {
"sync" => Ok(Self::Sync),
"async" => Ok(Self::Async),
_ => Err(anyhow::anyhow!(
"Invalid job mode '{s}'. Valid modes: sync, async"
)),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum AgentKind {
Analyst,
Reviewer,
Tester,
Engineer,
Coder,
Sanitation,
Diagnostics,
}
impl AgentKind {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Analyst => "analyst",
Self::Reviewer => "reviewer",
Self::Tester => "tester",
Self::Engineer => "engineer",
Self::Coder => "coder",
Self::Sanitation => "sanitation",
Self::Diagnostics => "diagnostics",
}
}
}
pub(crate) 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>,
pub caller_agent_id: Option<String>,
mode: JobMode,
}
#[derive(Debug, Clone)]
pub(crate) struct AgentRow {
pub idx: Option<i64>,
pub agent_id: String,
pub kind: 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,
}
crate::columns! {
JOB_COLUMNS [JOB] {
ID => "id",
KIND => "kind",
WORKSPACE_NAME => "workspace_name",
RETRY_COUNT => "retry_count",
TICKET_ID => "ticket_id",
CALLER_AGENT_ID => "caller_agent_id",
MODE => "mode",
}
}
fn job_row_from(row: &Row) -> anyhow::Result<JobRow> {
Ok(JobRow {
id: row.get(COL_JOB_ID)?,
kind: row.get(COL_JOB_KIND)?,
workspace_name: row.get(COL_JOB_WORKSPACE_NAME)?,
retry_count: row.get(COL_JOB_RETRY_COUNT)?,
ticket_id: row.get(COL_JOB_TICKET_ID)?,
caller_agent_id: row.get(COL_JOB_CALLER_AGENT_ID)?,
mode: match row.get::<Option<String>>(COL_JOB_MODE)? {
Some(s) => s.parse().unwrap_or_else(|e| {
warn!(mode = %s, error = %e, "Invalid jobs.mode — falling back to async");
JobMode::Async
}),
None => JobMode::Async,
},
})
}
crate::columns! {
AGENT_COLUMNS [AGENT] {
IDX => "idx",
AGENT_ID => "agent_id",
KIND => "kind",
STATUS => "status",
OUTCOME => "outcome",
TASK => "task",
}
}
fn agent_row_from(row: &Row) -> anyhow::Result<AgentRow> {
Ok(AgentRow {
idx: row.get(COL_AGENT_IDX)?,
agent_id: row.get(COL_AGENT_AGENT_ID)?,
kind: row.get(COL_AGENT_KIND)?,
status: row.get(COL_AGENT_STATUS)?,
outcome: row.get(COL_AGENT_OUTCOME)?,
task: row.get(COL_AGENT_TASK)?,
})
}
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: 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, .. } => {
assert!(
crate::pipeline::board::WORKING_PHASES.contains(phase),
"non-working phase as a job kind: {}",
phase.as_ref(),
);
<&'static str>::from(*phase)
}
}
}
}
#[must_use]
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,
caller_agent_id: Option<&str>,
) -> Result<()> {
let tx = conn.begin_tx().await?;
insert_job_tx(
&tx,
id,
task,
workspace_name,
user_name,
channel,
role,
agents,
child,
caller_agent_id,
)
.await?;
tx.commit().await?;
Ok(())
}
#[expect(clippy::too_many_arguments)]
async fn insert_job_tx(
tx: &TxGuard<'_>,
id: &str,
task: &str,
workspace_name: &str,
user_name: &str,
channel: &str,
role: Role,
agents: &[NewAgent],
child: &SpawnChild,
caller_agent_id: Option<&str>,
) -> Result<()> {
let kind = child.kind_str();
let ticket_id = child_ticket_id(child);
let mode = if caller_agent_id.is_some() {
JobMode::Sync
} else {
JobMode::Async
};
let now = db::now();
tx.execute(
"INSERT INTO jobs (id, kind, status, task, workspace_name, user_name, channel, role, \
ticket_id, retry_count, created_at, updated_at, caller_agent_id, mode) \
VALUES (?1, ?2, 'launched', ?3, ?4, ?5, ?6, ?7, ?8, 0, ?9, ?9, ?10, ?11)",
params![
id,
kind,
task,
workspace_name,
user_name,
channel,
role.as_str(),
ticket_id,
now.clone(),
caller_agent_id,
mode.as_str(),
],
)
.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}"))?;
}
}
Ok(())
}
pub(crate) struct NewAgent {
pub agent_id: String,
pub kind: AgentKind,
pub idx: Option<i64>,
pub task: String,
}
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()),
originating_workspace: None,
};
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(())
}
#[expect(clippy::cast_possible_truncation)] pub(crate) async fn abandon_session_jobs(caller_agent_id: &str) -> Result<usize> {
let conn = &crate::session::store().conn;
let deleted = conn
.execute(
"DELETE FROM jobs WHERE mode = 'sync' AND caller_agent_id = ?1",
params![caller_agent_id],
)
.await
.context("abandon session sync jobs")?;
info!(
caller_agent_id,
count = deleted,
"abandoned session sync jobs"
);
Ok(deleted as usize)
}
pub(crate) async fn transition_research_to_cleanup(
conn: &Connection,
run_id: &str,
cleanup_task: &str,
workspace_name: &str,
) -> Result<()> {
let tx = conn.begin_tx().await?;
tx.execute("DELETE FROM pending_jobs WHERE id = ?1", params![run_id])
.await
.with_context(|| format!("failed to delete pending envelope for run {run_id}"))?;
tx.execute("DELETE FROM jobs WHERE id = ?1", params![run_id])
.await
.with_context(|| format!("failed to delete research job row for run {run_id}"))?;
insert_job_tx(
&tx,
run_id,
cleanup_task,
workspace_name,
"",
"",
Role::Sanitation,
&[NewAgent {
agent_id: crate::research_cleanup::cleanup_agent_id(run_id),
kind: AgentKind::Sanitation,
idx: None,
task: cleanup_task.to_string(),
}],
&SpawnChild::ResearchCleanup,
None,
)
.await?;
tx.commit().await?;
Ok(())
}
pub(crate) async fn abandon_workspace_research_runs(workspace_name: &str) -> Result<usize> {
let conn = &crate::session::store().conn;
let research_rows = conn
.query(
"SELECT id FROM jobs WHERE workspace_name = ?1 AND kind IN ('research','research_cleanup')",
params![workspace_name],
)
.await
.context("list research jobs for workspace abandon")?;
let mut cancelled = 0usize;
for row in &research_rows {
let id: String = row.get(0)?;
crate::research_cancel::cancel_research_run(&id).await;
cancelled += 1;
}
Ok(cancelled)
}
pub(crate) async fn find_owned_launched_jobs(
conn: &Connection,
caller_agent_id: &str,
) -> Result<Vec<JobRow>> {
let rows = conn
.query(
&format!(
"SELECT {JOB_COLUMNS} FROM jobs j \
WHERE j.caller_agent_id = ?1 AND j.kind IN ('analyze','implement') AND j.status = 'launched' \
AND j.mode = 'sync' \
ORDER BY j.created_at DESC"
),
params![caller_agent_id],
)
.await
.context("find owned launched jobs")?;
rows.iter().map(job_row_from).collect()
}
#[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: 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.upsert_row(
"UPDATE agents SET status = ?1 WHERE job_id = ?2 AND agent_id = ?3",
|| params![status.as_str(), job_id, agent_id],
"INSERT INTO agents (job_id, agent_id, kind, idx, status, task) \
VALUES (?1, ?2, ?3, NULL, ?4, '') \
ON CONFLICT(job_id, agent_id) DO NOTHING",
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())
}
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) enum ResumePreamble {
Proceed(JobCaller, Role),
DrainCut,
Gone,
Unreadable(anyhow::Error),
}
pub(crate) enum SyncResumeOutcome {
Terminal(Role, JobCaller, anyhow::Result<String>),
DrainCut,
Gone,
}
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)> {
match resume_job_preamble_discrete(conn, job_id, abort_site, missing_site).await {
ResumePreamble::Proceed(caller, caller_role) => Some((caller, caller_role)),
ResumePreamble::DrainCut | ResumePreamble::Gone | ResumePreamble::Unreadable(_) => None,
}
}
pub(crate) async fn resume_job_preamble_discrete(
conn: &Connection,
job_id: &str,
abort_site: &str,
missing_site: &str,
) -> ResumePreamble {
if crate::shutdown::aborting() {
tracing::info!(job = %job_id, "{abort_site} aborted — drain/shutdown in progress");
return ResumePreamble::DrainCut;
}
match job_caller(conn, job_id).await {
Ok(Some(caller)) => {
let caller_role = std::str::FromStr::from_str(&caller.role).unwrap_or(Role::Manager);
ResumePreamble::Proceed(caller, caller_role)
}
Ok(None) => {
tracing::info!(job = %job_id, "{missing_site}: job row gone (explicitly abandoned) — skipping resume");
ResumePreamble::Gone
}
Err(e) => {
tracing::warn!(job = %job_id, error = %e, "{missing_site}: failed to read job row — skipping resume");
ResumePreamble::Unreadable(e)
}
}
}
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
}
}
}
async fn list_active_jobs(conn: &Connection) -> Result<Vec<JobRow>> {
let rows = conn
.query(
&format!(
"SELECT {JOB_COLUMNS} 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(
&format!("SELECT {AGENT_COLUMNS} 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;
}
let Some(workspace_name) =
crate::users::enforce_personal_pinning(job.role, &job.workspace_name, &job.user_name)
else {
error!(
pending_job = %row.id,
role = %job.role.as_str(),
workspace = %job.workspace_name,
"Pending job replay: pinned role with empty user — dropping poisoned envelope"
);
if let Err(e) = delete_pending_job(conn, &row.id).await {
warn!(pending_job = %row.id, error = %e, "Failed to delete refused pending job");
}
continue;
};
job.workspace_name = workspace_name;
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 phase_kinds = WORKING_PHASES
.iter()
.map(|phase| format!("'{}'", phase.as_ref()))
.collect::<Vec<_>>()
.join(", ");
let rows = conn
.query(
&format!(
"SELECT j.id, j.workspace_name \
FROM jobs j \
WHERE j.updated_at < ?1 \
AND j.kind IN ({phase_kinds}) \
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 workspace_name: String = row.get(1)?;
if 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)
}
#[must_use]
fn is_ticket_phase_kind(kind: &str) -> bool {
WORKING_PHASES.iter().any(|phase| phase.as_ref() == kind)
}
fn skip_sync_job(job: &JobRow) -> bool {
if job.mode == JobMode::Sync {
info!(
job = %job.id,
caller = ?job.caller_agent_id,
"sync job — caller resumes it; skipping boot resume",
);
true
} else {
false
}
}
async fn workspace_unresolvable(job: &JobRow) -> bool {
if crate::users::resolve_workspace(&job.workspace_name)
.await
.is_ok_and(|ws| ws.is_some())
{
return false;
}
warn!(
job = %job.id,
workspace = %job.workspace_name,
"workspace unresolvable — job stays launched, skipped (no boot resume, no checkpoint)",
);
true
}
#[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" {
if skip_sync_job(job) || workspace_unresolvable(job).await {
continue;
}
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" {
if skip_sync_job(job) || workspace_unresolvable(job).await {
continue;
}
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" {
if workspace_unresolvable(job).await {
continue;
}
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 cleaner row left over from a previous lifetime — terminalizing (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.upsert_row(
"UPDATE agents SET status = ?1, task = ?2, outcome = NULL \
WHERE agent_id = ?3 AND job_id IS NULL",
|| params![status.as_str(), task, pin_id.as_str()],
"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 NOTHING",
params![pin_id.as_str(), status.as_str(), task, role.as_str()],
)
.await
.with_context(|| format!("failed to upsert {role} session pin for ticket {ticket_id}"))?;
Ok(())
}
#[must_use]
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
}
#[cfg(test)]
mod tests {
use super::*;
async fn seed_ticket(conn: &crate::db::Connection, ticket_id: &str, ws_name: &str) {
let now = crate::db::now();
conn.execute(
"INSERT INTO tickets (id, title, description, workspace_name, created_at, updated_at) \
VALUES (?1, 'title', 'desc', ?2, ?3, ?3)",
crate::db::params![ticket_id, ws_name, now],
)
.await
.unwrap();
}
async fn status_retry(conn: &crate::db::Connection, job_id: &str) -> (String, i64) {
let rows = conn
.query(
"SELECT status, retry_count FROM jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert_eq!(rows.len(), 1, "job row exists for {job_id}");
(
rows[0].get::<String>(0).unwrap(),
rows[0].get::<i64>(1).unwrap(),
)
}
#[serial_test::serial(recover_from_restart)]
#[tokio::test]
async fn boot_recovery_skips_caller_owned_jobs() {
crate::util::test::init_management_test_stores().await;
let conn = &crate::session::store().conn;
let ws_name = "boot_recovery_ws";
crate::util::test::create_test_workspace("/tmp/boot_recovery_ws", ws_name).await;
for (job_id, kind, child, caller) in [
(
"owned_analyze",
AgentKind::Analyst,
SpawnChild::Analyze,
Some("pin_a"),
),
(
"owned_implement",
AgentKind::Coder,
SpawnChild::Implement,
Some("pin_b"),
),
(
"null_analyze",
AgentKind::Analyst,
SpawnChild::Analyze,
None,
),
] {
spawn_job(
conn,
job_id,
"task",
ws_name,
"caller",
"telegram",
crate::Role::Engineer,
&[NewAgent {
agent_id: format!("{job_id}_agent"),
kind,
idx: Some(0),
task: "task".to_string(),
}],
&child,
caller,
)
.await
.unwrap();
}
let modes: Vec<String> = conn
.query(
"SELECT mode FROM jobs WHERE id IN ('owned_analyze','owned_implement','null_analyze') ORDER BY id",
(),
)
.await
.unwrap()
.iter()
.map(|r| r.get::<String>(0).unwrap())
.collect();
assert_eq!(
modes,
["async".to_string(), "sync".to_string(), "sync".to_string()],
"spawn derives mode from the caller pin (sync ⇔ pin present)"
);
let resumable = recover_from_restart().await.unwrap();
assert!(
resumable.iter().any(
|r| matches!(r, ResumableJob::Analyze { job_id, .. } if job_id == "null_analyze")
),
"the async analyze is selected for boot resume"
);
assert!(
!resumable.iter().any(
|r| matches!(r, ResumableJob::Analyze { job_id, .. } if job_id == "owned_analyze")
),
"the sync caller-owned analyze is not boot-resumed"
);
assert!(
!resumable.iter().any(
|r| matches!(r, ResumableJob::Implement { job_id, .. } if job_id == "owned_implement")
),
"the sync caller-owned implement is not boot-resumed"
);
for job_id in ["owned_analyze", "owned_implement"] {
assert_eq!(
status_retry(conn, job_id).await,
("launched".to_string(), 0),
"{job_id} stays launched with retry_count unchanged"
);
}
assert_eq!(
status_retry(conn, "null_analyze").await,
("launched".to_string(), 1),
"the async analyze is check-pointed (retry_count bumped)"
);
}
#[serial_test::serial(recover_from_restart)] #[tokio::test]
async fn boot_scan_leaves_unresolvable_workspace_cleanup_row_in_place() {
crate::util::test::init_management_test_stores().await;
let conn = &crate::session::store().conn;
let job_id = "cleanup_missing_ws_1";
let ws_name = "ws_boot_scan_missing_1";
crate::util::test::JobRowBuilder::new(
conn,
job_id,
"research_cleanup",
"sanitation",
ws_name,
)
.task("cleanup prompt")
.timestamps(crate::db::now())
.insert()
.await
.unwrap();
let resumable = recover_from_restart().await.unwrap();
let hit = resumable
.iter()
.any(|r| matches!(r, ResumableJob::ResearchCleanup { job_id: j, .. } if j == job_id));
assert!(
!hit,
"unresolvable-workspace cleanup row is not resumed (binding-3 skip)"
);
assert_eq!(
status_retry(conn, job_id).await,
("launched".to_string(), 0),
"binding-3 skip leaves the cleanup row untouched (no checkpoint bump)"
);
}
#[tokio::test]
async fn clear_abandons_session_sync_jobs() {
crate::util::test::init_management_test_stores().await;
let conn = &crate::session::store().conn;
let ws_name = "clear_abandon_ws";
seed_ticket(conn, "t_clear", ws_name).await;
for (job_id, child, caller) in [
("clear_sync", SpawnChild::Analyze, Some("pin_clear")),
("clear_async", SpawnChild::Analyze, None),
(
"clear_phase",
SpawnChild::Phase {
phase: crate::pipeline::board::TicketPhase::Analysis,
ticket_id: "t_clear".to_string(),
},
None,
),
] {
spawn_job(
conn,
job_id,
"task",
ws_name,
"caller",
"telegram",
crate::Role::Engineer,
&[NewAgent {
agent_id: format!("{job_id}_agent"),
kind: AgentKind::Analyst,
idx: Some(0),
task: "task".to_string(),
}],
&child,
caller,
)
.await
.unwrap();
}
let abandoned = abandon_session_jobs("pin_clear").await.unwrap();
assert_eq!(abandoned, 1, "only the caller-owned sync job is abandoned");
let remaining: Vec<String> = conn
.query(
"SELECT id FROM jobs WHERE id IN ('clear_sync','clear_async','clear_phase') ORDER BY id",
(),
)
.await
.unwrap()
.iter()
.map(|r| r.get::<String>(0).unwrap())
.collect();
assert_eq!(
remaining,
["clear_async", "clear_phase"],
"async + phase jobs survive the session abandon"
);
let roster: Vec<String> = conn
.query(
"SELECT agent_id FROM agents WHERE job_id = 'clear_sync'",
(),
)
.await
.unwrap()
.iter()
.map(|r| r.get::<String>(0).unwrap())
.collect();
assert!(roster.is_empty(), "sync job roster cascaded away");
}
#[tokio::test]
async fn workspace_delete_abandons_all_jobs() {
crate::util::test::init_management_test_stores().await;
let conn = &crate::session::store().conn;
let ws_name = "ws_abandon_test";
seed_ticket(conn, "t_abandon", ws_name).await;
for (job_id, child, caller) in [
("abandon_sync", SpawnChild::Analyze, Some("pin_abandon")),
(
"abandon_phase",
SpawnChild::Phase {
phase: crate::pipeline::board::TicketPhase::Analysis,
ticket_id: "t_abandon".to_string(),
},
None,
),
("abandon_research", SpawnChild::Research, None),
("abandon_cleanup", SpawnChild::ResearchCleanup, None),
] {
spawn_job(
conn,
job_id,
"task",
ws_name,
"caller",
"telegram",
crate::Role::Engineer,
&[NewAgent {
agent_id: format!("{job_id}_agent"),
kind: AgentKind::Analyst,
idx: Some(0),
task: "task".to_string(),
}],
&child,
caller,
)
.await
.unwrap();
}
let ws_store = crate::workspace::WorkspaceStore { conn: conn.clone() };
ws_store.delete(ws_name).await.unwrap();
let jobs = conn
.query(
"SELECT id FROM jobs WHERE workspace_name = ?1",
crate::db::params![ws_name],
)
.await
.unwrap();
assert!(jobs.is_empty(), "every job row for the workspace is gone");
let research = conn
.query(
"SELECT id FROM research_jobs WHERE id IN ('abandon_research','abandon_cleanup')",
(),
)
.await
.unwrap();
assert!(
research.is_empty(),
"research_jobs child rows cascaded away"
);
let roster = conn
.query(
"SELECT agent_id FROM agents WHERE job_id IN ('abandon_sync','abandon_phase','abandon_research','abandon_cleanup')",
(),
)
.await
.unwrap();
assert!(roster.is_empty(), "every seeded agent roster cascaded away");
}
}