Skip to main content

SnapshotThenStream

Struct SnapshotThenStream 

Source
pub struct SnapshotThenStream<TSnapshot, TPublication> { /* private fields */ }
Expand description

Coordinates REST snapshot hydration with a live protobuf channel.

Implementations§

Source§

impl<TSnapshot, TPublication> SnapshotThenStream<TSnapshot, TPublication>
where TSnapshot: Send + 'static, TPublication: Send + 'static,

Source

pub fn new(cfg: SnapshotThenStreamConfig<TSnapshot, TPublication>) -> Self

Source

pub async fn start(&self) -> Result<()>

Begin websocket streaming and perform the initial snapshot refresh.

Source

pub async fn refresh_snapshot(&self) -> Result<()>

Fetch a REST snapshot and merge buffered publications.

On failure, readiness stays false, Self::err is set, and the pending buffer is retained so a successful retry merges each buffered publication exactly once. Success clears err.

Source

pub fn request_refresh(&self)

Request a snapshot refresh from a sync context (e.g. sequence gap handler).

Requests are coalesced behind one worker. A request arriving during a fetch schedules a follow-up, while repeated failures or persistent gaps fail closed after a bounded number of attempts.

Source

pub fn is_ready(&self) -> bool

Source

pub fn is_disposed(&self) -> bool

Source

pub fn err(&self) -> Option<Error>

Terminal stream error, if recovery failed closed.

Source

pub fn set_on_error<F>(&self, callback: F)
where F: Fn(Error) + Send + Sync + 'static,

Register a callback for transport, decode, snapshot, and terminal buffering errors.

If an error was already recorded, the callback is invoked immediately. Callback panics are isolated from the stream worker.

Source

pub fn close(&self)

Stop the stream.

Trait Implementations§

Source§

impl<TSnapshot, TPublication> Clone for SnapshotThenStream<TSnapshot, TPublication>

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<TSnapshot, TPublication> Drop for SnapshotThenStream<TSnapshot, TPublication>

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more

Auto Trait Implementations§

§

impl<TSnapshot, TPublication> !RefUnwindSafe for SnapshotThenStream<TSnapshot, TPublication>

§

impl<TSnapshot, TPublication> !UnwindSafe for SnapshotThenStream<TSnapshot, TPublication>

§

impl<TSnapshot, TPublication> Freeze for SnapshotThenStream<TSnapshot, TPublication>

§

impl<TSnapshot, TPublication> Send for SnapshotThenStream<TSnapshot, TPublication>
where TPublication: Send,

§

impl<TSnapshot, TPublication> Sync for SnapshotThenStream<TSnapshot, TPublication>
where TPublication: Send,

§

impl<TSnapshot, TPublication> Unpin for SnapshotThenStream<TSnapshot, TPublication>

§

impl<TSnapshot, TPublication> UnsafeUnpin for SnapshotThenStream<TSnapshot, TPublication>

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> 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 = 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<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

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