use crate::turso::{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,
db_name = "board",
schema = SCHEMA,
post_open = after_open,
expect = "BOARD not initialized — call init_global() 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_engineer_anchors().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"),
}
match board.clear_terminal_reservations().await {
Ok(n) if n > 0 => info!(count = n, "Cleared stale terminal pipeline reservations"),
Ok(_) => debug!("Reservation sweep: no stale terminal reservations"),
Err(e) => warn!(error = %e, "Reservation sweep failed"),
}
}
}
const SCHEMA: &str = "\
CREATE TABLE IF NOT EXISTS tickets (
id TEXT PRIMARY KEY,
title TEXT NOT NULL,
description TEXT NOT NULL,
phase TEXT NOT NULL DEFAULT 'backlog',
assigned_to TEXT,
workspace_name TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
prerequisites TEXT NOT NULL DEFAULT '[]',
supersedes TEXT,
superseded_by TEXT,
commit_hash TEXT,
lines_added INTEGER,
lines_removed INTEGER,
reporter TEXT NOT NULL DEFAULT '',
is_archived INTEGER NOT NULL DEFAULT 0,
embedding BLOB,
pipeline_reservation INTEGER NOT NULL DEFAULT 0,
priority INTEGER NOT NULL DEFAULT 1,
reviewed_head TEXT,
reviewed_tree TEXT,
done_at TEXT,
bounce_count INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS ticket_comments (
id TEXT PRIMARY KEY,
ticket_id TEXT NOT NULL,
role TEXT NOT NULL,
content TEXT NOT NULL,
created_at TEXT NOT NULL,
FOREIGN KEY (ticket_id) REFERENCES tickets(id)
);
CREATE INDEX IF NOT EXISTS idx_ticket_comments_ticket_id ON ticket_comments(ticket_id);
CREATE TABLE IF NOT EXISTS ticket_counters (
workspace_name TEXT PRIMARY KEY,
next_id INTEGER NOT NULL DEFAULT 1
);
";
const TICKETS_FTS_INDEX_NAME: &str = "idx_tickets_title_fts";
const TICKETS_FTS_INDEX_DDL: &str = "\
CREATE INDEX IF NOT EXISTS idx_tickets_title_fts ON tickets \
USING fts (title) WITH (tokenizer = 'ngram')";
const CANCELLED_ARCHIVE_HOURS: i64 = 1;
crate::columns! {
TICKET_COLUMNS [TICKET] {
ID => "id",
TITLE => "title",
DESCRIPTION => "description",
PHASE => "phase",
ASSIGNED_TO => "assigned_to",
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",
PIPELINE_RESERVATION => "pipeline_reservation",
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_BLOCKING_PHASES: &[TicketPhase] = &[
TicketPhase::InDevelopment,
TicketPhase::InDiagnostics,
TicketPhase::DiagnosticsDone,
TicketPhase::InReview,
TicketPhase::Reviewed,
TicketPhase::InQa,
TicketPhase::QaPassed,
TicketPhase::InSanitation,
TicketPhase::SanitationPassed,
];
const SANITATION_PIPELINE_PHASES: &[TicketPhase] =
&[TicketPhase::InSanitation, TicketPhase::SanitationPassed];
#[cfg(test)]
const TRANSITORY_HANDOFF_PHASES: &[TicketPhase] = &[
TicketPhase::DiagnosticsDone,
TicketPhase::SanitationPassed,
TicketPhase::Reviewed,
TicketPhase::QaPassed,
];
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 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 assigned_to: Option<String>,
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 pipeline_reservation: 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 = 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 = self.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.push_str(&self.format_comments());
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 {
let end = 19.min(c.created_at.len());
let ts = &c.created_at[..end];
let _ = write!(s, "\n [{}] ({}): {}", c.role, ts, c.content);
}
}
s
}
}
#[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,
ReadyForDevelopment,
InDevelopment,
InDiagnostics,
DiagnosticsDone,
InSanitation,
SanitationPassed,
InReview,
Reviewed,
InQa,
QaPassed,
Done,
Cancelled,
Failed,
}
impl TicketPhase {
#[cfg(test)]
#[must_use]
fn is_transitory_handoff(self) -> bool {
TRANSITORY_HANDOFF_PHASES.contains(&self)
}
#[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_blocking(&self) -> bool {
PIPELINE_BLOCKING_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<turso::Value>,
ticket_id: String,
}
impl PreparedUpdate {
async fn execute_tx(self, tx: &turso::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: &turso::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: &turso::Connection) -> Result<()> {
let rows = conn.execute(&self.sql, self.params).await?;
BoardStore::ensure_ticket_found(rows, &self.ticket_id)?;
crate::registry::AGENT_REGISTRY.cancel_by_ticket_id(&self.ticket_id);
Ok(())
}
async fn execute_tx_matched(self, tx: &turso::TxGuard<'_>) -> Result<bool> {
let rows = tx.execute(&self.sql, self.params).await?;
Ok(rows > 0)
}
}
#[derive(Debug, Clone, Copy)]
struct ResetTransition {
from: TicketPhase,
to: TicketPhase,
pipeline_reservation: bool,
}
#[derive(Copy, Clone, Debug, PartialEq, Eq)]
pub(crate) enum PipelineCheck {
Skip,
Enforce,
}
#[derive(Copy, Clone, Debug, PartialEq, Eq)]
pub(crate) enum LoadComments {
Yes,
No,
}
impl BoardStore {
async fn after_open(&self) -> anyhow::Result<()> {
crate::turso::ensure_fts_index(
&self.conn,
TICKETS_FTS_INDEX_NAME,
"ngram",
TICKETS_FTS_INDEX_DDL,
)
.await?;
Ok(())
}
async fn insert_ticket_tx(
tx: &TxGuard<'_>,
ticket_id: &str,
params: &TicketParams,
supersedes: Option<&str>,
) -> Result<()> {
let now = turso::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)",
turso::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",
turso::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",
turso::params![new_json, turso::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",
turso::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",
turso::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 = turso::now();
let cancelled_rows = tx
.execute(
"UPDATE tickets SET phase = ?1, updated_at = ?2, assigned_to = NULL, \
superseded_by = ?4, is_archived = 1, done_at = NULL, pipeline_reservation = 0 \
WHERE id = ?3",
turso::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::registry::AGENT_REGISTRY.cancel_by_ticket_id(supersede_id);
tx.commit().await?;
Ok(new_id)
}
async fn ticket_from_row(
&self,
row: &turso::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>()?,
assigned_to: row.get(COL_TICKET_ASSIGNED_TO)?,
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)?,
pipeline_reservation: row.get::<bool>(COL_TICKET_PIPELINE_RESERVATION)?,
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_ticket_in_workspace(
&self,
current_phase: TicketPhase,
target_phase: TicketPhase,
workspace_name: &str,
pipeline_check: PipelineCheck,
claim_grace: Option<Duration>,
) -> Result<Option<Ticket>> {
let now = turso::now();
let prereq_filter = format!(
"AND 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),
);
let pipeline_blocker_clause = if pipeline_check == PipelineCheck::Enforce {
let blocker_sql = phase_list_sql_fragment(PIPELINE_BLOCKING_PHASES);
format!(
"AND NOT EXISTS (SELECT 1 FROM tickets t2 \
WHERE t2.workspace_name = t1.workspace_name \
AND t2.phase IN ({blocker_sql}) \
AND t2.id != t1.id) "
)
} else {
String::new()
};
let grace_clause = if claim_grace.is_some() {
"AND t1.created_at <= ?5 "
} else {
""
};
let sql = format!(
"UPDATE tickets SET phase = ?1, assigned_to = NULL, updated_at = ?2, \
pipeline_reservation = 0 \
WHERE id = (SELECT t1.id FROM tickets t1 \
WHERE t1.phase = ?3 AND t1.assigned_to IS NULL AND t1.workspace_name = ?4 \
AND t1.is_archived = 0 \
{grace_clause}{pipeline_blocker_clause}{prereq_filter} \
ORDER BY t1.pipeline_reservation DESC, t1.priority ASC, t1.created_at ASC LIMIT 1) \
RETURNING {TICKET_COLUMNS}"
);
let mut params: Vec<Value> = vec![
Value::from(target_phase.as_ref()),
Value::from(now),
Value::from(current_phase.as_ref()),
Value::from(workspace_name),
];
if let Some(grace) = claim_grace {
params.push(Value::from((Utc::now() - grace).to_rfc3339()));
}
let rows = self.conn.query(&sql, params).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 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",
turso::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 ({})", turso::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",
turso::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, turso::params![ticket_id], |row| row.get::<i64>(0))
.await
}
fn build_ticket_update_with_updated_at(
set_clause: &str,
set_params: Vec<turso::Value>,
ticket_id: &str,
) -> PreparedUpdate {
let now = turso::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,
reservation: Option<bool>,
) -> PreparedUpdate {
let now = turso::now();
let guard: Option<&str> = expected_phase.as_ref().map(TicketPhase::as_ref);
let reservation = if target_phase.is_terminal() {
Some(false)
} else {
reservation
};
let sql = "UPDATE tickets SET phase = ?1, assigned_to = NULL, updated_at = ?2, \
pipeline_reservation = COALESCE(?5, pipeline_reservation), \
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<turso::Value> = vec![
Value::from(target_phase.as_ref()),
Value::from(now),
Value::from(ticket_id),
Value::from(guard),
Value::from(reservation),
];
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,
reservation: Option<bool>,
) -> Result<()> {
let prepared =
Self::build_transition_sql(ticket_id, expected_phase, target_phase, reservation);
prepared.execute_and_cancel(&self.conn).await
}
pub(crate) async fn transition_to_tx(
tx: &TxGuard<'_>,
ticket_id: &str,
expected_phase: Option<TicketPhase>,
target_phase: TicketPhase,
reservation: Option<bool>,
) -> Result<bool> {
let prepared =
Self::build_transition_sql(ticket_id, expected_phase, target_phase, reservation);
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 async fn set_assigned_to_no_cancel(
&self,
ticket_id: &str,
assigned_to: Option<&str>,
) -> Result<()> {
let prepared = Self::build_ticket_update_with_updated_at(
"assigned_to = ?",
vec![Value::from(assigned_to)],
ticket_id,
);
prepared.execute_no_cancel(&self.conn).await
}
pub(crate) async fn set_assigned_to_tx(
tx: &TxGuard<'_>,
ticket_id: &str,
assigned_to: Option<&str>,
) -> Result<()> {
let prepared = Self::build_ticket_update_with_updated_at(
"assigned_to = ?",
vec![Value::from(assigned_to)],
ticket_id,
);
prepared.execute_tx(tx).await
}
pub async fn claim_diagnostics(&self, ticket_id: &str, assigned_to: &str) -> Result<bool> {
let now = turso::now();
let rows = self
.conn
.execute(
"UPDATE tickets \
SET assigned_to = ?1, updated_at = ?2 \
WHERE id = ?3 \
AND assigned_to IS NULL \
AND phase = ?4 \
AND is_archived = 0",
turso::params![
assigned_to,
now,
ticket_id,
TicketPhase::InDiagnostics.as_ref()
],
)
.await?;
if rows > 0 {
crate::registry::AGENT_REGISTRY.cancel_by_ticket_id(ticket_id);
}
Ok(rows > 0)
}
pub async fn claim_sanitation(&self, ticket_id: &str, assigned_to: &str) -> Result<bool> {
let now = turso::now();
let blocker = phase_list_sql_fragment(SANITATION_PIPELINE_PHASES);
let sql = format!(
"UPDATE tickets SET phase = ?1, assigned_to = ?2, updated_at = ?3 \
WHERE id = ?4 AND phase = ?5 AND is_archived = 0 \
AND NOT EXISTS (SELECT 1 FROM tickets t2 \
WHERE t2.workspace_name = \
(SELECT workspace_name FROM tickets WHERE id = ?4) \
AND t2.id != ?4 \
AND t2.phase IN ({blocker}))"
);
let rows = self
.conn
.execute(
&sql,
turso::params![
TicketPhase::InSanitation.as_ref(),
assigned_to,
now,
ticket_id,
TicketPhase::QaPassed.as_ref(),
],
)
.await?;
Ok(rows > 0)
}
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",
turso::params![turso::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"
)),
}
}
pub(crate) async fn bounce_back_to_dev(&self, ticket_id: &str) -> Result<bool> {
let applied = crate::turso::with_tx_outcome(
&self.conn,
ticket_id,
"bounce back to dev",
async |tx| {
if Self::transition_to_tx(
tx,
ticket_id,
Some(TicketPhase::Reviewed),
TicketPhase::ReadyForDevelopment,
Some(true),
)
.await?
{
Self::increment_bounce_count_tx(tx, ticket_id).await?;
Ok(true)
} else {
Ok(false)
}
},
)
.await?;
if applied {
crate::registry::AGENT_REGISTRY.cancel_by_ticket_id(ticket_id);
}
Ok(applied)
}
const RESET_TRANSITIONS: &[ResetTransition] = &[
ResetTransition {
from: TicketPhase::InDevelopment,
to: TicketPhase::ReadyForDevelopment,
pipeline_reservation: true,
},
ResetTransition {
from: TicketPhase::InDiagnostics,
to: TicketPhase::ReadyForDevelopment,
pipeline_reservation: true,
},
ResetTransition {
from: TicketPhase::InSanitation,
to: TicketPhase::QaPassed,
pipeline_reservation: true,
},
ResetTransition {
from: TicketPhase::InQa,
to: TicketPhase::Reviewed,
pipeline_reservation: false,
},
ResetTransition {
from: TicketPhase::InReview,
to: TicketPhase::DiagnosticsDone,
pipeline_reservation: false,
},
ResetTransition {
from: TicketPhase::Analysis,
to: TicketPhase::Backlog,
pipeline_reservation: false,
},
];
pub(crate) fn reset_transition(from: TicketPhase) -> Option<(TicketPhase, bool)> {
Self::RESET_TRANSITIONS
.iter()
.find(|t| t.from == from)
.map(|t| (t.to, t.pipeline_reservation))
}
pub(crate) const RESET_TICKET_SET_CLAUSE: &str =
"phase = ?1, assigned_to = NULL, updated_at = ?2, pipeline_reservation = ?4";
pub async fn reset_inflight_tickets(&self, exclude_ticket_ids: &[String]) -> Result<()> {
let tx = self.conn.begin_tx().await?;
let now = turso::now();
for transition in Self::RESET_TRANSITIONS {
let mut values: Vec<turso::Value> = vec![
turso::Value::Text(transition.to.as_ref().to_string()),
turso::Value::Text(now.clone()),
turso::Value::Text(transition.from.as_ref().to_string()),
turso::Value::Integer(i64::from(transition.pipeline_reservation)),
];
let clause = if exclude_ticket_ids.is_empty() {
String::new()
} else {
values.extend(
exclude_ticket_ids
.iter()
.map(|s| turso::Value::Text(s.clone())),
);
format!(
" AND id NOT IN ({})",
turso::sql_in_placeholders(exclude_ticket_ids.len())
)
};
let sql = format!(
"UPDATE tickets SET {} WHERE phase = ?3{clause}",
Self::RESET_TICKET_SET_CLAUSE
);
tx.execute(&sql, values).await?;
}
tx.commit().await?;
Ok(())
}
async fn has_active_tickets_internal(
&self,
workspace_name: &str,
exclude_ticket_id: Option<&str>,
require_rfd_reservation: bool,
) -> Result<bool> {
let blocker_sql = phase_list_sql_fragment(PIPELINE_BLOCKING_PHASES);
let rfd_condition = if require_rfd_reservation {
"(phase = ?3 AND pipeline_reservation = 1)".to_string()
} else {
"phase = ?3".to_string()
};
let sql = format!(
"SELECT 1 FROM tickets WHERE \
(phase IN ({blocker_sql}) OR {rfd_condition}) \
AND workspace_name = ?1 AND is_archived = 0 \
AND (?2 IS NULL OR id != ?2) LIMIT 1",
);
let rfd = TicketPhase::ReadyForDevelopment.as_ref();
let rows = self
.conn
.query(&sql, turso::params![workspace_name, exclude_ticket_id, rfd])
.await?;
Ok(!rows.is_empty())
}
#[cfg(test)]
pub(crate) async fn has_pipeline_blocker_for_workspace(
&self,
workspace_name: &str,
) -> Result<bool> {
self.has_active_tickets_internal(workspace_name, None, true)
.await
}
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), false)
.await
}
pub async fn add_comment(&self, ticket_id: &str, role: &str, content: &str) -> Result<()> {
crate::turso::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 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 Some(assigned_to) = ticket.assigned_to.as_ref() else {
return; };
for agent_id in assigned_to.split(',') {
let agent_id = agent_id.trim();
if agent_id.is_empty() {
continue;
}
let commenter_role = role.parse::<crate::Role>().unwrap_or(crate::Role::Manager);
let job = crate::message_router::AgentJob {
content: content.to_string(),
workspace_name: ticket.workspace_name.clone(),
user_name: role.to_string(),
channel: String::new(),
kind: crate::message_router::JobKind::TicketComment,
role: commenter_role,
reply_target: None,
pending_job_id: None,
};
if crate::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 = turso::now();
tx.execute(
"INSERT INTO ticket_comments (id, ticket_id, role, content, created_at) \
VALUES (?1, ?2, ?3, ?4, ?5)",
turso::params![comment_id, ticket_id, role, content, now.as_str()],
)
.await?;
tx.execute(
"UPDATE tickets SET updated_at = ?1 WHERE id = ?2",
turso::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, turso::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",
turso::params![workspace_name, phase_str],
LoadComments::No,
)
.await
}
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",
turso::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, assigned_to = NULL",
vec![],
ticket_id,
);
prepared.execute_and_cancel(&self.conn).await
}
pub(crate) async fn drain_ready_for_development_to_planning(
&self,
workspace_name: &str,
) -> Result<u64> {
let now = turso::now();
let updated = self
.conn
.execute(
"UPDATE tickets SET phase = ?1, assigned_to = NULL, updated_at = ?2 \
WHERE phase = ?3 AND workspace_name = ?4 AND is_archived = 0",
turso::params![
TicketPhase::Planning.as_ref(),
now,
TicketPhase::ReadyForDevelopment.as_ref(),
workspace_name,
],
)
.await?;
Ok(updated)
}
pub async fn archive_stale_cancelled(&self, hours: i64) -> Result<u64> {
let now = turso::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 assigned_to IS NULL \
AND is_archived = 0",
turso::params![now, TicketPhase::Cancelled.as_ref(), cutoff],
)
.await
.context("Failed to archive stale cancelled tickets")?;
Ok(updated)
}
pub(crate) async fn clear_terminal_reservations(&self) -> Result<u64> {
let sql = format!(
"UPDATE tickets SET pipeline_reservation = 0 \
WHERE phase IN ({}) AND pipeline_reservation = 1",
phase_list_sql_fragment(TERMINAL_PHASES),
);
let updated = self
.conn
.execute(&sql, turso::params![])
.await
.context("Failed to clear stale terminal pipeline reservations")?;
Ok(updated)
}
pub async fn archive_all_done_and_cancelled(
&self,
workspace_name: Option<&str>,
) -> Result<u64> {
let now = turso::now();
let sql = format!(
"UPDATE tickets SET is_archived = 1, updated_at = ?1 \
WHERE phase IN ({}) AND assigned_to IS NULL AND is_archived = 0 \
AND (?2 IS NULL OR workspace_name = ?2)",
phase_list_sql_fragment(UNBLOCKING_PHASES),
);
let updated = self
.conn
.execute(&sql, turso::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_blocking()
|| ticket.phase == TicketPhase::ReadyForDevelopment
{
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::ReadyForDevelopment)
.copied()
.collect();
let ready = pipeline
.iter()
.filter(|t| t.phase == TicketPhase::ReadyForDevelopment)
.copied()
.collect();
[in_progress, ready, 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::turso::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,
turso::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::turso::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, turso::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",
turso::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)
}
}
#[cfg(test)]
async fn open_test_store() -> (BoardStore, tempfile::TempDir) {
crate::open_test_store!(BoardStore, "board")
}
#[cfg(test)]
#[path = "board_tests.rs"]
mod tests;