Skip to main content

Consumer

Struct Consumer 

Source
pub struct Consumer<T> { /* private fields */ }
Expand description

Consumes a JSON value from a track, reconstructing it from snapshots and deltas.

Implementations§

Source§

impl<T: DeserializeOwned> Consumer<T>

Source

pub fn new(track: Subscriber, config: ConsumerConfig) -> Self

Create a consumer reading from the given track subscriber.

Set ConsumerConfig::compression to read a track written by a producer with ProducerConfig::compression on.

Source

pub async fn next(&mut self) -> Result<Option<T>>
where T: Unpin,

Get the next reconstructed value, or None once the track ends.

Source

pub fn poll_next(&mut self, waiter: &Waiter) -> Poll<Result<Option<T>>>

Poll for the next reconstructed value, without blocking.

Jumps to the newest group, reads its snapshot, and applies deltas in order. All frames already buffered in the group are applied in one poll but only the resulting latest value is yielded: the intermediate reconstructions are stale, so a late joiner (or any consumer that has fallen behind) catches up to the head in a single step instead of replaying every superseded state. Frames must still be decoded in order (the DEFLATE window and merge patches are sequential); only the per-frame deserialize and yield are skipped. Switching to a newer group discards the older one.

Auto Trait Implementations§

§

impl<T> !Freeze for Consumer<T>

§

impl<T> RefUnwindSafe for Consumer<T>

§

impl<T> Send for Consumer<T>

§

impl<T> Sync for Consumer<T>

§

impl<T> Unpin for Consumer<T>

§

impl<T> UnsafeUnpin for Consumer<T>

§

impl<T> UnwindSafe for Consumer<T>

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<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> MaybeSend for T
where T: Send,

Source§

impl<T> MaybeSend for T
where T: Send,

Source§

impl<T> MaybeSync for T
where T: Sync,

Source§

impl<T> MaybeSync for T
where T: Sync,

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