pub struct ServerLedger { /* private fields */ }Implementations§
Source§impl ServerLedger
impl ServerLedger
pub fn open(path: impl AsRef<Path>, retention_days: u32) -> StoreResult<Self>
pub fn path(&self) -> &Path
pub fn retention_days(&self) -> i64
pub fn upsert_role(&self, role: &RoleRow) -> StoreResult<bool>
pub fn list_roles(&self) -> StoreResult<Vec<RoleRow>>
pub fn remove_role_missing_from(&self, names: &[String]) -> StoreResult<usize>
Sourcepub fn project_session(&self, write: &SessionWrite) -> StoreResult<bool>
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.
Sourcepub fn rebind_session(
&self,
from_session_id: &str,
write: &SessionWrite,
) -> StoreResult<bool>
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.
Sourcepub fn open_binding(
&self,
session_id: &str,
task_id: &str,
bound_at: i64,
) -> StoreResult<bool>
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.
Sourcepub fn release_binding(
&self,
session_id: &str,
task_id: &str,
released_at: i64,
) -> StoreResult<bool>
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.
Sourcepub fn open_binding_of(
&self,
session_id: &str,
) -> StoreResult<Option<SessionBindingRow>>
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.
Sourcepub fn publish_mirror_outcome(
&self,
session_id: &str,
observed_json: &str,
expected_observed_json: &str,
updated_at: i64,
) -> StoreResult<bool>
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).
Sourcepub fn beat_session(
&self,
session_id: &str,
at: i64,
flush_after_secs: i64,
) -> StoreResult<bool>
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.
Sourcepub fn get_session_row(
&self,
session_id: &str,
) -> StoreResult<Option<ServerSessionRow>>
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.
Sourcepub fn session_row_for_task(
&self,
task_id: &str,
) -> StoreResult<Option<ServerSessionRow>>
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.
pub fn list_sessions( &self, filter: QuerySessionsArgs, ) -> StoreResult<Vec<ServerSessionRow>>
pub fn append_ledger(&self, row: &LedgerRow) -> StoreResult<Append>
pub fn mark_in_flight(&self, msg_id: &str) -> StoreResult<bool>
pub fn mark_acked(&self, msg_id: &str, at: DateTime<Utc>) -> StoreResult<bool>
pub fn mark_rejected(&self, msg_id: &str, reason: &str) -> StoreResult<bool>
pub fn expire_queued_before(&self, now: DateTime<Utc>) -> StoreResult<usize>
pub fn queued_for(&self, role: &str, limit: u32) -> StoreResult<Vec<LedgerRow>>
Sourcepub fn queued_count_for(&self, role: &str) -> StoreResult<u32>
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).
pub fn in_flight_for(&self, role: &str) -> StoreResult<Vec<LedgerRow>>
Sourcepub fn pending_expiries(&self) -> StoreResult<Vec<(String, DateTime<Utc>)>>
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.
pub fn requeue_in_flight(&self, role: &str) -> StoreResult<usize>
Sourcepub fn requeue_one(&self, msg_id: &str) -> StoreResult<LedgerRow>
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.
Sourcepub fn expire_one(&self, msg_id: &str, reason: &str) -> StoreResult<LedgerRow>
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.
Sourcepub fn fail_one(&self, msg_id: &str, reason: &str) -> StoreResult<LedgerRow>
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.
Sourcepub fn update_fault_state(
&self,
task_id: &str,
next_state: &str,
reason: &str,
) -> StoreResult<Vec<ServerFaultRow>>
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.
pub fn ledger_query(&self, query: LedgerQuery) -> StoreResult<Vec<LedgerRow>>
Sourcepub fn ledger_task(&self, task: &str, limit: u32) -> StoreResult<Vec<LedgerRow>>
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.
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 prune(&self, older_than: DateTime<Utc>) -> StoreResult<usize>
pub fn record_fault(&self, fault: &ServerFaultRow) -> StoreResult<i64>
pub fn open_faults(&self) -> StoreResult<Vec<ServerFaultRow>>
pub fn ack_fault(&self, fault_id: i64) -> StoreResult<bool>
pub fn faults_query( &self, query: FaultQuery, ) -> StoreResult<Vec<ServerFaultRow>>
pub fn faults_query_proto( &self, query: QueryFaultsArgs, ) -> StoreResult<Vec<ServerFaultRow>>
Sourcepub fn record_ghost_sweep(&self, sweep: &GhostSweepRow) -> StoreResult<i64>
pub fn record_ghost_sweep(&self, sweep: &GhostSweepRow) -> StoreResult<i64>
Persist one ghost-sweep audit row and return its id.
Sourcepub fn list_ghost_sweeps(&self, limit: u32) -> StoreResult<Vec<GhostSweepRow>>
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.
pub fn cursor_for(&self, role: &str) -> StoreResult<Option<CursorRow>>
pub fn set_cursor( &self, role: &str, msg_id: Option<&str>, seq: i64, ) -> StoreResult<bool>
Sourcepub fn hook_cursor(&self, hook: &str) -> StoreResult<Option<i64>>
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”).
Sourcepub fn set_hook_cursor(&self, hook: &str, seq: i64) -> StoreResult<bool>
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.
Sourcepub fn session_row_writes(&self) -> u64
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.