use std::sync::{Arc, OnceLock};
use tokio::sync::Mutex;
use super::types::{
CardPatch, TaskBoard, TaskBoardCard, TaskCardStatus, TodosSnapshot, non_empty, normalise_board,
now_stamp, render_markdown,
};
use crate::error::{Result, TinyAgentsError};
use crate::graph::thread_locks::ThreadLockMap;
use crate::harness::store::Store;
pub const TODOS_NAMESPACE: &str = "graph.todos";
fn thread_lock(thread_id: &str) -> Arc<Mutex<()>> {
static LOCKS: OnceLock<ThreadLockMap> = OnceLock::new();
LOCKS
.get_or_init(|| ThreadLockMap::new("todo lock map"))
.lock_for(thread_id)
}
fn key(thread_id: &str) -> String {
thread_id
.as_bytes()
.iter()
.map(|b| format!("{b:02x}"))
.collect()
}
fn validate_thread_id(thread_id: &str) -> Result<String> {
let trimmed = thread_id.trim();
if trimmed.is_empty() {
return Err(TinyAgentsError::Validation(
"task board thread_id must not be empty or whitespace".to_string(),
));
}
Ok(trimmed.to_string())
}
async fn load_cards(store: &Arc<dyn Store>, thread_id: &str) -> Result<Vec<TaskBoardCard>> {
match store.get(TODOS_NAMESPACE, &key(thread_id)).await? {
Some(value) => {
let board: TaskBoard = serde_json::from_value(value)?;
Ok(board.cards)
}
None => Ok(Vec::new()),
}
}
async fn save_cards(
store: &Arc<dyn Store>,
thread_id: &str,
cards: Vec<TaskBoardCard>,
) -> Result<Vec<TaskBoardCard>> {
let mut board = TaskBoard {
thread_id: thread_id.to_string(),
cards,
updated_at: now_stamp(),
};
normalise_board(&mut board);
let value = serde_json::to_value(&board)?;
store.put(TODOS_NAMESPACE, &key(thread_id), value).await?;
Ok(board.cards)
}
fn snapshot(thread_id: &str, cards: Vec<TaskBoardCard>) -> TodosSnapshot {
let markdown = render_markdown(&cards);
TodosSnapshot {
thread_id: thread_id.to_string(),
cards,
markdown,
}
}
fn enforce_single_in_progress(cards: &[TaskBoardCard]) -> Result<()> {
let in_progress = cards
.iter()
.filter(|c| matches!(c.status, TaskCardStatus::InProgress))
.count();
if in_progress > 1 {
return Err(TinyAgentsError::Validation(format!(
"only one todo may be `in_progress` at a time (got {in_progress})"
)));
}
Ok(())
}
pub async fn list(store: &Arc<dyn Store>, thread_id: &str) -> Result<TodosSnapshot> {
let thread_id = validate_thread_id(thread_id)?;
let lock = thread_lock(&thread_id);
let _guard = lock.lock().await;
let cards = load_cards(store, &thread_id).await?;
Ok(snapshot(&thread_id, cards))
}
pub async fn add(
store: &Arc<dyn Store>,
thread_id: &str,
content: &str,
patch: CardPatch,
) -> Result<TodosSnapshot> {
let thread_id = validate_thread_id(thread_id)?;
let content = content.trim();
if content.is_empty() {
return Err(TinyAgentsError::Validation(
"todo content must not be empty".to_string(),
));
}
let lock = thread_lock(&thread_id);
let _guard = lock.lock().await;
let mut cards = load_cards(store, &thread_id).await?;
let order = cards.len() as u32;
cards.push(TaskBoardCard {
title: content.to_string(),
status: patch.status.unwrap_or(TaskCardStatus::Todo),
objective: patch.objective.and_then(non_empty),
plan: patch.plan.unwrap_or_default(),
assigned_agent: patch.assigned_agent.and_then(non_empty),
allowed_tools: patch.allowed_tools.unwrap_or_default(),
approval_mode: patch.approval_mode.flatten(),
acceptance_criteria: patch.acceptance_criteria.unwrap_or_default(),
evidence: patch.evidence.unwrap_or_default(),
notes: patch.notes.and_then(non_empty),
blocker: patch.blocker.and_then(non_empty),
source_metadata: patch.source_metadata,
order,
..TaskBoardCard::new(content)
});
enforce_single_in_progress(&cards)?;
let cards = save_cards(store, &thread_id, cards).await?;
Ok(snapshot(&thread_id, cards))
}
pub async fn edit(
store: &Arc<dyn Store>,
thread_id: &str,
id: &str,
patch: CardPatch,
) -> Result<TodosSnapshot> {
let thread_id = validate_thread_id(thread_id)?;
let lock = thread_lock(&thread_id);
let _guard = lock.lock().await;
let mut cards = load_cards(store, &thread_id).await?;
let card = cards
.iter_mut()
.find(|c| c.id == id)
.ok_or_else(|| TinyAgentsError::Validation(format!("todo id '{id}' not found")))?;
if let Some(content) = patch.content {
let trimmed = content.trim().to_string();
if trimmed.is_empty() {
return Err(TinyAgentsError::Validation(
"todo content must not be empty".to_string(),
));
}
card.title = trimmed;
}
if let Some(status) = patch.status {
card.status = status;
}
if let Some(objective) = patch.objective {
card.objective = non_empty(objective);
}
if let Some(plan) = patch.plan {
card.plan = plan;
}
if let Some(assigned_agent) = patch.assigned_agent {
card.assigned_agent = non_empty(assigned_agent);
}
if let Some(allowed_tools) = patch.allowed_tools {
card.allowed_tools = allowed_tools;
}
if let Some(approval_mode) = patch.approval_mode {
card.approval_mode = approval_mode;
}
if let Some(acceptance_criteria) = patch.acceptance_criteria {
card.acceptance_criteria = acceptance_criteria;
}
if let Some(evidence) = patch.evidence {
card.evidence = evidence;
}
if let Some(notes) = patch.notes {
card.notes = non_empty(notes);
}
if let Some(blocker) = patch.blocker {
card.blocker = non_empty(blocker);
}
if let Some(source_metadata) = patch.source_metadata {
card.source_metadata = Some(source_metadata);
}
card.updated_at = now_stamp();
enforce_single_in_progress(&cards)?;
let cards = save_cards(store, &thread_id, cards).await?;
Ok(snapshot(&thread_id, cards))
}
pub async fn update_status(
store: &Arc<dyn Store>,
thread_id: &str,
id: &str,
status: TaskCardStatus,
) -> Result<TodosSnapshot> {
edit(
store,
thread_id,
id,
CardPatch {
status: Some(status),
..Default::default()
},
)
.await
}
pub async fn set_session_thread(
store: &Arc<dyn Store>,
thread_id: &str,
id: &str,
session_thread_id: Option<String>,
) -> Result<TodosSnapshot> {
let thread_id = validate_thread_id(thread_id)?;
let lock = thread_lock(&thread_id);
let _guard = lock.lock().await;
let mut cards = load_cards(store, &thread_id).await?;
let card = cards
.iter_mut()
.find(|c| c.id == id)
.ok_or_else(|| TinyAgentsError::Validation(format!("todo id '{id}' not found")))?;
card.session_thread_id = session_thread_id.and_then(non_empty);
card.updated_at = now_stamp();
let cards = save_cards(store, &thread_id, cards).await?;
Ok(snapshot(&thread_id, cards))
}
pub async fn decide_plan(
store: &Arc<dyn Store>,
thread_id: &str,
id: &str,
approve: bool,
) -> Result<TodosSnapshot> {
let thread_id = validate_thread_id(thread_id)?;
let lock = thread_lock(&thread_id);
let _guard = lock.lock().await;
let mut cards = load_cards(store, &thread_id).await?;
let card = cards
.iter_mut()
.find(|c| c.id == id)
.ok_or_else(|| TinyAgentsError::Validation(format!("todo id '{id}' not found")))?;
if card.status != TaskCardStatus::AwaitingApproval {
return Err(TinyAgentsError::Validation(format!(
"card '{id}' is not awaiting approval (status: {})",
card.status.as_str()
)));
}
card.status = if approve {
TaskCardStatus::Ready
} else {
TaskCardStatus::Rejected
};
card.updated_at = now_stamp();
let cards = save_cards(store, &thread_id, cards).await?;
Ok(snapshot(&thread_id, cards))
}
pub async fn revise_plan(store: &Arc<dyn Store>, thread_id: &str) -> Result<TodosSnapshot> {
let thread_id = validate_thread_id(thread_id)?;
let lock = thread_lock(&thread_id);
let _guard = lock.lock().await;
let mut cards = load_cards(store, &thread_id).await?;
for card in cards.iter_mut() {
if card.status == TaskCardStatus::AwaitingApproval {
card.status = TaskCardStatus::Rejected;
card.updated_at = now_stamp();
}
}
let cards = save_cards(store, &thread_id, cards).await?;
Ok(snapshot(&thread_id, cards))
}
pub async fn remove(store: &Arc<dyn Store>, thread_id: &str, id: &str) -> Result<TodosSnapshot> {
let thread_id = validate_thread_id(thread_id)?;
let lock = thread_lock(&thread_id);
let _guard = lock.lock().await;
let mut cards = load_cards(store, &thread_id).await?;
let before = cards.len();
cards.retain(|c| c.id != id);
if cards.len() == before {
return Err(TinyAgentsError::Validation(format!(
"todo id '{id}' not found"
)));
}
let cards = save_cards(store, &thread_id, cards).await?;
Ok(snapshot(&thread_id, cards))
}
pub async fn replace(
store: &Arc<dyn Store>,
thread_id: &str,
cards: Vec<TaskBoardCard>,
) -> Result<TodosSnapshot> {
let thread_id = validate_thread_id(thread_id)?;
let lock = thread_lock(&thread_id);
let _guard = lock.lock().await;
enforce_single_in_progress(&cards)?;
let cards = save_cards(store, &thread_id, cards).await?;
Ok(snapshot(&thread_id, cards))
}
pub async fn clear(store: &Arc<dyn Store>, thread_id: &str) -> Result<TodosSnapshot> {
let thread_id = validate_thread_id(thread_id)?;
let lock = thread_lock(&thread_id);
let _guard = lock.lock().await;
let cards = save_cards(store, &thread_id, Vec::new()).await?;
Ok(snapshot(&thread_id, cards))
}
pub async fn claim_card(
store: &Arc<dyn Store>,
thread_id: &str,
card_id: &str,
expected: &[TaskCardStatus],
target: TaskCardStatus,
) -> Result<TaskBoardCard> {
let thread_id = validate_thread_id(thread_id)?;
let lock = thread_lock(&thread_id);
let _guard = lock.lock().await;
let mut cards = load_cards(store, &thread_id).await?;
let card = cards.iter_mut().find(|c| c.id == card_id).ok_or_else(|| {
TinyAgentsError::Validation(format!("claim_card: card '{card_id}' not found on board"))
})?;
if !expected.contains(&card.status) {
return Err(TinyAgentsError::Validation(format!(
"claim_card: card '{card_id}' status is '{}', expected one of [{}]; claim rejected",
card.status.as_str(),
expected
.iter()
.map(TaskCardStatus::as_str)
.collect::<Vec<_>>()
.join(", ")
)));
}
card.status = target;
card.updated_at = now_stamp();
let claimed_id = card.id.clone();
enforce_single_in_progress(&cards)?;
let cards = save_cards(store, &thread_id, cards).await?;
cards
.into_iter()
.find(|c| c.id == claimed_id)
.ok_or_else(|| {
TinyAgentsError::Graph(format!("claim_card: card '{claimed_id}' lost after save"))
})
}