use std::time::Duration;
use ruststream::SubscriptionSource;
use crate::broker::ConnectedSqsBroker;
use crate::error::SqsError;
use crate::subscriber::SqsSubscriber;
const MAX_WAIT: Duration = Duration::from_secs(20);
#[derive(Debug, Clone, PartialEq, Eq)]
#[must_use]
pub struct SqsQueue {
queue: String,
wait: Duration,
batch: i32,
visibility: Option<Duration>,
create_if_missing: bool,
}
impl SqsQueue {
pub fn new(queue: impl Into<String>) -> Self {
Self {
queue: queue.into(),
wait: MAX_WAIT,
batch: 10,
visibility: None,
create_if_missing: false,
}
}
pub fn wait(mut self, wait: Duration) -> Self {
self.wait = wait;
self
}
pub fn batch(mut self, batch: i32) -> Self {
self.batch = batch;
self
}
pub fn visibility(mut self, visibility: Duration) -> Self {
self.visibility = Some(visibility);
self
}
pub fn create_if_missing(mut self) -> Self {
self.create_if_missing = true;
self
}
#[must_use]
pub fn queue(&self) -> &str {
&self.queue
}
pub(crate) fn wait_value(&self) -> Duration {
self.wait
}
pub(crate) fn batch_value(&self) -> i32 {
self.batch
}
pub(crate) fn visibility_value(&self) -> Option<Duration> {
self.visibility
}
pub(crate) fn create_value(&self) -> bool {
self.create_if_missing
}
pub(crate) fn validate(&self) -> Result<(), SqsError> {
if self.queue.is_empty() {
return Err(SqsError::InvalidQueue("queue must be non-empty".into()));
}
if self.wait > MAX_WAIT {
return Err(SqsError::InvalidQueue(
"wait exceeds the 20 second long-polling cap".into(),
));
}
if !(1..=10).contains(&self.batch) {
return Err(SqsError::InvalidQueue(
"batch must be within 1..=10 (the receive cap)".into(),
));
}
if let Some(visibility) = self.visibility
&& (visibility.is_zero() || visibility > Duration::from_hours(12))
{
return Err(SqsError::InvalidQueue(
"visibility must be within 1s..=12h".into(),
));
}
Ok(())
}
}
impl SubscriptionSource<ConnectedSqsBroker> for SqsQueue {
type Subscriber = SqsSubscriber;
fn name(&self) -> &str {
self.queue()
}
async fn subscribe(self, connected: &ConnectedSqsBroker) -> Result<SqsSubscriber, SqsError> {
connected.subscribe_queue(self).await
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn empty_queue_is_rejected_before_io() {
assert!(matches!(
SqsQueue::new("").validate(),
Err(SqsError::InvalidQueue(_))
));
}
#[test]
fn overlong_wait_is_rejected_before_io() {
assert!(matches!(
SqsQueue::new("q").wait(Duration::from_secs(21)).validate(),
Err(SqsError::InvalidQueue(_))
));
}
#[test]
fn out_of_range_batch_is_rejected_before_io() {
assert!(matches!(
SqsQueue::new("q").batch(11).validate(),
Err(SqsError::InvalidQueue(_))
));
assert!(matches!(
SqsQueue::new("q").batch(0).validate(),
Err(SqsError::InvalidQueue(_))
));
}
}