Menu
akurai-tasks
publicLatest change 2c9c32aa6fc1bd2551a09c6a1fc1e71b2dd98c79 - feat: complete and deploy AkurAI Tasks by Ólafur Búi Ólafsson
use super::*;
impl Store {
pub fn verify_lease(
&self,
item_id: &str,
principal: &str,
token: &str,
generation: i64,
) -> Result<bool> {
let mut db = self
.db
.lock()
.map_err(|_| Error::Storage("lock poisoned".into()))?;
let item = require_item(&mut db, item_id)?;
Ok(lease_matches(&item, principal, token, generation))
}
pub fn renew_lease(
&self,
item_id: &str,
token: &str,
generation: i64,
expected_revision: i64,
actor: &str,
idempotency: &str,
) -> Result<Value> {
let request = object(vec![
("command", string("renew_lease")),
("item", string(item_id)),
("generation", Value::Int(generation)),
("expectedRevision", Value::Int(expected_revision)),
]);
self.mutate(actor, idempotency, &request, |db, _| {
let mut item = require_item(db, item_id)?;
check_revision(&item, expected_revision)?;
require_lease(&item, actor, token, generation)?;
update(&mut item, "leaseExpiresAt", Value::Int(now() + 900));
bump(&mut item);
save_item(db, &item)?;
update(&mut item, "leaseTokenHash", Value::Null);
Ok(item)
})
}
pub fn release_lease(
&self,
item_id: &str,
token: &str,
generation: i64,
expected_revision: i64,
actor: &str,
idempotency: &str,
) -> Result<Value> {
let request = object(vec![
("command", string("release_lease")),
("item", string(item_id)),
("generation", Value::Int(generation)),
("expectedRevision", Value::Int(expected_revision)),
]);
self.mutate(actor, idempotency, &request, |db, _| {
let mut item = require_item(db, item_id)?;
check_revision(&item, expected_revision)?;
require_lease(&item, actor, token, generation)?;
clear_lease(&mut item);
if field(&item, "state") == Some("In Progress") {
update(&mut item, "state", string("Ready"));
}
bump(&mut item);
save_item(db, &item)?;
Ok(item)
})
}
pub fn reclaim_expired_lease(
&self,
item_id: &str,
expected_revision: i64,
actor: &str,
idempotency: &str,
) -> Result<Value> {
let request = object(vec![
("command", string("reclaim_expired_lease")),
("item", string(item_id)),
("expectedRevision", Value::Int(expected_revision)),
]);
self.mutate(actor, idempotency, &request, |db, _| {
let mut item = require_item(db, item_id)?;
check_revision(&item, expected_revision)?;
let expiry = item
.get("leaseExpiresAt")
.and_then(Value::as_i64)
.unwrap_or(0);
if item.get("owner").and_then(Value::as_str).is_none() || expiry > now() {
return Err(Error::Conflict("lease is not expired".into()));
}
clear_lease(&mut item);
if field(&item, "state") == Some("In Progress") {
update(&mut item, "state", string("Ready"));
}
bump(&mut item);
save_item(db, &item)?;
Ok(item)
})
}
}
pub(crate) fn require_lease(
item: &Value,
principal: &str,
token: &str,
generation: i64,
) -> Result<()> {
if lease_matches(item, principal, token, generation) {
Ok(())
} else {
Err(Error::Conflict(
"current lease token and generation required".into(),
))
}
}
fn lease_matches(item: &Value, principal: &str, token: &str, generation: i64) -> bool {
field(item, "owner") == Some(principal)
&& item.get("leaseGeneration").and_then(Value::as_i64) == Some(generation)
&& item
.get("leaseExpiresAt")
.and_then(Value::as_i64)
.is_some_and(|expiry| expiry > now())
&& field(item, "leaseTokenHash").is_some_and(|stored| {
auth::constant_eq(stored.as_bytes(), auth::digest(token).as_bytes())
})
}
fn clear_lease(item: &mut Value) {
let generation = item
.get("leaseGeneration")
.and_then(Value::as_i64)
.unwrap_or(0)
+ 1;
update(item, "owner", Value::Null);
update(item, "leaseTokenHash", Value::Null);
update(item, "leaseExpiresAt", Value::Null);
update(item, "leaseGeneration", Value::Int(generation));
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn lease_token_and_generation_fence_owner_mutations() {
let path = std::env::temp_dir().join(format!(
"akurai-tasks-lease-{}-{}.db",
std::process::id(),
now()
));
let _ = std::fs::remove_file(&path);
let store = Store::open(&path).unwrap();
store
.create_project("LEASE", "Lease", vec![], "admin", "project")
.unwrap();
let mut item = store
.create_item(CreateItem {
project: "LEASE",
title: "Fence",
description: "",
repo: None,
priority: "normal",
actor: "admin",
idempotency_key: "item",
})
.unwrap();
for (index, state) in ["Triage", "Discovery", "Ready"].iter().enumerate() {
item = store
.transition(
"LEASE-1",
state,
item.get("revision").and_then(Value::as_i64).unwrap(),
"admin",
&format!("move-{index}"),
)
.unwrap();
}
item = store
.claim(
"LEASE-1",
"agent",
item.get("revision").and_then(Value::as_i64).unwrap(),
"admin",
"claim",
)
.unwrap();
let token = field(&item, "leaseToken").unwrap().to_string();
let generation = item.get("leaseGeneration").and_then(Value::as_i64).unwrap();
assert!(store
.verify_lease("LEASE-1", "agent", &token, generation)
.unwrap());
assert!(!store
.verify_lease("LEASE-1", "agent", "wrong", generation)
.unwrap());
item = store
.renew_lease(
"LEASE-1",
&token,
generation,
item.get("revision").and_then(Value::as_i64).unwrap(),
"agent",
"renew",
)
.unwrap();
let released = store
.release_lease(
"LEASE-1",
&token,
generation,
item.get("revision").and_then(Value::as_i64).unwrap(),
"agent",
"release",
)
.unwrap();
assert_eq!(released.get("owner"), Some(&Value::Null));
assert!(!store
.verify_lease("LEASE-1", "agent", &token, generation)
.unwrap());
drop(store);
let _ = std::fs::remove_file(path);
}
}