RSMQ port to async rust. RSMQ is a simple redis queue system that works in any redis v2.6+. It contains the same methods as the original one in https://github.com/smrchy/rsmq
This crate uses async in the implementation. If you want to use it in your sync code you can use tokio/asyncstd "blockon" method. Async was used in order to simplify the code and allow 1-to-1 port oft he JS code.
```rust,no_run
use rsmq_async::{Rsmq, RsmqError, RsmqConnection};
let mut rsmq = Rsmq::new(Default::default()).await?;
let message = rsmq.receive_message("myqueue", None).await?;
if let Some(message) = message { rsmq.delete_message("myqueue", &message.id).await?; }
```
Main object documentation are in: Rsmq and PooledRsmq and they both implement the trait RsmqConnection where you can see all the RSMQ methods. Make sure you always import the trait RsmqConnection.
Check https://crates.io/crates/rsmq_async
When initializing RSMQ you can enable the realtime PUBLISH for
new messages. On every new message that gets sent to RSQM via sendMessage
a
Redis PUBLISH will be issued to {rsmq.ns}:rt:{qname}
.
Besides the PUBLISH when a new message is sent to RSMQ nothing else will happen.
Your app could use the Redis SUBSCRIBE command to be notified of new messages
and issue a receiveMessage
then. However make sure not to listen with multiple
workers for new messages with SUBSCRIBE to prevent multiple simultaneous
receiveMessage
calls.
If you want to implement "at least one delivery" guarantee, you need to receive the messages using "receivemessage" and then, once the message is successfully processed, delete it with "deletemessage".
If you want to use a connection pool, just use PooledRsmq instad of Rsmq. It implements the RsmqConnection trait as the normal Rsmq.
Since version 0.16 where this pull request was merged
redis dependency supports tokio and async_std executors. By default it will guess what
you are using when creating the connection. You can
check redis Cargo.tolm
for
the flags async-std-comp
and tokio-comp
in order to choose one or the other.
```rust,norun use rsmqasync::{Rsmq, RsmqConnection};
async fn it_works() { let mut rsmq = Rsmq::new(Default::default()) .await .expect("connection failed");
rsmq.create_queue("myqueue", None, None, None)
.await
.expect("failed to create queue");
rsmq.send_message("myqueue", "testmessage", None)
.await
.expect("failed to send message");
let message = rsmq
.receive_message("myqueue", None)
.await
.expect("cannot receive message");
if let Some(message) = message {
rsmq.delete_message("myqueue", &message.id).await;
}
}
```