use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use sz_orm_queue::{
ActiveConfig, KafkaConfig, NatsConfig, PulsarConfig, RabbitConfig, RocketConfig,
};
use sz_orm_queue::{InMemoryQueue, MessageQueue, MqProvider, QueueWrapper};
#[tokio::test]
async fn stress_queue_100k_messages_single_thread() {
let queue = InMemoryQueue::new();
let total: u64 = 100_000;
for i in 0..total {
let payload = format!("msg-{}", i);
queue
.publish("bulk-topic", payload.as_bytes())
.await
.unwrap();
}
assert_eq!(queue.message_count("bulk-topic").await, total as usize);
for i in 0..total {
let msg = queue
.consume("bulk-topic")
.await
.unwrap()
.expect("msg must exist");
let expected = format!("msg-{}", i);
assert_eq!(
msg.payload,
expected.as_bytes(),
"FIFO broken at index {}",
i
);
queue.ack(&msg.id).await.unwrap();
}
assert_eq!(queue.message_count("bulk-topic").await, 0);
assert_eq!(
queue.in_flight_count().await,
0,
"in_flight must be 0 after ack"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 8)]
async fn stress_queue_concurrent_publish_consume() {
let queue = Arc::new(InMemoryQueue::new());
let total_per_task: u64 = 10_000;
let task_count: u64 = 8;
let total = total_per_task * task_count;
let mut handles = Vec::new();
for task_id in 0..task_count {
let q = queue.clone();
handles.push(tokio::spawn(async move {
let topic = format!("task-{}", task_id);
for i in 0..total_per_task {
let payload = format!("t{}-m{}", task_id, i);
q.publish(&topic, payload.as_bytes()).await.unwrap();
}
for i in 0..total_per_task {
let msg = q.consume(&topic).await.unwrap().expect("msg must exist");
let expected = format!("t{}-m{}", task_id, i);
assert_eq!(
msg.payload,
expected.as_bytes(),
"FIFO broken in task {}",
task_id
);
q.ack(&msg.id).await.unwrap();
}
}));
}
for h in handles {
h.await.unwrap();
}
for task_id in 0..task_count {
let topic = format!("task-{}", task_id);
assert_eq!(
queue.message_count(&topic).await,
0,
"topic {} not drained",
topic
);
}
assert_eq!(
queue.in_flight_count().await,
0,
"in_flight must be 0 after ack"
);
let _ = total;
}
#[tokio::test]
async fn stress_queue_many_topics() {
let queue = InMemoryQueue::new();
let topic_count: u64 = 1000;
let per_topic: u64 = 100;
for t in 0..topic_count {
let topic = format!("topic-{}", t);
for i in 0..per_topic {
queue
.publish(&topic, format!("m{}", i).as_bytes())
.await
.unwrap();
}
}
for t in 0..topic_count {
let topic = format!("topic-{}", t);
assert_eq!(
queue.message_count(&topic).await,
per_topic as usize,
"topic {} count mismatch",
topic
);
}
}
#[tokio::test]
async fn stress_queue_large_payload() {
let queue = InMemoryQueue::new();
let payload_size: usize = 1_000_000; let count: usize = 1000;
let payload = vec![0xABu8; payload_size];
for _ in 0..count {
queue.publish("large", &payload).await.unwrap();
}
assert_eq!(queue.message_count("large").await, count);
for _ in 0..count {
let msg = queue.consume("large").await.unwrap().unwrap();
assert_eq!(msg.payload.len(), payload_size);
queue.ack(&msg.id).await.unwrap();
}
assert_eq!(queue.in_flight_count().await, 0);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 6)]
async fn stress_queue_all_providers_concurrent() {
let providers = vec![
MqProvider::Kafka(KafkaConfig::default()),
MqProvider::RabbitMQ(RabbitConfig::default()),
MqProvider::RocketMQ(RocketConfig::default()),
MqProvider::ActiveMQ(ActiveConfig::default()),
MqProvider::Nats(NatsConfig::default()),
MqProvider::Pulsar(PulsarConfig::default()),
];
let success = Arc::new(AtomicU64::new(0));
let mut handles = Vec::new();
for provider in providers {
let s = success.clone();
handles.push(tokio::spawn(async move {
let wrapper = QueueWrapper::new(provider);
for i in 0..1000 {
let payload = format!("msg-{}", i);
wrapper.publish("test", payload.as_bytes()).await.unwrap();
let msg = wrapper.consume("test").await.unwrap().expect("msg");
assert_eq!(msg.payload, payload.as_bytes());
wrapper.ack(&msg.id).await.unwrap();
s.fetch_add(1, Ordering::Relaxed);
}
}));
}
for h in handles {
h.await.unwrap();
}
assert_eq!(success.load(Ordering::Relaxed), 6000);
}
#[tokio::test]
async fn stress_queue_subscriber_count() {
let queue = InMemoryQueue::new();
let n: usize = 10_000;
for _ in 0..n {
queue.subscribe("hot-topic").await.unwrap();
}
assert_eq!(queue.subscriber_count("hot-topic").await, n);
}
#[tokio::test]
async fn stress_queue_consume_empty_repeated() {
let queue = InMemoryQueue::new();
for _ in 0..10_000 {
let result = queue.consume("never-published").await.unwrap();
assert!(result.is_none());
}
}
#[tokio::test]
async fn stress_queue_ack_unknown_repeated() {
let queue = InMemoryQueue::new();
for i in 0..1000 {
let result = queue.ack(&format!("unknown-{}", i)).await;
assert!(result.is_err(), "ack unknown must fail at iter {}", i);
}
}