Skip to main content

ClientStore

Struct ClientStore 

Source
pub struct ClientStore { /* private fields */ }

Implementations§

Source§

impl ClientStore

Source

pub fn open(path: impl AsRef<Path>) -> StoreResult<Self>

Source

pub fn path(&self) -> &Path

Source

pub fn enqueue_intent(&self, op_id: &str, env_json: &Value) -> StoreResult<bool>

Source

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.

Source

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.

Source

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.

Source

pub fn accept_intent( &self, op_id: &str, receipt_json: &Value, ) -> StoreResult<bool>

Source

pub fn exhaust_intent(&self, op_id: &str, error: &str) -> StoreResult<bool>

Source

pub fn pending_intent_count(&self) -> StoreResult<i64>

Source

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.

Source

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.

Source

pub fn append_event(&self, kind: &str, data: &Value) -> StoreResult<i64>

Source

pub fn events_since( &self, seq: i64, limit: u32, ) -> StoreResult<Vec<EventRecord>>

Source

pub fn event_head(&self) -> StoreResult<i64>

Source

pub fn put_out_head(&self, task_id: &str, head: &str) -> StoreResult<bool>

Source

pub fn out_head(&self, task_id: &str) -> StoreResult<Option<String>>

Source

pub fn put_prose( &self, role: &str, prose: &str, spec_hash: &str, ) -> StoreResult<bool>

Source

pub fn prose(&self, role: &str) -> StoreResult<Option<(String, String)>>

Source

pub fn put_config(&self, key: &str, value: &str) -> StoreResult<bool>

Source

pub fn config(&self, key: &str) -> StoreResult<Option<String>>

Source

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.

Source

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.

Source

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.

Source

pub fn task(&self, task_id: &str) -> StoreResult<Option<TaskRow>>

The record of one task, when this role opened it.

Source

pub fn open_tasks(&self, limit: u32) -> StoreResult<Vec<TaskRow>>

Tasks this role has opened and not settled, oldest first.

Source

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.

Source

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.

Source

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

Source§

fn clone(&self) -> Self

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Debug for ClientStore

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl SessionLedger for ClientStore

Source§

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>

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>

Whether the task itself is tracked, so a stray report cannot conjure a row.
Source§

fn task_attempt(&self, task_id: &str) -> Result<i64>

Attempt count for a task, for the fault audit entry.
Source§

fn list_faults(&self, task_id: &str) -> Result<Vec<FaultRecord>>

Existing faults for one task, for the (task_id, kind, generation) dedupe.
Source§

fn insert_fault(&self, fault: &FaultRecord) -> Result<i64>

Append one fault row. Returns its id.
Source§

fn emit(&self, kind: &str, data: Value)

Observe an emitted event. The bridge never reads back from this hook.
Source§

fn note_alert(&self, line: String)

Observe an operator-visible alert line.

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> DynClone for T
where T: Clone,

Source§

fn __clone_box(&self, _: Private) -> *mut ()

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more