Skip to main content

SqliteProcessRegistry

Struct SqliteProcessRegistry 

Source
pub struct SqliteProcessRegistry { /* private fields */ }
Expand description

SQLite-backed process registry for one configured runtime deployment.

It is intentionally separate from Store: session databases persist one conversation, while this registry persists background process state and handle visibility across all sessions sharing the registry.

Implementations§

Source§

impl SqliteProcessRegistry

Source

pub async fn open( path: &Path, session_store_root: impl Into<PathBuf>, ) -> Result<Self>

Open a process registry whose terminal-retention prune removes the two process-owned session stores from session_store_root before the process row. The root is required and explicit; no sibling-directory convention is inferred.

Source

pub async fn open_with_clock( path: &Path, clock: Arc<dyn Clock>, session_store_root: impl Into<PathBuf>, ) -> Result<Self>

Source

pub async fn memory() -> Result<Self>

Source

pub async fn memory_with_clock(clock: Arc<dyn Clock>) -> Result<Self>

Trait Implementations§

Source§

impl ProcessRegistry for SqliteProcessRegistry

Source§

fn durability_tier(&self) -> DurabilityTier

Durability tier this process registry provides; defaults to DurabilityTier::Inline.
Source§

fn register_process<'life0, 'async_trait>( &'life0 self, registration: ProcessRegistration, ) -> Pin<Box<dyn Future<Output = Result<ProcessRecord, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Source§

fn put_segment_handover<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, handover: PersistedSegmentHandover, ) -> Pin<Box<dyn Future<Output = Result<(), PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Persist the bounded engine continuation for exactly one segment. Repeating an identical write is an idempotent no-op; conflicting data for the same (process_id, segment_ordinal) is rejected.
Source§

fn get_segment_handover<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, segment_ordinal: u64, ) -> Pin<Box<dyn Future<Output = Result<Option<PersistedSegmentHandover>, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Load the continuation for an exact process segment ordinal.
Source§

fn latest_segment_handover<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Option<PersistedSegmentHandover>, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Load the highest persisted segment ordinal for recovery.
Source§

fn delete_segment_handovers<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<(), PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Remove all cross-segment execution state once the process is terminal.
Source§

fn set_external_ref<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, external_ref: ProcessExternalRef, ) -> Pin<Box<dyn Future<Output = Result<ProcessRecord, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Attach a durable backend reference to a registered process. Read more
Source§

fn grant_handle<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, session_scope: &'life1 SessionScope, process_id: &'life2 str, descriptor: ProcessHandleDescriptor, ) -> Pin<Box<dyn Future<Output = Result<ProcessHandleGrant, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Source§

fn revoke_handle<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, session_scope: &'life1 SessionScope, process_id: &'life2 str, ) -> Pin<Box<dyn Future<Output = Result<(), PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Source§

fn transfer_handle_grants<'life0, 'life1, 'life2, 'life3, 'async_trait>( &'life0 self, from_scope: &'life1 SessionScope, to_scope: &'life2 SessionScope, process_ids: &'life3 [String], ) -> Pin<Box<dyn Future<Output = Result<(), PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait, 'life3: 'async_trait,

Source§

fn list_handle_grants<'life0, 'life1, 'async_trait>( &'life0 self, session_scope: &'life1 SessionScope, ) -> Pin<Box<dyn Future<Output = Result<Vec<ProcessHandleGrantEntry>, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Source§

fn list_live_handle_grants<'life0, 'life1, 'async_trait>( &'life0 self, session_scope: &'life1 SessionScope, ) -> Pin<Box<dyn Future<Output = Result<Vec<ProcessHandleGrantEntry>, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Source§

fn has_handle_grant<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, session_scope: &'life1 SessionScope, process_id: &'life2 str, ) -> Pin<Box<dyn Future<Output = Result<bool, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Source§

fn handle_grants_for_process<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Vec<ProcessHandleGrant>, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Source§

fn delete_session_process_state<'life0, 'life1, 'async_trait>( &'life0 self, session_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<ProcessSessionDeleteReport, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Source§

fn append_event<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, request: ProcessEventAppendRequest, ) -> Pin<Box<dyn Future<Output = Result<ProcessEventAppendResult, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Source§

fn events_after<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, after_sequence: u64, ) -> Pin<Box<dyn Future<Output = Result<Vec<ProcessEvent>, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Source§

fn count_events_through<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, process_id: &'life1 str, event_type: &'life2 str, up_to_sequence: u64, ) -> Pin<Box<dyn Future<Output = Result<u64, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Count events of event_type with sequence <= up_to_sequence. Read more
Source§

fn recent_events<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, limit: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<ProcessEvent>, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

The most recent limit events, in ascending sequence order. Read more
Source§

fn wake_events_after<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, after_sequence: u64, ) -> Pin<Box<dyn Future<Output = Result<Vec<ProcessEvent>, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Source§

fn complete_process<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, await_output: ProcessAwaitOutput, authority: ProcessCompletionAuthority, ) -> Pin<Box<dyn Future<Output = Result<ProcessRecord, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Complete a process without a Lash process lease, under an explicit, auditable completion authority. Read more
Source§

fn complete_process_with_lease<'life0, 'life1, 'async_trait>( &'life0 self, lease: &'life1 ProcessLease, await_output: ProcessAwaitOutput, ) -> Pin<Box<dyn Future<Output = Result<ProcessRecord, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Atomically append the terminal output while the supplied process lease is still current, then release that lease in the same transaction. Read more
Source§

fn record_first_started<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, started: ProcessStarted, ) -> Pin<Box<dyn Future<Output = Result<ProcessRecord, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Record the durable, lease-fenced “execution started” fact (ADR 0019). Read more
Source§

fn request_process_abandon<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, request: AbandonRequest, ) -> Pin<Box<dyn Future<Output = Result<ProcessRecord, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Set the durable, non-terminal Abandon Request marker (ADR 0019). Read more
Source§

fn set_process_wait<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, wait: WaitState, ) -> Pin<Box<dyn Future<Output = Result<ProcessRecord, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Source§

fn clear_process_wait<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<ProcessRecord, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Source§

fn get_process<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Option<ProcessRecord>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Source§

fn try_get_process<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Option<ProcessRecord>, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Fallible process lookup for correctness-critical execution paths where a transient store failure must not be mistaken for an absent row.
Source§

fn list_processes<'life0, 'life1, 'async_trait>( &'life0 self, filter: &'life1 ProcessListFilter, ) -> Pin<Box<dyn Future<Output = Result<Vec<ProcessRecord>, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Source§

fn processes_changed_since<'life0, 'async_trait>( &'life0 self, cursor: ProcessChangeCursor, limit: usize, ) -> Pin<Box<dyn Future<Output = Result<(Vec<ProcessRecord>, ProcessChangeCursor), PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Return process records whose persisted row changed strictly after cursor, ordered by the backend’s per-store change sequence. Read more
Source§

fn ack_wake<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, sequence: u64, ) -> Pin<Box<dyn Future<Output = Result<(), PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Source§

fn list_non_terminal<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<Vec<ProcessRecord>, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

All non-terminal process records, in stable process_id order. Read more
Source§

fn filter_unregistered_process_ids<'life0, 'life1, 'async_trait>( &'life0 self, process_ids: &'life1 [String], ) -> Pin<Box<dyn Future<Output = Result<Vec<String>, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Return the candidate ids that have no persisted process row, preserving input order. Durable backends override this with one NOT EXISTS anti-join so recovery does not issue one point read per candidate.
Source§

fn live_reference_summary<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<Vec<ProcessLiveReferenceSummary>, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Count non-terminal process rows by their captured definition and execution-environment references.
Source§

fn claim_process_lease<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, process_id: &'life1 str, owner: &'life2 LeaseOwnerIdentity, lease_ttl_ms: u64, ) -> Pin<Box<dyn Future<Output = Result<ProcessLeaseClaimOutcome, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Claim the durable single-owner lease over a non-terminal process. Read more
Source§

fn reclaim_process_lease<'life0, 'life1, 'life2, 'life3, 'async_trait>( &'life0 self, process_id: &'life1 str, owner: &'life2 LeaseOwnerIdentity, observed_holder: &'life3 ProcessLease, lease_ttl_ms: u64, ) -> Pin<Box<dyn Future<Output = Result<ProcessLeaseClaimOutcome, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait, 'life3: 'async_trait,

Reclaim an unexpired process lease whose observed holder is definitely dead according to persisted local-process liveness metadata. Read more
Source§

fn renew_process_lease<'life0, 'life1, 'async_trait>( &'life0 self, lease: &'life1 ProcessLease, lease_ttl_ms: u64, ) -> Pin<Box<dyn Future<Output = Result<ProcessLease, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Extend the expiry of a live lease the caller still owns. Read more
Source§

fn get_process_lease<'life0, 'life1, 'async_trait>( &'life0 self, process_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Option<ProcessLease>, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Read the current lease row for a process without claiming it. Read more
Source§

fn complete_process_lease<'life0, 'life1, 'async_trait>( &'life0 self, completion: &'life1 ProcessLeaseCompletion, ) -> Pin<Box<dyn Future<Output = Result<(), PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Release a lease the caller owns, fenced by the completion’s (process_id, lease_token). Read more
Source§

fn prune_terminal_processes<'life0, 'async_trait>( &'life0 self, cutoff_epoch_ms: u64, filter: Option<ProcessListFilter>, up_to_change_seq: Option<ProcessChangeCursor>, ) -> Pin<Box<dyn Future<Output = Result<ProcessPruneReport, PluginError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Physically delete terminal process rows whose updated_at_ms is older than cutoff_epoch_ms, match filter when one is supplied, and have a process change sequence no later than up_to_change_seq when supplied, together with their events, wake acks, handle grants, lease rows, and trigger-delivery reservations whose deterministic process id points at a pruned row. The same cutoff also prunes trigger-mutation idempotency receipts, bounding receipt retention under the host’s existing cleanup schedule. Durable backends also release attachment intents and delete the process-owned process-env:<id> and process-session-turn:<id> session stores before deleting the process row. Backends must fail toward retaining the terminal process if that cleanup cannot complete. Host-scheduled retention: hosts that project results/events into their own store call this to keep the registry bounded. Non-terminal rows are never touched. Callers must choose a retention window comfortably longer than any waiter lifetime — a pruned process id becomes “unknown process” to late awaits. Re-emitting the same trigger occurrence id after its process has aged out of retention may reserve a fresh delivery process id; occurrence-level idempotency still holds, and ordinary emit replays do not straddle a retention window in practice. 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> 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> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
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, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

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

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

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