Skip to main content

RedisMessage

Struct RedisMessage 

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

A Redis Streams delivery, read from a consumer group via XREADGROUP or XAUTOCLAIM.

Settlement follows the republish-retry model: ack is XACK; nack(requeue = true) re-appends a copy of the entry to the same stream and then acks the original (at-least-once, so a duplicate is possible if the process crashes between the two); nack(requeue = false) acks the original to drop it.

Implementations§

Source§

impl RedisMessage

Source

pub fn id(&self) -> Option<&str>

The stream entry ID (for example 1700000000000-0) this message was read at.

Source

pub const fn entry_id(&self) -> EntryId

The parsed entry id this message was read at.

Source

pub fn group(&self) -> Option<&str>

The consumer group this delivery was read through, or None once the message has settled.

Trait Implementations§

Source§

impl BuildContext<RedisMessage> for StreamContext

Source§

fn build(msg: &RedisMessage) -> Self

Builds the context value by reading fields out of msg.
Source§

impl Debug for RedisMessage

Source§

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

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

impl IncomingMessage for RedisMessage

Source§

fn supports_nack_after(&self) -> bool

Native delayed redelivery is available only when the subscription opted into a durable ZSET delay queue with RedisStream::delayed_retry; otherwise the runtime applies its broker-agnostic deferred-republish fallback.

Source§

async fn nack_after(self, delay: Duration) -> Result<(), AckError>

Schedules the message for redelivery no sooner than delay from now via the configured ZSET delay queue (ZADD the delayed copy, then XACK the original), with the retry-count header incremented. The subscriber’s sweeper re-XADDs it to the source stream once due.

§Errors

Returns AckError::Unsupported when the subscription did not opt into a delay queue, or AckError::Broker when the ZADD or XACK fails.

Source§

fn payload(&self) -> &[u8] ⓘ

Returns the raw payload of the message.
Source§

fn headers(&self) -> &Headers

Returns the headers attached to the message.
Source§

async fn ack(self) -> Result<(), AckError>

Acknowledges successful processing. Consumes the message handle. Read more
Source§

async fn nack(self, requeue: bool) -> Result<(), AckError>

Negatively acknowledges the message. When requeue is true the broker should redeliver according to its own retry policy; when false it should drop or dead-letter the message. Read more
Source§

fn partition_key(&self) -> Option<&[u8]>

Returns the routing key the broker partitioned this message by, or None when the message carries no key. Read more
Source§

impl Partitioned for RedisMessage

Source§

fn partition_key(&self) -> Option<&[u8]>

Returns the partition key for this item, or None if the broker should pick a partition.
Source§

impl Positioned for RedisMessage

The position of a delivery is the group cursor that redelivers it.

The cursor is exclusive (a group resumes after the id it holds), so the pinned position is the id immediately below this entry’s - seeking to it delivers this message again, followed by the entries after it. Repositioning is group-wide; see RedisGroupSeeker.

Source§

type Position = RedisGroupPosition

The position type, matching the subscription’s Seeker::Position.
Source§

fn position(&self) -> RedisGroupPosition

Returns the position of this delivery; seeking to it redelivers this message.

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