Skip to main content

LapinBroker

Struct LapinBroker 

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

A RabbitMQ broker backed by lapin: configuration captured, no I/O performed yet.

new is synchronous and records only the connection settings, so a RabbitMQ service is assembled with the synchronous #[ruststream::app] builder like any other broker. The runtime calls Broker::connect once at startup, which consumes this value and yields the ConnectedLapinBroker witness: subscriptions, publishers, and requesters exist only from there, so “not connected” is not representable.

By default the broker never creates infrastructure: descriptors describe the EXPECTED topology, and a missing queue is a subscribe error. Opt into declaration with declare_topology(true).

§Examples

use ruststream::nonzero;
use ruststream_lapin::{LapinBroker, QueueType};

let broker = LapinBroker::new("amqp://localhost:5672")
    .prefetch(nonzero!(64))
    .default_queue_type(QueueType::Quorum);

Implementations§

Source§

impl LapinBroker

Source

pub fn new(uri: impl Into<String>) -> Self

Records the connection URI; no I/O happens until Broker::connect.

The URI carries credentials, virtual host, and TLS scheme: amqp://user:pass@host:5672/vhost (or amqps:// with a TLS feature enabled).

Source

pub fn connection_name(self, name: impl Into<String>) -> Self

A connection name shown in the RabbitMQ management UI.

Source

pub fn prefetch(self, prefetch: NonZeroU16) -> Self

Caps unacknowledged deliveries in flight per subscription (basic.qos).

This is the back-pressure window for subscriber streams; individual queue descriptors can override it. Without it the server imposes no prefetch limit.

The count is a NonZeroU16 because AMQP reads basic.qos(0) as “no limit”, the exact opposite of a cap: the zero sentinel is unrepresentable here, and leaving the prefetch unset is how “unlimited” is spelled.

Source

pub fn declare_topology(self, declare: bool) -> Self

Whether subscribing declares the descriptor’s expected topology first. Defaults to false: managing infrastructure is the user’s job, so creation is a deliberate opt-in.

When enabled, subscribing declares the bound exchanges (except the built-in amq.* ones and the default exchange), the queue, and the bindings.

Source

pub fn default_queue_type(self, queue_type: QueueType) -> Self

The queue type declared for descriptors that do not set one.

Only consulted when declare_topology is enabled. Without a broker default or a per-queue type, no x-queue-type argument is sent and the server default applies.

Trait Implementations§

Source§

impl Broker for LapinBroker

Source§

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

Opens the connection and its shared publish channel, consuming the configuration.

§Errors

Returns AmqpError::Connect when the URI cannot be parsed or the connection fails.

Source§

type Error = AmqpError

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

type Connected = ConnectedLapinBroker

The connected form of this broker: the typed witness that connect succeeded.
Source§

fn bindable(self) -> Bindable<Self>
where Self: 'static,

Wraps the broker for publisher-token minting before registration: the returned Bindable hands out Bound tokens via its bind, and is itself what with_broker takes. Read more
Source§

impl Clone for LapinBroker

Source§

fn clone(&self) -> LapinBroker

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 LapinBroker

Source§

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

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

impl DescribeServer for LapinBroker

DescribeServer reports the configured AMQP address, which is what the AsyncAPI document records for the service.

Source§

fn describe_server(&self) -> ServerSpec

Returns the server coordinates for this broker.

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<'a, T, E> AsTaggedExplicit<'a, E> for T
where T: 'a,

Source§

fn explicit(self, class: Class, tag: u32) -> TaggedParser<'a, Explicit, Self, E>

Source§

impl<'a, T, E> AsTaggedImplicit<'a, E> for T
where T: 'a,

Source§

fn implicit( self, class: Class, constructed: bool, tag: u32, ) -> TaggedParser<'a, Implicit, Self, E>

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<B> BrokerRegistration for B
where B: Broker + 'static,

Source§

type Broker = B

The broker being registered.
Source§

fn into_parts(self) -> (B, Arc<Mutex<Option<Arc<<B as Broker>::Connected>>>>)

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

Source§

fn compat(self) -> Compat<T>
where T: Sized,

Applies the Compat adapter by value. Read more
Source§

fn compat_ref(&self) -> Compat<&T>

Applies the Compat adapter by shared reference. Read more
Source§

fn compat_mut(&mut self) -> Compat<&mut T>

Applies the Compat adapter by mutable reference. 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> 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.
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