use std::collections::HashMap;
use std::fmt::Write;
use std::sync::{Mutex, OnceLock};
use crate::agent::role::SYSTEM_ROLE;
use crate::db::{Connection, Row, params};
use crate::pipeline::board::{TicketPhase, store as board};
use crate::util::UnwrapPoison;
fn attributed_actor(raw: Option<&str>) -> &str {
raw.filter(|s| !s.trim().is_empty()).unwrap_or(SYSTEM_ROLE)
}
fn text_cell(row: &Row, idx: usize) -> String {
row.get_value(idx)
.ok()
.and_then(|v| v.as_text().cloned())
.unwrap_or_default()
}
type Cursor = i64;
static CURSORS: OnceLock<Mutex<HashMap<String, Cursor>>> = OnceLock::new();
static CURSOR_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
static SERVICE_START: OnceLock<String> = OnceLock::new();
pub fn init_global() {
let _ = SERVICE_START.get_or_init(crate::db::now);
CURSORS.get_or_init(|| Mutex::new(HashMap::new()));
}
fn cursor_for(workspace_name: &str) -> Cursor {
let mut map = CURSORS
.get()
.expect("chronicle not initialized — call init_global() first")
.lock()
.unwrap_poison();
*map.entry(workspace_name.to_string()).or_insert(-1)
}
fn advance_cursor(workspace_name: &str, last_id: i64) {
let mut map = CURSORS
.get()
.expect("chronicle not initialized — call init_global() first")
.lock()
.unwrap_poison();
if let Some(cursor) = map.get_mut(workspace_name) {
*cursor = last_id;
}
}
pub fn start_subscriber() {
static STARTED: OnceLock<()> = OnceLock::new();
if STARTED.get().is_some() {
return;
}
let _ = STARTED.set(());
let conn = board().conn.clone();
crate::db::cdc::register_ticket_materializer(move |event| {
let conn = conn.clone();
Box::pin(async move {
if apply_change(&conn, event).await? {
prune_acked(&conn).await?;
}
Ok(())
})
});
}
async fn apply_change(
conn: &Connection,
event: &crate::db::cdc::ChangeEvent,
) -> anyhow::Result<bool> {
if event.table != "tickets" || event.change_type != crate::db::cdc::ChangeType::Update {
return Ok(false);
}
let before_phase = event
.before
.as_ref()
.and_then(|r| r.get("phase"))
.and_then(crate::db::cdc::CdcValue::as_text);
let after_phase = event
.after
.as_ref()
.and_then(|r| r.get("phase"))
.and_then(crate::db::cdc::CdcValue::as_text);
let (Some(before), Some(after)) = (before_phase, after_phase) else {
return Ok(false);
};
if before == after {
return Ok(false);
}
let (Ok(source), Ok(target)) = (before.parse::<TicketPhase>(), after.parse::<TicketPhase>())
else {
tracing::warn!(
before = %before,
after = %after,
"ticket chronicle: unreadable phase pair — skipping the hop",
);
return Ok(false);
};
let ticket_id = event.ticket_id().unwrap_or_default();
let workspace = event
.after
.as_ref()
.and_then(|r| r.get("workspace_name"))
.and_then(crate::db::cdc::CdcValue::as_text)
.unwrap_or_default();
let at = event
.after
.as_ref()
.and_then(|r| r.get("updated_at"))
.and_then(crate::db::cdc::CdcValue::as_text)
.map_or_else(
|| {
chrono::DateTime::<chrono::Utc>::from_timestamp(event.change_time, 0)
.map_or_else(crate::db::now, |dt| dt.to_rfc3339())
},
str::to_string,
);
let actor = attributed_actor(
event
.after
.as_ref()
.and_then(|r| r.get("last_transition_actor"))
.and_then(crate::db::cdc::CdcValue::as_text),
);
let rows = conn
.execute(
"INSERT OR IGNORE INTO ticket_chronicle (ticket_id, workspace_name, source_phase, \
target_phase, at, actor) VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
params![
ticket_id,
workspace,
source.as_ref(),
target.as_ref(),
at,
actor
],
)
.await?;
Ok(rows > 0)
}
async fn prune_acked(conn: &Connection) -> anyhow::Result<()> {
let pairs: Vec<(String, i64)> = {
let map = CURSORS
.get()
.expect("chronicle not initialized — call init_global() first")
.lock()
.unwrap_poison();
map.iter()
.filter(|(_, c)| **c >= 0)
.map(|(ws, c)| (ws.clone(), *c))
.collect()
};
for (workspace, last_id) in pairs {
conn.execute(
"DELETE FROM ticket_chronicle WHERE workspace_name = ?1 AND id <= ?2",
params![workspace, last_id],
)
.await?;
}
Ok(())
}
#[must_use]
pub(crate) async fn drain(workspace_name: &str) -> String {
let conn = &board().conn;
let cursor_guard = CURSOR_LOCK.lock().await;
let mut cursor = cursor_for(workspace_name);
if cursor < 0 {
let service_start = SERVICE_START.get().cloned().unwrap_or_else(crate::db::now);
let seed: i64 = conn
.query_row(
"SELECT COALESCE(MAX(id), 0) FROM ticket_chronicle \
WHERE workspace_name = ?1 AND at <= ?2",
params![workspace_name, service_start],
|row| row.get(0),
)
.await
.unwrap_or(0);
cursor = seed;
advance_cursor(workspace_name, seed);
}
let rows = match conn
.query(
"SELECT id, ticket_id, source_phase, target_phase, at, actor \
FROM ticket_chronicle \
WHERE workspace_name = ?1 AND id > ?2 ORDER BY id",
params![workspace_name, cursor],
)
.await
{
Ok(rows) => rows,
Err(e) => {
tracing::warn!(workspace = workspace_name, error = %e, "chronicle drain failed");
return String::new();
}
};
if rows.is_empty() {
return String::new();
}
let hops: Vec<Hop> = rows
.iter()
.map(|row| {
let actor = text_cell(row, 5);
Hop {
id: text_cell(row, 1),
source: text_cell(row, 2),
target: text_cell(row, 3),
at: text_cell(row, 4),
actor: attributed_actor(Some(&actor)).to_string(),
}
})
.collect();
let last_id = rows
.iter()
.filter_map(|row| row.get_value(0).ok())
.filter_map(|v| v.as_integer().copied())
.max()
.unwrap_or(cursor);
advance_cursor(workspace_name, last_id);
drop(cursor_guard);
if let Err(e) = prune_acked(conn).await {
tracing::warn!(error = %e, "chronicle prune after drain failed");
}
format_chronicle(&hops)
}
#[derive(Debug, Clone)]
struct Hop {
id: String,
source: String,
target: String,
at: String,
actor: String,
}
fn format_chronicle(hops: &[Hop]) -> String {
if hops.is_empty() {
return String::new();
}
let mut out = String::from("<ticket-updates>\n");
let mut current: Option<&str> = None;
for hop in hops {
if current != Some(hop.id.as_str()) {
let _ = writeln!(out, "• {}:", hop.id);
current = Some(&hop.id);
}
let _ = writeln!(
out,
" {} → {} ({}) [{}]",
hop.source, hop.target, hop.at, hop.actor
);
}
let _ = writeln!(out, "</ticket-updates>");
out
}
#[cfg(test)]
pub fn reset() {
init_global();
CURSORS.get().unwrap().lock().unwrap_poison().clear();
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn format_groups_consecutive_same_ticket_hops_under_one_header() {
let hops = vec![
Hop {
id: "mahbot-1736".into(),
source: "in_development".into(),
target: "verification".into(),
at: "2026-08-17T08:11:34.225709+00:00".into(),
actor: "engineer".into(),
},
Hop {
id: "mahbot-1736".into(),
source: "verification".into(),
target: "in_sanitation".into(),
at: "2026-08-17T08:21:19.225709+00:00".into(),
actor: "system".into(),
},
];
let result = format_chronicle(&hops);
assert!(result.starts_with("<ticket-updates>\n"));
assert!(result.ends_with("</ticket-updates>\n"));
assert_eq!(result.matches("• mahbot-1736:").count(), 1);
assert!(result.contains(
" in_development → verification (2026-08-17T08:11:34.225709+00:00) [engineer]"
));
assert!(result.contains(
" verification → in_sanitation (2026-08-17T08:21:19.225709+00:00) [system]"
));
}
#[test]
fn format_interleaved_tickets_are_run_grouped() {
let hops = vec![
Hop {
id: "mahbot-A".into(),
source: "backlog".into(),
target: "analysis".into(),
at: "2026-08-17T08:00:00+00:00".into(),
actor: "system".into(),
},
Hop {
id: "mahbot-B".into(),
source: "analysis".into(),
target: "planning".into(),
at: "2026-08-17T08:01:00+00:00".into(),
actor: "analyst".into(),
},
Hop {
id: "mahbot-A".into(),
source: "planning".into(),
target: "queued".into(),
at: "2026-08-17T08:02:00+00:00".into(),
actor: "system".into(),
},
];
let result = format_chronicle(&hops);
assert_eq!(result.matches("• mahbot-A:").count(), 2);
assert_eq!(result.matches("• mahbot-B:").count(), 1);
let hop_a1 = result.find(" backlog → analysis").unwrap();
let hop_b = result.find(" analysis → planning").unwrap();
let hop_a2 = result.find(" planning → queued").unwrap();
assert!(hop_a1 < hop_b && hop_b < hop_a2);
}
#[test]
fn format_empty_returns_empty() {
assert!(format_chronicle(&[]).is_empty());
}
async fn chronicle_rows_for(board: &crate::pipeline::BoardStore, ticket_id: &str) -> i64 {
board
.conn
.query(
"SELECT COUNT(*) FROM ticket_chronicle WHERE ticket_id = ?1",
crate::db::params![ticket_id],
)
.await
.unwrap()[0]
.get_value(0)
.unwrap()
.as_integer()
.copied()
.unwrap_or(0)
}
#[tokio::test]
async fn cdc_materializes_chronicle_and_drain_delivers_grouped_block() {
crate::util::test::init_test_stores().await;
let tmp = tempfile::tempdir().unwrap();
let ws =
crate::util::test::create_test_workspace(tmp.path().to_str().unwrap(), "cdc-ws").await;
let board = crate::pipeline::board::store();
let ticket_id =
crate::util::test::make_ticket(board, &ws, "CDC ticket", TicketPhase::Backlog).await;
board
.transition_to(
&ticket_id,
Some(TicketPhase::Backlog),
TicketPhase::Analysis,
"pipeline",
)
.await
.unwrap();
let mut count = 0;
for _ in 0..20 {
crate::db::cdc::drain_once(&board.conn).await.unwrap();
count = chronicle_rows_for(board, &ticket_id).await;
if count > 0 {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
}
assert_eq!(
count, 1,
"a phase change materializes exactly one chronicle row"
);
board
.conn
.execute(
"UPDATE tickets SET last_transition_actor = '' WHERE id = ?1",
crate::db::params![ticket_id.as_str()],
)
.await
.unwrap();
board
.conn
.execute(
"UPDATE tickets SET phase = ?1 WHERE id = ?2",
crate::db::params![TicketPhase::Planning.as_ref(), ticket_id.as_str()],
)
.await
.unwrap();
for _ in 0..20 {
crate::db::cdc::drain_once(&board.conn).await.unwrap();
count = chronicle_rows_for(board, &ticket_id).await;
if count >= 2 {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
}
assert_eq!(
count, 2,
"both transitions materialize one chronicle row each"
);
for phase in ["in_review", TicketPhase::Cancelled.as_ref()] {
board
.conn
.execute(
"UPDATE tickets SET phase = ?1 WHERE id = ?2",
crate::db::params![phase, ticket_id.as_str()],
)
.await
.unwrap();
}
board
.conn
.execute(
"UPDATE tickets SET phase = ?1 WHERE id = ?2",
crate::db::params![TicketPhase::Done.as_ref(), ticket_id.as_str()],
)
.await
.unwrap();
for _ in 0..20 {
crate::db::cdc::drain_once(&board.conn).await.unwrap();
count = chronicle_rows_for(board, &ticket_id).await;
if count >= 3 {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
}
assert_eq!(
count, 3,
"unreadable phase pairs are skipped and pruned without wedging the drainer"
);
let block = drain(ws.name.as_str()).await;
assert!(block.contains("<ticket-updates>"));
assert!(block.contains(&format!("• {ticket_id}:")));
assert!(block.contains("backlog → analysis"));
assert!(block.contains("[pipeline]"));
assert!(block.contains("analysis → planning"));
assert!(block.contains("[system]"));
}
}