pub struct LapinBroker { /* private fields */ }Expand description
A RabbitMQ broker backed by lapin.
Follows the RustStream lazy startup contract: new is synchronous and does no
I/O; the network work happens in the idempotent async Broker::connect, which the runtime
calls once at startup. Publishers handed out earlier share the connection cell and resolve it
on first use.
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_lapin::{LapinBroker, QueueType};
let broker = LapinBroker::new("amqp://localhost:5672")
.prefetch(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 async fn connect(uri: impl Into<String>) -> Result<Self, AmqpError>
pub async fn connect(uri: impl Into<String>) -> Result<Self, AmqpError>
Connects eagerly: new followed by Broker::connect.
§Errors
Returns AmqpError::Connect when the connection cannot be established.
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: u16) -> Self
pub fn prefetch(self, prefetch: u16) -> 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.
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.
Sourcepub async fn subscribe(
&self,
def: RabbitQueue,
) -> Result<LapinSubscriber, AmqpError>
pub async fn subscribe( &self, def: RabbitQueue, ) -> Result<LapinSubscriber, AmqpError>
Opens a subscription for def, declaring its topology first when the broker opted in.
§Errors
Returns AmqpError::NotConnected before Broker::connect, AmqpError::Declare when
opted-in declaration fails, AmqpError::InvalidOptions for contradictory descriptor
options, and AmqpError::Subscribe when the channel or consumer cannot be opened (for
example the queue does not exist and declaration was not opted into).
Sourcepub fn publisher(&self) -> LapinPublisher
pub fn publisher(&self) -> LapinPublisher
Sourcepub fn requester(&self) -> LapinRequester
pub fn requester(&self) -> LapinRequester
A request/reply client over RabbitMQ direct reply-to.
Trait Implementations§
Source§impl Broker for LapinBroker
impl Broker for LapinBroker
Source§async fn connect(&self) -> Result<(), Self::Error>
async fn connect(&self) -> Result<(), Self::Error>
Establishes the connection and the shared publish channel; idempotent.
§Errors
Returns AmqpError::Connect when the URI cannot be parsed or the connection fails.
Source§async fn shutdown(&self) -> Result<(), Self::Error>
async fn shutdown(&self) -> Result<(), Self::Error>
Closes the connection; further operations fail with AmqpError::NotConnected or a
channel error. Idempotent: closing an already-closed connection succeeds.
§Errors
Returns AmqpError::Connect when the close handshake fails.
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
impl DescribeServer for LapinBroker
Source§fn describe_server(&self) -> ServerSpec
fn describe_server(&self) -> ServerSpec
Source§impl Subscribe for LapinBroker
impl Subscribe for LapinBroker
Source§async fn subscribe(&self, name: &str) -> Result<Self::Subscriber, Self::Error>
async fn subscribe(&self, name: &str) -> Result<Self::Subscriber, Self::Error>
Subscribes to the queue name with descriptor defaults (durable, shared).