Skip to main content

RedisPublisher

Struct RedisPublisher 

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

The live stream publisher: RedisPublish paired with a connection. Cheap to clone.

Publisher::publish appends the message to the stream named by OutgoingMessage::name with XADD <name> * .... The payload and headers are encoded as entry fields (see crate::RedisStream for the consuming side).

A publisher may outlive the connected broker it came from (it is a handle aliasing the connection), so every operation after shutdown reports RedisError::ShutDown rather than running against a closed pool.

§Transactions

Both framework transaction kinds are available on standalone and sentinel topologies, and both commit the same way: the buffer is held client-side while the transaction is open and flushed as one MULTI / EXEC block, in publish order, so subscribers see the whole batch or none of it. They differ only in where that buffer lives.

  • Borrowed (TransactionalPublisher): the handle carries one transaction. begin_transaction claims it and starts buffering published messages, commit flushes them, and abort discards them. Clones of a handle share the same open transaction, and a second begin_transaction while one is open is rejected.
  • Owned (OwnedTransactions): every transaction call returns a RedisTransaction owning its own buffer, so any number can be open on one handle concurrently and the handle keeps publishing directly meanwhile.

Two Redis properties apply to both kinds. Cluster supports neither, because a MULTI block cannot span hash slots, so opening a transaction there returns RedisError::InvalidOptions. And Redis has no rollback: a command that fails at runtime inside EXEC does not undo the commands before it. For a block of XADDs against stream keys that is practically limited to out-of-memory and wrong-type keys; a command the server refuses to queue discards the whole block.

Trait Implementations§

Source§

impl Clone for RedisPublisher

Source§

fn clone(&self) -> RedisPublisher

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 RedisPublisher

Source§

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

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

impl OwnedTransactions for RedisPublisher

Owned transactions: every transaction call opens an independent buffer-owning RedisTransaction, so any number can be open concurrently on one handle, next to (and unaffected by) the handle-level TransactionalPublisher transaction.

Source§

async fn transaction(&self) -> Result<RedisTransaction, RedisError>

§Errors

Returns RedisError::InvalidOptions on a cluster topology, which cannot offer multi-key transactions.

Source§

type Transaction = RedisTransaction

The buffer-owning transaction opened by transaction.
Source§

impl Publisher for RedisPublisher

Source§

type Error = RedisError

The error type returned by publish.
Source§

async fn publish(&self, msg: OutgoingMessage<'_>) -> Result<(), Self::Error>

Publishes a message to the broker. Read more
Source§

impl TransactionalPublisher for RedisPublisher

Source§

async fn begin_transaction(&self) -> Result<(), Self::Error>

Starts buffering published messages on this handle.

§Errors

Returns RedisError::InvalidOptions on a cluster topology, which cannot offer multi-key transactions, or RedisError::TransactionBusy when a transaction is already open on this handle (the open one is left untouched).

Source§

async fn commit(&self) -> Result<(), Self::Error>

Flushes the buffered XADDs as one MULTI / EXEC block, in publish order, then clears the transaction.

§Errors

Returns RedisError::NoTransaction when no transaction is open on this handle, RedisError::ShutDown when the connection is gone, or RedisError::Publish if the block is rejected. On failure the transaction is already closed: the buffer is lost, and recovery is redelivery of the inputs rather than resubmission of the buffer.

Source§

async fn abort(&self) -> Result<(), Self::Error>

Discards the buffered messages.

§Errors

Returns RedisError::NoTransaction when no transaction is open on this handle.

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<P> ErasedPublisher for P
where P: Publisher,

Source§

fn publish_bytes<'a>( &'a self, name: &'a str, payload: &'a [u8], ) -> Pin<Box<dyn Future<Output = Result<(), Box<dyn Error + Sync + Send>>> + Send + 'a>>

Publishes payload to name, with no headers. Read more
Source§

fn publish_message<'a>( &'a self, name: &'a str, payload: &'a [u8], headers: &'a Headers, ) -> Pin<Box<dyn Future<Output = Result<(), Box<dyn Error + Sync + Send>>> + Send + 'a>>

Publishes payload to name with headers. 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<Reply, DeliveryCx, Pipeline, Bare> ReplySink<Reply, DeliveryCx, Pipeline> for Bare
where Reply: AsRef<[u8]> + Sync, DeliveryCx: Sync, Pipeline: Send + Sync, Bare: Publisher,

Source§

type Error = <Bare as Publisher>::Error

The error surfaced when the reply cannot be published.
Source§

async fn deliver( &self, name: &str, reply: &Reply, _pipeline: &Pipeline, _cx: &PublishContext<'_, DeliveryCx>, ) -> Result<(), <Bare as ReplySink<Reply, DeliveryCx, Pipeline>>::Error>

Publishes reply to name.
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