Skip to main content

RedisSubscriber

Struct RedisSubscriber 

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

A Redis Streams subscription bound to a consumer group.

Constructed by crate::ConnectedRedisBroker::subscribe from a crate::RedisStream descriptor. The read mode (fresh tail vs reclaim) is fixed at construction.

Trait Implementations§

Source§

impl BatchSubscriber for RedisSubscriber

Source§

fn batches( &mut self, ) -> impl Stream<Item = Result<Self::Batch, Self::Error>> + Send + '_

Yields one batch per non-empty read (XREADGROUP COUNT / XAUTOCLAIM), up to RedisStream::count entries. Never yields an empty batch.

§Cancel safety

Same as Subscriber::stream: dropping the stream mid-read leaves fetched-but-unacked entries in the pending list.

Source§

type Batch = Vec<RedisMessage>

Container yielded by batches. Implementations choose between Vec, custom iterators, or anything else that yields the underlying Subscriber::Message.
Source§

impl Debug for RedisSubscriber

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl Seekable for RedisSubscriber

Repositioning the subscription moves its consumer group’s cursor, which is shared by every consumer of that group; see RedisGroupSeeker for the full contract.

Source§

type Seeker = RedisGroupSeeker

The handle usable while this subscriber’s stream is running.
Source§

fn seeker(&self) -> RedisGroupSeeker

Mints a handle for repositioning this subscription.
Source§

impl Subscriber for RedisSubscriber

Source§

fn stream( &mut self, ) -> impl Stream<Item = Result<Self::Message, Self::Error>> + Send + '_

Yields one message per entry, refilling from Redis when the local buffer drains.

§Cancel safety

Dropping the returned stream between items is safe. Dropping it while a read is in flight drops the read future; entries already delivered to this consumer but not yet acked stay in the group’s pending list and are redelivered (fresh mode) or reclaimable (reclaim mode).

Source§

type Message = RedisMessage

The message type yielded by this subscriber.
Source§

type Error = RedisError

The error type yielded by the stream when delivery fails.

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> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

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