use lapin;
use lapin::types::AMQPValue;
use lapin::types::FieldTable;
use lapin::types::LongString;
use lapin::types::ShortString;
use log::info;
use crate::Result;
use crate::rmq_primitive;
use crate::rmq_primitive::constant::RMQ_QUEUE_BIND_OPTIONS;
use crate::rmq_primitive::constant::RMQ_QUEUE_DECLARE_OPTIONS;
pub async fn init_work_queue<S: AsRef<str>>(
rmq_uri: S,
work_direct_exchange: S,
work_queue: S,
retry_direct_exchange: S,
retry_queue: S,
retry_interval_in_seconds: u32,
) -> Result<()> {
let rmq_uri = rmq_uri.as_ref();
let work_direct_exchange = work_direct_exchange.as_ref();
let retry_direct_exchange = retry_direct_exchange.as_ref();
let work_queue = work_queue.as_ref();
let retry_queue = retry_queue.as_ref();
let channel = rmq_primitive::create_channel(rmq_uri).await?;
let exchanges: [&str; 2] = [work_direct_exchange, retry_direct_exchange];
for exchange in exchanges.iter() {
info!("declaring exchange {}", exchange);
channel
.exchange_declare(
exchange,
lapin::ExchangeKind::Direct,
lapin::options::ExchangeDeclareOptions {
passive: false,
durable: true,
auto_delete: false,
internal: false,
nowait: false,
},
FieldTable::default(),
)
.await?;
}
info!("declaring work queue {}", work_queue);
let mut queue_args = FieldTable::default();
queue_args.insert(
ShortString::from("x-dead-letter-exchange"),
AMQPValue::from(LongString::from(retry_direct_exchange)),
);
queue_args.insert(
ShortString::from("x-dead-letter-routing-key"),
AMQPValue::from(LongString::from(retry_queue)),
);
channel
.queue_declare(work_queue, RMQ_QUEUE_DECLARE_OPTIONS, queue_args)
.await?;
info!(
"binding work queue {} to exchange {}",
work_queue, work_direct_exchange
);
channel
.queue_bind(
work_queue,
work_direct_exchange,
work_queue,
RMQ_QUEUE_BIND_OPTIONS,
FieldTable::default(),
)
.await?;
info!(
"declaring retry queue {} for work queue {}",
retry_queue, work_queue
);
let mut queue_args = FieldTable::default();
queue_args.insert(
ShortString::from("x-dead-letter-exchange"),
AMQPValue::from(LongString::from(work_direct_exchange)),
);
queue_args.insert(
ShortString::from("x-dead-letter-routing-key"),
AMQPValue::from(LongString::from(work_queue)),
);
queue_args.insert(
ShortString::from("x-message-ttl"),
AMQPValue::from(retry_interval_in_seconds * 1000),
);
channel
.queue_declare(retry_queue, RMQ_QUEUE_DECLARE_OPTIONS, queue_args)
.await?;
info!(
"binding retry queue {} to exchange {}",
retry_queue, retry_direct_exchange
);
channel
.queue_bind(
retry_queue,
retry_direct_exchange,
retry_queue,
RMQ_QUEUE_BIND_OPTIONS,
FieldTable::default(),
)
.await?;
channel
.close(
lapin::protocol::constants::REPLY_SUCCESS as lapin::types::ShortUInt,
"normal close of channel",
)
.await?;
info!(
"done creating exchanges and queues for work queue {}",
work_queue
);
Ok(())
}
pub async fn init_exchanges_and_queues<S: AsRef<str>>(
rmq_uri: S,
work_direct_exchange: S,
retry_direct_exchange: S,
queues: Vec<(S, S, u32)>,
) -> Result<()> {
for (work_queue, retry_queue, retry_interval_in_seconds) in queues.into_iter() {
let () = init_work_queue(
rmq_uri.as_ref(),
work_direct_exchange.as_ref(),
work_queue.as_ref(),
retry_direct_exchange.as_ref(),
retry_queue.as_ref(),
retry_interval_in_seconds,
).await?;
}
Ok(())
}