Skip to main content

RedisGroupSeeker

Struct RedisGroupSeeker 

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

Moves a consumer group’s cursor.

The handle behind the core Seekable capability, minted with Seekable::seeker or injected into a handler as a Seek(seeker): Seek<RedisGroupSeeker> parameter.

§A seek is group-wide

Redis keeps one cursor per consumer group, so XGROUP SETID moves the read position for every consumer of that group, not only for the subscription this seeker came from. That is unlike a partitioned log, where a seek is scoped to one consumer instance. Treat a seek as an operation on the group: replaying a range replays it for the whole worker pool, and skipping forward skips for all of them.

§What a seek does not do

  • It does not clear the pending entries list. Entries already delivered and not acknowledged stay pending and remain reachable through the reclaim path (RedisStream::reclaim), whichever way the cursor moved.
  • It does not cancel delayed retries. Copies already scheduled in a DelayedRetry::DurableZset queue are keyed by their due time, not by the cursor, so they are appended to the stream when they fall due regardless of where the group is reading.
  • It does not reset delivery counts. A replayed entry is delivered again, so its native delivery count (the one the reclaim path reads and RedisStream::max_deliveries caps) grows with each replay; the framework retry-count header only moves on an actual nack.

§Timing

The cursor changes as soon as seek returns, but a subscription parked in a blocking XREADGROUP observes it only on its next read - within one RedisStream::block interval. Entries selected under the old cursor are discarded rather than delivered, so a seek never yields a message from the position it moved away from.

§Examples

use ruststream::{Broker, Seekable, Seeker};
use ruststream_fred::{RedisBroker, RedisGroupPosition, RedisStream};

let connected = RedisBroker::standalone("redis://localhost:6379").connect().await?;
let subscriber = connected
    .subscribe(RedisStream::new("orders").group("workers"))
    .await?;

// Minted before the stream opens; usable while it runs.
let seeker = subscriber.seeker();
seeker.seek(RedisGroupPosition::beginning()).await?;

Trait Implementations§

Source§

impl Clone for RedisGroupSeeker

Source§

fn clone(&self) -> RedisGroupSeeker

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 Debug for RedisGroupSeeker

Source§

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

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

impl Seeker for RedisGroupSeeker

Source§

async fn seek(&self, to: RedisGroupPosition) -> Result<(), RedisError>

Moves the group cursor with XGROUP SETID.

§Errors

Returns RedisError::Stream when the group or the stream does not exist, or the command fails.

Source§

type Position = RedisGroupPosition

The broker’s own position type (Kafka partition offsets, a Redis entry id, a byte offset), constructed by the broker crate or captured from a delivered message via Positioned::position.
Source§

type Error = RedisError

The error returned when the broker rejects the reposition.

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