AkurAI Build
Menu

akurai-tasks

public

Latest 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);
    }
}