1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
//! The worklist contract.
use std::fmt::Debug;
use async_trait::async_trait;
use crate::core::{CaseId, RunId, 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>;
/// Settle a pending task, returning whether this call settled it.
///
/// A compare-and-set from the pending states: a task already completed,
/// expired or withdrawn is left as it stands and the answer is `false`.
/// The expiry sweep and a reviewer's decision race to settle one task,
/// and whichever lost must not overwrite the state the winner wrote —
/// the worklist would then contradict the answer the run consumed.
///
/// # Errors
///
/// [`StoreError::NotFound`] when no task has this id.
async fn set_state(&self, id: TaskId, state: TaskState) -> Result<bool, StoreError>;
/// Withdraw the tasks in `awaited` that belong to `run` and are still
/// pending, returning how many.
///
/// Called when the run concludes closed, naming the tasks it was still
/// waiting on. Nobody can answer a task whose run is sealed, so left
/// pending it is a decision the worklist offers and the backlog counts
/// that no answer can reach. A task the run opened *beside* its answer is
/// not awaited, and outlives the run by design; a task whose answer the
/// run consumed is settled by whoever delivered it.
async fn withdraw_run(&self, run: RunId, awaited: &[TaskId]) -> Result<usize, 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>;
}