pub struct ClientStore { /* private fields */ }Implementations§
Source§impl ClientStore
impl ClientStore
pub fn open(path: impl AsRef<Path>) -> StoreResult<Self>
pub fn path(&self) -> &Path
pub fn enqueue_intent(&self, op_id: &str, env_json: &Value) -> StoreResult<bool>
Sourcepub fn due_intents(
&self,
now: DateTime<Utc>,
limit: u32,
) -> StoreResult<Vec<IntentRow>>
pub fn due_intents( &self, now: DateTime<Utc>, limit: u32, ) -> StoreResult<Vec<IntentRow>>
The rows due by now, oldest deadline first, at most limit of them.
This is ClientStore::flush_order with the clock supplied by the
caller, which is how a test reads a retry the machine has parked in the
future without waiting on it.
Sourcepub fn bump_intent(
&self,
op_id: &str,
next_attempt_at: DateTime<Utc>,
error: &str,
) -> StoreResult<bool>
pub fn bump_intent( &self, op_id: &str, next_attempt_at: DateTime<Utc>, error: &str, ) -> StoreResult<bool>
Store a caller-supplied retry time and move the intent to retrying.
The caller owns the attempt ceiling; the client’s IntentMachine checks attempts and calls exhaust_intent when the ceiling is reached.
Sourcepub fn defer_intent(
&self,
op_id: &str,
next_attempt_at: DateTime<Utc>,
error: &str,
) -> StoreResult<bool>
pub fn defer_intent( &self, op_id: &str, next_attempt_at: DateTime<Utc>, error: &str, ) -> StoreResult<bool>
Push an intent’s next attempt forward without consuming its budget.
A transport failure is not the server refusing the message, so the plan’s
disconnect rule keeps the queue intact and the reconnect flushes it
(plan §6 line 289). Counting those against intent.attempts would drop a
completion that was written while the link was down.
pub fn accept_intent( &self, op_id: &str, receipt_json: &Value, ) -> StoreResult<bool>
pub fn exhaust_intent(&self, op_id: &str, error: &str) -> StoreResult<bool>
pub fn pending_intent_count(&self) -> StoreResult<i64>
Sourcepub fn delete_intent(&self, op_id: &str) -> StoreResult<bool>
pub fn delete_intent(&self, op_id: &str) -> StoreResult<bool>
Remove a row the server refused for good.
The row is gone rather than parked in a terminal state because a permanent refusal will never turn into an acceptance, and the queue must not carry it forward on every pass. Like every write here it runs on the store’s one connection, under the store’s lock and its busy timeout: a caller that opened a second handle to the file to drop this row wrote around that serialization, and whether the row went depended on how that handle happened to treat a busy database.
Sourcepub fn flush_order(&self) -> StoreResult<Vec<IntentRow>>
pub fn flush_order(&self) -> StoreResult<Vec<IntentRow>>
The rows one flush pass should send, in the order it should send them.
Three bounds, each of which the last pass lacked. Only a row whose
next_attempt_at has arrived is returned, so a retry the machine pushed
into the future stays put until it is due instead of being sent again on
the strength of being queued; the deadline orders the batch, so the row
that has waited longest is first and cannot be starved by later arrivals;
and the batch is capped at INTENT_FLUSH_BATCH_SIZE rows, so reading a
queue an offline client filled costs one bounded read.
A caller that means to look past the current deadline — a test asking what
a run left in the queue — reads ClientStore::due_intents with its own
horizon instead, which is the same query with the clock it supplies.
pub fn append_event(&self, kind: &str, data: &Value) -> StoreResult<i64>
pub fn events_since( &self, seq: i64, limit: u32, ) -> StoreResult<Vec<EventRecord>>
pub fn event_head(&self) -> StoreResult<i64>
pub fn put_out_head(&self, task_id: &str, head: &str) -> StoreResult<bool>
pub fn out_head(&self, task_id: &str) -> StoreResult<Option<String>>
pub fn put_prose( &self, role: &str, prose: &str, spec_hash: &str, ) -> StoreResult<bool>
pub fn prose(&self, role: &str) -> StoreResult<Option<(String, String)>>
pub fn put_config(&self, key: &str, value: &str) -> StoreResult<bool>
pub fn config(&self, key: &str) -> StoreResult<Option<String>>
Sourcepub fn bump_session_version(
&self,
task_id: &str,
generation: u64,
seq: u64,
) -> StoreResult<bool>
pub fn bump_session_version( &self, task_id: &str, generation: u64, seq: u64, ) -> StoreResult<bool>
Write a heartbeat’s (generation, seq) onto the stored row when that pair is strictly newer. The observation tuple stays as stored. Returns true when the row changed.
Sourcepub fn open_task(&self, causality: &Causality, kind: &str) -> StoreResult<bool>
pub fn open_task(&self, causality: &Causality, kind: &str) -> StoreResult<bool>
Open the task record for one delivery.
The causality of the envelope that became a session is the record: the
task id, its parent, its depth, and the redelivery count it arrived on.
An already-open task keeps its settle, its opened_at, and its verdict —
a re-dispatch refreshes only the chain it arrived on, and takes the
larger attempt so the count never runs backwards. kind is how the
delivery reached this role — root or relay — never the envelope’s
message kind, which says nothing a reader can act on. Answers true when
the record was written, which covers both the fresh open and a chain
refresh on a task that is already there.
Sourcepub fn settle_task(
&self,
task_id: &str,
task_state: TaskState,
) -> StoreResult<bool>
pub fn settle_task( &self, task_id: &str, task_state: TaskState, ) -> StoreResult<bool>
Settle one task with the verdict its completion carries.
The first terminal verdict wins. A task sitting open takes the verdict
and stamps settled_at; one already settled keeps its record and answers
false, which is what stops a zombie session that came back after a
newer one answered from rewriting the answer it already gave.
The write never depends on the open having run. A task this process never
dispatched — a completion reported for work an earlier process opened, or
a delivery that found its slot already serving and returned before the
open — gets its record here, with no kind, no parent, and both clocks at
this verdict. Without that, the verdict would be dropped on the floor and
the session could never project its way out of working.
Sourcepub fn task(&self, task_id: &str) -> StoreResult<Option<TaskRow>>
pub fn task(&self, task_id: &str) -> StoreResult<Option<TaskRow>>
The record of one task, when this role opened it.
Sourcepub fn open_tasks(&self, limit: u32) -> StoreResult<Vec<TaskRow>>
pub fn open_tasks(&self, limit: u32) -> StoreResult<Vec<TaskRow>>
Tasks this role has opened and not settled, oldest first.
Sourcepub fn active_sessions(&self) -> StoreResult<Vec<LiveSession>>
pub fn active_sessions(&self) -> StoreResult<Vec<LiveSession>>
The sessions this client holds and can recover, for the hello claim:
each one’s id, the delivery it is on now, and whether its process has
been released.
A fresh process after a crash has empty slots and this table is the only witness left that the sessions exist, so it must declare them to prevent duplicate dispatch.
Only claims sessions whose runtime may still be there: the ones that have
not mounted yet or are between turns (booting, ready, idle). A
session that was mid-turn when this process died is not claimed — the
plugin that was running it is gone — so the server requeues its
deliveries instead.
A suspended session — one whose resource is closed while its agent has
not gone — is claimed with suspended: true. The work it holds is still
owed, so the server holds its row rather than requeueing it, and the
session resumes that delivery instead of opening a second one. An exited
session (agent_state gone) is never claimed: it holds nothing.
Sourcepub fn release_binding(
&self,
session_id: &str,
task_id: &str,
) -> StoreResult<usize>
pub fn release_binding( &self, session_id: &str, task_id: &str, ) -> StoreResult<usize>
Stop serving one delivery: its session_tasks row gets released_at.
A task or role scope session outlives the delivery it served, so
settling that delivery is not the end of the binding. Releasing it here
is what stops the session from reading as bound to settled work: the next
delivery’s binding would close the other open one anyway
([crate::server::open_binding_conn]), but a session that goes idle with
no next delivery keeps answering OPEN_BINDING_TASK with a task that is
over, and hello then claims it for work nobody owes. Only an open
binding is released, so the first call is the one the row keeps and a
second changes nothing; the count is how many rows moved.
Sourcepub fn bind_task(&self, session_id: &str, task_id: &str) -> StoreResult<usize>
pub fn bind_task(&self, session_id: &str, task_id: &str) -> StoreResult<usize>
Take one delivery for a session: its session_tasks row opens now.
A scoped session serves its next delivery without its tuple moving, and
the binding has to open first: SESSION_ID_FOR_TASK reads it, so a write
for a delivery whose binding is not open yet lands on no row at all.
Opening a pair also closes the session’s other open binding, so the
session is on one delivery at a time either way. The count is how many
rows the insert moved — zero for a pair already open.
Trait Implementations§
Source§impl Clone for ClientStore
impl Clone for ClientStore
Source§impl Debug for ClientStore
impl Debug for ClientStore
Source§impl SessionLedger for ClientStore
impl SessionLedger for ClientStore
Source§fn get_session(&self, task_id: &str) -> Result<Option<SessionRecord>>
fn get_session(&self, task_id: &str) -> Result<Option<SessionRecord>>
The session serving one delivery, read through its binding.
The record answers for the delivery the caller named: a session that served one delivery and then another still answers for the first, which is what a caller holding that delivery’s tuple asks for.
Source§fn upsert_session(
&self,
task_id: &str,
version: &VersionedSession,
) -> Result<bool>
fn upsert_session( &self, task_id: &str, version: &VersionedSession, ) -> Result<bool>
Write one session tuple, and bind the delivery it is about.
The row is addressed by the session’s own id, which the backend reference names; the delivery the tuple serves is the binding this write opens, and a session serves one delivery at a time.
Source§fn task_is_known(&self, task_id: &str) -> Result<bool>
fn task_is_known(&self, task_id: &str) -> Result<bool>
Source§fn task_attempt(&self, task_id: &str) -> Result<i64>
fn task_attempt(&self, task_id: &str) -> Result<i64>
Source§fn list_faults(&self, task_id: &str) -> Result<Vec<FaultRecord>>
fn list_faults(&self, task_id: &str) -> Result<Vec<FaultRecord>>
(task_id, kind, generation) dedupe.