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
impl LapinBroker
Sourcepub fn new(uri: impl Into<String>) -> Self
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).
Sourcepub fn connection_name(self, name: impl Into<String>) -> Self
pub fn connection_name(self, name: impl Into<String>) -> Self
A connection name shown in the RabbitMQ management UI.
Sourcepub fn prefetch(self, prefetch: NonZeroU16) -> Self
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.
Sourcepub fn declare_topology(self, declare: bool) -> Self
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.
Sourcepub fn default_queue_type(self, queue_type: QueueType) -> Self
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
impl Broker for LapinBroker
Source§async fn connect(self) -> Result<Self::Connected, Self::Error>
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 Connected = ConnectedLapinBroker
type Connected = ConnectedLapinBroker
connect
succeeded.Source§impl Clone for LapinBroker
impl Clone for LapinBroker
Source§fn clone(&self) -> LapinBroker
fn clone(&self) -> LapinBroker
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read moreSource§impl Debug for LapinBroker
impl Debug for LapinBroker
Source§impl DescribeServer for LapinBroker
DescribeServer reports the configured AMQP address, which is what the AsyncAPI document
records for the service.
impl DescribeServer for LapinBroker
DescribeServer reports the configured AMQP address, which is what the AsyncAPI document
records for the service.