Menu
akurai-tasks
publicLatest change bfff2c5d7515a248f06ff8fe7f0cec15982b3deb - Remove lease gate from mutation operations by Ólafur Búi Ólafsson
#![forbid(unsafe_code)]
mod auth;
mod import;
mod lease;
mod maintenance;
mod metrics;
mod outcomes;
mod query;
mod webhook;
pub use auth::Permission;
pub use import::{discover_sources, plan_file_board};
pub use lease::LeaseProof;
pub use maintenance::ExclusiveLock;
pub use metrics::MetricsQuery;
pub use outcomes::{CreateOutcome, OutcomePatch};
pub use query::Query;
pub use webhook::TransitionWebhook;
use akurai_json::{parse, Value};
use akurai_storage::BTree;
#[cfg(test)]
use std::cell::Cell;
use std::collections::BTreeSet;
use std::fmt;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use std::time::{SystemTime, UNIX_EPOCH};
#[cfg(test)]
thread_local! {
static FAIL_NEXT_COMMIT: Cell<bool> = const { Cell::new(false) };
}
pub const STATES: &[&str] = &[
"Inbox",
"Triage",
"Discovery",
"Ready",
"In Progress",
"Review",
"Verification",
"Release Ready",
"Done",
];
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Error {
Invalid(String),
Forbidden(String),
NotFound(String),
Conflict(String),
Storage(String),
}
impl fmt::Display for Error {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Invalid(message) => write!(f, "invalid: {message}"),
Self::Forbidden(message) => write!(f, "forbidden: {message}"),
Self::NotFound(message) => write!(f, "not found: {message}"),
Self::Conflict(message) => write!(f, "conflict: {message}"),
Self::Storage(message) => write!(f, "storage: {message}"),
}
}
}
impl std::error::Error for Error {}
pub type Result<T, E = Error> = std::result::Result<T, E>;
#[derive(Clone)]
pub struct Store {
db: Arc<Mutex<BTree>>,
path: Arc<PathBuf>,
}
impl Store {
pub fn open(path: impl AsRef<Path>) -> Result<Self> {
let path = path.as_ref().to_path_buf();
let db = BTree::open(&path).map_err(storage)?;
Ok(Self {
db: Arc::new(Mutex::new(db)),
path: Arc::new(path),
})
}
pub fn create_project(
&self,
key: &str,
name: &str,
repos: Vec<String>,
actor: &str,
idempotency_key: &str,
) -> Result<Value> {
validate_key(key)?;
required("project name", name)?;
required("actor", actor)?;
let repos = normalize_repos(repos)?;
let request = object(vec![
("command", string("create_project")),
("key", string(key)),
("name", string(name.trim())),
("repos", strings(&repos)),
]);
self.mutate(actor, idempotency_key, &request, |db, sequence| {
let project_key = format!("project/{key}");
if db.get(project_key.as_bytes()).map_err(storage)?.is_some() {
return Err(Error::Conflict(format!("project {key} already exists")));
}
let now = now();
let project = object(vec![
("id", string(key)),
("key", string(key)),
("name", string(name.trim())),
("repos", strings(&repos)),
("outcomes", Value::Array(Vec::new())),
("archived", Value::Bool(false)),
("revision", Value::Int(1)),
("createdAt", Value::Int(now)),
("updatedAt", Value::Int(now)),
]);
put(db, &project_key, &project)?;
put(db, &format!("project-order/{sequence:020}/{key}"), &project)?;
Ok(project)
})
}
pub fn list_projects(&self) -> Result<Value> {
self.scan("project/")
}
pub fn get_project(&self, key: &str) -> Result<Value> {
self.get(&format!("project/{key}"), format!("project {key}"))
}
pub fn update_project(
&self,
key: &str,
repos: Vec<String>,
actor: &str,
idempotency_key: &str,
) -> Result<Value> {
validate_key(key)?;
required("actor", actor)?;
let repos = normalize_repos(repos)?;
let request = object(vec![
("command", string("update_project")),
("key", string(key)),
("repos", strings(&repos)),
]);
self.mutate(actor, idempotency_key, &request, |db, _| {
let mut project = get_value(db, &format!("project/{key}"))?
.ok_or_else(|| Error::NotFound(format!("project {key}")))?;
update(&mut project, "repos", strings(&repos));
let revision = project.get("revision").and_then(Value::as_i64).unwrap_or(0) + 1;
update(&mut project, "revision", Value::Int(revision));
update(&mut project, "updatedAt", Value::Int(now()));
put(db, &format!("project/{key}"), &project)?;
Ok(project)
})
}
pub fn archive_project(
&self,
key: &str,
archived: bool,
actor: &str,
idempotency_key: &str,
) -> Result<Value> {
validate_key(key)?;
required("actor", actor)?;
let request = object(vec![
("command", string("archive_project")),
("key", string(key)),
("archived", Value::Bool(archived)),
]);
self.mutate(actor, idempotency_key, &request, |db, _| {
let mut project = get_value(db, &format!("project/{key}"))?
.ok_or_else(|| Error::NotFound(format!("project {key}")))?;
update(&mut project, "archived", Value::Bool(archived));
let revision = project.get("revision").and_then(Value::as_i64).unwrap_or(0) + 1;
update(&mut project, "revision", Value::Int(revision));
update(&mut project, "updatedAt", Value::Int(now()));
put(db, &format!("project/{key}"), &project)?;
Ok(project)
})
}
pub fn delete_project(&self, key: &str, actor: &str, idempotency_key: &str) -> Result<Value> {
validate_key(key)?;
required("actor", actor)?;
let request = object(vec![
("command", string("delete_project")),
("key", string(key)),
]);
self.mutate(actor, idempotency_key, &request, |db, _| {
get_value(db, &format!("project/{key}"))?
.ok_or_else(|| Error::NotFound(format!("project {key}")))?;
let item_prefix = format!("item-by-project/{key}/");
let items = db
.range(item_prefix.as_bytes(), &upper_bound(item_prefix.as_bytes()))
.map_err(storage)?;
if !items.is_empty() {
return Err(Error::Conflict(format!(
"project {key} has {} work item(s); delete or reassign them before removing the project",
items.len()
)));
}
db.delete(format!("project/{key}").as_bytes())
.map_err(storage)?;
let suffix = format!("/{key}");
let order_entries = db
.range(b"project-order/", &upper_bound(b"project-order/"))
.map_err(storage)?;
for (order_key, _) in order_entries {
if order_key.ends_with(suffix.as_bytes()) {
db.delete(&order_key).map_err(storage)?;
}
}
Ok(object(vec![
("key", string(key)),
("deleted", Value::Bool(true)),
]))
})
}
/// Delete a project together with all its work items and their per-item
/// sub-records (dependencies, handoffs), the per-project item/event indexes,
/// and the item counter. Destructive and irreversible.
pub fn purge_project(&self, key: &str, actor: &str, idempotency_key: &str) -> Result<Value> {
validate_key(key)?;
required("actor", actor)?;
let request = object(vec![
("command", string("purge_project")),
("key", string(key)),
]);
self.mutate(actor, idempotency_key, &request, |db, _| {
get_value(db, &format!("project/{key}"))?
.ok_or_else(|| Error::NotFound(format!("project {key}")))?;
let item_prefix = format!("item-by-project/{key}/");
let index = db
.range(item_prefix.as_bytes(), &upper_bound(item_prefix.as_bytes()))
.map_err(storage)?;
let mut item_ids: Vec<String> = Vec::new();
for (_, bytes) in &index {
if let Some(id) = field(&decode(bytes)?, "id") {
item_ids.push(id.to_string());
}
}
let mut removed = 0i64;
for id in &item_ids {
db.delete(format!("item/{id}").as_bytes())
.map_err(storage)?;
for prefix in [format!("dependency/{id}/"), format!("handoff/{id}/")] {
let entries = db
.range(prefix.as_bytes(), &upper_bound(prefix.as_bytes()))
.map_err(storage)?;
for (edge_key, _) in entries {
db.delete(&edge_key).map_err(storage)?;
}
}
let pending = format!("handoff-pending/{id}");
if db.get(pending.as_bytes()).map_err(storage)?.is_some() {
db.delete(pending.as_bytes()).map_err(storage)?;
}
removed += 1;
}
for prefix in [
format!("item-by-project/{key}/"),
format!("event-by-project/{key}/"),
] {
let entries = db
.range(prefix.as_bytes(), &upper_bound(prefix.as_bytes()))
.map_err(storage)?;
for (scoped_key, _) in entries {
db.delete(&scoped_key).map_err(storage)?;
}
}
let counter = format!("counter/item/{key}");
if db.get(counter.as_bytes()).map_err(storage)?.is_some() {
db.delete(counter.as_bytes()).map_err(storage)?;
}
db.delete(format!("project/{key}").as_bytes())
.map_err(storage)?;
let suffix = format!("/{key}");
let order_entries = db
.range(b"project-order/", &upper_bound(b"project-order/"))
.map_err(storage)?;
for (order_key, _) in order_entries {
if order_key.ends_with(suffix.as_bytes()) {
db.delete(&order_key).map_err(storage)?;
}
}
Ok(object(vec![
("key", string(key)),
("purged", Value::Bool(true)),
("itemsRemoved", Value::Int(removed)),
]))
})
}
pub fn set_item_repo(
&self,
id: &str,
repo: Option<&str>,
_lease: Option<LeaseProof<'_>>,
actor: &str,
idempotency_key: &str,
) -> Result<Value> {
required("actor", actor)?;
let request = object(vec![
("command", string("set_item_repo")),
("id", string(id)),
("repo", optional_string(repo)),
]);
self.mutate(actor, idempotency_key, &request, |db, _| {
let mut item = require_item(db, id)?;
let project_key = field(&item, "project")
.ok_or_else(|| Error::Storage("item project missing".into()))?
.to_string();
let project = get_value(db, &format!("project/{project_key}"))?
.ok_or_else(|| Error::NotFound(format!("project {project_key}")))?;
let resolved = resolve_repo(&project, repo)?;
update(&mut item, "repo", optional_string(resolved.as_deref()));
bump(&mut item);
save_item(db, &item)?;
Ok(item)
})
}
pub fn create_item(&self, input: CreateItem<'_>) -> Result<Value> {
validate_key(input.project)?;
required("title", input.title)?;
required("actor", input.actor)?;
let priority = validate_priority(input.priority)?;
let request = object(vec![
("command", string("create_item")),
("project", string(input.project)),
("title", string(input.title.trim())),
("repo", optional_string(input.repo)),
("priority", string(input.priority)),
]);
self.mutate(
input.actor,
input.idempotency_key,
&request,
|db, sequence| {
let project = get_value(db, &format!("project/{}", input.project))?
.ok_or_else(|| Error::NotFound(format!("project {}", input.project)))?;
let resolved_repo = resolve_repo(&project, input.repo)?;
let item_sequence = next_counter(db, &format!("counter/item/{}", input.project))?;
let id = format!("{}-{item_sequence}", input.project);
let now = now();
let item = object(vec![
("id", string(&id)),
("project", string(input.project)),
("sequence", Value::Int(item_sequence)),
("title", string(input.title.trim())),
("description", string(input.description.trim())),
("repo", optional_string(resolved_repo.as_deref())),
("state", string("Inbox")),
("priority", string(priority)),
("type", string("task")),
("labels", Value::Array(Vec::new())),
("parent", Value::Null),
("leaseGeneration", Value::Int(0)),
("blocked", Value::Bool(false)),
("blockReason", Value::Null),
("blockedAt", Value::Null),
("outcomes", Value::Array(Vec::new())),
("revision", Value::Int(1)),
("createdAt", Value::Int(now)),
("updatedAt", Value::Int(now)),
]);
put(db, &format!("item/{id}"), &item)?;
put(
db,
&format!("item-by-project/{}/{item_sequence:020}/{id}", input.project),
&item,
)?;
let _ = sequence;
Ok(item)
},
)
}
pub fn get_item(&self, id: &str) -> Result<Value> {
self.get(&format!("item/{id}"), format!("item {id}"))
}
pub fn list_items(&self, project: &str) -> Result<Value> {
validate_key(project)?;
self.scan(&format!("item-by-project/{project}/"))
}
pub fn board(&self, project: &str) -> Result<Value> {
let Value::Array(items) = self.list_items(project)? else {
unreachable!();
};
let columns = STATES
.iter()
.map(|state| {
let cards = items
.iter()
.filter(|item| field(item, "state") == Some(*state))
.cloned()
.collect();
object(vec![
("state", string(state)),
("items", Value::Array(cards)),
])
})
.collect();
Ok(object(vec![
("project", self.get_project(project)?),
("columns", Value::Array(columns)),
(
"workflow",
object(vec![
(
"states",
Value::Array(STATES.iter().map(|state| string(state)).collect()),
),
("claimState", string("Ready")),
("claimTarget", string("In Progress")),
]),
),
]))
}
pub fn transition(
&self,
id: &str,
target: &str,
expected_revision: i64,
lease: Option<LeaseProof<'_>>,
actor: &str,
idempotency_key: &str,
) -> Result<Value> {
state_index(target)?; // validate the target is a known state
let request = object(vec![
("command", string("transition")),
("id", string(id)),
("target", string(target)),
("expectedRevision", Value::Int(expected_revision)),
(
"leaseGeneration",
lease
.map(|proof| Value::Int(proof.generation))
.unwrap_or(Value::Null),
),
]);
self.mutate(actor, idempotency_key, &request, |db, _| {
let mut item = require_item(db, id)?;
// Free movement: any known state → any known state. Operators and
// agents may move items backward, skip stages, or set In Progress
// directly. Leases are advisory (for worker coordination via claim);
// mutation operations no longer require a lease proof.
// gates (evidence, dependencies) still apply via enforce_transition.
enforce_transition(db, &item, target)?;
update(&mut item, "state", string(target));
bump(&mut item);
save_item(db, &item)?;
Ok(item)
})
}
pub fn claim(
&self,
id: &str,
owner: &str,
expected_revision: i64,
actor: &str,
idempotency_key: &str,
) -> Result<Value> {
required("owner", owner)?;
let token = auth::stable_token("claim-lease", actor, idempotency_key, &[id, owner], 32)?;
let token_hash = auth::digest(&token);
let request = object(vec![
("command", string("claim")),
("id", string(id)),
("owner", string(owner)),
("expectedRevision", Value::Int(expected_revision)),
]);
let mut result = self.mutate(actor, idempotency_key, &request, |db, _| {
let mut item = require_item(db, id)?;
check_revision(&item, expected_revision)?;
let owner_present = item.get("owner").and_then(Value::as_str).is_some();
let lease_active = item
.get("leaseExpiresAt")
.and_then(Value::as_i64)
.is_some_and(|expiry| expiry > now());
if owner_present && lease_active {
return Err(Error::Conflict(format!("item {id} is already claimed")));
}
// Any item without an ACTIVE lease may be (re)claimed from any
// state — the already-claimed guard above still protects a live
// lease. Lets a stranded mid-pipeline item be picked back up
// without first walking it back to Ready.
ensure_dependencies_done(db, id)?;
let generation = item
.get("leaseGeneration")
.and_then(Value::as_i64)
.unwrap_or(0)
+ 1;
update(&mut item, "owner", string(owner));
update(&mut item, "leaseGeneration", Value::Int(generation));
update(&mut item, "leaseTokenHash", string(&token_hash));
update(&mut item, "leaseExpiresAt", Value::Int(now() + 900));
update(&mut item, "state", string("In Progress"));
bump(&mut item);
save_item(db, &item)?;
Ok(item)
})?;
if field(&result, "leaseTokenHash") == Some(&token_hash) {
update(&mut result, "leaseToken", string(&token));
}
update(&mut result, "leaseTokenHash", Value::Null);
Ok(result)
}
pub fn add_dependency(
&self,
item_id: &str,
depends_on: &str,
actor: &str,
idempotency_key: &str,
) -> Result<Value> {
if item_id == depends_on {
return Err(Error::Invalid("an item cannot depend on itself".into()));
}
let request = object(vec![
("command", string("add_dependency")),
("item", string(item_id)),
("dependsOn", string(depends_on)),
]);
self.mutate(actor, idempotency_key, &request, |db, _| {
let item = require_item(db, item_id)?;
let dependency = require_item(db, depends_on)?;
if field(&item, "project") != field(&dependency, "project") {
return Err(Error::Invalid("dependencies must be in one project".into()));
}
if dependency_reaches(db, depends_on, item_id, &mut BTreeSet::new())? {
return Err(Error::Invalid("dependency would create a cycle".into()));
}
let edge = object(vec![
("item", string(item_id)),
("dependsOn", string(depends_on)),
("createdAt", Value::Int(now())),
]);
put(db, &format!("dependency/{item_id}/{depends_on}"), &edge)?;
Ok(edge)
})
}
pub fn add_record(
&self,
kind: &str,
item_id: &str,
body: &str,
lease: Option<LeaseProof<'_>>,
actor: &str,
idempotency_key: &str,
) -> Result<Value> {
if !matches!(kind, "comment" | "evidence") {
return Err(Error::Invalid(
"record kind must be comment or evidence".into(),
));
}
required("body", body)?;
let request = object(vec![
("command", string("add_record")),
("kind", string(kind)),
("item", string(item_id)),
("body", string(body.trim())),
(
"leaseGeneration",
lease
.map(|proof| Value::Int(proof.generation))
.unwrap_or(Value::Null),
),
]);
self.mutate(actor, idempotency_key, &request, |db, sequence| {
require_item(db, item_id)?;
let record = object(vec![
("id", Value::Int(sequence)),
("kind", string(kind)),
("item", string(item_id)),
("body", string(body.trim())),
("actor", string(actor)),
("createdAt", Value::Int(now())),
]);
put(db, &format!("{kind}/{item_id}/{sequence:020}"), &record)?;
Ok(record)
})
}
pub fn handoff(
&self,
item_id: &str,
to: &str,
summary: &str,
lease: Option<LeaseProof<'_>>,
actor: &str,
idempotency_key: &str,
) -> Result<Value> {
required("receiver", to)?;
required("summary", summary)?;
let request = object(vec![
("command", string("handoff")),
("item", string(item_id)),
("to", string(to)),
("summary", string(summary.trim())),
(
"leaseGeneration",
lease
.map(|proof| Value::Int(proof.generation))
.unwrap_or(Value::Null),
),
]);
self.mutate(actor, idempotency_key, &request, |db, sequence| {
let item = require_item(db, item_id)?;
if field(&item, "owner") != Some(actor) {
return Err(Error::Conflict(
"only the current owner can hand off".into(),
));
}
let record = object(vec![
("id", Value::Int(sequence)),
("item", string(item_id)),
("from", string(actor)),
("to", string(to)),
("summary", string(summary.trim())),
("status", string("pending")),
("createdAt", Value::Int(now())),
]);
put(db, &format!("handoff/{item_id}/{sequence:020}"), &record)?;
put(db, &format!("handoff-pending/{item_id}"), &record)?;
Ok(record)
})
}
pub fn accept_handoff(
&self,
item_id: &str,
actor: &str,
idempotency_key: &str,
) -> Result<Value> {
let token = auth::stable_token("handoff-lease", actor, idempotency_key, &[item_id], 32)?;
let token_hash = auth::digest(&token);
let request = object(vec![
("command", string("accept_handoff")),
("item", string(item_id)),
]);
let mut result = self.mutate(actor, idempotency_key, &request, |db, _| {
let key = format!("handoff-pending/{item_id}");
let mut handoff = get_value(db, &key)?
.ok_or_else(|| Error::NotFound(format!("pending handoff for {item_id}")))?;
if field(&handoff, "to") != Some(actor) {
return Err(Error::Conflict(
"handoff receiver does not match actor".into(),
));
}
let mut item = require_item(db, item_id)?;
update(&mut item, "owner", string(actor));
let generation = item
.get("leaseGeneration")
.and_then(Value::as_i64)
.unwrap_or(0)
+ 1;
update(&mut item, "leaseGeneration", Value::Int(generation));
update(&mut item, "leaseTokenHash", string(&token_hash));
update(&mut item, "leaseExpiresAt", Value::Int(now() + 900));
bump(&mut item);
save_item(db, &item)?;
update(&mut handoff, "status", string("accepted"));
update(&mut handoff, "acceptedAt", Value::Int(now()));
put(
db,
&format!(
"handoff/{item_id}/{:020}",
handoff.get("id").and_then(Value::as_i64).unwrap_or(0)
),
&handoff,
)?;
db.delete(key.as_bytes()).map_err(storage)?;
Ok(item)
})?;
if field(&result, "leaseTokenHash") == Some(&token_hash) {
update(&mut result, "leaseToken", string(&token));
}
update(&mut result, "leaseTokenHash", Value::Null);
Ok(result)
}
pub fn set_blocked(
&self,
item_id: &str,
reason: Option<&str>,
expected_revision: i64,
lease: Option<LeaseProof<'_>>,
actor: &str,
idempotency_key: &str,
) -> Result<Value> {
if let Some(reason) = reason {
required("block reason", reason)?;
}
let request = object(vec![
("command", string("set_blocked")),
("item", string(item_id)),
("reason", optional_string(reason)),
("expectedRevision", Value::Int(expected_revision)),
(
"leaseGeneration",
lease
.map(|proof| Value::Int(proof.generation))
.unwrap_or(Value::Null),
),
]);
self.mutate(actor, idempotency_key, &request, |db, _| {
let mut item = require_item(db, item_id)?;
check_revision(&item, expected_revision)?;
update(&mut item, "blocked", Value::Bool(reason.is_some()));
update(&mut item, "blockReason", optional_string(reason));
update(
&mut item,
"blockedAt",
if reason.is_some() {
Value::Int(now())
} else {
Value::Null
},
);
bump(&mut item);
save_item(db, &item)?;
Ok(item)
})
}
pub fn ready_frontier(&self, project: &str) -> Result<Value> {
let Value::Array(items) = self.list_items(project)? else {
unreachable!();
};
let db = self
.db
.lock()
.map_err(|_| Error::Storage("lock poisoned".into()))?;
let mut db = db;
let mut ready = Vec::new();
for item in items {
if field(&item, "state") == Some("Ready")
&& item.get("blocked").and_then(Value::as_bool) != Some(true)
&& ensure_dependencies_done(&mut db, field(&item, "id").unwrap_or("")).is_ok()
{
ready.push(item);
}
}
Ok(Value::Array(ready))
}
pub fn events(&self, project: Option<&str>) -> Result<Value> {
match project {
Some(project) => self.scan(&format!("event-by-project/{project}/")),
None => self.scan("event/"),
}
}
fn get(&self, key: &str, label: String) -> Result<Value> {
let mut db = self
.db
.lock()
.map_err(|_| Error::Storage("lock poisoned".into()))?;
get_value(&mut db, key)?.ok_or(Error::NotFound(label))
}
fn scan(&self, prefix: &str) -> Result<Value> {
let mut db = self
.db
.lock()
.map_err(|_| Error::Storage("lock poisoned".into()))?;
let entries = db
.range(prefix.as_bytes(), &upper_bound(prefix.as_bytes()))
.map_err(storage)?;
let mut values = Vec::with_capacity(entries.len());
for (_, bytes) in entries {
values.push(decode(&bytes)?);
}
Ok(Value::Array(values))
}
fn mutate<F>(
&self,
actor: &str,
idempotency_key: &str,
request: &Value,
apply: F,
) -> Result<Value>
where
F: FnOnce(&mut BTree, i64) -> Result<Value>,
{
required("actor", actor)?;
required("idempotency key", idempotency_key)?;
let request_json = request.to_json();
let idem_key = format!("idempotency/{actor}/{idempotency_key}");
let mut db = self
.db
.lock()
.map_err(|_| Error::Storage("lock poisoned".into()))?;
auth::authorize_mutation(&mut db, actor, request)?;
if let Some(existing) = get_value(&mut db, &idem_key)? {
if field(&existing, "request") != Some(&request_json) {
return Err(Error::Conflict(
"idempotency key reused for another request".into(),
));
}
return existing
.get("result")
.cloned()
.ok_or_else(|| Error::Storage("idempotency result missing".into()));
}
let attempt = (|| {
let sequence = next_counter(&mut db, "counter/event")?;
let result = apply(&mut db, sequence)?;
let project = result
.get("project")
.and_then(Value::as_str)
.or_else(|| request.get("project").and_then(Value::as_str))
.unwrap_or("system");
let event = object(vec![
("sequence", Value::Int(sequence)),
("project", string(project)),
("actor", string(actor)),
(
"command",
request.get("command").cloned().unwrap_or(Value::Null),
),
("request", request.clone()),
("result", result.clone()),
("createdAt", Value::Int(now())),
]);
put(&mut db, &format!("event/{sequence:020}"), &event)?;
put(
&mut db,
&format!("event-by-project/{project}/{sequence:020}"),
&event,
)?;
put(
&mut db,
&idem_key,
&object(vec![
("request", string(&request_json)),
("result", result.clone()),
("event", Value::Int(sequence)),
]),
)?;
#[cfg(test)]
if FAIL_NEXT_COMMIT.with(|fail| fail.replace(false)) {
return Err(Error::Storage("injected commit failure".into()));
}
db.commit().map_err(storage)?;
Ok(result)
})();
match attempt {
Ok(result) => Ok(result),
Err(error) => {
*db = BTree::open(self.path.as_path()).map_err(storage)?;
Err(error)
}
}
}
}
pub struct CreateItem<'a> {
pub project: &'a str,
pub title: &'a str,
pub description: &'a str,
pub repo: Option<&'a str>,
pub priority: &'a str,
pub actor: &'a str,
pub idempotency_key: &'a str,
}
fn enforce_transition(db: &mut BTree, item: &Value, target: &str) -> Result<()> {
let id = field(item, "id").unwrap_or("");
// `blocked` is advisory metadata, not a movement gate: the autonomous
// worker already skips blocked items at intake, and operators/agents must
// stay free to move a blocked item wherever they need it.
if target == "In Progress" {
ensure_dependencies_done(db, id)?;
}
if target == "Review" && !has_records(db, "evidence", id)? {
return Err(Error::Invalid(
"implementation evidence is required for Review".into(),
));
}
if target == "Done" && !has_records(db, "evidence", id)? {
return Err(Error::Invalid("evidence is required for Done".into()));
}
Ok(())
}
fn has_records(db: &mut BTree, kind: &str, item: &str) -> Result<bool> {
let prefix = format!("{kind}/{item}/");
Ok(!db
.range(prefix.as_bytes(), &upper_bound(prefix.as_bytes()))
.map_err(storage)?
.is_empty())
}
fn ensure_dependencies_done(db: &mut BTree, item: &str) -> Result<()> {
let prefix = format!("dependency/{item}/");
for (_, bytes) in db
.range(prefix.as_bytes(), &upper_bound(prefix.as_bytes()))
.map_err(storage)?
{
let edge = decode(&bytes)?;
let dependency = field(&edge, "dependsOn").unwrap_or("");
let dependency_item = require_item(db, dependency)?;
if field(&dependency_item, "state") != Some("Done") {
return Err(Error::Invalid(format!(
"dependency {dependency} is not Done"
)));
}
}
Ok(())
}
fn dependency_reaches(
db: &mut BTree,
start: &str,
target: &str,
visited: &mut BTreeSet<String>,
) -> Result<bool> {
if start == target {
return Ok(true);
}
if !visited.insert(start.to_string()) {
return Ok(false);
}
let prefix = format!("dependency/{start}/");
for (_, bytes) in db
.range(prefix.as_bytes(), &upper_bound(prefix.as_bytes()))
.map_err(storage)?
{
let edge = decode(&bytes)?;
if dependency_reaches(db, field(&edge, "dependsOn").unwrap_or(""), target, visited)? {
return Ok(true);
}
}
Ok(false)
}
fn resolve_repo(project: &Value, repo: Option<&str>) -> Result<Option<String>> {
let repos: Vec<&str> = match project.get("repos") {
Some(Value::Array(items)) => items.iter().filter_map(Value::as_str).collect(),
_ => Vec::new(),
};
match repo {
Some(candidate) => {
let candidate = candidate.trim();
required("repo", candidate)?;
if repos.contains(&candidate) {
Ok(Some(candidate.to_string()))
} else {
Err(Error::Invalid(format!(
"repo {candidate} is not registered to the project"
)))
}
}
None => match repos.as_slice() {
[] => Ok(None),
[single] => Ok(Some((*single).to_string())),
_ => Err(Error::Invalid(
"repo is required: project registers multiple repos".into(),
)),
},
}
}
fn normalize_repos(repos: Vec<String>) -> Result<Vec<String>> {
let mut unique = BTreeSet::new();
for repo in repos {
let repo = repo.trim();
required("repo", repo)?;
if repo.contains(char::is_whitespace) {
return Err(Error::Invalid(format!("repo contains whitespace: {repo}")));
}
unique.insert(repo.to_string());
}
Ok(unique.into_iter().collect())
}
fn validate_key(key: &str) -> Result<()> {
if key.is_empty()
|| key.len() > 24
|| !key
.bytes()
.all(|byte| byte.is_ascii_uppercase() || byte.is_ascii_digit() || byte == b'-')
{
return Err(Error::Invalid(
"project key must be 1-24 uppercase letters, digits, or '-'".into(),
));
}
Ok(())
}
fn validate_priority(priority: &str) -> Result<&str> {
if matches!(priority, "low" | "normal" | "high" | "critical") {
Ok(priority)
} else {
Err(Error::Invalid(
"priority must be low, normal, high, or critical".into(),
))
}
}
fn required(label: &str, value: &str) -> Result<()> {
if value.trim().is_empty() {
Err(Error::Invalid(format!("{label} is required")))
} else {
Ok(())
}
}
fn state_index(state: &str) -> Result<usize> {
STATES
.iter()
.position(|candidate| *candidate == state)
.ok_or_else(|| Error::Invalid(format!("unknown state {state}")))
}
fn check_revision(item: &Value, expected: i64) -> Result<()> {
let current = item.get("revision").and_then(Value::as_i64).unwrap_or(0);
if current == expected {
Ok(())
} else {
Err(Error::Conflict(format!(
"stale revision {expected}; current revision is {current}"
)))
}
}
fn require_item(db: &mut BTree, id: &str) -> Result<Value> {
get_value(db, &format!("item/{id}"))?.ok_or_else(|| Error::NotFound(format!("item {id}")))
}
fn save_item(db: &mut BTree, item: &Value) -> Result<()> {
let id = field(item, "id").ok_or_else(|| Error::Storage("item id missing".into()))?;
let project =
field(item, "project").ok_or_else(|| Error::Storage("item project missing".into()))?;
let sequence = item
.get("sequence")
.and_then(Value::as_i64)
.ok_or_else(|| Error::Storage("item sequence missing".into()))?;
put(db, &format!("item/{id}"), item)?;
put(
db,
&format!("item-by-project/{project}/{sequence:020}/{id}"),
item,
)
}
fn bump(value: &mut Value) {
let revision = value.get("revision").and_then(Value::as_i64).unwrap_or(0) + 1;
update(value, "revision", Value::Int(revision));
update(value, "updatedAt", Value::Int(now()));
}
fn update(object: &mut Value, key: &str, value: Value) {
let Value::Object(fields) = object else {
return;
};
if let Some((_, existing)) = fields.iter_mut().find(|(name, _)| name == key) {
*existing = value;
} else {
fields.push((key.to_string(), value));
}
}
fn next_counter(db: &mut BTree, key: &str) -> Result<i64> {
let current = db
.get(key.as_bytes())
.map_err(storage)?
.map(|bytes| String::from_utf8_lossy(&bytes).parse::<i64>())
.transpose()
.map_err(|_| Error::Storage(format!("invalid counter {key}")))?
.unwrap_or(0);
let next = current
.checked_add(1)
.ok_or_else(|| Error::Storage(format!("counter overflow {key}")))?;
db.insert(key.as_bytes(), next.to_string().as_bytes())
.map_err(storage)?;
Ok(next)
}
fn put(db: &mut BTree, key: &str, value: &Value) -> Result<()> {
db.insert(key.as_bytes(), value.to_json().as_bytes())
.map_err(storage)
}
fn get_value(db: &mut BTree, key: &str) -> Result<Option<Value>> {
db.get(key.as_bytes())
.map_err(storage)?
.map(|bytes| decode(&bytes))
.transpose()
}
fn decode(bytes: &[u8]) -> Result<Value> {
parse(&String::from_utf8_lossy(bytes)).map_err(|error| Error::Storage(error.to_string()))
}
fn storage(error: std::io::Error) -> Error {
Error::Storage(error.to_string())
}
fn upper_bound(prefix: &[u8]) -> Vec<u8> {
let mut bound = prefix.to_vec();
bound.push(0xff);
bound
}
fn now() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs() as i64
}
pub fn object(fields: Vec<(&str, Value)>) -> Value {
Value::Object(
fields
.into_iter()
.map(|(key, value)| (key.to_string(), value))
.collect(),
)
}
pub fn string(value: &str) -> Value {
Value::Str(value.to_string())
}
pub fn optional_string(value: Option<&str>) -> Value {
value.map(string).unwrap_or(Value::Null)
}
pub fn strings(values: &[String]) -> Value {
Value::Array(values.iter().map(|value| string(value)).collect())
}
pub fn field<'a>(value: &'a Value, key: &str) -> Option<&'a str> {
value.get(key).and_then(Value::as_str)
}
#[cfg(test)]
mod tests {
use super::*;
fn store(name: &str) -> (Store, std::path::PathBuf) {
let path = std::env::temp_dir().join(format!(
"akurai-tasks-{name}-{}-{}.db",
std::process::id(),
now()
));
let _ = std::fs::remove_file(&path);
(Store::open(&path).unwrap(), path)
}
#[test]
fn supports_multi_project_repo_boards_and_idempotency() {
let (store, path) = store("multi");
let project = store
.create_project(
"OPS",
"Operations",
vec!["olibuijr/api".into(), "olibuijr/web".into()],
"admin",
"project-1",
)
.unwrap();
assert_eq!(field(&project, "key"), Some("OPS"));
let item = store
.create_item(CreateItem {
project: "OPS",
title: "Ship API",
description: "",
repo: Some("olibuijr/api"),
priority: "high",
actor: "admin",
idempotency_key: "item-1",
})
.unwrap();
let replay = store
.create_item(CreateItem {
project: "OPS",
title: "Ship API",
description: "",
repo: Some("olibuijr/api"),
priority: "high",
actor: "admin",
idempotency_key: "item-1",
})
.unwrap();
assert_eq!(item, replay);
assert_eq!(store.list_items("OPS").unwrap(), Value::Array(vec![item]));
drop(store);
let _ = std::fs::remove_file(path);
}
#[test]
fn repo_is_resolved_validated_and_backfillable() {
let (store, path) = store("repos");
store
.create_project("SVC", "Service", vec!["owner/api".into()], "admin", "p")
.unwrap();
// Single-repo project auto-fills the item repo when none is given.
let auto = store
.create_item(CreateItem {
project: "SVC",
title: "Auto",
description: "",
repo: None,
priority: "normal",
actor: "admin",
idempotency_key: "auto",
})
.unwrap();
assert_eq!(field(&auto, "repo"), Some("owner/api"));
// An unregistered repo is rejected.
assert!(store
.create_item(CreateItem {
project: "SVC",
title: "Bad",
description: "",
repo: Some("owner/ghost"),
priority: "normal",
actor: "admin",
idempotency_key: "bad",
})
.is_err());
// Registering a second repo makes an omitted repo ambiguous.
store
.update_project(
"SVC",
vec!["owner/api".into(), "owner/web".into()],
"admin",
"u",
)
.unwrap();
assert!(store
.create_item(CreateItem {
project: "SVC",
title: "Ambiguous",
description: "",
repo: None,
priority: "normal",
actor: "admin",
idempotency_key: "ambig",
})
.is_err());
// set_item_repo moves an item onto a newly registered repo.
let moved = store
.set_item_repo("SVC-1", Some("owner/web"), None, "admin", "move")
.unwrap();
assert_eq!(field(&moved, "repo"), Some("owner/web"));
assert!(store
.set_item_repo("SVC-1", Some("owner/ghost"), None, "admin", "move2")
.is_err());
drop(store);
let _ = std::fs::remove_file(path);
}
#[test]
fn rejects_stale_revision_and_dependency_cycles() {
let (store, path) = store("rules");
store
.create_project("APP", "App", vec![], "admin", "p")
.unwrap();
for (key, title) in [("a", "A"), ("b", "B")] {
store
.create_item(CreateItem {
project: "APP",
title,
description: "",
repo: None,
priority: "normal",
actor: "admin",
idempotency_key: key,
})
.unwrap();
}
store
.add_dependency("APP-1", "APP-2", "admin", "dep-1")
.unwrap();
assert!(store
.add_dependency("APP-2", "APP-1", "admin", "dep-2")
.is_err());
assert!(store
.transition("APP-1", "Triage", 9, None, "admin", "move")
.is_err());
drop(store);
let _ = std::fs::remove_file(path);
}
#[test]
fn ready_claim_evidence_and_handoff_follow_one_command_path() {
let (store, path) = store("coordination");
store
.create_project(
"FLOW",
"Flow",
vec!["owner/repo".into()],
"admin",
"project",
)
.unwrap();
let mut item = store
.create_item(CreateItem {
project: "FLOW",
title: "Coordinated task",
description: "",
repo: Some("owner/repo"),
priority: "normal",
actor: "admin",
idempotency_key: "item",
})
.unwrap();
for (index, state) in ["Triage", "Discovery", "Ready"].iter().enumerate() {
item = store
.transition(
"FLOW-1",
state,
item.get("revision").and_then(Value::as_i64).unwrap(),
None,
"admin",
&format!("move-{index}"),
)
.unwrap();
}
assert_eq!(
store.ready_frontier("FLOW").unwrap(),
Value::Array(vec![item.clone()])
);
item = store
.claim(
"FLOW-1",
"agent-a",
item.get("revision").and_then(Value::as_i64).unwrap(),
"admin",
"claim",
)
.unwrap();
assert_eq!(field(&item, "state"), Some("In Progress"));
let proof = LeaseProof {
token: field(&item, "leaseToken").unwrap(),
generation: item.get("leaseGeneration").and_then(Value::as_i64).unwrap(),
};
store
.add_record(
"evidence",
"FLOW-1",
"cargo test: pass",
Some(proof),
"agent-a",
"evidence",
)
.unwrap();
let handoff = store
.handoff(
"FLOW-1",
"agent-b",
"Implementation complete",
Some(proof),
"agent-a",
"handoff",
)
.unwrap();
assert_eq!(field(&handoff, "status"), Some("pending"));
let accepted = store.accept_handoff("FLOW-1", "agent-b", "accept").unwrap();
assert_eq!(field(&accepted, "owner"), Some("agent-b"));
assert!(
accepted
.get("leaseGeneration")
.and_then(Value::as_i64)
.unwrap()
> item.get("leaseGeneration").and_then(Value::as_i64).unwrap()
);
drop(store);
let _ = std::fs::remove_file(path);
}
#[test]
fn archive_and_delete_project_lifecycle() {
let (store, path) = store("archive");
store
.create_project(
"OPS",
"Operations",
vec!["olibuijr/api".into()],
"admin",
"p1",
)
.unwrap();
// Archive is reversible and flips the flag + bumps revision.
let archived = store.archive_project("OPS", true, "admin", "a1").unwrap();
assert_eq!(
archived.get("archived").and_then(Value::as_bool),
Some(true)
);
let restored = store.archive_project("OPS", false, "admin", "a2").unwrap();
assert_eq!(
restored.get("archived").and_then(Value::as_bool),
Some(false)
);
// A project with work items refuses deletion.
store
.create_item(CreateItem {
project: "OPS",
title: "Ship API",
description: "",
repo: Some("olibuijr/api"),
priority: "high",
actor: "admin",
idempotency_key: "i1",
})
.unwrap();
assert!(store.delete_project("OPS", "admin", "d1").is_err());
// An empty project deletes and then reads as not-found.
store
.create_project("TMP", "Temp", vec![], "admin", "p2")
.unwrap();
let deleted = store.delete_project("TMP", "admin", "d2").unwrap();
assert_eq!(deleted.get("deleted").and_then(Value::as_bool), Some(true));
assert!(store.get_project("TMP").is_err());
assert!(store.delete_project("TMP", "admin", "d3").is_err());
drop(store);
let _ = std::fs::remove_file(path);
}
#[test]
fn purge_project_removes_items_and_project() {
let (store, path) = store("purge");
store
.create_project("PRG", "Purge", vec!["olibuijr/api".into()], "admin", "p1")
.unwrap();
let item = store
.create_item(CreateItem {
project: "PRG",
title: "Scaffold",
description: "",
repo: Some("olibuijr/api"),
priority: "normal",
actor: "admin",
idempotency_key: "i1",
})
.unwrap();
let item_id = field(&item, "id").unwrap().to_string();
// Plain delete refuses while items exist; purge cascades.
assert!(store.delete_project("PRG", "admin", "d1").is_err());
let purged = store.purge_project("PRG", "admin", "pg1").unwrap();
assert_eq!(purged.get("purged").and_then(Value::as_bool), Some(true));
assert_eq!(purged.get("itemsRemoved").and_then(Value::as_i64), Some(1));
assert!(store.get_project("PRG").is_err());
assert!(store.get_item(&item_id).is_err());
// A fresh project reusing the key starts empty (indexes cleared).
store
.create_project("PRG", "Purge", vec![], "admin", "p2")
.unwrap();
assert_eq!(store.list_items("PRG").unwrap(), Value::Array(vec![]));
drop(store);
let _ = std::fs::remove_file(path);
}
#[test]
fn failed_commit_reopens_the_durable_root_before_the_next_mutation() {
let (store, path) = store("commit-recovery");
FAIL_NEXT_COMMIT.with(|fail| fail.set(true));
let failed = store.create_project("LOST", "Lost", vec![], "admin", "failed");
assert!(matches!(failed, Err(Error::Storage(_))));
assert!(store.get_project("LOST").is_err());
store
.create_project("KEPT", "Kept", vec![], "admin", "successful")
.unwrap();
assert!(store.get_project("LOST").is_err());
assert_eq!(
field(&store.get_project("KEPT").unwrap(), "name"),
Some("Kept")
);
drop(store);
let reopened = Store::open(&path).unwrap();
assert!(reopened.get_project("LOST").is_err());
assert_eq!(
field(&reopened.get_project("KEPT").unwrap(), "name"),
Some("Kept")
);
drop(reopened);
let _ = std::fs::remove_file(path);
}
}