Skip to main content

lapin_futures_native_tls/
lib.rs

1#![deny(missing_docs)]
2#![warn(rust_2018_idioms)]
3#![doc(html_root_url = "https://docs.rs/lapin-futures-native-tls/0.10.1/")]
4
5//! lapin-futures-native-tls
6//!
7//! This library offers a nice integration of `native-tls` with the `lapin-futures` library.
8//! It uses `amq-protocol` URI parsing feature and adds the `connect` and `connect_cancellable`
9//! methods to `AMQPUri` which will provide you with a `lapin_futures::client::Client` and
10//! optionally a `lapin_futures::client::HeartbeatHandle` wrapped in a `Future`.
11//!
12//! It autodetects whether you're using `amqp` or `amqps` and opens either a raw `TcpStream`
13//! or a `TlsStream` using `native-tls` as the SSL engine.
14//!
15//! ## Connecting and opening a channel
16//!
17//! ```rust,no_run
18//! use env_logger;
19//! use failure::Error;
20//! use futures::future::Future;
21//! use lapin_futures_native_tls::{AMQPConnectionNativeTlsExt, lapin};
22//! use lapin::channel::ConfirmSelectOptions;
23//! use tokio;
24//!
25//! fn main() {
26//!     env_logger::init();
27//!
28//!     tokio::run(
29//!         "amqps://user:pass@host/vhost?heartbeat=10".connect_cancellable(|err| {
30//!             eprintln!("heartbeat error: {:?}", err);
31//!         }).map_err(Error::from).and_then(|(client, heartbeat_handle)| {
32//!             println!("Connected!");
33//!             client.create_confirm_channel(ConfirmSelectOptions::default()).map(|channel| (channel, heartbeat_handle)).and_then(|(channel, heartbeat_handle)| {
34//!                 println!("Stopping heartbeat.");
35//!                 heartbeat_handle.stop();
36//!                 println!("Closing channel.");
37//!                 channel.close(200, "Bye")
38//!             }).map_err(Error::from)
39//!         }).map_err(|err| {
40//!             eprintln!("amqp error: {:?}", err);
41//!         })
42//!     );
43//! }
44//! ```
45
46/// Reexport of the `lapin_futures_tls_internal` errors
47#[deprecated(note = "use lapin directly instead")]
48pub mod error;
49/// Reexport of the `lapin_futures` crate
50#[deprecated(note = "use lapin directly instead")]
51pub mod lapin;
52/// Reexport of the `uri` module from the `amq_protocol` crate
53#[deprecated(note = "use lapin directly instead")]
54pub mod uri;
55
56/// Reexport of `AMQPStream`
57#[deprecated(note = "use lapin directly instead")]
58pub type AMQPStream = lapin_futures_tls_internal::AMQPStream<TlsStream<TcpStream>>;
59
60use futures::{self, future::Future};
61use lapin_futures_tls_internal::{self, AMQPConnectionTlsExt, error::Error, lapin::client::ConnectionProperties, TcpStream};
62use native_tls;
63use tokio_tls::{TlsConnector, TlsStream};
64
65use std::io;
66
67use uri::AMQPUri;
68
69fn connector(host: String, stream: TcpStream) -> Box<dyn Future<Item = Box<TlsStream<TcpStream>>, Error = io::Error> + Send + 'static> {
70    Box::new(futures::future::result(native_tls::TlsConnector::builder().build().map_err(|_| io::Error::new(io::ErrorKind::Other, "Failed to create connector"))).and_then(move |connector| {
71        TlsConnector::from(connector).connect(&host, stream).map_err(|_| io::Error::new(io::ErrorKind::Other, "Failed to connect")).map(Box::new)
72    }))
73}
74
75/// Add a connect method providing a `lapin_futures::client::Client` wrapped in a `Future`.
76#[deprecated(note = "use lapin directly instead")]
77pub trait AMQPConnectionNativeTlsExt: AMQPConnectionTlsExt<TlsStream<TcpStream>> where Self: Sized {
78    /// Method providing a `lapin_futures::client::Client`, a `lapin_futures::client::HeartbeatHandle` and a `lapin::client::Heartbeat` pulse wrapped in a `Future`
79    fn connect(self) -> Box<dyn Future<Item = (lapin::client::Client<AMQPStream>, lapin::client::HeartbeatHandle, Box<dyn Future<Item = (), Error = Error> + Send + 'static>), Error = Error> + Send + 'static> {
80        AMQPConnectionTlsExt::connect(self, connector)
81    }
82    /// Method providing a `lapin_futures::client::Client` and `lapin_futures::client::HeartbeatHandle` wrapped in a `Future`
83    fn connect_cancellable<F: FnOnce(Error) + Send + 'static>(self, heartbeat_error_handler: F) -> Box<dyn Future<Item = (lapin::client::Client<AMQPStream>, lapin::client::HeartbeatHandle), Error = Error> + Send + 'static> {
84        AMQPConnectionTlsExt::connect_cancellable(self, heartbeat_error_handler, connector)
85    }
86    /// Method providing a `lapin_futures::client::Client`, a `lapin_futures::client::HeartbeatHandle` and a `lapin::client::Heartbeat` pulse wrapped in a `Future`
87    fn connect_full(self, properties: ConnectionProperties) -> Box<dyn Future<Item = (lapin::client::Client<AMQPStream>, lapin::client::HeartbeatHandle, Box<dyn Future<Item = (), Error = Error> + Send + 'static>), Error = Error> + Send + 'static> {
88        AMQPConnectionTlsExt::connect_full(self, connector, properties)
89    }
90    /// Method providing a `lapin_futures::client::Client` and `lapin_futures::client::HeartbeatHandle` wrapped in a `Future`
91    fn connect_cancellable_full<F: FnOnce(Error) + Send + 'static>(self, heartbeat_error_handler: F, properties: ConnectionProperties) -> Box<dyn Future<Item = (lapin::client::Client<AMQPStream>, lapin::client::HeartbeatHandle), Error = Error> + Send + 'static> {
92        AMQPConnectionTlsExt::connect_cancellable_full(self, heartbeat_error_handler, connector, properties)
93    }
94}
95
96impl AMQPConnectionNativeTlsExt for AMQPUri {}
97impl<'a> AMQPConnectionNativeTlsExt for &'a str {}