Skip to main content

ConnectedPulsarBroker

Struct ConnectedPulsarBroker 

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

The typed witness that connect succeeded: holds the live client directly.

Implementations§

Source§

impl ConnectedPulsarBroker

Source

pub fn publisher(&self) -> PulsarPublisher

A publisher from the connected form. It rides the same cell-backed publisher type as the early path; by now connect has filled the cell, so it resolves immediately.

Source

pub async fn subscribe_descriptor( &self, descriptor: PulsarSubscription, ) -> Result<PulsarSubscriber, PulsarError>

Opens the subscription described by descriptor.

§Errors

Returns PulsarError when the descriptor is invalid, the consumer cannot be created, or the broker is shut down.

Trait Implementations§

Source§

impl ConnectedBroker for ConnectedPulsarBroker

Source§

type Error = PulsarError

The error type returned by connected-broker operations.
Source§

type Closed = ()

The terminal witness returned by shutdown. Read more
Source§

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

Closes the broker connection, flushing in-flight publishes and stopping background tasks, consuming the connected form. Read more
Source§

impl Debug for ConnectedPulsarBroker

Source§

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

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

impl DefaultPublish for ConnectedPulsarBroker

Source§

type Policy = PulsarPublish

The broker’s plain publish policy, constructible with its defaults.
Source§

impl PublishPolicy<ConnectedPulsarBroker> for PulsarPublish

Source§

type Live = PulsarPublisher

The live form this policy pairs into: a Publisher for a leaf policy, or the live wiring form for a combinator stack (a typed publisher over a policy pairs into the same typed publisher over the live leaf).
Source§

async fn pair( self, connected: &ConnectedPulsarBroker, ) -> Result<Self::Live, PairError>

Pairs the policy with a connected broker, producing the live publisher. Read more
Source§

impl Subscribe for ConnectedPulsarBroker

Source§

type Subscriber = PulsarSubscriber

The subscriber type opened by a by-name subscription.
Source§

async fn subscribe(&self, name: &str) -> Result<Self::Subscriber, Self::Error>

Opens a subscription to name, producing this broker’s Subscriber. Read more
Source§

impl SubscriptionSource<ConnectedPulsarBroker> for PulsarSubscription

Source§

type Subscriber = PulsarSubscriber

The subscriber type this source opens.
Source§

fn name(&self) -> &str

The name (subject / channel) this subscription binds to. Read more
Source§

async fn subscribe( self, connected: &ConnectedPulsarBroker, ) -> Result<PulsarSubscriber, PulsarError>

Opens the subscription against the connected broker. Called once at startup. Read more

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