use super::*;
#[derive(Clone, Debug, Default)]
pub(crate) struct QueueStats {
pub visible: i64,
pub in_flight: i64,
pub delayed: i64,
}
#[derive(Clone, Debug)]
pub(crate) struct QueueMessage {
pub id: String,
pub receipt_handle: String,
pub body: String,
pub receive_count: i64,
pub sent_at: Option<DateTime<Utc>>,
pub task: Option<SqsdTask>,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub(crate) struct SqsdTask {
pub name: Option<String>,
pub path: Option<String>,
pub scheduled_time_raw: Option<String>,
pub scheduled_at: Option<DateTime<Utc>>,
}
pub(crate) fn sqsd_task_from(
attrs: Option<&std::collections::HashMap<String, aws_sdk_sqs::types::MessageAttributeValue>>,
) -> Option<SqsdTask> {
let attrs = attrs?;
let get = |k: &str| -> Option<String> {
attrs
.get(k)
.and_then(|v| v.string_value.clone())
.filter(|s| !s.is_empty())
};
let name = get("beanstalk.sqsd.task_name");
let path = get("beanstalk.sqsd.path");
let scheduled_time_raw = get("beanstalk.sqsd.scheduled_time");
if name.is_none() && path.is_none() && scheduled_time_raw.is_none() {
return None;
}
let scheduled_at = scheduled_time_raw
.as_deref()
.and_then(parse_sqsd_scheduled_time);
Some(SqsdTask {
name,
path,
scheduled_time_raw,
scheduled_at,
})
}
pub(crate) fn parse_sqsd_scheduled_time(raw: &str) -> Option<DateTime<Utc>> {
chrono::NaiveDateTime::parse_from_str(raw.trim(), "%Y-%m-%d %H:%M:%S UTC")
.ok()
.map(|naive| naive.and_utc())
}
pub(crate) fn derive_dlq_url(main: &str) -> Option<String> {
let trimmed = main.trim_end_matches('/');
if trimmed.ends_with("-dlq") {
return None;
}
Some(format!("{trimmed}-dlq"))
}
impl AwsClient {
pub(crate) async fn queue_stats(&self, queue_url: &str) -> Result<QueueStats> {
use aws_sdk_sqs::types::QueueAttributeName as Q;
let resp = self
.sqs
.get_queue_attributes()
.queue_url(queue_url)
.attribute_names(Q::ApproximateNumberOfMessages)
.attribute_names(Q::ApproximateNumberOfMessagesNotVisible)
.attribute_names(Q::ApproximateNumberOfMessagesDelayed)
.send()
.await?;
let attrs = resp.attributes.unwrap_or_default();
let parse = |k: Q| -> i64 {
attrs
.get(&k)
.and_then(|v| v.parse::<i64>().ok())
.unwrap_or(0)
};
Ok(QueueStats {
visible: parse(Q::ApproximateNumberOfMessages),
in_flight: parse(Q::ApproximateNumberOfMessagesNotVisible),
delayed: parse(Q::ApproximateNumberOfMessagesDelayed),
})
}
pub(crate) async fn peek_messages(
&self,
queue_url: &str,
max: i32,
) -> Result<Vec<QueueMessage>> {
use aws_sdk_sqs::types::MessageSystemAttributeName as M;
let target = max.clamp(1, 100) as usize;
let mut out: Vec<QueueMessage> = Vec::new();
let mut seen: std::collections::HashSet<String> = std::collections::HashSet::new();
let mut empty_in_a_row = 0;
for _ in 0..((target / 10).max(1) + 4) {
if out.len() >= target {
break;
}
let resp = self
.sqs
.receive_message()
.queue_url(queue_url)
.max_number_of_messages(((target - out.len()).clamp(1, 10)) as i32)
.visibility_timeout(5)
.wait_time_seconds(1)
.message_system_attribute_names(M::ApproximateReceiveCount)
.message_system_attribute_names(M::SentTimestamp)
.message_attribute_names("All")
.send()
.await
.wrap_err("ReceiveMessage failed")?;
let batch = resp.messages.unwrap_or_default();
if batch.is_empty() {
empty_in_a_row += 1;
if empty_in_a_row >= 2 {
break;
}
continue;
}
empty_in_a_row = 0;
for m in batch {
let id = m.message_id.clone().unwrap_or_default();
if !id.is_empty() && !seen.insert(id.clone()) {
continue;
}
let attrs = m.attributes.unwrap_or_default();
let receive_count = attrs
.get(&M::ApproximateReceiveCount)
.and_then(|v| v.parse::<i64>().ok())
.unwrap_or(0);
let sent_at = attrs
.get(&M::SentTimestamp)
.and_then(|v| v.parse::<i64>().ok())
.and_then(DateTime::from_timestamp_millis);
let task = sqsd_task_from(m.message_attributes.as_ref());
out.push(QueueMessage {
id,
receipt_handle: m.receipt_handle.unwrap_or_default(),
body: m.body.unwrap_or_default(),
receive_count,
sent_at,
task,
});
if out.len() >= target {
break;
}
}
}
Ok(out)
}
pub(crate) async fn send_message(&self, queue_url: &str, body: &str) -> Result<()> {
self.sqs
.send_message()
.queue_url(queue_url)
.message_body(body)
.send()
.await?;
Ok(())
}
pub(crate) async fn delete_message(&self, queue_url: &str, receipt_handle: &str) -> Result<()> {
self.sqs
.delete_message()
.queue_url(queue_url)
.receipt_handle(receipt_handle)
.send()
.await?;
Ok(())
}
pub(crate) async fn purge_queue(&self, queue_url: &str) -> Result<()> {
self.sqs.purge_queue().queue_url(queue_url).send().await?;
Ok(())
}
}
#[cfg(test)]
mod sqsd_task_tests {
use super::{parse_sqsd_scheduled_time, sqsd_task_from, SqsdTask};
use aws_sdk_sqs::types::MessageAttributeValue;
use std::collections::HashMap;
fn attr(v: &str) -> MessageAttributeValue {
MessageAttributeValue::builder()
.data_type("String")
.string_value(v)
.build()
.expect("valid attribute")
}
fn fixture() -> HashMap<String, MessageAttributeValue> {
HashMap::from([
(
"beanstalk.sqsd.task_name".to_string(),
attr("Remove unattended jobs"),
),
(
"beanstalk.sqsd.path".to_string(),
attr("/STCleanupUnattendedJobs.do"),
),
(
"beanstalk.sqsd.scheduled_time".to_string(),
attr("2026-09-17 06:04:00 UTC"),
),
])
}
#[test]
fn a_real_worker_task_is_extracted_whole() {
let t = sqsd_task_from(Some(&fixture())).expect("an EB task");
assert_eq!(t.name.as_deref(), Some("Remove unattended jobs"));
assert_eq!(t.path.as_deref(), Some("/STCleanupUnattendedJobs.do"));
assert_eq!(
t.scheduled_time_raw.as_deref(),
Some("2026-09-17 06:04:00 UTC"),
"the raw string must be kept exactly as sent"
);
assert_eq!(
t.scheduled_at.map(|d| d.to_rfc3339()),
Some("2026-09-17T06:04:00+00:00".to_string())
);
}
#[test]
fn the_scheduled_time_format_is_ebs_own() {
assert!(
parse_sqsd_scheduled_time("2026-09-17 06:04:00 UTC").is_some(),
"the format EB actually sends must parse"
);
for not_it in [
"2026-09-17T06:04:00Z",
"2026-09-17T06:04:00+00:00",
"1789625040",
"1789625040068",
"",
"not a time at all",
] {
assert!(
parse_sqsd_scheduled_time(not_it).is_none(),
"{not_it:?} is not EB's format and must not silently parse"
);
}
}
#[test]
fn a_plain_message_is_not_a_task() {
assert_eq!(sqsd_task_from(None), None);
assert_eq!(sqsd_task_from(Some(&HashMap::new())), None);
let other = HashMap::from([("my.app.attribute".to_string(), attr("something"))]);
assert_eq!(
sqsd_task_from(Some(&other)),
None,
"a non-sqsd attribute must not make this look like an EB task"
);
}
#[test]
fn a_partial_task_keeps_what_it_has() {
let partial = HashMap::from([
(
"beanstalk.sqsd.task_name".to_string(),
attr("Nightly sweep"),
),
(
"beanstalk.sqsd.scheduled_time".to_string(),
attr("whenever EB feels like it"),
),
]);
let t = sqsd_task_from(Some(&partial)).expect("still a task");
assert_eq!(t.name.as_deref(), Some("Nightly sweep"));
assert_eq!(t.path, None);
assert_eq!(
t.scheduled_time_raw.as_deref(),
Some("whenever EB feels like it"),
"an unparseable time must still be shown verbatim"
);
assert_eq!(t.scheduled_at, None);
assert_ne!(t, SqsdTask::default(), "and must not be an empty task");
}
#[test]
fn the_peek_requests_custom_message_attributes() {
let src = std::fs::read_to_string("src/aws/sqs.rs").expect("read own source");
let prod = src.split("#[cfg(test)]").next().expect("production half");
assert!(
prod.contains("receive_message()"),
"the scan is not finding the peek at all"
);
assert!(
prod.contains(".message_attribute_names("),
"the peek must request custom attributes, or no message ever \
carries a task"
);
}
}