Skip to main content

ServerLedger

Struct ServerLedger 

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

Implementations§

Source§

impl ServerLedger

Source

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

Source

pub fn path(&self) -> &Path

Source

pub fn retention_days(&self) -> i64

Source

pub fn upsert_role(&self, role: &RoleRow) -> StoreResult<bool>

Source

pub fn list_roles(&self) -> StoreResult<Vec<RoleRow>>

Source

pub fn remove_role_missing_from(&self, names: &[String]) -> StoreResult<usize>

Source

pub fn project_session(&self, write: &SessionWrite) -> StoreResult<bool>

Write one mirror row, and bind the delivery the write carries.

The row is addressed by session id and the (generation, seq) gate is the row’s own, as it was when the row answered for a task. A write that lands binds its delivery in the same transaction: a session serves one delivery at a time, so taking this one releases whatever the session was on before. A write the gate refuses changes nothing, the binding included — the row it was refused by is the newer word.

A write that lands carries the session’s last_seen forward, so the beats it has already outlived are spent: [crate::liveness] is told what the row now holds and drops the ones that are no longer newer.

Source

pub fn rebind_session( &self, from_session_id: &str, write: &SessionWrite, ) -> StoreResult<bool>

Move one mirror row to the session the write names, and write it there.

This is what repair rebind means once a row is addressed by its session: the operator says the delivery is now carried by another session, so the row moves to that address instead of a second row appearing beside it. The bindings move with it, since they name the same session. A row already sitting at the new address is the one being replaced, so it goes first: the operator’s word is the newer fact.

Source

pub fn open_binding( &self, session_id: &str, task_id: &str, bound_at: i64, ) -> StoreResult<bool>

Take one delivery for a session: a session_tasks row with bound_at.

Answers whether the binding moved, which a pair already open does not.

Source

pub fn release_binding( &self, session_id: &str, task_id: &str, released_at: i64, ) -> StoreResult<bool>

Stop serving one delivery: its session_tasks row gets released_at.

Only an open binding is released, so the first release is the one the row keeps and a second call changes nothing.

Source

pub fn open_binding_of( &self, session_id: &str, ) -> StoreResult<Option<SessionBindingRow>>

The delivery a session is on now, when it is on one.

At most one binding of a session is open, so the order here only decides which row a hand-written pair of open bindings answers with.

Source

pub fn publish_mirror_outcome( &self, session_id: &str, observed_json: &str, expected_observed_json: &str, updated_at: i64, ) -> StoreResult<bool>

Publish a late mirror verdict when the stored projection bytes still match the bytes the server compared. The stored tuple remains the session version while observed_json carries the task verdict.

last_seen is left where it stands: this write says the projection gained a verdict, not that the session was heard from, and the beats that were newer than the row keep answering for it (ServerLedger::beat_session).

Source

pub fn beat_session( &self, session_id: &str, at: i64, flush_after_secs: i64, ) -> StoreResult<bool>

Take one beat for a session whose projection did not change.

This is the whole of v2’s liveness rule, and it is deliberately not a projection write: (generation, seq), updated_at, and the event stream are all left exactly as they were, and the beat reaches the row only when the row is at least flush_after_secs behind — the interval the server states beside its call, which is also the reader’s worst-case staleness. Between those flushes the beat lives in memory, where every read of the row picks it up ([crate::liveness]).

The caller supplies the interval rather than this store: it is a promise about what a reader of last_seen is owed, and the server is the layer that knows the cluster’s presence window. The interval is floored at one second, so a spec that asks for a window below it gets one flush per second rather than one per beat — the v1 write rate this slice removes.

Answers whether the beat reached the table.

Source

pub fn get_session_row( &self, session_id: &str, ) -> StoreResult<Option<ServerSessionRow>>

One mirror row, addressed by its session id: the live last_seen this process holds when it has a newer one, the persisted value otherwise.

Source

pub fn session_row_for_task( &self, task_id: &str, ) -> StoreResult<Option<ServerSessionRow>>

The mirror row serving one delivery, read through its binding.

A reader asks “the session serving this task”; the binding is where that fact lives, and the row’s own task_id is derived from it either way.

Source

pub fn list_sessions( &self, filter: QuerySessionsArgs, ) -> StoreResult<Vec<ServerSessionRow>>

Source

pub fn append_ledger(&self, row: &LedgerRow) -> StoreResult<Append>

Source

pub fn mark_in_flight(&self, msg_id: &str) -> StoreResult<bool>

Source

pub fn mark_acked(&self, msg_id: &str, at: DateTime<Utc>) -> StoreResult<bool>

Source

pub fn mark_rejected(&self, msg_id: &str, reason: &str) -> StoreResult<bool>

Source

pub fn expire_queued_before(&self, now: DateTime<Utc>) -> StoreResult<usize>

Source

pub fn queued_for(&self, role: &str, limit: u32) -> StoreResult<Vec<LedgerRow>>

Source

pub fn queued_count_for(&self, role: &str) -> StoreResult<u32>

How many deliveries are queued for one role’s inbox, counted exactly.

ServerLedger::queued_for takes a limit, so a caller that counted its rows would report the cap as the depth: an operator reading “512” when the truth is 900 reads a saturated role as a full one. The rows are the ones pull would hand this role, so note rows are left out exactly as relay::pull leaves them out (crates/onlyne-server/src/relay.rs).

Source

pub fn in_flight_for(&self, role: &str) -> StoreResult<Vec<LedgerRow>>

Source

pub fn pending_expiries(&self) -> StoreResult<Vec<(String, DateTime<Utc>)>>

msg_id plus parsed deadline for every queued or in-flight row whose expires_at is set, oldest deadline first.

Source

pub fn requeue_in_flight(&self, role: &str) -> StoreResult<usize>

Source

pub fn requeue_one(&self, msg_id: &str) -> StoreResult<LedgerRow>

Move one in-flight row back to queued and publish its ledger_state event in the same transaction. A row already queued is returned unchanged, so a disconnect path that fires twice settles once.

Source

pub fn expire_one(&self, msg_id: &str, reason: &str) -> StoreResult<LedgerRow>

Move one queued row to expired and publish its ledger_state event in the same transaction. Settle one row whose deadline passed.

A row reaches this from queued while it waits for its role, and from in_flight when the role held it past the deadline: the sender asked for a deadline, so the sweep answers with expired in both states.

Source

pub fn fail_one(&self, msg_id: &str, reason: &str) -> StoreResult<LedgerRow>

Move one row to rejected and publish its ledger_state event in the same transaction. The automatic requeue gate uses this when the row has used its requeue budget.

Source

pub fn update_fault_state( &self, task_id: &str, next_state: &str, reason: &str, ) -> StoreResult<Vec<ServerFaultRow>>

Move every open fault of a task to next_state, publishing one fault event per moved row in the same transaction. Returns the moved rows.

Source

pub fn ledger_query(&self, query: LedgerQuery) -> StoreResult<Vec<LedgerRow>>

Source

pub fn ledger_task(&self, task: &str, limit: u32) -> StoreResult<Vec<LedgerRow>>

The plan’s task-keyed ledger read: every row of one task across roles, in insertion order. docs/v1-PLAN.md line 501 reads three rows back in order after a reconnect and line 508 reads one task’s rows after a relocation. The ledger carries no monotonic column of its own, so the order is enqueued_at with SQLite’s rowid breaking ties between rows written in one second: rowid is the insertion counter, and it makes the order the ledger’s write order without a second copy of that fact.

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 prune(&self, older_than: DateTime<Utc>) -> StoreResult<usize>

Source

pub fn record_fault(&self, fault: &ServerFaultRow) -> StoreResult<i64>

Source

pub fn open_faults(&self) -> StoreResult<Vec<ServerFaultRow>>

Source

pub fn ack_fault(&self, fault_id: i64) -> StoreResult<bool>

Source

pub fn faults_query( &self, query: FaultQuery, ) -> StoreResult<Vec<ServerFaultRow>>

Source

pub fn faults_query_proto( &self, query: QueryFaultsArgs, ) -> StoreResult<Vec<ServerFaultRow>>

Source

pub fn record_ghost_sweep(&self, sweep: &GhostSweepRow) -> StoreResult<i64>

Persist one ghost-sweep audit row and return its id.

Source

pub fn list_ghost_sweeps(&self, limit: u32) -> StoreResult<Vec<GhostSweepRow>>

The recorded sweeps, newest first.

The order is the sessions listing’s order: swept_at DESC, rowid DESC against an ascending index, which SQLite walks backwards. rowid breaks the ties inside one second, and it breaks them in the order the pass wrote them.

Source

pub fn cursor_for(&self, role: &str) -> StoreResult<Option<CursorRow>>

Source

pub fn set_cursor( &self, role: &str, msg_id: Option<&str>, seq: i64, ) -> StoreResult<bool>

Source

pub fn hook_cursor(&self, hook: &str) -> StoreResult<Option<i64>>

The last event seq one hook handled successfully, None for a hook this cluster has never run.

This is the resume point of the at-least-once rule: a hook that restarted picks up after the last event it handled, so nothing between the cursor and the head is skipped (docs/v2-CONTRACT.md §“Slice 7”).

Source

pub fn set_hook_cursor(&self, hook: &str, seq: i64) -> StoreResult<bool>

Record the last event seq one hook handled successfully. A seq at or below the recorded one is left alone: a worker that resumed from an older row after a failure must not walk the cursor backwards.

Source

pub fn session_row_writes(&self) -> u64

Statements this store has executed against the sessions table.

The projection writes, the mirror-outcome publishes, the liveness flushes, and the address moves a rebind makes all count here. It is the observation a liveness claim is made against: “a beat that changed nothing wrote no row” is answered by this number rather than by reading the code, and a test that reads it does not have to be a party to the write path.

Trait Implementations§

Source§

impl Clone for ServerLedger

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 ServerLedger

Source§

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

Formats the value using the given formatter. Read more

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