Skip to main content

RedisEventBus

Struct RedisEventBus 

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

Redis-backed transient event publisher and reconnecting subscriber.

Implementations§

Source§

impl RedisEventBus

Source

pub fn from_url( url: &str, topic: impl Into<String>, ) -> Result<Self, RedisEventBusError>

Creates a bus from a Redis URL and logical topic.

Source

pub fn from_client(client: Client, topic: impl Into<String>) -> Self

Creates a bus from an existing Redis client.

Source

pub fn with_reconnect_policy(self, policy: RedisEventReconnectPolicy) -> Self

Overrides the reconnect policy.

Source

pub fn topic(&self) -> &str

Returns the configured logical topic.

Source

pub async fn publish( &self, payload: impl Into<String>, ) -> Result<(), RedisEventBusError>

Publishes one opaque payload. Payload interpretation belongs to the product layer.

Source

pub async fn subscribe( &self, ) -> Result<RedisEventSubscription, RedisEventBusError>

Opens one Redis subscription attempt.

Source

pub async fn run_subscription<F, Fut>( &self, shutdown: CancellationToken, observer: Option<&dyn EventConnectionObserver>, on_payload: F, )
where F: FnMut(String) -> Fut, Fut: Future<Output = ()>,

Runs a reconnecting subscription until shutdown is cancelled.

Malformed Redis payloads are logged and skipped. The callback is responsible for decoding the product payload and deciding whether an event belongs to the current runtime.

Trait Implementations§

Source§

impl Clone for RedisEventBus

Source§

fn clone(&self) -> RedisEventBus

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 EventSubscriptionSource for RedisEventBus

Source§

type Item = String

Item emitted by the transport.
Source§

type Subscription = RedisEventSubscription

One active transport subscription.
Source§

type Error = RedisEventBusError

Transport error returned by subscribe or receive.
Source§

fn subscribe<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<Self::Subscription, Self::Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Opens one subscription attempt.
Source§

fn receive<'life0, 'life1, 'async_trait>( &'life0 self, subscription: &'life1 mut Self::Subscription, ) -> Pin<Box<dyn Future<Output = Result<Self::Item, Self::Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Receives one item from an active subscription.

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

§

impl<T> Instrument for T

§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided [Span], returning an Instrumented wrapper. Read more
§

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

Source§

type Output = T

Should always be Self
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.
§

impl<T> WithSubscriber for T

§

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
§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a [WithDispatch] wrapper. Read more