fix: resolve compilation errors in rabbitmq module by correcting async function return types and method calls

This commit is contained in:
Per Stark" (aider)
2024-09-18 12:24:32 +02:00
parent 0704a011b9
commit 303ddc5811

View File

@@ -37,10 +37,10 @@ impl RabbitMQ {
} }
pub async fn declare_queue(&self, queue_name: &str) -> Result<(String, u32, u32), Box<dyn std::error::Error>> { pub async fn declare_queue(&self, queue_name: &str) -> Result<(String, u32, u32), Box<dyn std::error::Error>> {
Ok(self.channel self.channel
.queue_declare(QueueDeclareArguments::durable_client_named(queue_name)) .queue_declare(QueueDeclareArguments::durable_client_named(queue_name))
.await?? .await?
) .ok_or_else(|| Box::new(std::io::Error::new(std::io::ErrorKind::Other, "Failed to declare queue")) as Box<dyn std::error::Error>)
} }
pub async fn bind_queue(&self, queue_name: &str, exchange_name: &str, routing_key: &str) -> Result<(), Box<dyn std::error::Error>> { pub async fn bind_queue(&self, queue_name: &str, exchange_name: &str, routing_key: &str) -> Result<(), Box<dyn std::error::Error>> {
@@ -62,17 +62,21 @@ impl RabbitMQ {
let (tx, rx) = mpsc::channel(100); let (tx, rx) = mpsc::channel(100);
let args = BasicConsumeArguments::new(queue_name, consumer_tag); let args = BasicConsumeArguments::new(queue_name, consumer_tag);
let consumer = DefaultConsumer::new(args.no_ack).with_callback(move |_deliver, _basic_properties, content| { let consumer = DefaultConsumer::new(args.no_ack);
let content_str = String::from_utf8_lossy(&content).to_string(); let tx_clone = tx.clone();
let tx = tx.clone();
tokio::spawn(async move { self.channel.basic_consume(consumer, args)
if let Err(e) = tx.send(content_str).await { .await?
eprintln!("Failed to send message: {}", e); .set_delegate(move |deliver: amqprs::consumer::DeliverEvent| {
} let content_str = String::from_utf8_lossy(&deliver.content).to_string();
let tx = tx_clone.clone();
tokio::spawn(async move {
if let Err(e) = tx.send(content_str).await {
eprintln!("Failed to send message: {}", e);
}
});
}); });
});
self.channel.basic_consume(consumer, args).await?;
Ok(rx) Ok(rx)
} }
} }