use std::collections::HashMap;
use std::path::{Path, PathBuf};
use rusqlite::{Connection, OptionalExtension, params};
use crate::error::{Error, Result};
use crate::model::{
Dep, DepKind, DepRef, Doc, Event, LivenessSource, Priority, Project, RefineOutcome, Repo,
RunningRow, Session, SessionOutcome, Task, TaskState, location_is_url,
};
const MIGRATIONS: &[&str] = &[
include_str!("../migrations/0001_init.sql"),
include_str!("../migrations/0002_rename_backlog_to_parked.sql"),
include_str!("../migrations/0003_track_pr.sql"),
include_str!("../migrations/0004_add_session_ref.sql"),
include_str!("../migrations/0005_add_branch.sql"),
include_str!("../migrations/0006_one_open_session_per_task.sql"),
include_str!("../migrations/0007_add_human.sql"),
include_str!("../migrations/0008_add_stalled_state.sql"),
include_str!("../migrations/0009_add_review_action.sql"),
include_str!("../migrations/0010_add_waiting_state.sql"),
include_str!("../migrations/0011_add_archived.sql"),
include_str!("../migrations/0012_repos.sql"),
include_str!("../migrations/0013_add_deep.sql"),
include_str!("../migrations/0014_docs.sql"),
include_str!("../migrations/0015_dep_kind_in_key.sql"),
include_str!("../migrations/0016_add_refining_state.sql"),
include_str!("../migrations/0017_schema_migrations.sql"),
include_str!("../migrations/0018_project_viewer.sql"),
include_str!("../migrations/0019_session_liveness_source.sql"),
include_str!("../migrations/0020_store_meta.sql"),
];
fn path_is_cargo_target(path: &Path) -> bool {
let parts: Vec<_> = path
.components()
.map(|c| c.as_os_str().to_string_lossy().into_owned())
.collect();
parts
.windows(2)
.any(|pair| pair[0] == "target" && (pair[1] == "debug" || pair[1] == "release"))
}
fn record_in_journal(tx: &Connection, from_version: usize, consent: Option<&str>) -> Result<()> {
let found: i64 = tx.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = 'schema_migrations'",
[],
|row| row.get(0),
)?;
if found == 0 {
return Ok(());
}
for idx in 1..=from_version {
tx.execute(
"INSERT OR IGNORE INTO schema_migrations (idx, sql, applied_at, applied_by)
VALUES (?1, NULL, datetime('now'), NULL)",
params![idx as i64],
)?;
}
let by = applied_by(consent);
for idx in (from_version + 1)..=MIGRATIONS.len() {
tx.execute(
"INSERT OR REPLACE INTO schema_migrations (idx, sql, applied_at, applied_by)
VALUES (?1, ?2, datetime('now'), ?3)",
params![idx as i64, MIGRATIONS[idx - 1], by],
)?;
}
Ok(())
}
fn applied_by(consent: Option<&str>) -> String {
let exe = std::env::current_exe()
.map(|p| p.display().to_string())
.unwrap_or_else(|_| "an unknown executable".to_string());
let via = consent.map(|c| format!(", {c}")).unwrap_or_default();
format!("voro {} at {exe}{via}", env!("CARGO_PKG_VERSION"))
}
fn remedy_for_divergence(path: Option<&Path>) -> String {
if path.is_some_and(|p| p == Store::dev_db_path()) {
format!(
"It is the dev store, which is disposable — rebuild it at this build's schema with \
`voro seed --force`, or delete {} and it will be reseeded on the next run.",
Store::dev_db_path().display()
)
} else {
format!(
"Restore the snapshot taken before that migration from {}; failing that, run the \
build named above, which is the only one whose schema matches this database.",
Store::backup_dir_for(path.unwrap_or(&Store::production_db_path())).display()
)
}
}
fn remedy_for_schema_ahead(path: Option<&Path>) -> String {
if path.is_some_and(|p| p == Store::dev_db_path()) {
format!(
"It is the dev store, which is disposable — rebuild it at this build's schema with \
`voro seed --force`, or delete {} and it will be reseeded on the next run.",
Store::dev_db_path().display()
)
} else {
format!(
"Restore a pre-migration snapshot from {}; failing that, run the build that migrated \
it — though if that build was never released, doing so entrenches a schema no other \
build can open.",
Store::backup_dir_for(path.unwrap_or(&Store::production_db_path())).display()
)
}
}
pub struct Store {
pub(crate) conn: Connection,
}
impl Store {
pub fn truncate_all(&mut self) -> Result<()> {
let tx = self.conn.transaction()?;
tx.pragma_update(None, "foreign_keys", false)?;
for table in [
"task_docs",
"docs",
"deps",
"events",
"sessions",
"tasks",
"repos",
"projects",
] {
tx.execute(&format!("DELETE FROM {table}"), [])?;
}
tx.execute("DELETE FROM sqlite_sequence", []).ok();
tx.commit()?;
self.conn.pragma_update(None, "foreign_keys", true)?;
Ok(())
}
}
#[derive(Debug, Clone)]
pub struct NewTask {
pub project_id: i64,
pub repo_id: Option<i64>,
pub title: String,
pub body: String,
pub priority: Priority,
pub state: TaskState,
pub agent: Option<String>,
pub human: bool,
pub deep: bool,
}
#[derive(Debug, Clone)]
pub struct TaskEdit {
pub title: String,
pub body: String,
pub priority: Priority,
pub agent: Option<String>,
pub human: bool,
pub deep: bool,
}
impl Store {
pub fn open(path: &Path) -> Result<Store> {
Store::open_with_consent(path, None)
}
pub fn open_migrate(path: &Path, consent: &str) -> Result<Store> {
Store::open_with_consent(path, Some(consent))
}
fn open_with_consent(path: &Path, consent: Option<&str>) -> Result<Store> {
if let Some(dir) = path.parent() {
std::fs::create_dir_all(dir)
.map_err(|e| Error::Invalid(format!("cannot create {}: {e}", dir.display())))?;
}
Store::open_at(
Connection::open(path)?,
path,
&Store::production_db_path(),
consent,
)
}
pub fn open_in_memory() -> Result<Store> {
Store::from_connection_at(Connection::open_in_memory()?, None)
}
pub fn data_dir() -> PathBuf {
let data_home = std::env::var_os("XDG_DATA_HOME")
.map(PathBuf::from)
.filter(|p| p.is_absolute())
.unwrap_or_else(|| {
let home = std::env::var_os("HOME")
.map(PathBuf::from)
.unwrap_or_default();
home.join(".local/share")
});
data_home.join("voro")
}
pub fn production_db_path() -> PathBuf {
Store::data_dir().join("voro.db")
}
pub fn dev_db_path() -> PathBuf {
Store::data_dir().join("dev.db")
}
pub fn backup_dir_for(path: &Path) -> PathBuf {
path.parent()
.filter(|p| !p.as_os_str().is_empty())
.unwrap_or(Path::new("."))
.join("backups")
}
pub fn is_dev_build() -> bool {
std::env::current_exe().is_ok_and(|exe| path_is_cargo_target(&exe))
}
pub fn default_db_path() -> PathBuf {
if Store::is_dev_build() {
Store::dev_db_path()
} else {
Store::production_db_path()
}
}
#[cfg(test)]
fn from_connection(conn: Connection) -> Result<Store> {
Store::from_connection_at(conn, None)
}
fn from_connection_at(conn: Connection, path: Option<&Path>) -> Result<Store> {
Store::open_at_opt(conn, path, &Store::production_db_path(), None)
}
fn open_at(
conn: Connection,
path: &Path,
production: &Path,
consent: Option<&str>,
) -> Result<Store> {
Store::open_at_opt(conn, Some(path), production, consent)
}
fn open_at_opt(
conn: Connection,
path: Option<&Path>,
production: &Path,
consent: Option<&str>,
) -> Result<Store> {
conn.pragma_update(None, "foreign_keys", true)?;
let mut store = Store { conn };
let version = store.schema_version()?;
store.verify_journal(path)?;
if version > MIGRATIONS.len() {
return Err(Error::SchemaAhead {
version,
known: MIGRATIONS.len(),
remedy: remedy_for_schema_ahead(path),
});
}
if version < MIGRATIONS.len()
&& let Some(path) = path
{
if version > 0 && consent.is_none() && store.is_protected(path, production)? {
return Err(Error::MigrationsPending {
path: path.to_path_buf(),
pending: MIGRATIONS.len() - version,
version,
known: MIGRATIONS.len(),
});
}
store.snapshot(path, version)?;
}
store.migrate(consent)?;
if path == Some(production) {
store.mark_protected()?;
}
Ok(store)
}
fn is_protected(&self, path: &Path, production: &Path) -> Result<bool> {
if path == production {
return Ok(true);
}
let has_meta: i64 = self.conn.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = 'store_meta'",
[],
|row| row.get(0),
)?;
if has_meta == 0 {
return Ok(false);
}
let marked: Option<String> = self
.conn
.query_row(
"SELECT value FROM store_meta WHERE key = 'protected'",
[],
|row| row.get(0),
)
.optional()?;
Ok(marked.as_deref() == Some("1"))
}
fn mark_protected(&self) -> Result<()> {
self.conn.execute(
"INSERT OR IGNORE INTO store_meta (key, value) VALUES ('protected', '1')",
[],
)?;
Ok(())
}
fn verify_journal(&self, path: Option<&Path>) -> Result<()> {
if !self.has_journal()? {
return Ok(());
}
let mut stmt = self.conn.prepare(
"SELECT idx, sql, applied_at, applied_by FROM schema_migrations
WHERE sql IS NOT NULL ORDER BY idx",
)?;
let rows = stmt.query_map([], |row| {
Ok((
row.get::<_, i64>(0)? as usize,
row.get::<_, String>(1)?,
row.get::<_, String>(2)?,
row.get::<_, Option<String>>(3)?,
))
})?;
for row in rows {
let (idx, applied, applied_at, applied_by) = row?;
let Some(carried) = MIGRATIONS.get(idx - 1) else {
continue;
};
if applied != *carried {
return Err(Error::SchemaDiverged {
idx,
applied_at,
applied_by: applied_by.unwrap_or_else(|| "an unrecorded build".to_string()),
remedy: remedy_for_divergence(path),
});
}
}
Ok(())
}
fn has_journal(&self) -> Result<bool> {
let found: i64 = self.conn.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = 'schema_migrations'",
[],
|row| row.get(0),
)?;
Ok(found > 0)
}
pub fn schema_version(&self) -> Result<usize> {
Ok(self
.conn
.query_row("PRAGMA user_version", [], |row| row.get::<_, i64>(0))? as usize)
}
fn snapshot(&self, path: &Path, version: usize) -> Result<()> {
if version == 0 || !path.exists() {
return Ok(());
}
let stamp: String =
self.conn
.query_row("SELECT strftime('%Y%m%d-%H%M%S', 'now')", [], |row| {
row.get(0)
})?;
let stem = path
.file_stem()
.map(|s| s.to_string_lossy().into_owned())
.unwrap_or_else(|| "voro".to_string());
let dir = Store::backup_dir_for(path);
let target = dir.join(format!("{stem}-v{version}-{stamp}.db"));
let copied = std::fs::create_dir_all(&dir).and_then(|()| std::fs::copy(path, &target));
if let Err(e) = copied {
eprintln!(
"voro: could not snapshot {} before migrating to schema {}: {e}",
path.display(),
MIGRATIONS.len()
);
}
Ok(())
}
pub fn data_version(&self) -> Result<i64> {
Ok(self
.conn
.query_row("PRAGMA data_version", [], |r| r.get(0))?)
}
fn migrate(&mut self, consent: Option<&str>) -> Result<()> {
self.conn.pragma_update(None, "foreign_keys", false)?;
let applied = self.apply_migrations(consent);
let restored = self.conn.pragma_update(None, "foreign_keys", true);
applied?;
restored?;
let violations: i64 =
self.conn
.query_row("SELECT COUNT(*) FROM pragma_foreign_key_check", [], |r| {
r.get(0)
})?;
if violations > 0 {
return Err(Error::Invalid(format!(
"{violations} foreign key violation(s) after migration"
)));
}
Ok(())
}
fn apply_migrations(&mut self, consent: Option<&str>) -> Result<()> {
let tx = self.conn.transaction()?;
let version: usize =
tx.query_row("PRAGMA user_version", [], |row| row.get::<_, i64>(0))? as usize;
for (i, sql) in MIGRATIONS.iter().enumerate().skip(version) {
tx.execute_batch(sql)?;
tx.pragma_update(None, "user_version", (i + 1) as i64)?;
}
record_in_journal(&tx, version, consent)?;
tx.commit()?;
Ok(())
}
pub fn create_project(&mut self, name: &str, path: &str) -> Result<Project> {
let tx = self.conn.transaction()?;
tx.execute("INSERT INTO projects (name) VALUES (?1)", params![name])?;
let id = tx.last_insert_rowid();
tx.execute(
"INSERT INTO repos (project_id, name, path, is_default) VALUES (?1, ?2, ?3, 1)",
params![id, name, path],
)?;
tx.commit()?;
self.project(id)
}
pub fn project(&self, id: i64) -> Result<Project> {
self.conn
.query_row(
&format!("SELECT {PROJECT_COLUMNS} FROM projects WHERE id = ?1"),
[id],
project_from_row,
)
.optional()?
.ok_or(Error::ProjectNotFound(id))
}
pub fn projects(&self) -> Result<Vec<Project>> {
let mut stmt = self.conn.prepare(&format!(
"SELECT {PROJECT_COLUMNS} FROM projects ORDER BY name"
))?;
let rows = stmt.query_map([], project_from_row)?;
Ok(rows.collect::<rusqlite::Result<_>>()?)
}
pub fn set_weight(&mut self, project_id: i64, weight: i64) -> Result<()> {
if !(0..=5).contains(&weight) {
return Err(Error::Invalid(format!("weight {weight} out of range 0-5")));
}
let changed = self.conn.execute(
"UPDATE projects SET weight = ?1 WHERE id = ?2",
params![weight, project_id],
)?;
if changed == 0 {
return Err(Error::ProjectNotFound(project_id));
}
Ok(())
}
pub fn set_viewer(&mut self, project_id: i64, viewer: Option<&str>) -> Result<Project> {
let viewer = match viewer.map(str::trim) {
Some("") => {
return Err(Error::Invalid(
"viewer name is required — name no viewer to use the default one".into(),
));
}
named => named,
};
let changed = self.conn.execute(
"UPDATE projects SET viewer = ?1 WHERE id = ?2",
params![viewer, project_id],
)?;
if changed == 0 {
return Err(Error::ProjectNotFound(project_id));
}
self.project(project_id)
}
pub fn set_archived(&mut self, project_id: i64, archived: bool) -> Result<Project> {
let project = self.project(project_id)?;
if project.archived == archived {
return Err(Error::Invalid(format!(
"project '{}' is {} archived",
project.name,
if archived { "already" } else { "not" }
)));
}
self.conn.execute(
"UPDATE projects SET archived = ?1 WHERE id = ?2",
params![archived, project_id],
)?;
self.project(project_id)
}
pub fn rename_project(&mut self, project_id: i64, name: &str) -> Result<Project> {
let changed = self.conn.execute(
"UPDATE projects SET name = ?1 WHERE id = ?2",
params![name, project_id],
)?;
if changed == 0 {
return Err(Error::ProjectNotFound(project_id));
}
self.project(project_id)
}
pub fn set_default_repo_path(&mut self, project_id: i64, path: &str) -> Result<Repo> {
let repo = self.default_repo(project_id)?;
self.set_repo_path(repo.id, path)
}
pub fn delete_project(&mut self, project_id: i64) -> Result<()> {
self.project(project_id)?;
let task_count: i64 = self.conn.query_row(
"SELECT COUNT(*) FROM tasks WHERE project_id = ?1",
[project_id],
|r| r.get(0),
)?;
if task_count > 0 {
return Err(Error::ProjectHasTasks {
id: project_id,
count: task_count,
});
}
let tx = self.conn.transaction()?;
tx.execute("DELETE FROM repos WHERE project_id = ?1", [project_id])?;
tx.execute("DELETE FROM projects WHERE id = ?1", [project_id])?;
tx.commit()?;
Ok(())
}
pub fn repos(&self, project_id: i64) -> Result<Vec<Repo>> {
let mut stmt = self.conn.prepare(&format!(
"SELECT {REPO_COLUMNS} FROM repos WHERE project_id = ?1
ORDER BY is_default DESC, name"
))?;
let rows = stmt.query_map([project_id], repo_from_row)?;
Ok(rows.collect::<rusqlite::Result<_>>()?)
}
pub fn repo(&self, id: i64) -> Result<Repo> {
self.conn
.query_row(
&format!("SELECT {REPO_COLUMNS} FROM repos WHERE id = ?1"),
[id],
repo_from_row,
)
.optional()?
.ok_or(Error::RepoIdNotFound(id))
}
pub fn repo_by_name(&self, project_id: i64, name: &str) -> Result<Repo> {
let repos = self.repos(project_id)?;
repos
.into_iter()
.find(|r| r.name == name)
.ok_or_else(|| Error::RepoNotFound {
project: self.project(project_id).map(|p| p.name).unwrap_or_default(),
name: name.to_string(),
known: self.repo_names(project_id),
})
}
pub fn default_repo(&self, project_id: i64) -> Result<Repo> {
self.conn
.query_row(
&format!(
"SELECT {REPO_COLUMNS} FROM repos WHERE project_id = ?1 AND is_default = 1"
),
[project_id],
repo_from_row,
)
.optional()?
.ok_or_else(|| match self.project(project_id) {
Err(e) => e,
Ok(_) => Error::Invalid(format!("project {project_id} has no default repo")),
})
}
pub fn repo_for_task(&self, task: &Task) -> Result<Repo> {
match task.repo_id {
Some(id) => self.repo(id),
None => self.default_repo(task.project_id),
}
}
pub fn add_repo(&mut self, project_id: i64, name: &str, path: &str) -> Result<Repo> {
self.project(project_id)?;
if name.trim().is_empty() {
return Err(Error::Invalid("a repo name is required".into()));
}
if self.repos(project_id)?.iter().any(|r| r.name == name) {
return Err(Error::Invalid(format!(
"project already has a repo named '{name}'"
)));
}
self.conn.execute(
"INSERT INTO repos (project_id, name, path, is_default) VALUES (?1, ?2, ?3, 0)",
params![project_id, name, path],
)?;
self.repo(self.conn.last_insert_rowid())
}
pub fn set_repo_path(&mut self, repo_id: i64, path: &str) -> Result<Repo> {
let changed = self.conn.execute(
"UPDATE repos SET path = ?1 WHERE id = ?2",
params![path, repo_id],
)?;
if changed == 0 {
return Err(Error::RepoIdNotFound(repo_id));
}
self.repo(repo_id)
}
pub fn set_default_repo(&mut self, repo_id: i64) -> Result<Repo> {
let repo = self.repo(repo_id)?;
let tx = self.conn.transaction()?;
tx.execute(
"UPDATE repos SET is_default = 0 WHERE project_id = ?1",
[repo.project_id],
)?;
tx.execute("UPDATE repos SET is_default = 1 WHERE id = ?1", [repo_id])?;
tx.commit()?;
self.repo(repo_id)
}
pub fn delete_repo(&mut self, repo_id: i64) -> Result<()> {
let repo = self.repo(repo_id)?;
let project = self.project(repo.project_id)?;
if self.repos(repo.project_id)?.len() == 1 {
return Err(Error::LastRepo {
project: project.name,
name: repo.name,
});
}
if repo.is_default {
return Err(Error::DefaultRepo {
project: project.name,
name: repo.name,
});
}
let used: i64 = self.conn.query_row(
"SELECT COUNT(*) FROM tasks WHERE repo_id = ?1",
[repo_id],
|r| r.get(0),
)?;
if used > 0 {
return Err(Error::RepoInUse {
name: repo.name,
count: used,
});
}
self.conn
.execute("DELETE FROM repos WHERE id = ?1", [repo_id])?;
Ok(())
}
pub fn set_task_repo(&mut self, task_id: i64, repo_id: Option<i64>) -> Result<Task> {
let task = self.task(task_id)?;
if let Some(id) = repo_id {
let repo = self.repo(id)?;
if repo.project_id != task.project_id {
return Err(Error::Invalid(format!(
"repo '{}' belongs to another project",
repo.name
)));
}
}
self.conn.execute(
"UPDATE tasks SET repo_id = ?1 WHERE id = ?2",
params![repo_id, task_id],
)?;
self.task(task_id)
}
fn repo_names(&self, project_id: i64) -> String {
self.repos(project_id)
.map(|repos| {
repos
.into_iter()
.map(|r| r.name)
.collect::<Vec<_>>()
.join(", ")
})
.unwrap_or_default()
}
pub fn create_doc(
&mut self,
project_id: i64,
repo_id: Option<i64>,
location: &str,
title: Option<&str>,
) -> Result<Doc> {
self.project(project_id)?;
let (location, repo_id) = self.normalise_location(project_id, repo_id, location)?;
if self
.docs(project_id)?
.iter()
.any(|d| d.location == location)
{
return Err(Error::Invalid(format!(
"this project already has a document at '{location}'"
)));
}
self.conn.execute(
"INSERT INTO docs (project_id, repo_id, title, location, created_at)
VALUES (?1, ?2, ?3, ?4, datetime('now'))",
params![project_id, repo_id, title, location],
)?;
let doc = self.doc(self.conn.last_insert_rowid())?;
log_global_event(&self.conn, "doc-added", Some(&doc.location))?;
Ok(doc)
}
fn normalise_location(
&self,
project_id: i64,
repo_id: Option<i64>,
location: &str,
) -> Result<(String, Option<i64>)> {
let location = location.trim();
if location.is_empty() {
return Err(Error::Invalid("a document path or URL is required".into()));
}
if let Some(id) = repo_id {
let repo = self.repo(id)?;
if repo.project_id != project_id {
return Err(Error::Invalid(format!(
"repo '{}' belongs to another project",
repo.name
)));
}
}
if location_is_url(location) {
if repo_id.is_some() {
return Err(Error::Invalid(
"a URL resolves on its own — drop --repo, which only picks the checkout a \
relative path is read from"
.into(),
));
}
return Ok((location.to_string(), None));
}
if !Path::new(location).is_absolute() {
return Ok((location.to_string(), repo_id));
}
let mut repos = match repo_id {
Some(id) => vec![self.repo(id)?],
None => self.repos(project_id)?,
};
repos.sort_by_key(|r| std::cmp::Reverse(r.path.len()));
for repo in &repos {
if let Ok(rel) = Path::new(location).strip_prefix(&repo.path) {
return Ok((rel.to_string_lossy().into_owned(), Some(repo.id)));
}
}
match repo_id {
Some(id) => Err(Error::Invalid(format!(
"'{location}' is not inside repo '{}' ({})",
self.repo(id)?.name,
self.repo(id)?.path
))),
None => Ok((location.to_string(), None)),
}
}
pub fn doc(&self, id: i64) -> Result<Doc> {
self.conn
.query_row(
&format!("SELECT {DOC_COLUMNS} FROM docs WHERE id = ?1"),
[id],
doc_from_row,
)
.optional()?
.ok_or(Error::DocNotFound(id))
}
pub fn docs(&self, project_id: i64) -> Result<Vec<Doc>> {
let mut stmt = self.conn.prepare(&format!(
"SELECT {DOC_COLUMNS} FROM docs WHERE project_id = ?1 ORDER BY id"
))?;
let rows = stmt.query_map([project_id], doc_from_row)?;
Ok(rows.collect::<rusqlite::Result<_>>()?)
}
pub fn all_docs(&self) -> Result<Vec<Doc>> {
let mut stmt = self
.conn
.prepare(&format!("SELECT {DOC_COLUMNS} FROM docs ORDER BY id"))?;
let rows = stmt.query_map([], doc_from_row)?;
Ok(rows.collect::<rusqlite::Result<_>>()?)
}
pub fn docs_at(&self, location: &str) -> Result<Vec<Doc>> {
let location = location.trim();
Ok(self
.all_docs()?
.into_iter()
.filter(|d| d.location == location)
.collect())
}
pub fn resolve_doc(&self, doc: &Doc) -> Result<String> {
if doc.is_url() || Path::new(&doc.location).is_absolute() {
return Ok(doc.location.clone());
}
let repo = match doc.repo_id {
Some(id) => self.repo(id)?,
None => self.default_repo(doc.project_id)?,
};
Ok(Path::new(&repo.path)
.join(&doc.location)
.to_string_lossy()
.into_owned())
}
pub fn delete_doc(&mut self, doc_id: i64) -> Result<Vec<i64>> {
let doc = self.doc(doc_id)?;
let linked = self.tasks_for_doc(doc_id)?;
let tx = self.conn.transaction()?;
for task in &linked {
log_event(&tx, task.id, "doc-unlinked", Some(doc.label()))?;
}
tx.execute("DELETE FROM task_docs WHERE doc_id = ?1", [doc_id])?;
tx.execute("DELETE FROM docs WHERE id = ?1", [doc_id])?;
log_global_event(&tx, "doc-removed", Some(&doc.location))?;
tx.commit()?;
Ok(linked.into_iter().map(|t| t.id).collect())
}
pub fn link_doc(&mut self, task_id: i64, doc_id: i64) -> Result<bool> {
self.task(task_id)?;
let doc = self.doc(doc_id)?;
let changed = self.conn.execute(
"INSERT OR IGNORE INTO task_docs (task_id, doc_id) VALUES (?1, ?2)",
params![task_id, doc_id],
)?;
if changed > 0 {
log_event(&self.conn, task_id, "doc-linked", Some(doc.label()))?;
}
Ok(changed > 0)
}
pub fn unlink_doc(&mut self, task_id: i64, doc_id: i64) -> Result<bool> {
self.task(task_id)?;
let doc = self.doc(doc_id)?;
let changed = self.conn.execute(
"DELETE FROM task_docs WHERE task_id = ?1 AND doc_id = ?2",
params![task_id, doc_id],
)?;
if changed > 0 {
log_event(&self.conn, task_id, "doc-unlinked", Some(doc.label()))?;
}
Ok(changed > 0)
}
pub fn set_task_docs(&mut self, task_id: i64, doc_ids: &[i64]) -> Result<Vec<Doc>> {
self.task(task_id)?;
let wanted: Vec<Doc> = doc_ids
.iter()
.map(|id| self.doc(*id))
.collect::<Result<_>>()?;
let current = self.docs_for_task(task_id)?;
let tx = self.conn.transaction()?;
for doc in ¤t {
if !wanted.iter().any(|d| d.id == doc.id) {
tx.execute(
"DELETE FROM task_docs WHERE task_id = ?1 AND doc_id = ?2",
params![task_id, doc.id],
)?;
log_event(&tx, task_id, "doc-unlinked", Some(doc.label()))?;
}
}
for doc in &wanted {
if !current.iter().any(|d| d.id == doc.id) {
tx.execute(
"INSERT INTO task_docs (task_id, doc_id) VALUES (?1, ?2)",
params![task_id, doc.id],
)?;
log_event(&tx, task_id, "doc-linked", Some(doc.label()))?;
}
}
tx.commit()?;
self.docs_for_task(task_id)
}
pub fn docs_for_task(&self, task_id: i64) -> Result<Vec<Doc>> {
let mut stmt = self.conn.prepare(&format!(
"SELECT {} FROM docs d JOIN task_docs td ON td.doc_id = d.id
WHERE td.task_id = ?1 ORDER BY d.id",
prefixed(DOC_COLUMNS, "d")
))?;
let rows = stmt.query_map([task_id], doc_from_row)?;
Ok(rows.collect::<rusqlite::Result<_>>()?)
}
pub fn docs_by_task(&self) -> Result<HashMap<i64, Vec<Doc>>> {
let mut stmt = self.conn.prepare(&format!(
"SELECT td.task_id, {} FROM docs d JOIN task_docs td ON td.doc_id = d.id
ORDER BY td.task_id, d.id",
prefixed(DOC_COLUMNS, "d")
))?;
let rows = stmt.query_map([], |row| {
Ok((row.get::<_, i64>(0)?, doc_from_row_at(row, 1)?))
})?;
let mut map: HashMap<i64, Vec<Doc>> = HashMap::new();
for row in rows {
let (task_id, doc) = row?;
map.entry(task_id).or_default().push(doc);
}
Ok(map)
}
pub fn tasks_for_doc(&self, doc_id: i64) -> Result<Vec<Task>> {
let mut stmt = self.conn.prepare(&format!(
"SELECT {} FROM tasks t JOIN task_docs td ON td.task_id = t.id
WHERE td.doc_id = ?1 ORDER BY t.id",
prefixed(TASK_COLUMNS, "t")
))?;
let rows = stmt.query_map([doc_id], task_from_row)?;
Ok(rows.collect::<rusqlite::Result<_>>()?)
}
pub fn create_task(&mut self, new: NewTask) -> Result<Task> {
if !matches!(
new.state,
TaskState::Proposed | TaskState::Parked | TaskState::Ready
) {
return Err(Error::Invalid(format!(
"a task cannot be created in state '{}'",
new.state
)));
}
if new.human && new.agent.is_some() {
return Err(Error::Invalid(
"a human-only task cannot carry an agent override — the override only \
selects a dispatch agent, and no agent can execute the task"
.into(),
));
}
if new.human && new.deep {
return Err(Error::Invalid(
"a human-only task cannot be deep — deep only selects a dispatch model, \
and no agent can execute the task"
.into(),
));
}
let project = self.project(new.project_id)?;
if project.archived {
return Err(Error::ProjectArchived { name: project.name });
}
if let Some(repo_id) = new.repo_id {
let repo = self.repo(repo_id)?;
if repo.project_id != new.project_id {
return Err(Error::Invalid(format!(
"repo '{}' belongs to another project",
repo.name
)));
}
}
let tx = self.conn.transaction()?;
tx.execute(
"INSERT INTO tasks (project_id, repo_id, title, body, priority, state, agent, human,
deep, state_since, created_at)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, datetime('now'), datetime('now'))",
params![
new.project_id,
new.repo_id,
new.title,
new.body,
new.priority,
new.state,
new.agent,
new.human,
new.deep
],
)?;
let id = tx.last_insert_rowid();
log_event(&tx, id, "created", Some(new.state.as_str()))?;
tx.commit()?;
self.task(id)
}
pub fn task(&self, id: i64) -> Result<Task> {
get_task(&self.conn, id)?.ok_or(Error::TaskNotFound(id))
}
pub fn tasks(&self) -> Result<Vec<Task>> {
let mut stmt = self
.conn
.prepare(&format!("SELECT {TASK_COLUMNS} FROM tasks ORDER BY id"))?;
let rows = stmt.query_map([], task_from_row)?;
Ok(rows.collect::<rusqlite::Result<_>>()?)
}
pub fn update_task(&mut self, id: i64, edit: TaskEdit) -> Result<Task> {
let current = self.task(id)?;
if edit.human && edit.agent.is_some() {
return Err(Error::HumanTask {
id,
reason: "an agent override is meaningless on a task no agent can execute — \
clear one or the other"
.into(),
});
}
if edit.human && edit.deep {
return Err(Error::HumanTask {
id,
reason: "the deep flag is meaningless on a task no agent can execute — it \
only selects a dispatch model; clear one or the other"
.into(),
});
}
if edit.human && !current.human {
if matches!(
current.state,
TaskState::NeedsInput | TaskState::Review | TaskState::Stalled
) {
return Err(Error::HumanTask {
id,
reason: format!(
"a task in state '{}' was executed by an agent; resolve it first",
current.state
),
});
}
let open_sessions: i64 = self.conn.query_row(
"SELECT COUNT(*) FROM sessions WHERE task_id = ?1 AND ended_at IS NULL",
[id],
|r| r.get(0),
)?;
if open_sessions > 0 {
return Err(Error::HumanTask {
id,
reason: "an agent session is still open on it; complete or abort it first"
.into(),
});
}
}
let tx = self.conn.transaction()?;
tx.execute(
"UPDATE tasks SET title = ?1, body = ?2, priority = ?3, agent = ?4, human = ?5,
deep = ?6
WHERE id = ?7",
params![
edit.title,
edit.body,
edit.priority,
edit.agent,
edit.human,
edit.deep,
id
],
)?;
if edit.body != current.body && !current.body.is_empty() {
log_event(&tx, id, "body", Some(¤t.body))?;
}
tx.commit()?;
self.task(id)
}
pub fn set_priority(&mut self, id: i64, priority: Priority) -> Result<Task> {
let changed = self.conn.execute(
"UPDATE tasks SET priority = ?1 WHERE id = ?2",
params![priority, id],
)?;
if changed == 0 {
return Err(Error::TaskNotFound(id));
}
log_event(&self.conn, id, "priority", Some(&priority.to_string()))?;
self.task(id)
}
pub fn set_deep(&mut self, id: i64, deep: bool) -> Result<Task> {
let task = self.task(id)?;
if deep && task.human {
return Err(Error::HumanTask {
id,
reason: "the deep flag only selects a dispatch model, and no agent can \
execute the task"
.into(),
});
}
self.conn.execute(
"UPDATE tasks SET deep = ?1 WHERE id = ?2",
params![deep, id],
)?;
log_event(
&self.conn,
id,
"deep",
Some(if deep { "set" } else { "cleared" }),
)?;
self.task(id)
}
pub fn set_pr(&mut self, id: i64, pr_url: Option<&str>) -> Result<Task> {
let changed = self.conn.execute(
"UPDATE tasks SET pr_url = ?1 WHERE id = ?2",
params![pr_url, id],
)?;
if changed == 0 {
return Err(Error::TaskNotFound(id));
}
log_event(&self.conn, id, "pr", pr_url.or(Some("cleared")))?;
self.task(id)
}
pub fn set_branch(&mut self, id: i64, branch: Option<&str>) -> Result<Task> {
let changed = self.conn.execute(
"UPDATE tasks SET branch = ?1 WHERE id = ?2",
params![branch, id],
)?;
if changed == 0 {
return Err(Error::TaskNotFound(id));
}
log_event(&self.conn, id, "branch", branch.or(Some("cleared")))?;
self.task(id)
}
pub fn set_summary(&mut self, id: i64, summary: &str) -> Result<Task> {
if summary.trim().is_empty() {
return Err(Error::Invalid("a summary is required".into()));
}
let task = self.task(id)?;
if !matches!(task.state, TaskState::Running | TaskState::Review) {
return Err(Error::Invalid(format!(
"a summary can only be set on a running or review task; task {} is {}",
id, task.state
)));
}
log_event(&self.conn, id, "summary", Some(summary.trim()))?;
self.task(id)
}
pub fn latest_refine_outcome(&self, task_id: i64) -> Result<Option<RefineOutcome>> {
let detail: Option<String> = self
.conn
.query_row(
"SELECT detail FROM events WHERE task_id = ?1 AND kind = 'refine'
ORDER BY id DESC LIMIT 1",
[task_id],
|r| r.get::<_, Option<String>>(0),
)
.optional()?
.flatten();
detail.map(|d| RefineOutcome::parse(&d)).transpose()
}
pub fn refined_flag(&self, task_id: i64) -> Result<bool> {
self.refine_marker(task_id, RefineOutcome::Applied)
}
pub fn refine_failed_flag(&self, task_id: i64) -> Result<bool> {
self.refine_marker(task_id, RefineOutcome::Failed)
}
pub fn correct_late_refine(&mut self, task_id: i64) -> Result<bool> {
if !self.refine_failed_flag(task_id)? {
return Ok(false);
}
log_event(
&self.conn,
task_id,
"refine",
Some(RefineOutcome::Applied.as_str()),
)?;
Ok(true)
}
fn refine_marker(&self, task_id: i64, wanted: RefineOutcome) -> Result<bool> {
let state: Option<TaskState> = self
.conn
.query_row("SELECT state FROM tasks WHERE id = ?1", [task_id], |r| {
r.get(0)
})
.optional()?;
if state != Some(TaskState::Proposed) {
return Ok(false);
}
Ok(self.latest_refine_outcome(task_id)? == Some(wanted))
}
pub fn latest_refine_note(&self, task_id: i64) -> Result<Option<String>> {
Ok(self
.conn
.query_row(
"SELECT detail FROM events WHERE task_id = ?1 AND kind = 'refined'
ORDER BY id DESC LIMIT 1",
[task_id],
|r| r.get::<_, Option<String>>(0),
)
.optional()?
.flatten())
}
pub fn discovered_from(&self, task_id: i64) -> Result<Option<Task>> {
let parent: Option<i64> = self
.conn
.query_row(
"SELECT depends_on FROM deps
WHERE task_id = ?1 AND kind = 'discovered-from'
ORDER BY depends_on DESC LIMIT 1",
[task_id],
|r| r.get(0),
)
.optional()?;
parent.map(|id| self.task(id)).transpose()
}
pub fn add_dep(&mut self, task_id: i64, depends_on: i64, kind: DepKind) -> Result<()> {
if kind != DepKind::Blocks && task_id == depends_on {
return Err(Error::Invalid("a task cannot depend on itself".into()));
}
let tx = self.conn.transaction()?;
if kind == DepKind::Blocks {
crate::transition::reject_blocks_cycle(&tx, task_id, depends_on)?;
}
let inserted = tx.execute(
"INSERT INTO deps (task_id, depends_on, kind) VALUES (?1, ?2, ?3)
ON CONFLICT (task_id, depends_on, kind) DO NOTHING",
params![task_id, depends_on, kind],
)?;
if inserted == 0 {
return Err(Error::Invalid(format!(
"#{task_id} already has a {kind} dependency on #{depends_on}"
)));
}
if kind == DepKind::Blocks {
crate::transition::reconcile_readiness(&tx, task_id)?;
}
tx.commit()?;
Ok(())
}
pub fn remove_dep(&mut self, task_id: i64, depends_on: i64, kind: DepKind) -> Result<()> {
let tx = self.conn.transaction()?;
let removed = tx.execute(
"DELETE FROM deps WHERE task_id = ?1 AND depends_on = ?2 AND kind = ?3",
params![task_id, depends_on, kind],
)?;
if removed == 0 {
return Err(Error::Invalid(format!(
"#{task_id} has no {kind} dependency on #{depends_on}"
)));
}
if kind == DepKind::Blocks {
crate::transition::reconcile_readiness(&tx, task_id)?;
}
tx.commit()?;
Ok(())
}
pub fn deps_by_task(&self) -> Result<HashMap<i64, Vec<DepRef>>> {
self.dep_refs(
"SELECT d.task_id, t.id, t.title, t.state, d.kind
FROM deps d JOIN tasks t ON t.id = d.depends_on
ORDER BY d.task_id, t.id, d.kind",
)
}
pub fn dependents_by_task(&self) -> Result<HashMap<i64, Vec<DepRef>>> {
self.dep_refs(
"SELECT d.depends_on, t.id, t.title, t.state, d.kind
FROM deps d JOIN tasks t ON t.id = d.task_id
ORDER BY d.depends_on, t.id, d.kind",
)
}
fn dep_refs(&self, sql: &str) -> Result<HashMap<i64, Vec<DepRef>>> {
let mut stmt = self.conn.prepare(sql)?;
let rows = stmt.query_map([], |row| {
let key: i64 = row.get(0)?;
let dep = DepRef {
id: row.get(1)?,
title: row.get(2)?,
state: row.get(3)?,
kind: row.get(4)?,
};
Ok((key, dep))
})?;
let mut map: HashMap<i64, Vec<DepRef>> = HashMap::new();
for row in rows {
let (key, dep) = row?;
map.entry(key).or_default().push(dep);
}
Ok(map)
}
pub fn deps_of(&self, task_id: i64) -> Result<Vec<Dep>> {
let mut stmt = self.conn.prepare(
"SELECT task_id, depends_on, kind FROM deps WHERE task_id = ?1
ORDER BY depends_on, kind",
)?;
let rows = stmt.query_map([task_id], |row| {
Ok(Dep {
task_id: row.get(0)?,
depends_on: row.get(1)?,
kind: row.get(2)?,
})
})?;
Ok(rows.collect::<rusqlite::Result<_>>()?)
}
pub fn create_session(
&mut self,
task_id: i64,
agent: &str,
pid: Option<i64>,
liveness_source: LivenessSource,
log_path: Option<&str>,
) -> Result<Session> {
let id = insert_session(&self.conn, task_id, agent, pid, liveness_source, log_path)?;
self.session(id)
}
pub fn set_session_ref(&mut self, id: i64, session_ref: &str) -> Result<Session> {
let changed = self.conn.execute(
"UPDATE sessions SET session_ref = ?1 WHERE id = ?2",
params![session_ref, id],
)?;
if changed == 0 {
return Err(Error::SessionNotFound(id));
}
self.session(id)
}
pub fn record_session_send(
&mut self,
id: i64,
session_ref: Option<&str>,
pid: i64,
) -> Result<Session> {
let changed = self.conn.execute(
"UPDATE sessions SET pid = ?1, session_ref = COALESCE(?2, session_ref)
WHERE id = ?3",
params![pid, session_ref, id],
)?;
if changed == 0 {
return Err(Error::SessionNotFound(id));
}
self.session(id)
}
pub fn end_session(&mut self, id: i64, outcome: SessionOutcome) -> Result<Session> {
if set_session_outcome(&self.conn, id, outcome)? == 0 {
return Err(Error::SessionNotFound(id));
}
self.session(id)
}
pub fn session(&self, id: i64) -> Result<Session> {
self.conn
.query_row(
&format!("SELECT {SESSION_COLUMNS} FROM sessions WHERE id = ?1"),
[id],
session_from_row,
)
.optional()?
.ok_or(Error::SessionNotFound(id))
}
pub fn sessions_for(&self, task_id: i64) -> Result<Vec<Session>> {
let mut stmt = self.conn.prepare(&format!(
"SELECT {SESSION_COLUMNS} FROM sessions WHERE task_id = ?1 ORDER BY id DESC"
))?;
let rows = stmt.query_map([task_id], session_from_row)?;
Ok(rows.collect::<rusqlite::Result<_>>()?)
}
pub fn latest_sessions(&self) -> Result<std::collections::HashMap<i64, Session>> {
let mut stmt = self.conn.prepare(&format!(
"SELECT {SESSION_COLUMNS} FROM sessions s
WHERE s.id = (SELECT max(id) FROM sessions WHERE task_id = s.task_id)"
))?;
let rows = stmt.query_map([], session_from_row)?;
rows.map(|r| r.map(|s| (s.task_id, s)))
.collect::<rusqlite::Result<_>>()
.map_err(Into::into)
}
pub fn live_sessions(&self) -> Result<Vec<Session>> {
let mut stmt = self.conn.prepare(&format!(
"SELECT {SESSION_COLUMNS} FROM sessions WHERE ended_at IS NULL ORDER BY id DESC"
))?;
let rows = stmt.query_map([], session_from_row)?;
Ok(rows.collect::<rusqlite::Result<_>>()?)
}
pub fn incomplete_report_flag(&self, task_id: i64) -> Result<bool> {
let row: Option<(TaskState, Option<String>)> = self
.conn
.query_row(
"SELECT state, branch FROM tasks WHERE id = ?1",
[task_id],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.optional()?;
let Some((state, branch)) = row else {
return Ok(false);
};
if state != TaskState::Review {
return Ok(false);
}
let has_branch = branch.is_some();
let has_summary = self.latest_summary(task_id)?.is_some();
Ok(has_branch && !has_summary)
}
pub fn running_rows(&self) -> Result<Vec<RunningRow>> {
let mut stmt = self.conn.prepare(
"WITH strip AS (
SELECT s.id AS session_id, t.id AS task_id, t.title, t.state,
s.agent, t.pr_url,
CASE WHEN t.state = 'waiting' THEN t.state_since
ELSE COALESCE(s.started_at, t.state_since) END AS since
FROM tasks t
JOIN projects p ON p.id = t.project_id
LEFT JOIN sessions s ON s.task_id = t.id AND s.ended_at IS NULL
WHERE t.state IN ('running','refining','waiting') AND p.archived = 0
)
SELECT session_id, task_id, title, state, agent, pr_url, since,
CAST(strftime('%s', 'now') - strftime('%s', since) AS INTEGER)
FROM strip
ORDER BY (state = 'waiting'), (session_id IS NULL),
session_id DESC, task_id DESC",
)?;
let rows = stmt.query_map([], |row| {
Ok(RunningRow {
session_id: row.get(0)?,
task_id: row.get(1)?,
task_title: row.get(2)?,
task_state: row.get(3)?,
agent: row.get(4)?,
pr_url: row.get(5)?,
started_at: row.get(6)?,
elapsed_secs: row.get(7)?,
})
})?;
Ok(rows.collect::<rusqlite::Result<_>>()?)
}
pub fn latest_summary(&self, task_id: i64) -> Result<Option<String>> {
Ok(self
.conn
.query_row(
"SELECT detail FROM events WHERE task_id = ?1 AND kind = 'summary'
ORDER BY id DESC LIMIT 1",
[task_id],
|r| r.get::<_, Option<String>>(0),
)
.optional()?
.flatten())
}
pub fn record_reviewed(&mut self, id: i64, sha: &str) -> Result<()> {
let sha = sha.trim();
if sha.is_empty() {
return Err(Error::Invalid("a reviewed revision is required".into()));
}
self.task(id)?;
log_event(&self.conn, id, crate::review::REVIEWED_EVENT, Some(sha))
}
pub fn last_reviewed(&self, task_id: i64) -> Result<Option<String>> {
Ok(self
.conn
.query_row(
"SELECT detail FROM events WHERE task_id = ?1 AND kind = ?2
ORDER BY id DESC LIMIT 1",
params![task_id, crate::review::REVIEWED_EVENT],
|r| r.get::<_, Option<String>>(0),
)
.optional()?
.flatten())
}
pub fn events_for(&self, task_id: i64) -> Result<Vec<Event>> {
let mut stmt = self.conn.prepare(
"SELECT id, task_id, at, kind, detail FROM events WHERE task_id = ?1 ORDER BY id",
)?;
let rows = stmt.query_map([task_id], |row| {
Ok(Event {
id: row.get(0)?,
task_id: row.get(1)?,
at: row.get(2)?,
kind: row.get(3)?,
detail: row.get(4)?,
})
})?;
Ok(rows.collect::<rusqlite::Result<_>>()?)
}
}
pub(crate) const TASK_COLUMNS: &str = "id, project_id, title, body, priority, state, agent, \
question, pr_url, branch, state_since, created_at, \
closed_at, human, repo_id, deep";
pub(crate) fn get_task(conn: &Connection, id: i64) -> Result<Option<Task>> {
Ok(conn
.query_row(
&format!("SELECT {TASK_COLUMNS} FROM tasks WHERE id = ?1"),
[id],
task_from_row,
)
.optional()?)
}
pub(crate) fn task_from_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<Task> {
Ok(Task {
id: row.get(0)?,
project_id: row.get(1)?,
title: row.get(2)?,
body: row.get(3)?,
priority: row.get(4)?,
state: row.get(5)?,
agent: row.get(6)?,
question: row.get(7)?,
pr_url: row.get(8)?,
branch: row.get(9)?,
state_since: row.get(10)?,
created_at: row.get(11)?,
closed_at: row.get(12)?,
human: row.get(13)?,
repo_id: row.get(14)?,
deep: row.get(15)?,
})
}
pub(crate) const SESSION_COLUMNS: &str = "id, task_id, agent, pid, session_ref, liveness_source, log_path, started_at, ended_at, outcome";
pub(crate) fn get_open_session(conn: &Connection, task_id: i64) -> Result<Option<Session>> {
Ok(conn
.query_row(
&format!(
"SELECT {SESSION_COLUMNS} FROM sessions
WHERE task_id = ?1 AND ended_at IS NULL"
),
[task_id],
session_from_row,
)
.optional()?)
}
pub(crate) fn get_session(conn: &Connection, id: i64) -> Result<Option<Session>> {
Ok(conn
.query_row(
&format!("SELECT {SESSION_COLUMNS} FROM sessions WHERE id = ?1"),
[id],
session_from_row,
)
.optional()?)
}
fn session_from_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<Session> {
Ok(Session {
id: row.get(0)?,
task_id: row.get(1)?,
agent: row.get(2)?,
pid: row.get(3)?,
session_ref: row.get(4)?,
liveness_source: row.get(5)?,
log_path: row.get(6)?,
started_at: row.get(7)?,
ended_at: row.get(8)?,
outcome: row.get(9)?,
})
}
pub(crate) const PROJECT_COLUMNS: &str = "id, name, weight, viewer, archived";
fn project_from_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<Project> {
Ok(Project {
id: row.get(0)?,
name: row.get(1)?,
weight: row.get(2)?,
viewer: row.get(3)?,
archived: row.get(4)?,
})
}
pub(crate) const DOC_COLUMNS: &str = "id, project_id, repo_id, title, location, created_at";
fn doc_from_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<Doc> {
doc_from_row_at(row, 0)
}
fn doc_from_row_at(row: &rusqlite::Row<'_>, at: usize) -> rusqlite::Result<Doc> {
Ok(Doc {
id: row.get(at)?,
project_id: row.get(at + 1)?,
repo_id: row.get(at + 2)?,
title: row.get(at + 3)?,
location: row.get(at + 4)?,
created_at: row.get(at + 5)?,
})
}
fn prefixed(columns: &str, alias: &str) -> String {
columns
.split(", ")
.map(|c| format!("{alias}.{c}"))
.collect::<Vec<_>>()
.join(", ")
}
pub(crate) const REPO_COLUMNS: &str = "id, project_id, name, path, is_default";
fn repo_from_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<Repo> {
Ok(Repo {
id: row.get(0)?,
project_id: row.get(1)?,
name: row.get(2)?,
path: row.get(3)?,
is_default: row.get(4)?,
})
}
pub(crate) fn insert_session(
conn: &Connection,
task_id: i64,
agent: &str,
pid: Option<i64>,
liveness_source: LivenessSource,
log_path: Option<&str>,
) -> Result<i64> {
close_open_session(conn, task_id, SessionOutcome::Aborted)?;
conn.execute(
"INSERT INTO sessions (task_id, agent, pid, liveness_source, log_path, started_at)
VALUES (?1, ?2, ?3, ?4, ?5, datetime('now'))",
params![task_id, agent, pid, liveness_source, log_path],
)?;
Ok(conn.last_insert_rowid())
}
pub(crate) fn close_open_session(
conn: &Connection,
task_id: i64,
outcome: SessionOutcome,
) -> Result<usize> {
Ok(conn.execute(
"UPDATE sessions SET ended_at = datetime('now'), outcome = ?1
WHERE task_id = ?2 AND ended_at IS NULL",
params![outcome, task_id],
)?)
}
pub(crate) fn set_session_outcome(
conn: &Connection,
id: i64,
outcome: SessionOutcome,
) -> Result<usize> {
Ok(conn.execute(
"UPDATE sessions SET ended_at = datetime('now'), outcome = ?1 WHERE id = ?2",
params![outcome, id],
)?)
}
pub(crate) fn log_event(
conn: &Connection,
task_id: i64,
kind: &str,
detail: Option<&str>,
) -> Result<()> {
conn.execute(
"INSERT INTO events (task_id, at, kind, detail) VALUES (?1, datetime('now'), ?2, ?3)",
params![task_id, kind, detail],
)?;
Ok(())
}
pub(crate) fn log_global_event(conn: &Connection, kind: &str, detail: Option<&str>) -> Result<()> {
conn.execute(
"INSERT INTO events (task_id, at, kind, detail)
VALUES (NULL, datetime('now'), ?1, ?2)",
params![kind, detail],
)?;
Ok(())
}
#[cfg(test)]
mod schema_guard_tests {
use super::*;
fn scratch(tag: &str) -> PathBuf {
tempfile::Builder::new()
.prefix(&format!("voro-store-{tag}-"))
.tempdir()
.unwrap()
.keep()
}
#[test]
fn a_cargo_build_directory_is_recognised_in_every_profile() {
for exe in [
"/home/u/proj/target/debug/voro",
"/home/u/proj/target/release/voro",
"/home/u/proj/target/debug/deps/voro-1a2b3c",
"/home/u/proj/.claude/worktrees/feature/target/debug/voro",
] {
assert!(
path_is_cargo_target(Path::new(exe)),
"{exe} should be a dev build"
);
}
for exe in [
"/home/u/.cargo/bin/voro",
"/usr/local/bin/voro",
"/opt/target-practice/voro",
] {
assert!(
!path_is_cargo_target(Path::new(exe)),
"{exe} should not be a dev build"
);
}
}
#[test]
fn a_database_from_the_future_is_refused_with_a_way_out() {
let dir = scratch("future");
let path = dir.join("voro.db");
Store::open(&path).unwrap();
Connection::open(&path)
.unwrap()
.pragma_update(None, "user_version", (MIGRATIONS.len() + 1) as i64)
.unwrap();
let message = match Store::open(&path) {
Ok(_) => panic!("a store from the future should not open"),
Err(e) => e.to_string(),
};
assert!(message.contains("schema version"), "{message}");
assert!(
message.contains("Restore a pre-migration snapshot"),
"{message}"
);
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn the_dev_store_is_told_to_reseed_rather_than_reinstall() {
let remedy = remedy_for_schema_ahead(Some(&Store::dev_db_path()));
assert!(remedy.contains("voro seed --force"), "{remedy}");
assert!(!remedy.contains("cargo install"), "{remedy}");
}
#[test]
fn a_migration_snapshots_the_database_beside_it_first() {
let dir = scratch("snapshot");
let path = dir.join("voro.db");
let conn = Connection::open(&path).unwrap();
conn.execute_batch(MIGRATIONS[0]).unwrap();
conn.pragma_update(None, "user_version", 1i64).unwrap();
drop(conn);
Store::open(&path).unwrap();
let backups: Vec<_> = std::fs::read_dir(Store::backup_dir_for(&path))
.unwrap()
.map(|e| e.unwrap().file_name().to_string_lossy().into_owned())
.collect();
assert_eq!(backups.len(), 1, "{backups:?}");
assert!(backups[0].starts_with("voro-v1-"), "{backups:?}");
let saved = Connection::open(Store::backup_dir_for(&path).join(&backups[0])).unwrap();
let version: i64 = saved
.query_row("PRAGMA user_version", [], |r| r.get(0))
.unwrap();
assert_eq!(version, 1);
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn a_migration_applied_from_a_different_branch_is_refused_at_the_same_version() {
let dir = scratch("diverged");
let path = dir.join("voro.db");
Store::open(&path).unwrap();
Connection::open(&path)
.unwrap()
.execute(
"UPDATE schema_migrations SET sql = ?1, applied_by = ?2 WHERE idx = ?3",
params![
"ALTER TABLE projects RENAME COLUMN review_action TO viewer;",
"voro 0.1.0 at /home/u/.claude/worktrees/project-viewer/target/debug/voro",
MIGRATIONS.len() as i64
],
)
.unwrap();
let message = match Store::open(&path) {
Ok(_) => panic!("a divergent schema should not open"),
Err(e) => e.to_string(),
};
let version: i64 = Connection::open(&path)
.unwrap()
.query_row("PRAGMA user_version", [], |r| r.get(0))
.unwrap();
assert_eq!(version, MIGRATIONS.len() as i64);
assert!(message.contains("project-viewer"), "{message}");
assert!(message.contains("Restore"), "{message}");
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn the_journal_records_what_was_applied_and_by_whom() {
let dir = scratch("journal");
let path = dir.join("voro.db");
Store::open(&path).unwrap();
let conn = Connection::open(&path).unwrap();
let rows: i64 = conn
.query_row("SELECT COUNT(*) FROM schema_migrations", [], |r| r.get(0))
.unwrap();
assert_eq!(rows, MIGRATIONS.len() as i64);
let (sql, by): (String, String) = conn
.query_row(
"SELECT sql, applied_by FROM schema_migrations WHERE idx = 1",
[],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.unwrap();
assert_eq!(sql, MIGRATIONS[0]);
assert!(by.starts_with("voro "), "{by}");
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn pre_journal_history_is_backfilled_unverifiable_and_opens_cleanly() {
let dir = scratch("backfill");
let path = dir.join("voro.db");
let conn = Connection::open(&path).unwrap();
for sql in &MIGRATIONS[..MIGRATIONS.len() - 1] {
conn.execute_batch(sql).unwrap();
}
conn.pragma_update(None, "user_version", (MIGRATIONS.len() - 1) as i64)
.unwrap();
drop(conn);
Store::open(&path).unwrap();
Store::open(&path).unwrap();
let conn = Connection::open(&path).unwrap();
let unverifiable: i64 = conn
.query_row(
"SELECT COUNT(*) FROM schema_migrations WHERE sql IS NULL",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(unverifiable, (MIGRATIONS.len() - 1) as i64);
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn a_fresh_database_is_not_snapshotted() {
let dir = scratch("fresh");
let path = dir.join("voro.db");
Store::open(&path).unwrap();
assert!(!Store::backup_dir_for(&path).exists());
std::fs::remove_dir_all(&dir).ok();
}
fn store_at_previous_version(path: &Path) {
let conn = Connection::open(path).unwrap();
for sql in &MIGRATIONS[..MIGRATIONS.len() - 1] {
conn.execute_batch(sql).unwrap();
}
conn.pragma_update(None, "user_version", (MIGRATIONS.len() - 1) as i64)
.unwrap();
}
fn open_as_production(path: &Path, consent: Option<&str>) -> Result<Store> {
Store::open_at(Connection::open(path).unwrap(), path, path, consent)
}
#[test]
fn the_production_store_refuses_to_migrate_without_consent() {
let dir = scratch("gate-refuse");
let path = dir.join("voro.db");
store_at_previous_version(&path);
let message = match open_as_production(&path, None) {
Ok(_) => panic!("a protected store with pending migrations should not open"),
Err(e) => e.to_string(),
};
assert!(message.contains("pending migration"), "{message}");
assert!(message.contains("voro migrate"), "{message}");
let version: i64 = Connection::open(&path)
.unwrap()
.query_row("PRAGMA user_version", [], |r| r.get(0))
.unwrap();
assert_eq!(version, (MIGRATIONS.len() - 1) as i64);
assert!(!Store::backup_dir_for(&path).exists());
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn consent_migrates_the_production_store_and_is_journalled() {
let dir = scratch("gate-consent");
let path = dir.join("voro.db");
store_at_previous_version(&path);
open_as_production(&path, Some("via voro migrate --yes")).unwrap();
let conn = Connection::open(&path).unwrap();
let version: i64 = conn
.query_row("PRAGMA user_version", [], |r| r.get(0))
.unwrap();
assert_eq!(version, MIGRATIONS.len() as i64);
let by: String = conn
.query_row(
"SELECT applied_by FROM schema_migrations WHERE idx = ?1",
[MIGRATIONS.len() as i64],
|r| r.get(0),
)
.unwrap();
assert!(by.contains("via voro migrate --yes"), "{by}");
assert!(Store::backup_dir_for(&path).exists());
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn the_protected_marker_travels_with_the_file() {
let dir = scratch("gate-marker");
let path = dir.join("voro.db");
open_as_production(&path, None).unwrap();
let moved = dir.join("restored-copy.db");
std::fs::rename(&path, &moved).unwrap();
Connection::open(&moved)
.unwrap()
.pragma_update(None, "user_version", (MIGRATIONS.len() - 1) as i64)
.unwrap();
assert!(matches!(
Store::open(&moved),
Err(Error::MigrationsPending { .. })
));
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn a_fresh_production_store_is_created_silently_and_marked() {
let dir = scratch("gate-fresh");
let path = dir.join("voro.db");
let store = open_as_production(&path, None).unwrap();
let marked: String = store
.conn
.query_row(
"SELECT value FROM store_meta WHERE key = 'protected'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(marked, "1");
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn an_unprotected_store_still_migrates_silently() {
let dir = scratch("gate-scratch");
let path = dir.join("scratch.db");
store_at_previous_version(&path);
let store = Store::open(&path).unwrap();
assert_eq!(store.schema_version().unwrap(), MIGRATIONS.len());
std::fs::remove_dir_all(&dir).ok();
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::transition::{Action, Triage};
fn new_ready(project_id: i64) -> NewTask {
NewTask {
project_id,
repo_id: None,
title: "t".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
}
}
#[test]
fn rename_project_updates_name_and_leaves_task_references_intact() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("old-name", "/tmp/old").unwrap();
let task = s
.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "t".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
})
.unwrap();
let renamed = s.rename_project(p.id, "new-name").unwrap();
assert_eq!(renamed.id, p.id);
assert_eq!(renamed.name, "new-name");
let reloaded = s.task(task.id).unwrap();
assert_eq!(reloaded.project_id, p.id);
assert_eq!(s.project(reloaded.project_id).unwrap().name, "new-name");
}
#[test]
fn project_viewer_defaults_to_none_and_round_trips() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("proj", "/tmp/proj").unwrap();
assert_eq!(p.viewer, None);
let updated = s.set_viewer(p.id, Some("zed")).unwrap();
assert_eq!(updated.viewer.as_deref(), Some("zed"));
assert_eq!(s.project(p.id).unwrap().viewer.as_deref(), Some("zed"));
assert_eq!(s.projects().unwrap()[0].viewer.as_deref(), Some("zed"));
s.set_viewer(p.id, None).unwrap();
assert_eq!(s.project(p.id).unwrap().viewer, None);
let raw: Option<String> = s
.conn
.query_row("SELECT viewer FROM projects WHERE id = ?1", [p.id], |r| {
r.get(0)
})
.unwrap();
assert_eq!(raw, None);
assert!(matches!(
s.set_viewer(p.id, Some(" ")),
Err(Error::Invalid(_))
));
assert!(matches!(
s.set_viewer(999, Some("zed")),
Err(Error::ProjectNotFound(999))
));
}
#[test]
fn migration_0018_reads_review_actions_as_viewer_names() {
let conn = Connection::open_in_memory().unwrap();
for sql in &MIGRATIONS[..17] {
conn.execute_batch(sql).unwrap();
}
conn.pragma_update(None, "user_version", 17).unwrap();
conn.execute(
"INSERT INTO projects (name, review_action) VALUES
('named', 'viewer:zed'),
('bare-viewer', 'viewer'),
('auto', 'auto'),
('pinned-to-pr', 'pr'),
('unset', NULL)",
[],
)
.unwrap();
let mut store = Store::from_connection(conn).unwrap();
let viewer_of = |store: &mut Store, name: &str| {
store
.projects()
.unwrap()
.into_iter()
.find(|p| p.name == name)
.unwrap()
.viewer
};
assert_eq!(viewer_of(&mut store, "named").as_deref(), Some("zed"));
for named_none in ["bare-viewer", "auto", "pinned-to-pr", "unset"] {
assert_eq!(viewer_of(&mut store, named_none), None, "{named_none}");
}
}
#[test]
fn rename_project_rejects_unknown_id() {
let mut s = Store::open_in_memory().unwrap();
assert!(matches!(
s.rename_project(999, "x"),
Err(Error::ProjectNotFound(999))
));
}
#[test]
fn set_pr_tracks_clears_and_logs() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let t = s
.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "review me".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
})
.unwrap();
assert!(s.task(t.id).unwrap().pr_url.is_none());
let tracked = s
.set_pr(t.id, Some("https://github.com/acme/widget/pull/42"))
.unwrap();
assert_eq!(
tracked.pr_url.as_deref(),
Some("https://github.com/acme/widget/pull/42")
);
assert_eq!(tracked.state, TaskState::Ready);
let cleared = s.set_pr(t.id, None).unwrap();
assert!(cleared.pr_url.is_none());
let events = s.events_for(t.id).unwrap();
let kinds: Vec<&str> = events.iter().map(|e| e.kind.as_str()).collect();
assert_eq!(kinds, vec!["created", "pr", "pr"]);
assert!(matches!(s.set_pr(999, None), Err(Error::TaskNotFound(999))));
}
#[test]
fn set_priority_updates_leaves_state_and_logs() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let t = s
.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "reprioritise me".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
})
.unwrap();
let raised = s.set_priority(t.id, Priority::P0).unwrap();
assert_eq!(raised.priority, Priority::P0);
assert_eq!(raised.state, TaskState::Ready);
let events = s.events_for(t.id).unwrap();
let last = events.last().unwrap();
assert_eq!(last.kind, "priority");
assert_eq!(last.detail.as_deref(), Some("P0"));
assert!(matches!(
s.set_priority(999, Priority::P1),
Err(Error::TaskNotFound(999))
));
}
#[test]
fn set_branch_records_clears_and_logs() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let t = s
.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "branch me".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
})
.unwrap();
assert!(s.task(t.id).unwrap().branch.is_none());
let named = s.set_branch(t.id, Some("feat/parser")).unwrap();
assert_eq!(named.branch.as_deref(), Some("feat/parser"));
assert_eq!(named.state, TaskState::Ready);
let renamed = s.set_branch(t.id, Some("feat/parser-v2")).unwrap();
assert_eq!(renamed.branch.as_deref(), Some("feat/parser-v2"));
let cleared = s.set_branch(t.id, None).unwrap();
assert!(cleared.branch.is_none());
let events = s.events_for(t.id).unwrap();
let kinds: Vec<&str> = events.iter().map(|e| e.kind.as_str()).collect();
assert_eq!(kinds, vec!["created", "branch", "branch", "branch"]);
assert!(matches!(
s.set_branch(999, None),
Err(Error::TaskNotFound(999))
));
}
fn human_fixture() -> (Store, i64) {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
(s, p.id)
}
fn new_with(project_id: i64, agent: Option<&str>, human: bool) -> NewTask {
NewTask {
project_id,
repo_id: None,
title: "hands-on".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: agent.map(str::to_string),
human,
deep: false,
}
}
fn edit_of(task: &Task, agent: Option<&str>, human: bool) -> TaskEdit {
TaskEdit {
title: task.title.clone(),
body: task.body.clone(),
priority: task.priority,
agent: agent.map(str::to_string),
human,
deep: false,
}
}
#[test]
fn create_task_refuses_a_human_task_with_an_agent_override() {
let (mut s, p) = human_fixture();
let err = s.create_task(new_with(p, Some("codex"), true)).unwrap_err();
assert!(err.to_string().contains("agent override"), "{err}");
assert!(s.tasks().unwrap().is_empty());
assert!(s.create_task(new_with(p, Some("codex"), false)).is_ok());
let human = s.create_task(new_with(p, None, true)).unwrap();
assert!(human.human);
}
#[test]
fn update_task_logs_the_body_it_replaced_and_nothing_else() {
let (mut s, p) = human_fixture();
let task = s.create_task(new_with(p, None, false)).unwrap();
let write = TaskEdit {
body: "the brief".into(),
..edit_of(&task, None, false)
};
let task = s.update_task(task.id, write).unwrap();
assert!(
!s.events_for(task.id)
.unwrap()
.iter()
.any(|e| e.kind == "body")
);
let retitle = TaskEdit {
title: "renamed".into(),
..edit_of(&task, None, false)
};
let task = s.update_task(task.id, retitle).unwrap();
assert!(
!s.events_for(task.id)
.unwrap()
.iter()
.any(|e| e.kind == "body")
);
let rewrite = TaskEdit {
body: "a rewrite".into(),
..edit_of(&task, None, false)
};
let task = s.update_task(task.id, rewrite).unwrap();
assert_eq!(task.body, "a rewrite");
let logged: Vec<String> = s
.events_for(task.id)
.unwrap()
.into_iter()
.filter(|e| e.kind == "body")
.map(|e| e.detail.unwrap_or_default())
.collect();
assert_eq!(logged, vec!["the brief".to_string()]);
}
#[test]
fn update_task_guards_the_agent_human_exclusivity_both_ways() {
let (mut s, p) = human_fixture();
let human = s.create_task(new_with(p, None, true)).unwrap();
let err = s
.update_task(human.id, edit_of(&human, Some("codex"), true))
.unwrap_err();
assert!(matches!(err, Error::HumanTask { id, .. } if id == human.id));
let agented = s.create_task(new_with(p, Some("codex"), false)).unwrap();
let err = s
.update_task(agented.id, edit_of(&agented, Some("codex"), true))
.unwrap_err();
assert!(matches!(err, Error::HumanTask { id, .. } if id == agented.id));
let flipped = s
.update_task(agented.id, edit_of(&agented, None, true))
.unwrap();
assert!(flipped.human);
assert!(flipped.agent.is_none());
}
#[test]
fn update_task_refuses_flagging_human_in_agent_executed_states() {
use crate::transition::Action;
for walk in [TaskState::NeedsInput, TaskState::Review, TaskState::Stalled] {
let (mut s, p) = human_fixture();
let t = s.create_task(new_with(p, None, false)).unwrap();
match walk {
TaskState::NeedsInput => {
s.apply(t.id, Action::Start).unwrap();
s.apply(t.id, Action::Ask("A or B?".into())).unwrap();
}
TaskState::Stalled => {
let (_, session) = s
.record_dispatch(t.id, "claude", Some(1), LivenessSource::Pid, None)
.unwrap();
s.reconcile_session(session.id, false, false).unwrap();
}
_ => {
s.apply(t.id, Action::Start).unwrap();
s.apply(t.id, Action::Complete(None)).unwrap();
}
}
assert_eq!(s.task(t.id).unwrap().state, walk);
let err = s.update_task(t.id, edit_of(&t, None, true)).unwrap_err();
assert!(
matches!(err, Error::HumanTask { id, .. } if id == t.id),
"{walk}: {err}"
);
assert!(!s.task(t.id).unwrap().human);
}
}
#[test]
fn update_task_refuses_flagging_human_while_a_session_is_open() {
use crate::transition::Action;
let (mut s, p) = human_fixture();
let t = s.create_task(new_with(p, None, false)).unwrap();
s.record_dispatch(t.id, "claude", Some(1), LivenessSource::Pid, None)
.unwrap();
let err = s.update_task(t.id, edit_of(&t, None, true)).unwrap_err();
assert!(matches!(err, Error::HumanTask { id, .. } if id == t.id));
s.apply(t.id, Action::Abort).unwrap();
let flipped = s.update_task(t.id, edit_of(&t, None, true)).unwrap();
assert!(flipped.human);
let by_hand = s.create_task(new_with(p, None, false)).unwrap();
s.apply(by_hand.id, Action::Start).unwrap();
assert!(
s.update_task(by_hand.id, edit_of(&by_hand, None, true))
.unwrap()
.human
);
}
fn deep_new(project_id: i64, human: bool, deep: bool) -> NewTask {
NewTask {
deep,
..new_with(project_id, None, human)
}
}
#[test]
fn set_deep_toggles_the_flag_and_logs_it() {
let (mut s, p) = human_fixture();
let t = s.create_task(deep_new(p, false, false)).unwrap();
assert!(!t.deep);
assert!(s.set_deep(t.id, true).unwrap().deep);
assert!(!s.set_deep(t.id, false).unwrap().deep);
let kinds: Vec<String> = s
.events_for(t.id)
.unwrap()
.into_iter()
.map(|e| e.kind)
.collect();
assert_eq!(kinds, vec!["created", "deep", "deep"]);
assert!(matches!(
s.set_deep(999, true),
Err(Error::TaskNotFound(999))
));
}
#[test]
fn a_human_task_cannot_be_deep() {
let (mut s, p) = human_fixture();
let err = s.create_task(deep_new(p, true, true)).unwrap_err();
assert!(err.to_string().contains("deep"), "{err}");
assert!(s.tasks().unwrap().is_empty());
let human = s.create_task(deep_new(p, true, false)).unwrap();
let err = s.set_deep(human.id, true).unwrap_err();
assert!(matches!(err, Error::HumanTask { id, .. } if id == human.id));
assert!(!s.task(human.id).unwrap().deep);
let edit = TaskEdit {
deep: true,
..edit_of(&human, None, true)
};
let err = s.update_task(human.id, edit).unwrap_err();
assert!(matches!(err, Error::HumanTask { id, .. } if id == human.id));
assert!(!s.set_deep(human.id, false).unwrap().deep);
}
#[test]
fn deep_defaults_off_and_is_constrained() {
let (mut s, p) = human_fixture();
let t = s.create_task(deep_new(p, false, false)).unwrap();
assert!(!t.deep);
assert!(
s.conn
.execute("UPDATE tasks SET deep = 2 WHERE id = ?1", [t.id])
.is_err()
);
}
#[test]
fn migration_0007_defaults_existing_tasks_to_dispatchable() {
let conn = Connection::open_in_memory().unwrap();
for sql in &MIGRATIONS[..6] {
conn.execute_batch(sql).unwrap();
}
conn.pragma_update(None, "user_version", 6).unwrap();
conn.execute("INSERT INTO projects (name, path) VALUES ('p', '/tmp')", [])
.unwrap();
conn.execute(
"INSERT INTO tasks (project_id, title, state, state_since, created_at)
VALUES (1, 'pre-flag', 'ready', datetime('now'), datetime('now'))",
[],
)
.unwrap();
let store = Store::from_connection(conn).unwrap();
assert!(!store.task(1).unwrap().human);
let junk = store
.conn
.execute("UPDATE tasks SET human = 2 WHERE id = 1", []);
assert!(junk.is_err(), "the CHECK must reject values outside 0/1");
}
#[test]
fn set_summary_appends_a_superseding_summary_event() {
use crate::transition::Action;
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let t = s
.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "summarise me".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
})
.unwrap();
s.apply(t.id, Action::Start).unwrap();
let updated = s.set_summary(t.id, " early account ").unwrap();
assert_eq!(updated.state, TaskState::Running);
assert_eq!(
s.latest_summary(t.id).unwrap().as_deref(),
Some("early account")
);
s.apply(t.id, Action::Complete(Some("done-time".into())))
.unwrap();
let updated = s.set_summary(t.id, "amended for the PR body").unwrap();
assert_eq!(updated.state, TaskState::Review);
assert_eq!(
s.latest_summary(t.id).unwrap().as_deref(),
Some("amended for the PR body")
);
let events = s.events_for(t.id).unwrap();
let summaries = events.iter().filter(|e| e.kind == "summary").count();
assert_eq!(summaries, 3);
}
#[test]
fn set_summary_clears_the_incomplete_report_flag() {
use crate::transition::Action;
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let t = s
.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "half a report".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
})
.unwrap();
s.apply(t.id, Action::Start).unwrap();
s.apply(t.id, Action::Complete(None)).unwrap();
s.set_branch(t.id, Some("feat/x")).unwrap();
assert!(s.incomplete_report_flag(t.id).unwrap());
s.set_summary(t.id, "the missing half").unwrap();
assert!(!s.incomplete_report_flag(t.id).unwrap());
}
#[test]
fn set_summary_is_refused_outside_running_and_review() {
use crate::transition::Action;
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let t = s
.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "not yet".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
})
.unwrap();
let err = s.set_summary(t.id, "too early").unwrap_err();
assert!(err.to_string().contains("ready"), "{err}");
s.apply(t.id, Action::Start).unwrap();
s.apply(t.id, Action::Complete(None)).unwrap();
s.apply(t.id, Action::Accept).unwrap();
let err = s.set_summary(t.id, "too late").unwrap_err();
assert!(err.to_string().contains("done"), "{err}");
assert!(s.set_summary(t.id, " ").is_err());
assert!(matches!(
s.set_summary(999, "x"),
Err(Error::TaskNotFound(999))
));
}
#[test]
fn creating_a_project_creates_its_default_repo() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let repos = s.repos(p.id).unwrap();
assert_eq!(repos.len(), 1);
assert_eq!(repos[0].name, "voro");
assert_eq!(repos[0].path, "/tmp/voro");
assert!(repos[0].is_default);
assert_eq!(s.default_repo(p.id).unwrap().id, repos[0].id);
}
#[test]
fn added_repos_are_not_default_and_names_are_unique_per_project() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("odm", "/tmp/odm").unwrap();
let oats = s.add_repo(p.id, "oats", "/tmp/oats").unwrap();
assert!(!oats.is_default);
assert!(s.add_repo(p.id, "oats", "/tmp/elsewhere").is_err());
assert!(s.add_repo(p.id, " ", "/tmp/blank").is_err());
let other = s.create_project("voro", "/tmp/voro").unwrap();
assert!(s.add_repo(other.id, "oats", "/tmp/oats").is_ok());
let names: Vec<_> = s.repos(p.id).unwrap().into_iter().map(|r| r.name).collect();
assert_eq!(names, vec!["odm", "oats"]);
}
#[test]
fn only_one_repo_is_ever_default() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("odm", "/tmp/odm").unwrap();
let oats = s.add_repo(p.id, "oats", "/tmp/oats").unwrap();
let promoted = s.set_default_repo(oats.id).unwrap();
assert!(promoted.is_default);
let defaults = s
.repos(p.id)
.unwrap()
.into_iter()
.filter(|r| r.is_default)
.count();
assert_eq!(defaults, 1);
assert_eq!(s.default_repo(p.id).unwrap().name, "oats");
}
#[test]
fn a_projects_last_repo_cannot_be_deleted() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("odm", "/tmp/odm").unwrap();
let only = s.default_repo(p.id).unwrap();
assert!(matches!(
s.delete_repo(only.id),
Err(Error::LastRepo { .. })
));
}
#[test]
fn the_default_repo_cannot_be_deleted_while_others_remain() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("odm", "/tmp/odm").unwrap();
let oats = s.add_repo(p.id, "oats", "/tmp/oats").unwrap();
let default = s.default_repo(p.id).unwrap();
assert!(matches!(
s.delete_repo(default.id),
Err(Error::DefaultRepo { .. })
));
s.set_default_repo(oats.id).unwrap();
s.delete_repo(default.id).unwrap();
assert_eq!(s.repos(p.id).unwrap().len(), 1);
}
#[test]
fn a_repo_a_task_names_cannot_be_deleted() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("odm", "/tmp/odm").unwrap();
let oats = s.add_repo(p.id, "oats", "/tmp/oats").unwrap();
let mut new = new_ready(p.id);
new.repo_id = Some(oats.id);
let t = s.create_task(new).unwrap();
assert!(matches!(
s.delete_repo(oats.id),
Err(Error::RepoInUse { count: 1, .. })
));
s.set_task_repo(t.id, None).unwrap();
s.delete_repo(oats.id).unwrap();
}
#[test]
fn a_task_resolves_its_own_repo_then_the_project_default() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("odm", "/tmp/odm").unwrap();
let oats = s.add_repo(p.id, "oats", "/tmp/oats").unwrap();
let plain = s.create_task(new_ready(p.id)).unwrap();
assert_eq!(s.repo_for_task(&plain).unwrap().path, "/tmp/odm");
let pointed = s.set_task_repo(plain.id, Some(oats.id)).unwrap();
assert_eq!(pointed.repo_id, Some(oats.id));
assert_eq!(s.repo_for_task(&pointed).unwrap().path, "/tmp/oats");
let back = s.set_task_repo(plain.id, None).unwrap();
s.set_default_repo(oats.id).unwrap();
assert_eq!(s.repo_for_task(&back).unwrap().path, "/tmp/oats");
}
#[test]
fn a_task_cannot_name_another_projects_repo() {
let mut s = Store::open_in_memory().unwrap();
let odm = s.create_project("odm", "/tmp/odm").unwrap();
let voro = s.create_project("voro", "/tmp/voro").unwrap();
let foreign = s.default_repo(voro.id).unwrap();
let mut new = new_ready(odm.id);
new.repo_id = Some(foreign.id);
assert!(s.create_task(new).is_err());
let t = s.create_task(new_ready(odm.id)).unwrap();
assert!(s.set_task_repo(t.id, Some(foreign.id)).is_err());
}
#[test]
fn an_unknown_repo_name_errors_listing_the_projects_repos() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("odm", "/tmp/odm").unwrap();
s.add_repo(p.id, "oats", "/tmp/oats").unwrap();
let err = s.repo_by_name(p.id, "nope").unwrap_err().to_string();
assert!(err.contains("odm"), "{err}");
assert!(err.contains("oats"), "{err}");
}
#[test]
fn deleting_a_project_takes_its_repos_with_it() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("odm", "/tmp/odm").unwrap();
s.add_repo(p.id, "oats", "/tmp/oats").unwrap();
s.delete_project(p.id).unwrap();
assert!(s.repos(p.id).unwrap().is_empty());
}
#[test]
fn set_path_updates_the_default_repo() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("proj", "/tmp/old").unwrap();
let updated = s.set_default_repo_path(p.id, "/tmp/new").unwrap();
assert_eq!(updated.path, "/tmp/new");
assert!(updated.is_default);
assert_eq!(s.default_repo(p.id).unwrap().path, "/tmp/new");
}
#[test]
fn set_path_rejects_unknown_id() {
let mut s = Store::open_in_memory().unwrap();
assert!(matches!(
s.set_default_repo_path(999, "/tmp"),
Err(Error::ProjectNotFound(999))
));
}
#[test]
fn delete_project_removes_a_taskless_project() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("empty", "/tmp/empty").unwrap();
s.delete_project(p.id).unwrap();
assert!(matches!(s.project(p.id), Err(Error::ProjectNotFound(_))));
assert!(s.projects().unwrap().is_empty());
}
#[test]
fn delete_project_rejects_unknown_id() {
let mut s = Store::open_in_memory().unwrap();
assert!(matches!(
s.delete_project(999),
Err(Error::ProjectNotFound(999))
));
}
fn task_in_state(s: &mut Store, project_id: i64, state: TaskState) -> i64 {
use TaskState::*;
let create = |s: &mut Store, state| {
s.create_task(NewTask {
project_id,
repo_id: None,
title: format!("task in {state}"),
body: String::new(),
priority: Priority::P1,
state,
agent: None,
human: false,
deep: false,
})
.unwrap()
.id
};
match state {
Proposed | Parked | Ready => create(s, state),
Refining => {
let id = create(s, Proposed);
s.record_refine_launch(
id,
"thin body",
"claude",
Some(1),
LivenessSource::Pid,
None,
)
.unwrap();
id
}
Running => {
let id = create(s, Ready);
s.apply(id, Action::Start).unwrap();
id
}
NeedsInput => {
let id = task_in_state(s, project_id, Running);
s.apply(id, Action::Ask("which schema?".into())).unwrap();
id
}
Review => {
let id = task_in_state(s, project_id, Running);
s.apply(id, Action::Complete(None)).unwrap();
id
}
Waiting => {
let id = task_in_state(s, project_id, Review);
s.apply(id, Action::HandOff).unwrap();
id
}
Stalled => {
let id = create(s, Ready);
let (_, session) = s
.record_dispatch(id, "claude", Some(1), LivenessSource::Pid, None)
.unwrap();
s.reconcile_session(session.id, false, false).unwrap();
id
}
Done => {
let id = task_in_state(s, project_id, Review);
s.apply(id, Action::Accept).unwrap();
id
}
Rejected => {
let id = create(s, Proposed);
s.apply(id, Action::Triage(Triage::Reject)).unwrap();
id
}
}
}
#[test]
fn delete_project_refuses_with_a_task_in_any_state() {
for state in TaskState::ALL {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("proj", "/tmp/proj").unwrap();
task_in_state(&mut s, p.id, state);
let err = s.delete_project(p.id).unwrap_err();
assert!(
matches!(err, Error::ProjectHasTasks { id, count } if id == p.id && count == 1),
"state {state}: expected ProjectHasTasks, got {err}"
);
assert!(s.project(p.id).is_ok());
}
}
#[test]
fn set_archived_round_trips_and_refuses_noops() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("retiring", "/tmp/retiring").unwrap();
assert!(!p.archived);
let archived = s.set_archived(p.id, true).unwrap();
assert!(archived.archived);
assert!(s.projects().unwrap()[0].archived);
let err = s.set_archived(p.id, true).unwrap_err();
assert!(err.to_string().contains("already archived"), "{err}");
let restored = s.set_archived(p.id, false).unwrap();
assert!(!restored.archived);
let err = s.set_archived(p.id, false).unwrap_err();
assert!(err.to_string().contains("not archived"), "{err}");
assert!(matches!(
s.set_archived(999, true),
Err(Error::ProjectNotFound(999))
));
}
#[test]
fn create_task_refuses_an_archived_project() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("retired", "/tmp/retired").unwrap();
s.set_archived(p.id, true).unwrap();
for state in [TaskState::Proposed, TaskState::Parked, TaskState::Ready] {
let err = s
.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "too late".into(),
body: String::new(),
priority: Priority::P2,
state,
agent: None,
human: false,
deep: false,
})
.unwrap_err();
assert!(
matches!(&err, Error::ProjectArchived { name } if name == "retired"),
"{state}: {err}"
);
}
assert!(s.tasks().unwrap().is_empty());
s.set_archived(p.id, false).unwrap();
assert!(
s.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "welcome back".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
})
.is_ok()
);
}
#[test]
fn archive_freezes_task_states_and_history_and_unarchive_restores_them() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("retiring", "/tmp/retiring").unwrap();
let tasks: Vec<i64> = TaskState::ALL
.iter()
.map(|state| task_in_state(&mut s, p.id, *state))
.collect();
let before: Vec<Task> = tasks.iter().map(|id| s.task(*id).unwrap()).collect();
let events_before: Vec<usize> = tasks
.iter()
.map(|id| s.events_for(*id).unwrap().len())
.collect();
s.set_archived(p.id, true).unwrap();
let frozen: Vec<Task> = tasks.iter().map(|id| s.task(*id).unwrap()).collect();
assert_eq!(frozen, before);
s.set_archived(p.id, false).unwrap();
let after: Vec<Task> = tasks.iter().map(|id| s.task(*id).unwrap()).collect();
assert_eq!(after, before);
let events_after: Vec<usize> = tasks
.iter()
.map(|id| s.events_for(*id).unwrap().len())
.collect();
assert_eq!(events_after, events_before);
}
#[test]
fn running_rows_exclude_archived_projects() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("retiring", "/tmp/retiring").unwrap();
let id = task_in_state(&mut s, p.id, TaskState::Running);
assert_eq!(s.running_rows().unwrap().len(), 1);
s.set_archived(p.id, true).unwrap();
assert!(s.running_rows().unwrap().is_empty());
assert_eq!(s.task(id).unwrap().state, TaskState::Running);
s.set_archived(p.id, false).unwrap();
assert_eq!(s.running_rows().unwrap()[0].task_id, id);
}
#[test]
fn migration_0011_defaults_existing_projects_to_active() {
let conn = Connection::open_in_memory().unwrap();
for sql in &MIGRATIONS[..10] {
conn.execute_batch(sql).unwrap();
}
conn.pragma_update(None, "user_version", 10).unwrap();
conn.execute("INSERT INTO projects (name, path) VALUES ('p', '/tmp')", [])
.unwrap();
let store = Store::from_connection(conn).unwrap();
assert!(!store.project(1).unwrap().archived);
let junk = store
.conn
.execute("UPDATE projects SET archived = 2 WHERE id = 1", []);
assert!(junk.is_err(), "the CHECK must reject values outside 0/1");
}
#[test]
fn migration_0012_turns_each_project_path_into_its_default_repo() {
let conn = Connection::open_in_memory().unwrap();
for sql in &MIGRATIONS[..11] {
conn.execute_batch(sql).unwrap();
}
conn.pragma_update(None, "user_version", 11).unwrap();
conn.execute(
"INSERT INTO projects (name, path) VALUES ('odm', '/tmp/odm'), ('voro', '/tmp/voro')",
[],
)
.unwrap();
conn.execute(
"INSERT INTO tasks (project_id, title, state, state_since, created_at)
VALUES (1, 'old', 'ready', datetime('now'), datetime('now'))",
[],
)
.unwrap();
let store = Store::from_connection(conn).unwrap();
for (id, name, path) in [(1, "odm", "/tmp/odm"), (2, "voro", "/tmp/voro")] {
let repos = store.repos(id).unwrap();
assert_eq!(repos.len(), 1);
assert_eq!(repos[0].name, name);
assert_eq!(repos[0].path, path);
assert!(repos[0].is_default);
}
let task = store.task(1).unwrap();
assert_eq!(task.repo_id, None);
assert_eq!(store.repo_for_task(&task).unwrap().path, "/tmp/odm");
assert!(
store
.conn
.query_row("SELECT path FROM projects WHERE id = 1", [], |r| r
.get::<_, String>(0))
.is_err()
);
assert!(
store
.conn
.execute(
"INSERT INTO repos (project_id, name, path, is_default)
VALUES (1, 'second', '/tmp/second', 1)",
[],
)
.is_err()
);
}
#[test]
fn migration_0015_widens_the_dep_key_without_losing_edges() {
let conn = Connection::open_in_memory().unwrap();
for sql in &MIGRATIONS[..14] {
conn.execute_batch(sql).unwrap();
}
conn.pragma_update(None, "user_version", 14).unwrap();
conn.execute("INSERT INTO projects (name) VALUES ('voro')", [])
.unwrap();
conn.execute(
"INSERT INTO repos (project_id, name, path, is_default)
VALUES (1, 'voro', '/tmp/voro', 1)",
[],
)
.unwrap();
conn.execute(
"INSERT INTO tasks (project_id, title, state, state_since, created_at)
VALUES (1, 'source', 'ready', datetime('now'), datetime('now')),
(1, 'spawned', 'ready', datetime('now'), datetime('now'))",
[],
)
.unwrap();
conn.execute(
"INSERT INTO deps (task_id, depends_on, kind) VALUES (2, 1, 'discovered-from')",
[],
)
.unwrap();
let mut store = Store::from_connection(conn).unwrap();
let carried = store.deps_of(2).unwrap();
assert_eq!(carried.len(), 1);
assert_eq!(carried[0].kind, DepKind::DiscoveredFrom);
store.set_blocks_deps(2, &[1]).unwrap();
let kinds: Vec<DepKind> = store.deps_of(2).unwrap().iter().map(|d| d.kind).collect();
assert_eq!(kinds, vec![DepKind::Blocks, DepKind::DiscoveredFrom]);
}
#[test]
fn migration_0002_converts_backlog_rows() {
let conn = Connection::open_in_memory().unwrap();
conn.execute_batch(MIGRATIONS[0]).unwrap();
conn.pragma_update(None, "user_version", 1).unwrap();
conn.execute("INSERT INTO projects (name, path) VALUES ('p', '/tmp')", [])
.unwrap();
conn.execute(
"INSERT INTO tasks (project_id, title, state, state_since, created_at)
VALUES (1, 'blocker', 'ready', datetime('now'), datetime('now')),
(1, 'waiting', 'backlog', datetime('now'), datetime('now'))",
[],
)
.unwrap();
conn.execute("INSERT INTO deps (task_id, depends_on) VALUES (2, 1)", [])
.unwrap();
conn.execute(
"INSERT INTO events (task_id, at, kind, detail)
VALUES (2, datetime('now'), 'created', 'backlog')",
[],
)
.unwrap();
let store = Store::from_connection(conn).unwrap();
assert_eq!(store.task(2).unwrap().state, TaskState::Parked);
assert_eq!(store.task(1).unwrap().state, TaskState::Ready);
assert_eq!(store.deps_of(2).unwrap().len(), 1);
assert_eq!(
store.events_for(2).unwrap()[0].detail.as_deref(),
Some("backlog")
);
let version: i64 = store
.conn
.query_row("PRAGMA user_version", [], |r| r.get(0))
.unwrap();
assert_eq!(version, MIGRATIONS.len() as i64);
let refs: i64 = store
.conn
.query_row("SELECT COUNT(session_ref) FROM sessions", [], |r| r.get(0))
.unwrap();
assert_eq!(refs, 0);
}
#[test]
fn migration_0006_dedupes_open_sessions_and_enforces_the_index() {
let conn = Connection::open_in_memory().unwrap();
for sql in &MIGRATIONS[..5] {
conn.execute_batch(sql).unwrap();
}
conn.pragma_update(None, "user_version", 5).unwrap();
conn.execute("INSERT INTO projects (name, path) VALUES ('p', '/tmp')", [])
.unwrap();
conn.execute(
"INSERT INTO tasks (project_id, title, state, state_since, created_at)
VALUES (1, 'run me', 'running', datetime('now'), datetime('now'))",
[],
)
.unwrap();
for _ in 0..3 {
conn.execute(
"INSERT INTO sessions (task_id, agent, started_at) VALUES (1, 'a', datetime('now'))",
[],
)
.unwrap();
}
let store = Store::from_connection(conn).unwrap();
let open: Vec<i64> = store
.sessions_for(1)
.unwrap()
.into_iter()
.filter(|s| s.ended_at.is_none())
.map(|s| s.id)
.collect();
assert_eq!(open, vec![3]);
assert_eq!(
store.session(1).unwrap().outcome,
Some(SessionOutcome::Aborted)
);
let second = store.conn.execute(
"INSERT INTO sessions (task_id, agent, started_at) VALUES (1, 'b', datetime('now'))",
[],
);
assert!(second.is_err());
}
#[test]
fn migration_0008_backfills_flagged_ready_tasks_to_stalled() {
let conn = Connection::open_in_memory().unwrap();
for sql in &MIGRATIONS[..7] {
conn.execute_batch(sql).unwrap();
}
conn.pragma_update(None, "user_version", 7).unwrap();
conn.execute("INSERT INTO projects (name, path) VALUES ('p', '/tmp')", [])
.unwrap();
conn.execute(
"INSERT INTO tasks (project_id, title, state, state_since, created_at)
VALUES (1, 't1', 'ready', datetime('now'), datetime('now')),
(1, 't2', 'ready', datetime('now'), datetime('now')),
(1, 't3', 'ready', datetime('now'), datetime('now')),
(1, 't4', 'ready', datetime('now'), datetime('now')),
(1, 't5', 'ready', datetime('now'), datetime('now')),
(1, 't6', 'running', datetime('now'), datetime('now'))",
[],
)
.unwrap();
conn.execute(
"INSERT INTO sessions (task_id, agent, started_at, ended_at, outcome)
VALUES (1, 'a', datetime('now'), datetime('now'), 'failed'),
(2, 'a', datetime('now'), datetime('now'), 'capped'),
(3, 'a', datetime('now'), datetime('now'), 'aborted'),
(4, 'a', datetime('now'), datetime('now'), 'failed'),
(4, 'a', datetime('now'), datetime('now'), 'aborted'),
(6, 'a', datetime('now'), datetime('now'), 'failed')",
[],
)
.unwrap();
conn.execute("UPDATE tasks SET human = 1 WHERE id = 5", [])
.unwrap();
let store = Store::from_connection(conn).unwrap();
assert_eq!(store.task(1).unwrap().state, TaskState::Stalled);
assert_eq!(store.task(2).unwrap().state, TaskState::Stalled);
assert_eq!(store.task(3).unwrap().state, TaskState::Ready);
assert_eq!(store.task(4).unwrap().state, TaskState::Ready);
assert_eq!(store.task(5).unwrap().state, TaskState::Ready);
assert_eq!(store.task(6).unwrap().state, TaskState::Running);
assert!(store.task(5).unwrap().human);
assert!(!store.task(1).unwrap().human);
let junk = store
.conn
.execute("UPDATE tasks SET human = 2 WHERE id = 5", []);
assert!(junk.is_err(), "the CHECK must reject values outside 0/1");
}
#[test]
fn migration_0010_admits_waiting_and_preserves_existing_tasks() {
let conn = Connection::open_in_memory().unwrap();
for sql in &MIGRATIONS[..9] {
conn.execute_batch(sql).unwrap();
}
conn.pragma_update(None, "user_version", 9).unwrap();
conn.execute("INSERT INTO projects (name, path) VALUES ('p', '/tmp')", [])
.unwrap();
conn.execute(
"INSERT INTO tasks (project_id, title, state, agent, pr_url, branch, human,
state_since, created_at)
VALUES (1, 'in review', 'review', 'claude', 'https://x/pull/1', 'feat/x', 1,
datetime('now'), datetime('now'))",
[],
)
.unwrap();
let store = Store::from_connection(conn).unwrap();
let task = store.task(1).unwrap();
assert_eq!(task.state, TaskState::Review);
assert_eq!(task.pr_url.as_deref(), Some("https://x/pull/1"));
assert_eq!(task.branch.as_deref(), Some("feat/x"));
assert!(task.human);
assert!(
store
.conn
.execute("UPDATE tasks SET state = 'waiting' WHERE id = 1", [])
.is_ok()
);
assert!(
store
.conn
.execute("UPDATE tasks SET state = 'bogus' WHERE id = 1", [])
.is_err()
);
let version: i64 = store
.conn
.query_row("PRAGMA user_version", [], |r| r.get(0))
.unwrap();
assert_eq!(version, MIGRATIONS.len() as i64);
}
#[test]
fn migration_0016_admits_refining_and_preserves_existing_tasks() {
let conn = Connection::open_in_memory().unwrap();
for sql in &MIGRATIONS[..15] {
conn.execute_batch(sql).unwrap();
}
conn.pragma_update(None, "user_version", 15).unwrap();
conn.execute("INSERT INTO projects (name) VALUES ('voro')", [])
.unwrap();
conn.execute(
"INSERT INTO repos (project_id, name, path, is_default)
VALUES (1, 'voro', '/tmp/voro', 1)",
[],
)
.unwrap();
conn.execute(
"INSERT INTO tasks (project_id, repo_id, title, state, agent, pr_url, branch,
human, deep, state_since, created_at)
VALUES (1, 1, 'in review', 'review', 'claude', 'https://x/pull/1', 'feat/x', 1, 1,
datetime('now'), datetime('now'))",
[],
)
.unwrap();
let store = Store::from_connection(conn).unwrap();
let task = store.task(1).unwrap();
assert_eq!(task.state, TaskState::Review);
assert_eq!(task.pr_url.as_deref(), Some("https://x/pull/1"));
assert_eq!(task.branch.as_deref(), Some("feat/x"));
assert_eq!(task.repo_id, Some(1));
assert!(task.human);
assert!(task.deep);
assert!(
store
.conn
.execute("UPDATE tasks SET state = 'refining' WHERE id = 1", [])
.is_ok()
);
assert!(
store
.conn
.execute("UPDATE tasks SET state = 'bogus' WHERE id = 1", [])
.is_err()
);
let version: i64 = store
.conn
.query_row("PRAGMA user_version", [], |r| r.get(0))
.unwrap();
assert_eq!(version, MIGRATIONS.len() as i64);
}
fn task_fixture(s: &mut Store) -> i64 {
s.conn
.execute("INSERT OR IGNORE INTO projects (name) VALUES ('voro')", [])
.unwrap();
let project_id: i64 = s
.conn
.query_row("SELECT id FROM projects WHERE name = 'voro'", [], |r| {
r.get(0)
})
.unwrap();
s.conn
.execute(
"INSERT INTO tasks (project_id, title, state, state_since, created_at)
VALUES (?1, 'run me', 'running', datetime('now'), datetime('now'))",
params![project_id],
)
.unwrap();
s.conn.last_insert_rowid()
}
#[test]
fn events_for_orders_oldest_first() {
use crate::transition::Action;
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let task = s
.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "trace me".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
})
.unwrap();
s.apply(task.id, Action::Start).unwrap();
s.apply(task.id, Action::Ask("A or B?".into())).unwrap();
s.apply(task.id, Action::Resume).unwrap();
let events = s.events_for(task.id).unwrap();
let kinds: Vec<&str> = events.iter().map(|e| e.kind.as_str()).collect();
assert_eq!(
kinds,
vec!["created", "transition", "transition", "transition"]
);
assert!(events.windows(2).all(|w| w[0].id < w[1].id));
}
#[test]
fn latest_summary_returns_the_newest_summary_event() {
use crate::transition::Action;
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let t = s
.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "summary me".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
})
.unwrap();
assert_eq!(s.latest_summary(t.id).unwrap(), None);
s.apply(t.id, Action::Start).unwrap();
s.apply(t.id, Action::Complete(Some("first pass".into())))
.unwrap();
assert_eq!(
s.latest_summary(t.id).unwrap().as_deref(),
Some("first pass")
);
s.apply(t.id, Action::RejectWork("redo".into())).unwrap();
s.apply(t.id, Action::Complete(Some("second pass".into())))
.unwrap();
assert_eq!(
s.latest_summary(t.id).unwrap().as_deref(),
Some("second pass")
);
}
#[test]
fn last_reviewed_supersedes_and_starts_absent() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let t = s
.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "review me".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
})
.unwrap();
assert_eq!(s.last_reviewed(t.id).unwrap(), None);
s.record_reviewed(t.id, " aaaa1111 ").unwrap();
assert_eq!(s.last_reviewed(t.id).unwrap().as_deref(), Some("aaaa1111"));
s.record_reviewed(t.id, "bbbb2222").unwrap();
assert_eq!(s.last_reviewed(t.id).unwrap().as_deref(), Some("bbbb2222"));
assert!(s.record_reviewed(t.id, " ").is_err());
assert!(s.record_reviewed(999, "aaaa1111").is_err());
assert_eq!(s.task(t.id).unwrap().state, TaskState::Ready);
}
#[test]
fn incomplete_report_flag_marks_a_review_task_with_a_branch_and_no_summary() {
use crate::transition::Action;
fn reviewed(branch: Option<&str>, summary: Option<&str>) -> (Store, i64) {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let t = s
.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "report me".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
})
.unwrap();
s.apply(t.id, Action::Start).unwrap();
s.apply(t.id, Action::Complete(summary.map(str::to_string)))
.unwrap();
if let Some(name) = branch {
s.set_branch(t.id, Some(name)).unwrap();
}
(s, t.id)
}
let (s, id) = reviewed(Some("feat/x"), None);
assert!(
s.incomplete_report_flag(id).unwrap(),
"half report: branch, no summary"
);
let (s, id) = reviewed(None, Some("already fixed by PR #96; nothing to do"));
assert!(
!s.incomplete_report_flag(id).unwrap(),
"no-code report: summary, no branch"
);
let (s, id) = reviewed(Some("feat/x"), Some("did the thing"));
assert!(!s.incomplete_report_flag(id).unwrap());
let (s, id) = reviewed(None, None);
assert!(!s.incomplete_report_flag(id).unwrap());
}
#[test]
fn incomplete_report_flag_is_gated_on_review() {
use crate::transition::Action;
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let t = s
.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "in flight".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
})
.unwrap();
s.set_branch(t.id, Some("feat/x")).unwrap();
assert!(!s.incomplete_report_flag(t.id).unwrap(), "ready");
s.apply(t.id, Action::Start).unwrap();
assert!(!s.incomplete_report_flag(t.id).unwrap(), "running");
s.apply(t.id, Action::Complete(None)).unwrap();
assert!(s.incomplete_report_flag(t.id).unwrap(), "review");
s.apply(t.id, Action::Accept).unwrap();
assert!(!s.incomplete_report_flag(t.id).unwrap(), "done");
}
fn proposal(s: &mut Store, title: &str) -> Task {
let p = s.projects().unwrap().first().cloned().unwrap_or_else(|| {
s.create_project("voro", "/tmp/voro").unwrap();
s.projects().unwrap().remove(0)
});
s.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: title.into(),
body: "thin body".into(),
priority: Priority::P2,
state: TaskState::Proposed,
agent: None,
human: false,
deep: false,
})
.unwrap()
}
#[test]
fn record_refine_launch_moves_the_task_and_logs_the_note() {
let mut s = Store::open_in_memory().unwrap();
let blocker = proposal(&mut s, "blocker");
let t = proposal(&mut s, "refine me");
s.add_dep(t.id, blocker.id, DepKind::Blocks).unwrap();
let before = s.task(t.id).unwrap();
let (after, session) = s
.record_refine_launch(
t.id,
" name the files it touches ",
"claude",
Some(4321),
LivenessSource::Pid,
Some("/var/log/refine.log"),
)
.unwrap();
assert_eq!(after.state, TaskState::Refining);
assert_eq!(after.priority, before.priority);
assert_eq!(after.body, before.body);
assert_eq!(s.deps_of(t.id).unwrap().len(), 1);
assert_eq!(session.pid, Some(4321));
assert_eq!(session.log_path.as_deref(), Some("/var/log/refine.log"));
assert!(session.ended_at.is_none());
assert_eq!(
s.latest_refine_note(t.id).unwrap().as_deref(),
Some("name the files it touches")
);
assert!(!s.refined_flag(t.id).unwrap());
assert!(!s.refine_failed_flag(t.id).unwrap());
}
#[test]
fn a_note_less_refine_launch_logs_no_note() {
let mut s = Store::open_in_memory().unwrap();
let t = proposal(&mut s, "refine me");
s.record_refine_launch(t.id, "", "claude", Some(1), LivenessSource::Pid, None)
.unwrap();
assert_eq!(s.latest_refine_note(t.id).unwrap(), None);
assert_eq!(s.task(t.id).unwrap().state, TaskState::Refining);
}
#[test]
fn the_markers_read_the_round_that_just_concluded() {
use crate::transition::{Action, Triage};
let mut s = Store::open_in_memory().unwrap();
let t = proposal(&mut s, "refine me");
for (outcome, refined, failed) in [
(RefineOutcome::Applied, true, false),
(RefineOutcome::Failed, false, true),
(RefineOutcome::Cancelled, false, false),
(RefineOutcome::Applied, true, false),
] {
s.record_refine_launch(t.id, "note", "claude", Some(1), LivenessSource::Pid, None)
.unwrap();
let after = s.conclude_refine(t.id, outcome).unwrap();
assert_eq!(after.state, TaskState::Proposed, "{outcome}");
assert_eq!(s.refined_flag(t.id).unwrap(), refined, "{outcome}");
assert_eq!(s.refine_failed_flag(t.id).unwrap(), failed, "{outcome}");
assert_eq!(s.latest_refine_outcome(t.id).unwrap(), Some(outcome));
}
s.apply(t.id, Action::Triage(Triage::Parked)).unwrap();
assert!(!s.refined_flag(t.id).unwrap());
assert!(!s.refine_failed_flag(t.id).unwrap());
}
#[test]
fn a_late_rewrite_corrects_a_failed_round_to_applied() {
let mut s = Store::open_in_memory().unwrap();
let t = proposal(&mut s, "refine me");
s.record_refine_launch(t.id, "note", "claude", Some(1), LivenessSource::Pid, None)
.unwrap();
s.conclude_refine(t.id, RefineOutcome::Failed).unwrap();
assert!(s.refine_failed_flag(t.id).unwrap());
assert!(s.correct_late_refine(t.id).unwrap());
assert!(s.refined_flag(t.id).unwrap());
assert!(!s.refine_failed_flag(t.id).unwrap());
assert_eq!(
s.latest_refine_outcome(t.id).unwrap(),
Some(RefineOutcome::Applied)
);
assert_eq!(s.task(t.id).unwrap().state, TaskState::Proposed);
assert_eq!(
s.sessions_for(t.id).unwrap()[0].outcome,
Some(SessionOutcome::Failed),
"the session keeps the outcome the reconciler observed"
);
assert!(!s.correct_late_refine(t.id).unwrap());
}
#[test]
fn correcting_a_round_is_a_no_op_off_the_failed_case() {
use crate::transition::{Action, Triage};
for outcome in [RefineOutcome::Applied, RefineOutcome::Cancelled] {
let mut s = Store::open_in_memory().unwrap();
let t = proposal(&mut s, "refine me");
s.record_refine_launch(t.id, "note", "claude", Some(1), LivenessSource::Pid, None)
.unwrap();
s.conclude_refine(t.id, outcome).unwrap();
assert!(!s.correct_late_refine(t.id).unwrap(), "{outcome}");
assert_eq!(s.latest_refine_outcome(t.id).unwrap(), Some(outcome));
}
let mut s = Store::open_in_memory().unwrap();
let t = proposal(&mut s, "never refined");
assert!(!s.correct_late_refine(t.id).unwrap());
assert_eq!(s.latest_refine_outcome(t.id).unwrap(), None);
s.record_refine_launch(t.id, "note", "claude", Some(1), LivenessSource::Pid, None)
.unwrap();
s.conclude_refine(t.id, RefineOutcome::Failed).unwrap();
s.apply(t.id, Action::Triage(Triage::Ready)).unwrap();
assert!(!s.correct_late_refine(t.id).unwrap());
assert_eq!(
s.latest_refine_outcome(t.id).unwrap(),
Some(RefineOutcome::Failed)
);
}
#[test]
fn concluding_a_round_closes_its_session_with_the_matching_outcome() {
for (outcome, session_outcome) in [
(RefineOutcome::Applied, SessionOutcome::Completed),
(RefineOutcome::Failed, SessionOutcome::Failed),
(RefineOutcome::Cancelled, SessionOutcome::Aborted),
] {
let mut s = Store::open_in_memory().unwrap();
let t = proposal(&mut s, "refine me");
let (_, session) = s
.record_refine_launch(t.id, "note", "claude", Some(1), LivenessSource::Pid, None)
.unwrap();
s.conclude_refine(t.id, outcome).unwrap();
let closed = s.session(session.id).unwrap();
assert!(closed.ended_at.is_some(), "{outcome}");
assert_eq!(closed.outcome, Some(session_outcome), "{outcome}");
}
}
#[test]
fn refine_transitions_are_refused_from_the_wrong_state() {
use crate::transition::{Action, Triage};
let mut s = Store::open_in_memory().unwrap();
let t = proposal(&mut s, "refine me");
assert!(matches!(
s.conclude_refine(t.id, RefineOutcome::Applied),
Err(Error::InvalidTransition { .. })
));
s.apply(t.id, Action::Triage(Triage::Parked)).unwrap();
assert!(matches!(
s.record_refine_launch(t.id, "too late", "claude", None, LivenessSource::Pid, None),
Err(Error::InvalidTransition { .. })
));
assert_eq!(s.task(t.id).unwrap().state, TaskState::Parked);
assert!(s.sessions_for(t.id).unwrap().is_empty());
}
#[test]
fn discovered_from_resolves_the_parent_proposal() {
let mut s = Store::open_in_memory().unwrap();
let parent = proposal(&mut s, "parent");
let child = proposal(&mut s, "child");
assert!(s.discovered_from(child.id).unwrap().is_none());
s.add_dep(child.id, parent.id, DepKind::DiscoveredFrom)
.unwrap();
assert_eq!(
s.discovered_from(child.id).unwrap().map(|t| t.id),
Some(parent.id)
);
let blocker = proposal(&mut s, "blocker");
assert!(s.discovered_from(blocker.id).unwrap().is_none());
}
#[test]
fn incomplete_report_flag_is_false_for_a_missing_task() {
let s = Store::open_in_memory().unwrap();
assert!(!s.incomplete_report_flag(999).unwrap());
}
#[test]
fn session_create_end_round_trip() {
let mut s = Store::open_in_memory().unwrap();
let task_id = task_fixture(&mut s);
let opened = s
.create_session(
task_id,
"claude",
Some(4321),
LivenessSource::Pid,
Some("/var/log/s.log"),
)
.unwrap();
assert_eq!(opened.task_id, task_id);
assert_eq!(opened.agent, "claude");
assert_eq!(opened.pid, Some(4321));
assert_eq!(opened.log_path.as_deref(), Some("/var/log/s.log"));
assert!(!opened.started_at.is_empty());
assert!(opened.ended_at.is_none());
assert!(opened.outcome.is_none());
let ended = s.end_session(opened.id, SessionOutcome::Completed).unwrap();
assert_eq!(ended.id, opened.id);
assert!(ended.ended_at.is_some());
assert_eq!(ended.outcome, Some(SessionOutcome::Completed));
assert_eq!(s.session(opened.id).unwrap(), ended);
}
#[test]
fn latest_sessions_keeps_only_the_newest_per_task() {
let mut s = Store::open_in_memory().unwrap();
let with_history = task_fixture(&mut s);
let sessionless = task_fixture(&mut s);
let first = s
.create_session(
with_history,
"claude",
None,
LivenessSource::Pid,
Some("/var/log/first.log"),
)
.unwrap();
s.end_session(first.id, SessionOutcome::Failed).unwrap();
let second = s
.create_session(
with_history,
"codex",
None,
LivenessSource::Pid,
Some("/var/log/second.log"),
)
.unwrap();
let latest = s.latest_sessions().unwrap();
assert_eq!(latest.len(), 1);
assert_eq!(latest[&with_history].id, second.id);
assert_eq!(
latest[&with_history].log_path.as_deref(),
Some("/var/log/second.log")
);
assert!(!latest.contains_key(&sessionless));
}
#[test]
fn session_optional_fields_are_null() {
let mut s = Store::open_in_memory().unwrap();
let task_id = task_fixture(&mut s);
let opened = s
.create_session(task_id, "codex", None, LivenessSource::Pid, None)
.unwrap();
assert!(opened.pid.is_none());
assert!(opened.session_ref.is_none());
assert!(opened.log_path.is_none());
}
#[test]
fn set_session_ref_records_and_rejects_unknown_ids() {
let mut s = Store::open_in_memory().unwrap();
let task_id = task_fixture(&mut s);
let opened = s
.create_session(task_id, "claude", None, LivenessSource::Pid, None)
.unwrap();
assert!(opened.session_ref.is_none());
let updated = s
.set_session_ref(opened.id, "3f6c0e6e-1111-2222-3333-444455556666")
.unwrap();
assert_eq!(
updated.session_ref.as_deref(),
Some("3f6c0e6e-1111-2222-3333-444455556666")
);
assert_eq!(s.session(opened.id).unwrap(), updated);
assert!(matches!(
s.set_session_ref(999, "x"),
Err(Error::SessionNotFound(999))
));
}
#[test]
fn a_session_records_the_liveness_source_it_was_launched_with() {
let mut s = Store::open_in_memory().unwrap();
let task_id = task_fixture(&mut s);
let listing = s
.create_session(task_id, "claude", Some(1), LivenessSource::Listing, None)
.unwrap();
assert_eq!(listing.liveness_source, LivenessSource::Listing);
assert_eq!(
s.session(listing.id).unwrap().liveness_source,
LivenessSource::Listing
);
let pid = s
.create_session(task_id, "manual", Some(1), LivenessSource::Pid, None)
.unwrap();
assert_eq!(pid.liveness_source, LivenessSource::Pid);
assert_eq!(
s.live_sessions().unwrap()[0].liveness_source,
pid.liveness_source
);
s.set_session_ref(pid.id, "uuid").unwrap();
assert_eq!(
s.session(pid.id).unwrap().liveness_source,
LivenessSource::Pid
);
}
#[test]
fn dispatch_and_refine_launches_carry_their_liveness_source() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("proj", "/tmp/proj").unwrap();
let ready = s
.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "run me".into(),
body: String::new(),
priority: Priority::P1,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
})
.unwrap();
let (_, dispatched) = s
.record_dispatch(ready.id, "claude", Some(1), LivenessSource::Listing, None)
.unwrap();
assert_eq!(dispatched.liveness_source, LivenessSource::Listing);
let proposal = s
.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "sloppy".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Proposed,
agent: None,
human: false,
deep: false,
})
.unwrap();
let (_, headless) = s
.record_refine_launch(
proposal.id,
"name the files",
"claude",
Some(2),
LivenessSource::Listing,
None,
)
.unwrap();
assert_eq!(headless.liveness_source, LivenessSource::Listing);
s.conclude_refine(proposal.id, RefineOutcome::Cancelled)
.unwrap();
let (_, interactive) = s
.record_refine_launch(
proposal.id,
"",
"claude",
Some(3),
LivenessSource::Pid,
None,
)
.unwrap();
assert_eq!(interactive.liveness_source, LivenessSource::Pid);
}
#[test]
fn migration_0017_defaults_existing_sessions_to_the_listing() {
let conn = Connection::open_in_memory().unwrap();
for sql in &MIGRATIONS[..16] {
conn.execute_batch(sql).unwrap();
}
conn.pragma_update(None, "user_version", 16).unwrap();
conn.execute("INSERT INTO projects (name) VALUES ('p')", [])
.unwrap();
conn.execute(
"INSERT INTO repos (project_id, name, path, is_default)
VALUES (1, 'p', '/tmp/p', 1)",
[],
)
.unwrap();
conn.execute(
"INSERT INTO tasks (project_id, title, state, state_since, created_at)
VALUES (1, 'dispatched', 'running', datetime('now'), datetime('now'))",
[],
)
.unwrap();
conn.execute(
"INSERT INTO sessions (task_id, agent, pid, started_at)
VALUES (1, 'claude', 4242, datetime('now'))",
[],
)
.unwrap();
let store = Store::from_connection(conn).unwrap();
assert_eq!(
store.session(1).unwrap().liveness_source,
LivenessSource::Listing
);
let junk = store.conn.execute(
"UPDATE sessions SET liveness_source = 'guess' WHERE id = 1",
[],
);
assert!(junk.is_err(), "the CHECK must reject an unknown source");
}
#[test]
fn record_session_send_moves_the_pid_and_follows_a_fork() {
let mut s = Store::open_in_memory().unwrap();
let task_id = task_fixture(&mut s);
let opened = s
.create_session(task_id, "claude", Some(1234), LivenessSource::Pid, None)
.unwrap();
s.set_session_ref(opened.id, "first-ref").unwrap();
let resumed = s.record_session_send(opened.id, None, 4321).unwrap();
assert_eq!(resumed.pid, Some(4321));
assert_eq!(resumed.session_ref.as_deref(), Some("first-ref"));
let forked = s
.record_session_send(opened.id, Some("forked-ref"), 5678)
.unwrap();
assert_eq!(forked.pid, Some(5678));
assert_eq!(forked.session_ref.as_deref(), Some("forked-ref"));
assert_eq!(s.session(opened.id).unwrap(), forked);
assert!(matches!(
s.record_session_send(999, None, 1),
Err(Error::SessionNotFound(999))
));
}
#[test]
fn end_session_rejects_unknown_id() {
let mut s = Store::open_in_memory().unwrap();
assert!(matches!(
s.end_session(999, SessionOutcome::Aborted),
Err(Error::SessionNotFound(999))
));
}
#[test]
fn sessions_for_returns_newest_first() {
let mut s = Store::open_in_memory().unwrap();
let task_id = task_fixture(&mut s);
let first = s
.create_session(task_id, "claude", None, LivenessSource::Pid, None)
.unwrap();
let second = s
.create_session(task_id, "claude", None, LivenessSource::Pid, None)
.unwrap();
let sessions = s.sessions_for(task_id).unwrap();
assert_eq!(
sessions.iter().map(|s| s.id).collect::<Vec<_>>(),
vec![second.id, first.id]
);
}
#[test]
fn live_sessions_excludes_ended() {
let mut s = Store::open_in_memory().unwrap();
let task_id = task_fixture(&mut s);
let done = s
.create_session(task_id, "claude", None, LivenessSource::Pid, None)
.unwrap();
let live = s
.create_session(task_id, "claude", None, LivenessSource::Pid, None)
.unwrap();
s.end_session(done.id, SessionOutcome::Failed).unwrap();
let ids = s.live_sessions().unwrap();
assert_eq!(ids.iter().map(|s| s.id).collect::<Vec<_>>(), vec![live.id]);
}
#[test]
fn running_rows_join_current_task_fields() {
let mut s = Store::open_in_memory().unwrap();
let task_id = task_fixture(&mut s);
let session = s
.create_session(task_id, "claude", None, LivenessSource::Pid, None)
.unwrap();
let rows = s.running_rows().unwrap();
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].session_id, Some(session.id));
assert_eq!(rows[0].task_id, task_id);
assert_eq!(rows[0].task_title, "run me");
assert_eq!(rows[0].task_state, TaskState::Running);
assert_eq!(rows[0].agent.as_deref(), Some("claude"));
assert!(rows[0].elapsed_secs >= 0);
}
#[test]
fn running_rows_exclude_ended_sessions_and_order_newest_first() {
let mut s = Store::open_in_memory().unwrap();
let task_id = task_fixture(&mut s);
let done = s
.create_session(task_id, "claude", None, LivenessSource::Pid, None)
.unwrap();
let live = s
.create_session(task_id, "codex", None, LivenessSource::Pid, None)
.unwrap();
s.end_session(done.id, SessionOutcome::Completed).unwrap();
let rows = s.running_rows().unwrap();
assert_eq!(
rows.iter().map(|r| r.session_id).collect::<Vec<_>>(),
vec![Some(live.id)]
);
assert_eq!(rows[0].agent.as_deref(), Some("codex"));
}
#[test]
fn running_rows_compute_elapsed_from_started_at() {
let mut s = Store::open_in_memory().unwrap();
let task_id = task_fixture(&mut s);
let session = s
.create_session(task_id, "claude", None, LivenessSource::Pid, None)
.unwrap();
s.conn
.execute(
"UPDATE sessions SET started_at = datetime('now', '-90 seconds') WHERE id = ?1",
params![session.id],
)
.unwrap();
let rows = s.running_rows().unwrap();
assert_eq!(rows.len(), 1);
assert!(
(85..=95).contains(&rows[0].elapsed_secs),
"expected ~90s elapsed, got {}",
rows[0].elapsed_secs
);
}
#[test]
fn running_rows_include_running_task_without_live_session() {
let mut s = Store::open_in_memory().unwrap();
let task_id = task_fixture(&mut s);
s.conn
.execute(
"UPDATE tasks SET state_since = datetime('now', '-90 seconds') WHERE id = ?1",
params![task_id],
)
.unwrap();
let rows = s.running_rows().unwrap();
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].session_id, None);
assert_eq!(rows[0].agent, None);
assert_eq!(rows[0].task_id, task_id);
assert_eq!(rows[0].task_state, TaskState::Running);
assert!(
(85..=95).contains(&rows[0].elapsed_secs),
"expected ~90s in running, got {}",
rows[0].elapsed_secs
);
}
#[test]
fn running_rows_include_task_whose_sessions_all_ended() {
let mut s = Store::open_in_memory().unwrap();
let task_id = task_fixture(&mut s);
let done = s
.create_session(task_id, "claude", None, LivenessSource::Pid, None)
.unwrap();
s.end_session(done.id, SessionOutcome::Failed).unwrap();
let rows = s.running_rows().unwrap();
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].session_id, None);
assert_eq!(rows[0].task_id, task_id);
}
#[test]
fn running_rows_order_live_sessions_before_session_less_tasks() {
let mut s = Store::open_in_memory().unwrap();
let live_task = task_fixture(&mut s);
let session = s
.create_session(live_task, "claude", None, LivenessSource::Pid, None)
.unwrap();
let orphan_task = task_fixture(&mut s);
let rows = s.running_rows().unwrap();
assert_eq!(rows.len(), 2);
assert_eq!(rows[0].session_id, Some(session.id));
assert_eq!(rows[0].task_id, live_task);
assert_eq!(rows[1].session_id, None);
assert_eq!(rows[1].task_id, orphan_task);
}
#[test]
fn running_rows_include_a_refining_task() {
let mut s = Store::open_in_memory().unwrap();
let t = proposal(&mut s, "refine me");
let (_, session) = s
.record_refine_launch(
t.id,
"thin body",
"claude",
Some(1),
LivenessSource::Pid,
None,
)
.unwrap();
let rows = s.running_rows().unwrap();
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].task_id, t.id);
assert_eq!(rows[0].task_state, TaskState::Refining);
assert_eq!(rows[0].session_id, Some(session.id));
assert_eq!(rows[0].agent.as_deref(), Some("claude"));
s.conclude_refine(t.id, RefineOutcome::Applied).unwrap();
assert!(s.running_rows().unwrap().is_empty());
}
#[test]
fn running_rows_measure_a_waiting_task_from_the_hand_off() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let id = s
.create_task(NewTask {
project_id: p.id,
repo_id: None,
title: "handed off".into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
})
.unwrap()
.id;
let (_, opened) = s
.record_dispatch(id, "claude", Some(1), LivenessSource::Pid, None)
.unwrap();
s.apply(id, Action::Complete(None)).unwrap();
s.apply(id, Action::HandOff).unwrap();
s.set_pr(id, Some("https://github.com/o/r/pull/7")).unwrap();
let session = opened.id;
s.conn
.execute(
"UPDATE sessions SET started_at = datetime('now', '-2 hours') WHERE id = ?1",
params![session],
)
.unwrap();
s.conn
.execute(
"UPDATE tasks SET state_since = datetime('now', '-90 seconds') WHERE id = ?1",
params![id],
)
.unwrap();
let rows = s.running_rows().unwrap();
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].task_id, id);
assert_eq!(rows[0].task_state, TaskState::Waiting);
assert_eq!(rows[0].session_id, Some(session));
assert_eq!(
rows[0].pr_url.as_deref(),
Some("https://github.com/o/r/pull/7")
);
assert!(
(85..=95).contains(&rows[0].elapsed_secs),
"expected ~90s waiting, got {} (the session's own age is 2h)",
rows[0].elapsed_secs
);
}
#[test]
fn running_rows_sort_waiting_after_work_under_way() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let waiting = task_in_state(&mut s, p.id, TaskState::Waiting);
let running = task_in_state(&mut s, p.id, TaskState::Running);
let refining = task_in_state(&mut s, p.id, TaskState::Refining);
let rows = s.running_rows().unwrap();
assert_eq!(rows.len(), 3);
assert_eq!(rows.last().unwrap().task_id, waiting);
let ahead: Vec<i64> = rows[..2].iter().map(|r| r.task_id).collect();
assert!(
ahead.contains(&running) && ahead.contains(&refining),
"{ahead:?}"
);
}
#[test]
fn running_rows_exclude_a_waiting_task_in_an_archived_project() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("retiring", "/tmp/retiring").unwrap();
let id = task_in_state(&mut s, p.id, TaskState::Waiting);
assert_eq!(s.running_rows().unwrap().len(), 1);
s.set_archived(p.id, true).unwrap();
assert!(s.running_rows().unwrap().is_empty());
assert_eq!(s.task(id).unwrap().state, TaskState::Waiting);
}
#[test]
fn running_rows_exclude_tasks_that_left_running() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let new = |title: &str| NewTask {
project_id: p.id,
repo_id: None,
title: title.into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
};
let review = s.create_task(new("review")).unwrap().id;
s.record_dispatch(review, "claude", Some(1), LivenessSource::Pid, None)
.unwrap();
s.apply(review, Action::Complete(None)).unwrap();
assert!(s.sessions_for(review).unwrap()[0].ended_at.is_none());
let done = s.create_task(new("done")).unwrap().id;
s.record_dispatch(done, "claude", Some(2), LivenessSource::Pid, None)
.unwrap();
s.apply(done, Action::Complete(None)).unwrap();
s.apply(done, Action::Accept).unwrap();
let rejected = s.create_task(new("rejected")).unwrap().id;
s.record_dispatch(rejected, "claude", Some(3), LivenessSource::Pid, None)
.unwrap();
s.apply(rejected, Action::Abort).unwrap();
s.apply(rejected, Action::Abandon).unwrap();
let running = s.create_task(new("running")).unwrap().id;
s.record_dispatch(running, "claude", Some(4), LivenessSource::Pid, None)
.unwrap();
let rows = s.running_rows().unwrap();
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].task_id, running);
}
#[test]
fn running_rows_ignore_a_stale_open_session_on_a_closed_task() {
let mut s = Store::open_in_memory().unwrap();
let task_id = task_fixture(&mut s);
s.create_session(task_id, "claude", Some(1), LivenessSource::Pid, None)
.unwrap();
s.conn
.execute("UPDATE tasks SET state = 'done' WHERE id = ?1", [task_id])
.unwrap();
assert!(s.running_rows().unwrap().is_empty());
}
#[test]
fn session_outcome_serialises_for_all_variants() {
let mut s = Store::open_in_memory().unwrap();
let task_id = task_fixture(&mut s);
for outcome in SessionOutcome::ALL {
let opened = s
.create_session(task_id, "claude", None, LivenessSource::Pid, None)
.unwrap();
let ended = s.end_session(opened.id, outcome).unwrap();
assert_eq!(ended.outcome, Some(outcome));
}
}
fn scratch_db() -> PathBuf {
tempfile::Builder::new()
.prefix("voro-dataversion-")
.tempdir()
.unwrap()
.keep()
.join("voro.db")
}
#[test]
fn data_version_tracks_external_commits_only() {
let path = scratch_db();
let mut a = Store::open(&path).unwrap();
let mut b = Store::open(&path).unwrap();
let start = a.data_version().unwrap();
a.create_project("alpha", "/tmp/alpha").unwrap();
assert_eq!(a.data_version().unwrap(), start);
b.create_project("beta", "/tmp/beta").unwrap();
assert_ne!(a.data_version().unwrap(), start);
drop(a);
drop(b);
let _ = std::fs::remove_file(&path);
}
#[test]
fn dep_maps_resolve_both_directions_with_title_state_and_kind() {
use crate::model::{DepKind, DepRef, Priority};
use crate::transition::Action;
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let new = |title: &str| NewTask {
project_id: p.id,
repo_id: None,
title: title.into(),
body: String::new(),
priority: Priority::P2,
state: TaskState::Ready,
agent: None,
human: false,
deep: false,
};
let blocker = s.create_task(new("blocker")).unwrap();
s.apply(blocker.id, Action::Start).unwrap();
s.apply(blocker.id, Action::Complete(None)).unwrap();
s.apply(blocker.id, Action::Accept).unwrap();
let source = s.create_task(new("source")).unwrap();
let task = s.create_task(new("task")).unwrap();
s.add_dep(task.id, blocker.id, DepKind::Blocks).unwrap();
s.add_dep(task.id, source.id, DepKind::DiscoveredFrom)
.unwrap();
let deps = s.deps_by_task().unwrap();
assert_eq!(
deps[&task.id],
vec![
DepRef {
id: blocker.id,
title: "blocker".into(),
state: TaskState::Done,
kind: DepKind::Blocks,
},
DepRef {
id: source.id,
title: "source".into(),
state: TaskState::Ready,
kind: DepKind::DiscoveredFrom,
},
]
);
assert!(!deps[&task.id][0].is_open());
assert!(!deps.contains_key(&blocker.id));
let dependents = s.dependents_by_task().unwrap();
assert_eq!(
dependents[&blocker.id],
vec![DepRef {
id: task.id,
title: "task".into(),
state: TaskState::Ready,
kind: DepKind::Blocks,
}]
);
assert_eq!(dependents[&source.id].len(), 1);
assert_eq!(dependents[&source.id][0].kind, DepKind::DiscoveredFrom);
assert!(!dependents.contains_key(&task.id));
}
#[test]
fn a_doc_links_tasks_across_projects_and_answers_both_directions() {
let mut s = Store::open_in_memory().unwrap();
let plan = s.create_project("augere", "/tmp/augere").unwrap();
let other = s.create_project("mote", "/tmp/mote").unwrap();
let doc = s
.create_doc(plan.id, None, "docs/strategy.md", Some("Strategy"))
.unwrap();
let a = s.create_task(new_ready(plan.id)).unwrap();
let b = s.create_task(new_ready(other.id)).unwrap();
let c = s.create_task(new_ready(other.id)).unwrap();
for task in [&a, &b, &c] {
assert!(s.link_doc(task.id, doc.id).unwrap());
}
assert!(!s.link_doc(a.id, doc.id).unwrap());
let derived: Vec<i64> = s
.tasks_for_doc(doc.id)
.unwrap()
.into_iter()
.map(|t| t.id)
.collect();
assert_eq!(derived, vec![a.id, b.id, c.id]);
assert_eq!(s.docs_for_task(b.id).unwrap(), vec![doc.clone()]);
assert_eq!(s.docs_by_task().unwrap()[&c.id], vec![doc.clone()]);
assert!(s.unlink_doc(b.id, doc.id).unwrap());
assert!(!s.unlink_doc(b.id, doc.id).unwrap());
assert_eq!(s.tasks_for_doc(doc.id).unwrap().len(), 2);
assert!(s.docs_for_task(b.id).unwrap().is_empty());
}
#[test]
fn every_doc_link_and_unlink_lands_on_the_task_event_trail() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let doc = s.create_doc(p.id, None, "docs/DESIGN.md", None).unwrap();
let t = s.create_task(new_ready(p.id)).unwrap();
s.link_doc(t.id, doc.id).unwrap();
s.unlink_doc(t.id, doc.id).unwrap();
s.unlink_doc(t.id, doc.id).unwrap();
let kinds: Vec<String> = s
.events_for(t.id)
.unwrap()
.into_iter()
.map(|e| e.kind)
.collect();
assert_eq!(kinds, vec!["created", "doc-linked", "doc-unlinked"]);
}
#[test]
fn set_task_docs_replaces_the_whole_list_and_logs_both_directions() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let one = s.create_doc(p.id, None, "docs/a.md", None).unwrap();
let two = s.create_doc(p.id, None, "docs/b.md", None).unwrap();
let t = s.create_task(new_ready(p.id)).unwrap();
s.set_task_docs(t.id, &[one.id]).unwrap();
let now = s.set_task_docs(t.id, &[two.id]).unwrap();
assert_eq!(now, vec![two.clone()]);
let events: Vec<(String, Option<String>)> = s
.events_for(t.id)
.unwrap()
.into_iter()
.map(|e| (e.kind, e.detail))
.collect();
assert_eq!(
events,
vec![
("created".into(), Some("ready".into())),
("doc-linked".into(), Some("docs/a.md".into())),
("doc-unlinked".into(), Some("docs/a.md".into())),
("doc-linked".into(), Some("docs/b.md".into())),
]
);
assert!(s.set_task_docs(t.id, &[]).unwrap().is_empty());
}
#[test]
fn an_absolute_path_inside_a_checkout_is_stored_relative_to_it() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("augere", "/tmp/augere").unwrap();
let doc = s
.create_doc(p.id, None, "/tmp/augere/docs/strategy.md", None)
.unwrap();
assert_eq!(doc.location, "docs/strategy.md");
assert_eq!(s.resolve_doc(&doc).unwrap(), "/tmp/augere/docs/strategy.md");
s.set_default_repo_path(p.id, "/srv/augere").unwrap();
assert_eq!(
s.resolve_doc(&s.doc(doc.id).unwrap()).unwrap(),
"/srv/augere/docs/strategy.md"
);
}
#[test]
fn a_relative_doc_resolves_against_the_repo_it_names() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("odm", "/tmp/odm").unwrap();
let oats = s.add_repo(p.id, "oats", "/tmp/oats").unwrap();
let default = s.create_doc(p.id, None, "docs/plan.md", None).unwrap();
assert_eq!(s.resolve_doc(&default).unwrap(), "/tmp/odm/docs/plan.md");
let named = s
.create_doc(p.id, Some(oats.id), "notes/plan.md", None)
.unwrap();
assert_eq!(s.resolve_doc(&named).unwrap(), "/tmp/oats/notes/plan.md");
let nested = s.add_repo(p.id, "inner", "/tmp/odm/vendor").unwrap();
let doc = s
.create_doc(p.id, None, "/tmp/odm/vendor/docs/x.md", None)
.unwrap();
assert_eq!(doc.repo_id, Some(nested.id));
assert_eq!(doc.location, "docs/x.md");
}
#[test]
fn a_url_resolves_verbatim_and_takes_no_repo() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let repo = s.default_repo(p.id).unwrap();
let doc = s
.create_doc(p.id, None, "https://example.com/plan", Some("Plan"))
.unwrap();
assert!(doc.is_url());
assert!(doc.repo_id.is_none());
assert_eq!(s.resolve_doc(&doc).unwrap(), "https://example.com/plan");
assert_eq!(doc.label(), "Plan");
assert!(
s.create_doc(p.id, Some(repo.id), "https://example.com/other", None)
.is_err()
);
}
#[test]
fn a_doc_outside_every_checkout_stays_absolute() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let doc = s
.create_doc(p.id, None, "/etc/notes/plan.md", None)
.unwrap();
assert_eq!(doc.location, "/etc/notes/plan.md");
assert!(doc.repo_id.is_none());
assert_eq!(s.resolve_doc(&doc).unwrap(), "/etc/notes/plan.md");
let repo = s.default_repo(p.id).unwrap();
assert!(
s.create_doc(p.id, Some(repo.id), "/etc/notes/other.md", None)
.is_err()
);
}
#[test]
fn a_doc_is_registered_once_per_project_and_labels_itself() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let other = s.create_project("mote", "/tmp/mote").unwrap();
let doc = s.create_doc(p.id, None, "docs/plan.md", None).unwrap();
assert_eq!(doc.label(), "docs/plan.md");
assert!(s.create_doc(p.id, None, "docs/plan.md", None).is_err());
let twin = s.create_doc(other.id, None, "docs/plan.md", None).unwrap();
assert_eq!(s.docs_at("docs/plan.md").unwrap().len(), 2);
assert_ne!(doc.id, twin.id);
assert_eq!(s.docs(p.id).unwrap(), vec![doc]);
assert!(s.create_doc(p.id, None, " ", None).is_err());
}
#[test]
fn removing_a_doc_unlinks_its_tasks_rather_than_refusing() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let doc = s.create_doc(p.id, None, "docs/plan.md", None).unwrap();
let a = s.create_task(new_ready(p.id)).unwrap();
let b = s.create_task(new_ready(p.id)).unwrap();
s.link_doc(a.id, doc.id).unwrap();
s.link_doc(b.id, doc.id).unwrap();
let freed = s.delete_doc(doc.id).unwrap();
assert_eq!(freed, vec![a.id, b.id]);
assert!(s.doc(doc.id).is_err());
assert!(s.docs_for_task(a.id).unwrap().is_empty());
assert!(s.docs_by_task().unwrap().is_empty());
assert_eq!(
s.events_for(a.id).unwrap().last().unwrap().kind,
"doc-unlinked"
);
}
#[test]
fn linking_names_a_task_and_a_doc_that_exist() {
let mut s = Store::open_in_memory().unwrap();
let p = s.create_project("voro", "/tmp/voro").unwrap();
let doc = s.create_doc(p.id, None, "docs/plan.md", None).unwrap();
let t = s.create_task(new_ready(p.id)).unwrap();
assert!(s.link_doc(999, doc.id).is_err());
assert!(s.link_doc(t.id, 999).is_err());
assert!(s.set_task_docs(t.id, &[999]).is_err());
assert!(s.docs_for_task(t.id).unwrap().is_empty());
}
#[test]
fn migration_0014_leaves_existing_tasks_untouched() {
let conn = Connection::open_in_memory().unwrap();
for sql in &MIGRATIONS[..13] {
conn.execute_batch(sql).unwrap();
}
conn.pragma_update(None, "user_version", 13).unwrap();
conn.execute("INSERT INTO projects (name) VALUES ('legacy')", [])
.unwrap();
conn.execute(
"INSERT INTO repos (project_id, name, path, is_default)
VALUES (1, 'legacy', '/tmp/legacy', 1)",
[],
)
.unwrap();
conn.execute(
"INSERT INTO tasks (project_id, title, state, state_since, created_at)
VALUES (1, 'old work', 'ready', datetime('now'), datetime('now'))",
[],
)
.unwrap();
let store = Store::from_connection(conn).unwrap();
let task = store.task(1).unwrap();
assert_eq!(task.title, "old work");
assert_eq!(task.state, TaskState::Ready);
assert!(store.all_docs().unwrap().is_empty());
assert!(store.docs_for_task(1).unwrap().is_empty());
}
}