Skip to main content

Consumer

Struct Consumer 

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

Consumes a sliding window of JSON records from a track, yielding one event per change.

A Decoder that owns its track: it reads groups, starts a cold DEFLATE window at each boundary, and turns each group’s header into just the changes this reader has not been told about. When something else already owns the track, use the Decoder directly.

Group rolls never surface. A publisher rolls for compression’s sake, and a header restating the window yields nothing for records already delivered, so this reads as one continuous stream of Events regardless of how the publisher framed them.

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.

Source

pub fn range(&self) -> Range<u64>

Absolute index of the oldest record in the window, and of the next to arrive.

Source

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

Get the next event, or None once the track ends.

Source

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

Poll for the next event, without blocking.

Auto Trait Implementations§

§

impl<T> Freeze for Consumer<T>
where Decoder<T>: Freeze,

§

impl<T> RefUnwindSafe for Consumer<T>

§

impl<T> Send for Consumer<T>
where Decoder<T>: Send,

§

impl<T> Sync for Consumer<T>
where Decoder<T>: Sync,

§

impl<T> Unpin for Consumer<T>
where Decoder<T>: Unpin,

§

impl<T> UnsafeUnpin for Consumer<T>
where Decoder<T>: UnsafeUnpin,

§

impl<T> UnwindSafe for Consumer<T>
where Decoder<T>: UnwindSafe,

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

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 = !

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