agentplane 0.26.0

Durable, replayable agent runtime — the journal is the plan of record
Documentation
//! The worklist contract.

use std::fmt::Debug;

use async_trait::async_trait;

use crate::core::{CaseId, StoreError, Task, TaskId, TaskState, Timestamp};

pub use crate::core::ClaimError;

/// Pending human work.
#[async_trait]
pub trait TaskStore: Send + Sync + Debug {
    /// Whose rows this handle can reach.
    ///
    /// Defaults to [`TenantId::DEFAULT`](crate::core::TenantId::DEFAULT), the
    /// tenant a store serves until told otherwise. Override it with the tenant
    /// the handle is actually scoped to.
    ///
    /// This exists so a mismatch with the plane's tenant is a **startup
    /// refusal**. When a key ring is wired, `build()` seals this state under
    /// the plane's tenant while the store writes rows under its own; the two
    /// disagreeing is not a leak — the scopes simply differ — but it seals the
    /// state under a scope erasure will never destroy. That is an erasure that
    /// reports success and misses, which is the one failure a deletion
    /// guarantee cannot have.
    fn tenant(&self) -> &str {
        crate::core::TenantId::DEFAULT
    }

    /// Create a task, or return the existing one with this id.
    ///
    /// Idempotent because task ids are derived from the awaiting effect rather
    /// than minted: a resumed run addresses the same task instead of opening a
    /// second one for the same decision.
    async fn open(&self, task: &Task) -> Result<Task, StoreError>;

    /// Fetch one task. Named for what it returns — see `CaseStore::case`.
    async fn task(&self, id: TaskId) -> Result<Option<Task>, StoreError>;

    /// Reserve a task for one actor.
    ///
    /// Enforces the four-eyes exclusion and role eligibility, atomically, so two
    /// reviewers cannot both believe they hold it.
    ///
    /// **Eligibility is checked before availability**, and the order is part of
    /// the contract rather than an artefact of how it was written:
    ///
    /// `NotFound` → `Excluded` → `WrongRole` → `NotPending` → `AlreadyClaimed`
    ///
    /// Told "held by Bob", a barred reviewer waits for Bob to release it and
    /// tries again — and is refused, for a reason nobody has yet mentioned. The
    /// permanent answer has to win over the transient one, or the transient one
    /// hides it. It also keeps queue state — who is reviewing what — from
    /// anybody not eligible for that queue.
    ///
    /// **A same-holder re-claim is idempotent success.** A claim by the actor
    /// who already holds the task returns the task, not
    /// `AlreadyClaimed { holder: yourself }`. A claim's acknowledgement can be
    /// lost — a dropped response, a crashed client that retries on restart —
    /// and the retry must converge on "you hold it" rather than bounce off its
    /// own success; an error naming the caller as the obstacle is one every
    /// client would have to special-case back into an `Ok`. The exclusivity
    /// the verb exists for is untouched: a task held by anybody *else* is
    /// still refused with [`ClaimError::AlreadyClaimed`].
    ///
    /// # Errors
    ///
    /// [`ClaimError`], per the order above.
    async fn claim(&self, id: TaskId, actor: &str, roles: &[String]) -> Result<Task, ClaimError>;

    /// Release a claim without deciding.
    ///
    /// Only the holder may. A release by anybody else must report
    /// [`ClaimError::NotHeld`] rather than succeed silently: a caller who is
    /// told "released" and then sees the task still assigned has no way to tell
    /// which of the two is lying.
    ///
    /// # Errors
    ///
    /// [`ClaimError::NotHeld`] if `actor` does not hold it, or it is not
    /// claimed.
    async fn release(&self, id: TaskId, actor: &str) -> Result<(), ClaimError>;

    /// Take a claim over from a holder who is not coming back.
    ///
    /// The absent-holder case [`release`](Self::release) cannot reach: only
    /// the holder may release, so a task claimed by a reviewer who has left is
    /// parked until its deadline breaches — a routine handover turned into an
    /// escalation, or "an operator edits the database", which is the
    /// anti-pattern the release endpoint exists to prevent.
    ///
    /// `from` names the holder being displaced and is a compare-and-swap
    /// guard, not documentation: a take-over decided from a stale view must
    /// fail rather than displace whoever holds it *now* — the same rule a
    /// case write follows by naming the version it read. Eligibility is
    /// re-checked in full for the new actor; a take-over is a claim, and
    /// four-eyes exclusion does not thin because the previous reviewer left.
    ///
    /// A reservation, not a decision: like claim and release it lives in the
    /// store, and the decision eventually taken still records its decider.
    /// The API gates it under its own action, so policy can hand it to a
    /// queue lead without handing it to every reviewer.
    ///
    /// # Errors
    ///
    /// [`ClaimError`], in claim's eligibility-first order, with
    /// [`ClaimError::NotHeld`] naming `from` when the task is not currently
    /// held by them — including when it is not held at all, where the right
    /// verb is [`claim`](Self::claim).
    async fn take_over(
        &self,
        id: TaskId,
        from: &str,
        actor: &str,
        roles: &[String],
    ) -> Result<Task, ClaimError>;

    async fn set_state(&self, id: TaskId, state: TaskState) -> Result<(), StoreError>;

    /// Widen an unanswered task to its declared escalation audience.
    ///
    /// One verb rather than three writes because the three must not be
    /// separable: [`Task::escalate`] moves the state, clears the reservation
    /// and widens the audience, and this verb applies all of it in one store
    /// transaction. Spelled as three `set_*` calls, a crash between them
    /// leaves a task telling the wider audience it exists while the old
    /// holder's claim still bars them from taking it.
    ///
    /// Escalating a task that is no longer pending is a **no-op returning the
    /// task as it stands**: the sweep that escalates races the reviewer it is
    /// escalating past, and the decision winning that race is the outcome
    /// everybody wanted — resurrecting a completed task into `escalated`
    /// would un-decide it.
    ///
    /// # Errors
    ///
    /// [`StoreError::NotFound`] when no task has this id.
    async fn escalate(&self, id: TaskId) -> Result<Task, StoreError>;

    /// Open work, highest priority and oldest first.
    async fn queue(&self, roles: &[String], limit: usize) -> Result<Vec<Task>, StoreError>;

    /// Everything pending on one matter.
    async fn for_case(&self, case: CaseId) -> Result<Vec<Task>, StoreError>;

    /// How many decisions are waiting on a person, across every role.
    ///
    /// Separate from `queue` for the same reason `pending_count` is separate
    /// from `pending`: a gauge must not be read from a `limit`-bounded list.
    async fn open_count(&self) -> Result<u64, StoreError>;

    /// Tasks whose window has closed and whose expiry policy has not yet fired.
    ///
    /// `open` and `claimed` only — **never `escalated`**, although escalated
    /// tasks are pending and past due. This list drives the sweep that applies
    /// each task's declared expiry policy, and escalation *is* that policy
    /// having fired: an escalated task is answered by a person or it is
    /// answered never, and no further tick has anything to do to it. Included,
    /// escalated rows would accumulate at the head of a bounded, oldest-first
    /// scan until one batch is nothing but rows the sweep will no-op, and the
    /// declared `deny`/`proceed` policies behind them silently stop firing
    /// plane-wide — an oversight queue that can be flooded is an oversight
    /// control that can be switched off.
    async fn overdue(&self, now: Timestamp, limit: usize) -> Result<Vec<Task>, StoreError>;
}