pub struct RMQConnectionManager { /* private fields */ }Expand description
A mobc::Manager for lapin::Connections.
§Example
use mobc::Pool;
use mobc_lapin::RMQConnectionManager;
use tokio_amqp::*;
use futures::StreamExt;
use std::time::Duration;
use lapin::{
options::*, types::FieldTable, BasicProperties, publisher_confirm::Confirmation,
ConnectionProperties,
};
const PAYLOAD: &[u8;13] = b"Hello, World!";
const QUEUE_NAME: &str = "test";
#[tokio::main]
async fn main() {
let addr = "amqp://rmq:rmq@127.0.0.1:5672/%2f";
let manager = RMQConnectionManager::new(
addr.to_owned(),
ConnectionProperties::default().with_executor(tokio_executor_trait::Tokio::current()),
);
let pool = Pool::<RMQConnectionManager>::builder()
.max_open(5)
.build(manager);
let conn = pool.get().await.unwrap();
let channel = conn.create_channel().await.unwrap();
let _ = channel
.queue_declare(
QUEUE_NAME,
QueueDeclareOptions::default(),
FieldTable::default(),
)
.await.unwrap();
// send messages to the queue
println!("spawning senders...");
for i in 0..50 {
let send_pool = pool.clone();
let send_props = BasicProperties::default().with_kind(format!("Sender: {}", i).into());
tokio::spawn(async move {
let mut interval = tokio::time::interval(Duration::from_millis(200));
loop {
interval.tick().await;
let send_conn = send_pool.get().await.unwrap();
let send_channel = send_conn.create_channel().await.unwrap();
let confirm = send_channel
.basic_publish(
"",
QUEUE_NAME,
BasicPublishOptions::default(),
PAYLOAD.as_ref(),
send_props.clone(),
)
.await.unwrap()
.await.unwrap();
assert_eq!(confirm, Confirmation::NotRequested);
}
});
}
// listen for incoming messages from the queue
let mut consumer = channel
.basic_consume(
QUEUE_NAME,
"my_consumer",
BasicConsumeOptions::default(),
FieldTable::default(),
)
.await.unwrap();
println!("listening to messages...");
while let Some(delivery) = consumer.next().await {
let delivery = delivery.expect("error in consumer");
println!("incoming message from: {:?}", delivery.properties.kind());
channel
.basic_ack(delivery.delivery_tag, BasicAckOptions::default())
.await
.expect("ack");
}
}Implementations§
Source§impl RMQConnectionManager
impl RMQConnectionManager
pub fn new(addr: String, connection_properties: ConnectionProperties) -> Self
Trait Implementations§
Source§impl Clone for RMQConnectionManager
impl Clone for RMQConnectionManager
Source§fn clone(&self) -> RMQConnectionManager
fn clone(&self) -> RMQConnectionManager
Returns a duplicate of the value. Read more
1.0.0 · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
Performs copy-assignment from
source. Read moreSource§impl Manager for RMQConnectionManager
impl Manager for RMQConnectionManager
Source§type Connection = Connection
type Connection = Connection
The connection type this manager deals with.
Source§fn connect<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<Self::Connection, Self::Error>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn connect<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<Self::Connection, Self::Error>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Attempts to create a new connection.
Source§fn check<'life0, 'async_trait>(
&'life0 self,
conn: Self::Connection,
) -> Pin<Box<dyn Future<Output = Result<Self::Connection, Self::Error>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn check<'life0, 'async_trait>(
&'life0 self,
conn: Self::Connection,
) -> Pin<Box<dyn Future<Output = Result<Self::Connection, Self::Error>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Determines if the connection is still connected to the database when check-out. Read more
Source§fn spawn_task<T>(&self, task: T)
fn spawn_task<T>(&self, task: T)
Spawns a new asynchronous task.
Source§fn validate(&self, _conn: &mut Self::Connection) -> bool
fn validate(&self, _conn: &mut Self::Connection) -> bool
Quickly determines a connection is still valid when check-in.
Auto Trait Implementations§
impl Freeze for RMQConnectionManager
impl !RefUnwindSafe for RMQConnectionManager
impl Send for RMQConnectionManager
impl Sync for RMQConnectionManager
impl Unpin for RMQConnectionManager
impl !UnwindSafe for RMQConnectionManager
Blanket Implementations§
Source§impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
Source§impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Mutably borrows from an owned value. Read more