use crate::db::{self, IntoParams, TxGuard, Value};
use anyhow::{Context, Result};
use chrono::{Duration, Utc};
use serde::Serialize;
use std::collections::HashMap;
use std::fmt::Write;
use std::sync::LazyLock;
use tracing::{debug, info, warn};
crate::define_store! {
pub static BOARD: BoardStore,
post_open = after_open,
expect = "BOARD not initialized — call init_all_stores() first",
}
pub async fn run_archive_cancelled_loop() {
let interval = std::time::Duration::from_mins(5);
loop {
if !crate::shutdown::sleep_or_shutdown_or_drain(interval).await {
break;
}
crate::jobs::purge_terminal_session_pins().await;
let Some(board) = BOARD.get() else {
warn!("Archive cancelled loop: board not initialized");
continue;
};
match board.archive_stale_cancelled(CANCELLED_ARCHIVE_HOURS).await {
Ok(n) if n > 0 => info!(count = n, "Archived stale cancelled tickets"),
Ok(_) => debug!("Archive cancelled loop: no stale tickets"),
Err(e) => warn!(error = %e, "Archive cancelled loop failed"),
}
}
}
const CANCELLED_ARCHIVE_HOURS: i64 = 1;
crate::columns! {
TICKET_COLUMNS [TICKET] {
ID => "id",
TITLE => "title",
DESCRIPTION => "description",
PHASE => "phase",
WORKSPACE_NAME => "workspace_name",
CREATED_AT => "created_at",
UPDATED_AT => "updated_at",
PREREQUISITES => "prerequisites",
SUPERSEDES => "supersedes",
SUPERSEDED_BY => "superseded_by",
COMMIT_HASH => "commit_hash",
LINES_ADDED => "lines_added",
LINES_REMOVED => "lines_removed",
REPORTER => "reporter",
IS_ARCHIVED => "is_archived",
PRIORITY => "priority",
REVIEWED_HEAD => "reviewed_head",
REVIEWED_TREE => "reviewed_tree",
DONE_AT => "done_at",
BOUNCE_COUNT => "bounce_count",
}
}
crate::columns! {
COMMENT_COLUMNS [COMMENT] {
ROLE => "role",
CONTENT => "content",
CREATED_AT => "created_at",
}
}
const PIPELINE_OCCUPIED_PHASES: &[TicketPhase] = &[
TicketPhase::InDevelopment,
TicketPhase::InDiagnostics,
TicketPhase::InReview,
TicketPhase::InQa,
TicketPhase::InSanitation,
];
pub(crate) const WORKING_PHASES: &[TicketPhase] = &[
TicketPhase::Analysis,
TicketPhase::InDevelopment,
TicketPhase::InDiagnostics,
TicketPhase::InReview,
TicketPhase::InQa,
TicketPhase::InSanitation,
];
pub const UNBLOCKING_PHASES: &[TicketPhase] = &[TicketPhase::Done, TicketPhase::Cancelled];
pub const TERMINAL_PHASES: &[TicketPhase] = &[
TicketPhase::Done,
TicketPhase::Cancelled,
TicketPhase::Failed,
];
fn phase_list_sql_fragment(phases: &[TicketPhase]) -> String {
phases
.iter()
.map(|p| format!("'{}'", p.as_ref()))
.collect::<Vec<_>>()
.join(", ")
}
fn claim_candidate_where(base: usize, source: TicketPhase) -> String {
let mut parts = vec![
format!("t1.phase = ?{base}"),
format!("t1.workspace_name = ?{}", base + 1),
"t1.is_archived = 0".to_string(),
];
if source == TicketPhase::Backlog {
parts.push(format!("t1.created_at <= ?{}", base + 2));
}
if source == TicketPhase::Queued {
parts.push(format!(
"NOT EXISTS (SELECT 1 FROM tickets t2 \
WHERE t2.workspace_name = t1.workspace_name \
AND t2.phase IN ({}) \
AND t2.is_archived = 0 \
AND t2.id != t1.id)",
phase_list_sql_fragment(PIPELINE_OCCUPIED_PHASES),
));
}
parts.push(format!(
"NOT EXISTS ( \
SELECT 1 FROM json_each(t1.prerequisites) AS je \
JOIN tickets t_pre ON t_pre.id = je.value \
WHERE t_pre.phase NOT IN ({}) \
)",
phase_list_sql_fragment(UNBLOCKING_PHASES),
));
parts.join(" AND ")
}
fn parse_prereqs(raw: &str) -> Result<Vec<String>> {
serde_json::from_str(raw).with_context(|| {
let preview = if raw.len() > 200 {
format!("{}…", crate::util::truncate_bytes(raw, 200))
} else {
raw.to_string()
};
format!("Corrupt prerequisites JSON in database: {preview}")
})
}
#[derive(Debug, Clone)]
pub(crate) struct TicketParams {
pub title: String,
pub description: String,
pub workspace_name: String,
pub phase: TicketPhase,
pub prerequisites: Vec<String>,
pub reporter: String,
pub embedding: Option<Vec<u8>>,
pub priority: i64,
}
#[derive(Debug, Clone, PartialEq, Serialize)]
pub struct TicketComment {
pub role: String,
pub content: String,
pub created_at: String,
}
#[derive(Debug, Clone, PartialEq)]
pub struct Ticket {
pub id: String,
pub title: String,
pub description: String,
pub phase: TicketPhase,
pub workspace_name: String,
pub created_at: String,
pub updated_at: String,
pub comments: Vec<TicketComment>,
pub prerequisites: Vec<String>,
pub supersedes: Option<String>,
pub superseded_by: Option<String>,
pub commit_hash: Option<String>,
pub lines_added: Option<i64>,
pub lines_removed: Option<i64>,
pub reporter: String,
pub is_archived: bool,
pub priority: i64,
pub reviewed_head: Option<String>,
pub reviewed_tree: Option<String>,
pub done_at: Option<String>,
pub bounce_count: i64,
}
impl Ticket {
#[must_use]
pub fn short_display(&self) -> String {
format!(
" [{}] [{}] {}: {}",
self.reporter, self.phase, self.id, self.title
)
}
#[must_use]
pub fn detailed_display(&self) -> String {
let mut out = self.header(&self.description);
out.push_str(&self.format_comments());
out
}
#[must_use]
pub(crate) fn detailed_display_limited(
&self,
desc_max_chars: usize,
comment_max_chars: usize,
last_n_full: usize,
) -> String {
let description = crate::util::truncate(&self.description, desc_max_chars);
let mut out = self.header(&description);
out.push_str(&self.format_comments_limited(comment_max_chars, last_n_full));
out
}
#[must_use]
fn header(&self, description: &str) -> String {
let mut out = format!(
"Ticket: {id}\n\
Title: {title}\n\
Description: {description}\n\
Phase: {phase}\n\
Reporter: {reporter}\n\
Workspace: {workspace}\n\
Created: {created}\n\
Updated: {updated}\n\
Priority: P{priority}\n",
id = self.id,
title = self.title,
description = description,
phase = self.phase,
reporter = self.reporter,
workspace = self.workspace_name,
created = self.created_at,
updated = self.updated_at,
priority = self.priority,
);
if let Some(ref s) = self.supersedes {
let _ = writeln!(out, "Supersedes: {s}");
}
if let Some(ref s) = self.superseded_by {
let _ = writeln!(out, "Superseded by: {s}");
}
if !self.prerequisites.is_empty() {
let _ = writeln!(out, "Prerequisites: {}", self.prerequisites.join(", "));
}
if self.is_archived {
out.push_str("Archived: yes\n");
}
out
}
#[must_use]
fn format_comments(&self) -> String {
let mut s = String::from("Comments:");
if self.comments.is_empty() {
s.push_str("\n (no comments)");
} else {
for c in &self.comments {
s.push_str(&Self::comment_line(c, &c.content));
}
}
s
}
#[must_use]
fn format_comments_limited(&self, max_chars: usize, last_n_full: usize) -> String {
let mut s = String::from("Comments:");
if self.comments.is_empty() {
s.push_str("\n (no comments)");
} else {
let n = self.comments.len();
for (i, c) in self.comments.iter().enumerate() {
let content = if i + last_n_full < n {
crate::util::truncate(&c.content, max_chars)
} else {
c.content.clone()
};
s.push_str(&Self::comment_line(c, &content));
}
}
s
}
#[must_use]
fn comment_line(c: &TicketComment, content: &str) -> String {
let end = 19.min(c.created_at.len());
let ts = &c.created_at[..end];
format!("\n [{}] ({}): {}", c.role, ts, content)
}
}
#[derive(
Debug, Clone, Copy, PartialEq, Serialize, strum::Display, strum::AsRefStr, strum::EnumIter,
)]
#[serde(rename_all = "snake_case")]
#[strum(serialize_all = "snake_case")]
pub enum TicketPhase {
Backlog,
Analysis,
Planning,
Queued,
InDevelopment,
InDiagnostics,
InSanitation,
InReview,
InQa,
Done,
Cancelled,
Failed,
}
impl TicketPhase {
#[must_use]
pub fn is_unblocking(&self) -> bool {
UNBLOCKING_PHASES.contains(self)
}
#[must_use]
pub fn is_terminal(&self) -> bool {
TERMINAL_PHASES.contains(self)
}
#[must_use]
pub fn is_pipeline_occupied(&self) -> bool {
PIPELINE_OCCUPIED_PHASES.contains(self)
}
#[must_use]
pub fn display_name(&self) -> String {
self.as_ref().replace('_', " ")
}
}
static ALL_TICKET_PHASE_NAMES: LazyLock<String> = LazyLock::new(|| {
<TicketPhase as strum::IntoEnumIterator>::iter()
.map(|p| p.to_string())
.collect::<Vec<_>>()
.join(", ")
});
impl std::str::FromStr for TicketPhase {
type Err = anyhow::Error;
fn from_str(s: &str) -> Result<Self, Self::Err> {
<TicketPhase as strum::IntoEnumIterator>::iter()
.find(|p| p.as_ref() == s)
.ok_or_else(|| {
anyhow::anyhow!(
"Invalid phase '{s}'. Valid phases: {}",
*ALL_TICKET_PHASE_NAMES
)
})
}
}
struct PreparedUpdate {
sql: String,
params: Vec<db::Value>,
ticket_id: String,
}
impl PreparedUpdate {
async fn execute_tx(self, tx: &db::TxGuard<'_>) -> Result<()> {
let rows = tx.execute(&self.sql, self.params).await?;
BoardStore::ensure_ticket_found(rows, &self.ticket_id)?;
Ok(())
}
async fn execute_no_cancel(self, conn: &db::Connection) -> Result<()> {
let rows = conn.execute(&self.sql, self.params).await?;
BoardStore::ensure_ticket_found(rows, &self.ticket_id)?;
Ok(())
}
async fn execute_and_cancel(self, conn: &db::Connection) -> Result<()> {
let rows = conn.execute(&self.sql, self.params).await?;
BoardStore::ensure_ticket_found(rows, &self.ticket_id)?;
crate::agent::registry::AGENT_REGISTRY.cancel_by_ticket_id(&self.ticket_id);
Ok(())
}
async fn execute_tx_matched(self, tx: &db::TxGuard<'_>) -> Result<bool> {
let rows = tx.execute(&self.sql, self.params).await?;
Ok(rows > 0)
}
}
#[derive(Copy, Clone, Debug, PartialEq, Eq)]
pub(crate) enum LoadComments {
Yes,
No,
}
impl BoardStore {
pub(crate) async fn after_open(&self) -> anyhow::Result<()> {
crate::db::ensure_fts_index(
&self.conn,
crate::db::TICKETS_FTS_INDEX_NAME,
"ngram",
crate::db::TICKETS_FTS_INDEX_DDL,
)
.await?;
crate::db::repair_ticket_title_fts_if_corrupt(
&self.conn,
crate::db::TICKETS_FTS_INDEX_NAME,
crate::db::TICKETS_FTS_INDEX_DDL,
)
.await;
Ok(())
}
async fn insert_ticket_tx(
tx: &TxGuard<'_>,
ticket_id: &str,
params: &TicketParams,
supersedes: Option<&str>,
) -> Result<()> {
let now = db::now();
let prereqs_json = serde_json::to_string(¶ms.prerequisites)?;
tx.execute(
"INSERT INTO tickets (id, title, description, phase, workspace_name, \
created_at, updated_at, prerequisites, supersedes, reporter, embedding, priority) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
db::params![
ticket_id,
params.title.as_str(),
params.description.as_str(),
params.phase.as_ref(),
params.workspace_name.as_str(),
now.as_str(),
now.as_str(),
prereqs_json.as_str(),
supersedes,
params.reporter.as_str(),
params.embedding.as_deref(),
params.priority,
],
)
.await?;
Ok(())
}
async fn rewire_dependents_tx(
tx: &TxGuard<'_>,
supersede_id: &str,
new_id: &str,
workspace_name: &str,
) -> Result<()> {
let dep_rows = tx
.query(
"SELECT DISTINCT t.id, t.prerequisites \
FROM tickets t, json_each(t.prerequisites) AS je \
WHERE je.value = ?1 AND t.workspace_name = ?2",
db::params![supersede_id, workspace_name],
)
.await?;
for row in &dep_rows {
let dep_id: String = row.get(0)?;
let raw: String = row.get(1)?;
let mut prereqs: Vec<String> = parse_prereqs(&raw)
.with_context(|| format!("Failed to parse prerequisites for ticket {dep_id}"))?;
let mut changed = false;
for p in &mut prereqs {
if *p == supersede_id {
*p = new_id.to_string();
changed = true;
}
}
if changed {
let new_json = serde_json::to_string(&prereqs)?;
tx.execute(
"UPDATE tickets SET prerequisites = ?1, updated_at = ?2 WHERE id = ?3",
db::params![new_json, db::now(), dep_id],
)
.await?;
}
}
Ok(())
}
async fn begin_tx_and_validate_prerequisites(
&self,
workspace_name: &str,
prerequisites: &[String],
) -> Result<(TxGuard<'_>, String)> {
let tx = self.conn.begin_tx().await?;
let seq: i64 = tx
.query_row(
"INSERT INTO ticket_counters (workspace_name, next_id) VALUES (?1, 1) \
ON CONFLICT(workspace_name) DO UPDATE SET next_id = ticket_counters.next_id + 1 \
RETURNING next_id - 1",
db::params![workspace_name],
|row| row.get(0),
)
.await?;
let id = format!("{workspace_name}-{seq}");
anyhow::ensure!(
!prerequisites.contains(&id),
"Ticket cannot depend on itself: {id}"
);
Self::validate_prerequisites(&tx, prerequisites, workspace_name).await?;
Ok((tx, id))
}
pub(crate) async fn create_ticket(&self, params: &TicketParams) -> Result<String> {
let (tx, id) = self
.begin_tx_and_validate_prerequisites(¶ms.workspace_name, ¶ms.prerequisites)
.await?;
Self::insert_ticket_tx(&tx, &id, params, None).await?;
tx.commit().await?;
Ok(id)
}
pub(crate) async fn supersede_and_create(
&self,
supersede_id: &str,
params: &TicketParams,
) -> Result<String> {
anyhow::ensure!(
!params.prerequisites.iter().any(|p| p == supersede_id),
"Ticket cannot supersede and depend on the same ticket: {supersede_id}"
);
let (tx, new_id) = self
.begin_tx_and_validate_prerequisites(¶ms.workspace_name, ¶ms.prerequisites)
.await?;
let rows = tx
.query(
"SELECT workspace_name FROM tickets WHERE id = ?1",
db::params![supersede_id],
)
.await?;
let row = rows
.into_iter()
.next()
.ok_or_else(|| anyhow::anyhow!("Superseded ticket not found: {supersede_id}"))?;
let old_ws: String = row.get(0)?;
anyhow::ensure!(
old_ws == params.workspace_name,
"Superseded ticket {supersede_id} belongs to workspace '{old_ws}', \
not the current workspace '{}'. \
Cross-workspace supersede is not allowed.",
params.workspace_name,
);
let now = db::now();
let cancelled_rows = tx
.execute(
"UPDATE tickets SET phase = ?1, updated_at = ?2, \
superseded_by = ?4, is_archived = 1, done_at = NULL \
WHERE id = ?3",
db::params![
TicketPhase::Cancelled.as_ref(),
now,
supersede_id,
new_id.as_str(),
],
)
.await?;
Self::ensure_ticket_found(cancelled_rows, supersede_id)?;
Self::insert_ticket_tx(&tx, &new_id, params, Some(supersede_id)).await?;
Self::rewire_dependents_tx(&tx, supersede_id, &new_id, ¶ms.workspace_name).await?;
crate::agent::registry::AGENT_REGISTRY.cancel_by_ticket_id(supersede_id);
tx.commit().await?;
Ok(new_id)
}
async fn ticket_from_row(&self, row: &db::Row, load_comments: LoadComments) -> Result<Ticket> {
let id: String = row.get(COL_TICKET_ID)?;
let comments = if load_comments == LoadComments::Yes {
self.get_comments(&id).await?
} else {
Vec::new()
};
let prerequisites_raw: String = row.get(COL_TICKET_PREREQUISITES)?;
let prerequisites = parse_prereqs(&prerequisites_raw)
.with_context(|| format!("Failed to parse prerequisites for ticket {id}"))?;
Ok(Ticket {
id,
title: row.get(COL_TICKET_TITLE)?,
description: row.get(COL_TICKET_DESCRIPTION)?,
phase: row
.get::<String>(COL_TICKET_PHASE)?
.parse::<TicketPhase>()?,
workspace_name: row.get(COL_TICKET_WORKSPACE_NAME)?,
created_at: row.get(COL_TICKET_CREATED_AT)?,
updated_at: row.get(COL_TICKET_UPDATED_AT)?,
comments,
prerequisites,
supersedes: row.get(COL_TICKET_SUPERSEDES)?,
superseded_by: row.get(COL_TICKET_SUPERSEDED_BY)?,
commit_hash: row.get(COL_TICKET_COMMIT_HASH)?,
lines_added: row.get(COL_TICKET_LINES_ADDED)?,
lines_removed: row.get(COL_TICKET_LINES_REMOVED)?,
reporter: row.get::<String>(COL_TICKET_REPORTER)?,
is_archived: row.get::<bool>(COL_TICKET_IS_ARCHIVED)?,
priority: row.get::<i64>(COL_TICKET_PRIORITY)?,
reviewed_head: row.get(COL_TICKET_REVIEWED_HEAD)?,
reviewed_tree: row.get(COL_TICKET_REVIEWED_TREE)?,
done_at: row.get(COL_TICKET_DONE_AT)?,
bounce_count: row.get(COL_TICKET_BOUNCE_COUNT)?,
})
}
pub(crate) const BACKLOG_CLAIM_GRACE: Duration = Duration::seconds(5);
pub(crate) async fn claim_backlog_for_analysis(
&self,
workspace_name: &str,
now: String,
grace_cutoff: String,
) -> Result<Option<Ticket>> {
self.claim_for_phase(
workspace_name,
now,
TicketPhase::Backlog,
TicketPhase::Analysis,
Some(&grace_cutoff),
)
.await
}
pub(crate) async fn claim_queued_for_development(
&self,
workspace_name: &str,
now: String,
) -> Result<Option<Ticket>> {
self.claim_for_phase(
workspace_name,
now,
TicketPhase::Queued,
TicketPhase::InDevelopment,
None,
)
.await
}
async fn claim_for_phase(
&self,
workspace_name: &str,
now: String,
source: TicketPhase,
target: TicketPhase,
grace_cutoff: Option<&str>,
) -> Result<Option<Ticket>> {
let sql = format!(
"UPDATE tickets SET phase = ?1, updated_at = ?2 \
WHERE id = (SELECT t1.id FROM tickets t1 \
WHERE {} \
ORDER BY t1.priority ASC, t1.created_at ASC LIMIT 1) \
RETURNING {TICKET_COLUMNS}",
claim_candidate_where(3, source),
);
let rows = match grace_cutoff {
Some(cutoff) => {
self.conn
.query_cached(
&sql,
db::params![
target.as_ref(),
now,
source.as_ref(),
workspace_name,
cutoff
],
)
.await?
}
None => {
self.conn
.query_cached(
&sql,
db::params![target.as_ref(), now, source.as_ref(), workspace_name],
)
.await?
}
};
match rows.into_iter().next() {
Some(row) => Ok(Some(self.ticket_from_row(&row, LoadComments::Yes).await?)),
None => Ok(None),
}
}
pub(crate) async fn backlog_claim_candidate_exists(
&self,
workspace_name: &str,
grace_cutoff: &str,
) -> Result<bool> {
self.claim_candidate_exists(TicketPhase::Backlog, workspace_name, Some(grace_cutoff))
.await
}
pub(crate) async fn queued_claim_candidate_exists(&self, workspace_name: &str) -> Result<bool> {
self.claim_candidate_exists(TicketPhase::Queued, workspace_name, None)
.await
}
async fn claim_candidate_exists(
&self,
source: TicketPhase,
workspace_name: &str,
grace_cutoff: Option<&str>,
) -> Result<bool> {
let sql = format!(
"SELECT 1 FROM tickets t1 WHERE {} LIMIT 1",
claim_candidate_where(1, source),
);
let rows = match grace_cutoff {
Some(cutoff) => {
self.conn
.query_cached(&sql, db::params![source.as_ref(), workspace_name, cutoff])
.await?
}
None => {
self.conn
.query_cached(&sql, db::params![source.as_ref(), workspace_name])
.await?
}
};
Ok(!rows.is_empty())
}
pub(crate) async fn select_tickets(
&self,
suffix: &str,
params: impl IntoParams + Send + 'static,
load_comments: LoadComments,
) -> Result<Vec<Ticket>> {
let sql = format!("SELECT {TICKET_COLUMNS} FROM tickets {suffix}");
let rows = self.conn.query(&sql, params).await?;
let mut tickets = Vec::with_capacity(rows.len());
for row in rows {
tickets.push(self.ticket_from_row(&row, load_comments).await?);
}
Ok(tickets)
}
pub async fn get_ticket(&self, ticket_id: &str) -> Result<Option<Ticket>> {
Ok(self
.select_tickets("WHERE id = ?1", db::params![ticket_id], LoadComments::Yes)
.await?
.into_iter()
.next())
}
pub(crate) async fn get_tickets_by_ids(
&self,
ids: &[String],
load_comments: LoadComments,
) -> Result<Vec<Ticket>> {
if ids.is_empty() {
return Ok(Vec::new());
}
let (suffix, params) = Self::in_clause_for_ids(ids);
self.select_tickets(&suffix, params, load_comments).await
}
fn in_clause_for_ids(ids: &[String]) -> (String, Vec<Value>) {
let suffix = format!("WHERE id IN ({})", db::sql_in_placeholders(ids.len()));
let params: Vec<Value> = ids.iter().map(|id| Value::Text(id.clone())).collect();
(suffix, params)
}
pub async fn get_ticket_phase(&self, ticket_id: &str) -> Result<Option<TicketPhase>> {
self.conn
.query_optional(
"SELECT phase FROM tickets WHERE id = ?1",
db::params![ticket_id],
|row| {
let phase: String = row.get(0)?;
phase.parse()
},
)
.await
}
pub(crate) async fn get_ticket_priority(&self, ticket_id: &str) -> Result<Option<i64>> {
let sql = "SELECT priority FROM tickets WHERE id = ?1";
self.conn
.query_optional(sql, db::params![ticket_id], |row| row.get::<i64>(0))
.await
}
fn build_ticket_update_with_updated_at(
set_clause: &str,
set_params: Vec<db::Value>,
ticket_id: &str,
) -> PreparedUpdate {
let now = db::now();
let sql = format!("UPDATE tickets SET {set_clause}, updated_at = ? WHERE id = ?");
let mut params = set_params;
params.push(Value::from(now));
params.push(Value::from(ticket_id));
PreparedUpdate {
sql,
params,
ticket_id: ticket_id.to_string(),
}
}
fn build_transition_sql(
ticket_id: &str,
expected_phase: Option<TicketPhase>,
target_phase: TicketPhase,
) -> PreparedUpdate {
let now = db::now();
let guard: Option<&str> = expected_phase.as_ref().map(TicketPhase::as_ref);
let sql = "UPDATE tickets SET phase = ?1, updated_at = ?2, \
done_at = CASE WHEN ?1 = 'done' THEN ?2 \
WHEN phase = 'done' THEN NULL \
ELSE done_at END \
WHERE id = ?3 AND (?4 IS NULL OR phase = ?4)";
let params: Vec<db::Value> = vec![
Value::from(target_phase.as_ref()),
Value::from(now),
Value::from(ticket_id),
Value::from(guard),
];
PreparedUpdate {
sql: sql.to_string(),
params,
ticket_id: ticket_id.to_string(),
}
}
pub async fn transition_to(
&self,
ticket_id: &str,
expected_phase: Option<TicketPhase>,
target_phase: TicketPhase,
) -> Result<()> {
let prepared = Self::build_transition_sql(ticket_id, expected_phase, target_phase);
prepared.execute_and_cancel(&self.conn).await?;
if target_phase.is_terminal()
&& let Some(session) = crate::session::SESSIONS.get()
{
let _ = crate::jobs::complete_ticket_phase_jobs(&session.conn, ticket_id).await;
}
Ok(())
}
pub(crate) async fn transition_to_tx(
tx: &TxGuard<'_>,
ticket_id: &str,
expected_phase: Option<TicketPhase>,
target_phase: TicketPhase,
) -> Result<bool> {
let prepared = Self::build_transition_sql(ticket_id, expected_phase, target_phase);
prepared.execute_tx_matched(tx).await
}
fn ensure_ticket_found(rows: u64, ticket_id: &str) -> Result<()> {
anyhow::ensure!(rows > 0, "Ticket {ticket_id} not found");
Ok(())
}
pub(crate) async fn set_commit_info_tx(
tx: &TxGuard<'_>,
ticket_id: &str,
hash: &str,
lines_added: i64,
lines_removed: i64,
) -> Result<()> {
debug_assert!(
lines_added >= 0,
"lines_added must be non-negative: {lines_added}"
);
debug_assert!(
lines_removed >= 0,
"lines_removed must be non-negative: {lines_removed}"
);
let prepared = Self::build_ticket_update_with_updated_at(
"commit_hash = ?, lines_added = ?, lines_removed = ?",
vec![
Value::from(hash),
Value::from(lines_added),
Value::from(lines_removed),
],
ticket_id,
);
prepared.execute_tx(tx).await
}
pub(crate) async fn set_reviewed_base(
&self,
ticket_id: &str,
head: Option<&str>,
tree: Option<&str>,
) -> Result<()> {
let prepared = Self::build_ticket_update_with_updated_at(
"reviewed_head = ?, reviewed_tree = ?",
vec![Value::from(head), Value::from(tree)],
ticket_id,
);
prepared.execute_no_cancel(&self.conn).await
}
pub(crate) async fn increment_bounce_count_tx(
tx: &TxGuard<'_>,
ticket_id: &str,
) -> Result<i64> {
let rows = tx
.query(
"UPDATE tickets SET bounce_count = bounce_count + 1, updated_at = ?1 \
WHERE id = ?2 RETURNING bounce_count",
db::params![db::now(), ticket_id],
)
.await
.map_err(anyhow::Error::from)?;
match rows.into_iter().next() {
Some(row) => row.get::<i64>(0).map_err(anyhow::Error::from),
None => Err(anyhow::anyhow!(
"ticket {ticket_id} not found — bounce counter not incremented"
)),
}
}
async fn has_active_tickets_internal(
&self,
workspace_name: &str,
exclude_ticket_id: Option<&str>,
) -> Result<bool> {
let occupied_sql = phase_list_sql_fragment(PIPELINE_OCCUPIED_PHASES);
let sql = format!(
"SELECT 1 FROM tickets WHERE \
(phase IN ({occupied_sql}) OR phase = ?3) \
AND workspace_name = ?1 AND is_archived = 0 \
AND (?2 IS NULL OR id != ?2) LIMIT 1",
);
let queued = TicketPhase::Queued.as_ref();
let rows = self
.conn
.query(&sql, db::params![workspace_name, exclude_ticket_id, queued])
.await?;
Ok(!rows.is_empty())
}
pub async fn has_active_tickets_excluding(
&self,
workspace_name: &str,
exclude_ticket_id: &str,
) -> Result<bool> {
self.has_active_tickets_internal(workspace_name, Some(exclude_ticket_id))
.await
}
pub async fn add_comment(&self, ticket_id: &str, role: &str, content: &str) -> Result<()> {
crate::db::with_tx(&self.conn, ticket_id, "add comment", async |tx| {
Self::add_comment_tx(tx, ticket_id, role, content).await
})
.await?;
self.route_comment_to_agents(ticket_id, role, content).await;
Ok(())
}
async fn route_comment_to_agents(&self, ticket_id: &str, role: &str, content: &str) {
let Some(sessions) = crate::session::SESSIONS.get() else {
return;
};
let ticket = match self.get_ticket(ticket_id).await {
Ok(Some(t)) => t,
Ok(None) => {
warn!(ticket = %ticket_id, "Comment routing: ticket not found");
return;
}
Err(e) => {
warn!(ticket = %ticket_id, error = %e, "Comment routing: failed to fetch ticket");
return;
}
};
let active =
match crate::jobs::list_running_agents_for_ticket(&sessions.conn, ticket_id).await {
Ok(agents) if !agents.is_empty() => agents,
Ok(_) => return, Err(e) => {
warn!(
ticket = %ticket_id,
error = %e,
"Comment routing: failed to read running agents",
);
return;
}
};
for agent_id in active {
if agent_id.is_empty() {
continue;
}
let commenter_role = role.parse::<crate::Role>().unwrap_or(crate::Role::Manager);
let job = crate::agent::message_router::AgentJob {
content: content.to_string(),
workspace_name: ticket.workspace_name.clone(),
user_name: role.to_string(),
channel: String::new(),
kind: crate::agent::message_router::MessageKind::TicketComment,
role: commenter_role,
reply_target: None,
pending_job_id: None,
};
if crate::agent::message_router::try_route(&agent_id, job) {
debug!(
ticket = %ticket_id,
agent = %agent_id,
"Routed comment to running agent",
);
}
}
}
pub(crate) async fn add_comment_tx(
tx: &TxGuard<'_>,
ticket_id: &str,
role: &str,
content: &str,
) -> Result<()> {
let comment_id = crate::generate_id();
let now = db::now();
tx.execute(
"INSERT INTO ticket_comments (id, ticket_id, role, content, created_at) \
VALUES (?1, ?2, ?3, ?4, ?5)",
db::params![comment_id, ticket_id, role, content, now.as_str()],
)
.await?;
tx.execute(
"UPDATE tickets SET updated_at = ?1 WHERE id = ?2",
db::params![now.as_str(), ticket_id],
)
.await?;
Ok(())
}
pub async fn get_comments(&self, ticket_id: &str) -> Result<Vec<TicketComment>> {
let sql = format!(
"SELECT {COMMENT_COLUMNS} FROM ticket_comments WHERE ticket_id = ?1 ORDER BY created_at ASC"
);
let rows = self.conn.query(&sql, db::params![ticket_id]).await?;
let mut comments = Vec::new();
for row in rows {
comments.push(TicketComment {
role: row.get(COL_COMMENT_ROLE)?,
content: row.get(COL_COMMENT_CONTENT)?,
created_at: row.get(COL_COMMENT_CREATED_AT)?,
});
}
Ok(comments)
}
async fn validate_prerequisites(
tx: &TxGuard<'_>,
prerequisite_ids: &[String],
workspace_name: &str,
) -> Result<()> {
if prerequisite_ids.is_empty() {
return Ok(());
}
let (suffix, params) = Self::in_clause_for_ids(prerequisite_ids);
let sql = format!("SELECT id, workspace_name FROM tickets {suffix}");
let rows = tx.query(&sql, params).await?;
let mut found: HashMap<String, String> = HashMap::new();
for row in rows {
let id: String = row.get(0)?;
let ws_name: String = row.get(1)?;
found.insert(id, ws_name);
}
for pid in prerequisite_ids {
let ws_name = found
.get(pid)
.ok_or_else(|| anyhow::anyhow!("Prerequisite ticket not found: {pid}"))?;
anyhow::ensure!(
ws_name == workspace_name,
"Prerequisite {pid} belongs to workspace '{ws_name}', \
not the ticket's workspace '{workspace_name}'. \
Cross-workspace prerequisites are not allowed.",
);
}
Ok(())
}
pub async fn list_all_tickets(
&self,
workspace_name: Option<&str>,
phase_filter: Option<TicketPhase>,
) -> Result<Vec<Ticket>> {
let phase_str: Option<&str> = phase_filter.as_ref().map(TicketPhase::as_ref);
self.select_tickets(
"WHERE (?1 IS NULL OR workspace_name = ?1) \
AND (?2 IS NULL OR phase = ?2) \
AND is_archived = 0 \
ORDER BY priority ASC, created_at DESC",
db::params![workspace_name, phase_str],
LoadComments::No,
)
.await
}
pub(crate) async fn list_working_tickets(
&self,
workspace_name: &str,
) -> Result<Vec<(TicketPhase, Ticket)>> {
let phase_placeholders = (2..=WORKING_PHASES.len() + 1)
.map(|i| format!("?{i}"))
.collect::<Vec<_>>()
.join(", ");
let sql = format!(
"SELECT {TICKET_COLUMNS} FROM tickets \
WHERE workspace_name = ?1 \
AND phase IN ({phase_placeholders}) \
AND is_archived = 0 \
ORDER BY priority ASC, created_at DESC"
);
let mut params: Vec<Value> = vec![Value::from(workspace_name)];
params.extend(WORKING_PHASES.iter().map(|p| Value::from(p.as_ref())));
let rows = self.conn.query_cached(&sql, params).await?;
let mut buckets: Vec<Vec<Ticket>> = vec![Vec::new(); WORKING_PHASES.len()];
for row in rows {
let ticket = self.ticket_from_row(&row, LoadComments::No).await?;
let Some(idx) = WORKING_PHASES.iter().position(|p| *p == ticket.phase) else {
continue; };
buckets[idx].push(ticket);
}
let mut out = Vec::with_capacity(buckets.iter().map(Vec::len).sum());
for (idx, phase) in WORKING_PHASES.iter().enumerate() {
out.extend(buckets[idx].drain(..).map(|ticket| (*phase, ticket)));
}
Ok(out)
}
pub async fn count_by_phase(
&self,
phase: TicketPhase,
workspace_name: Option<&str>,
) -> Result<i64> {
self.conn
.query_row(
"SELECT COUNT(*) FROM tickets \
WHERE phase = ?1 \
AND (?2 IS NULL OR workspace_name = ?2) \
AND is_archived = 0",
db::params![phase.as_ref(), workspace_name],
|row| row.get(0),
)
.await
.map_err(Into::into)
}
pub async fn set_archived(&self, ticket_id: &str) -> Result<()> {
let prepared =
Self::build_ticket_update_with_updated_at("is_archived = 1", vec![], ticket_id);
prepared.execute_and_cancel(&self.conn).await
}
pub(crate) async fn drain_queued_to_planning(&self, workspace_name: &str) -> Result<u64> {
let now = db::now();
let updated = self
.conn
.execute(
"UPDATE tickets SET phase = ?1, updated_at = ?2 \
WHERE phase = ?3 AND workspace_name = ?4 AND is_archived = 0",
db::params![
TicketPhase::Planning.as_ref(),
now,
TicketPhase::Queued.as_ref(),
workspace_name,
],
)
.await?;
Ok(updated)
}
pub async fn archive_stale_cancelled(&self, hours: i64) -> Result<u64> {
let now = db::now();
let cutoff = (Utc::now() - Duration::hours(hours)).to_rfc3339();
let updated = self
.conn
.execute(
"UPDATE tickets SET is_archived = 1, updated_at = ?1 \
WHERE phase = ?2 AND updated_at < ?3 \
AND is_archived = 0",
db::params![now, TicketPhase::Cancelled.as_ref(), cutoff],
)
.await
.context("Failed to archive stale cancelled tickets")?;
Ok(updated)
}
pub async fn archive_all_done_and_cancelled(
&self,
workspace_name: Option<&str>,
) -> Result<u64> {
let now = db::now();
let sql = format!(
"UPDATE tickets SET is_archived = 1, updated_at = ?1 \
WHERE phase IN ({}) AND is_archived = 0 \
AND (?2 IS NULL OR workspace_name = ?2)",
phase_list_sql_fragment(UNBLOCKING_PHASES),
);
let updated = self
.conn
.execute(&sql, db::params![now, workspace_name])
.await
.context("Failed to archive done/cancelled tickets")?;
Ok(updated)
}
#[must_use]
pub fn partition_board_tickets(
tickets: &[Ticket],
) -> (Vec<&Ticket>, Vec<&Ticket>, Vec<&Ticket>) {
let mut pending = Vec::new();
let mut pipeline = Vec::new();
let mut completed = Vec::new();
for ticket in tickets {
if ticket.is_archived {
continue; }
if ticket.phase.is_unblocking() {
completed.push(ticket);
} else if ticket.phase.is_pipeline_occupied() || ticket.phase == TicketPhase::Queued {
pipeline.push(ticket);
} else {
pending.push(ticket);
}
}
pending.sort_by(|a, b| {
a.priority
.cmp(&b.priority)
.then(a.created_at.cmp(&b.created_at))
});
pipeline.sort_by(|a, b| {
a.priority
.cmp(&b.priority)
.then(a.created_at.cmp(&b.created_at))
});
completed.sort_by(|a, b| {
let (a_done, b_done) = (a.phase == TicketPhase::Done, b.phase == TicketPhase::Done);
match (a_done, b_done) {
(true, true) => Self::board_done_sort_key(b).cmp(Self::board_done_sort_key(a)),
(true, false) => std::cmp::Ordering::Less,
(false, true) => std::cmp::Ordering::Greater,
(false, false) => b.created_at.cmp(&a.created_at),
}
});
(pending, pipeline, completed)
}
#[must_use]
pub fn board_sections(tickets: &[Ticket]) -> [Vec<&Ticket>; 4] {
let (pending, pipeline, completed) = Self::partition_board_tickets(tickets);
let in_progress = pipeline
.iter()
.filter(|t| t.phase != TicketPhase::Queued)
.copied()
.collect();
let queued = pipeline
.iter()
.filter(|t| t.phase == TicketPhase::Queued)
.copied()
.collect();
[in_progress, queued, pending, completed]
}
#[must_use]
pub fn board_display_order(tickets: &[Ticket]) -> Vec<&Ticket> {
Self::board_sections(tickets)
.into_iter()
.flatten()
.collect()
}
fn board_done_sort_key(ticket: &Ticket) -> &str {
ticket.done_at.as_deref().unwrap_or(&ticket.created_at)
}
pub async fn search_archived_by_fts(
&self,
query: &str,
limit: usize,
workspace_name: &str,
) -> Result<Vec<(String, f64)>> {
let sanitized = crate::db::sanitize_fts_query(query);
if sanitized.is_empty() {
return Ok(Vec::new());
}
let sql = format!(
"SELECT t.id, fts_score(t.title, ?2) AS score \
FROM tickets t \
WHERE t.is_archived = 1 \
AND t.workspace_name = ?1 \
AND t.title MATCH ?2 \
ORDER BY score DESC LIMIT {limit}"
);
match self
.conn
.query_map(
&sql,
db::params![workspace_name, sanitized.clone()],
|row| {
let id: String = row.get(0)?;
let score: f64 = row.get(1)?;
Ok::<_, anyhow::Error>((id, score))
},
)
.await
{
Ok(items) => Ok(items
.into_iter()
.collect::<std::result::Result<Vec<_>, _>>()?),
Err(e) => {
tracing::warn!(
query = %sanitized,
error = %e,
"FTS search for archived tickets failed"
);
Ok(Vec::new())
}
}
}
pub async fn search_by_fts(
&self,
query: &str,
limit: usize,
workspace_name: Option<&str>,
) -> Result<Vec<Ticket>> {
let sanitized = crate::db::sanitize_fts_query(query);
if sanitized.is_empty() {
return Ok(Vec::new());
}
let sql = format!(
"SELECT {TICKET_COLUMNS} \
FROM tickets t \
WHERE (?1 IS NULL OR t.workspace_name = ?1) \
AND t.title MATCH ?2 \
ORDER BY fts_score(t.title, ?2) DESC \
LIMIT {limit}"
);
match self
.conn
.query(&sql, db::params![workspace_name, sanitized.clone()])
.await
{
Ok(rows) => {
let mut tickets = Vec::with_capacity(rows.len());
for row in rows {
match self.ticket_from_row(&row, LoadComments::No).await {
Ok(t) => tickets.push(t),
Err(e) => {
tracing::warn!(
error = %e,
"Failed to parse ticket row from FTS search"
);
}
}
}
Ok(tickets)
}
Err(e) => {
tracing::warn!(
query = %sanitized,
error = %e,
"FTS search failed"
);
Ok(Vec::new())
}
}
}
pub async fn list_archived_with_embeddings(
&self,
workspace_name: &str,
) -> Result<Vec<(String, Vec<f32>)>> {
let rows = self
.conn
.query(
"SELECT id, embedding FROM tickets \
WHERE is_archived = 1 AND workspace_name = ?1 AND embedding IS NOT NULL",
db::params![workspace_name],
)
.await?;
let mut candidates: Vec<(String, Vec<f32>)> = Vec::new();
for row in &rows {
let id: String = row.get(0)?;
let stored: Vec<u8> = row.get(1)?;
let emb = crate::vector::bytes_to_vec(&stored);
candidates.push((id, emb));
}
Ok(candidates)
}
}