Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
128 changes: 71 additions & 57 deletions queue/src/adapters/rabbitmq/adapter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -117,14 +117,48 @@ struct SubscriptionInfo {
impl RabbitMQAdapter {
/// Resolve the DLQ queue name for a given topic. Handles both function
/// queues (`__fn_queue::` prefix) and topic-based queues.
fn resolve_dlq_name(topic: &str) -> (String, bool) {
fn resolve_dlq_name(topic: &str, is_function_queue: bool) -> (String, bool) {
if let Some(queue_name) = topic.strip_prefix("__fn_queue::") {
(FnQueueNames::new(queue_name).dlq(), true)
} else if is_function_queue {
(FnQueueNames::new(topic).dlq(), true)
} else {
(RabbitNames::new(topic).dlq(), false)
}
}

async fn resolve_dlq_name_for_topic(&self, topic: &str) -> (String, bool) {
let is_function_queue = self.function_queue_configs.read().await.contains_key(topic);
Self::resolve_dlq_name(topic, is_function_queue)
}

async fn operation_channel(&self, operation: &str) -> anyhow::Result<Channel> {
self.connection
.create_channel()
.await
.map_err(|e| anyhow::anyhow!("Failed to create channel for {}: {}", operation, e))
}

async fn passive_queue_count(&self, queue_name: &str) -> u64 {
let Ok(channel) = self.operation_channel("queue inspection").await else {
return 0;
};
match channel
.queue_declare(
queue_name,
lapin::options::QueueDeclareOptions {
passive: true,
..Default::default()
},
lapin::types::FieldTable::default(),
)
.await
{
Ok(info) => info.message_count() as u64,
Err(_) => 0,
}
}

/// Scan a DLQ for a specific message by ID, applying `on_found` when the
/// target is located. Non-target messages are nacked back to the queue.
///
Expand All @@ -143,10 +177,10 @@ impl RabbitMQAdapter {
F: FnOnce(&Delivery, bool, &str) -> Fut,
Fut: std::future::Future<Output = anyhow::Result<()>>,
{
let (dlq_name, is_fn_queue) = Self::resolve_dlq_name(topic);
let (dlq_name, is_fn_queue) = self.resolve_dlq_name_for_topic(topic).await;
let channel = self.operation_channel("DLQ message lookup").await?;

let queue_info = self
.channel
let queue_info = channel
.queue_declare(
&dlq_name,
lapin::options::QueueDeclareOptions {
Expand All @@ -161,8 +195,7 @@ impl RabbitMQAdapter {
let count = queue_info.message_count();

for _ in 0..count {
let get_result = self
.channel
let get_result = channel
.basic_get(&dlq_name, BasicGetOptions { no_ack: false })
.await
.map_err(|e| anyhow::anyhow!("Failed to get message from DLQ: {}", e))?;
Expand Down Expand Up @@ -489,12 +522,12 @@ impl QueueAdapter for RabbitMQAdapter {
}

async fn redrive_dlq(&self, topic: &str) -> anyhow::Result<u64> {
let (dlq_name, is_fn_queue) = Self::resolve_dlq_name(topic);
let (dlq_name, is_fn_queue) = self.resolve_dlq_name_for_topic(topic).await;
let channel = self.operation_channel("DLQ redrive").await?;
let mut count: u64 = 0;

loop {
let get_result = self
.channel
let get_result = channel
.basic_get(&dlq_name, BasicGetOptions { no_ack: false })
.await
.map_err(|e| anyhow::anyhow!("Failed to get message from DLQ: {}", e))?;
Expand All @@ -503,7 +536,7 @@ impl QueueAdapter for RabbitMQAdapter {
Some(delivery) => {
let republish_result: anyhow::Result<()> = if is_fn_queue {
// Function queue DLQ: raw data payload, republish directly
let queue_name = topic.strip_prefix("__fn_queue::").unwrap();
let queue_name = topic.strip_prefix("__fn_queue::").unwrap_or(topic);
let names = FnQueueNames::new(queue_name);

let mut headers = delivery
Expand All @@ -526,7 +559,7 @@ impl QueueAdapter for RabbitMQAdapter {
properties
};

self.channel
channel
.basic_publish(
&names.exchange(),
queue_name,
Expand Down Expand Up @@ -674,10 +707,10 @@ impl QueueAdapter for RabbitMQAdapter {
async fn dlq_count(&self, topic: &str) -> anyhow::Result<u64> {
// Function queues use FnQueueNames (e.g., __fn_queue::orders -> ::dlq.queue),
// while topic-based queues use RabbitNames (e.g., user.created -> .dlq).
let (dlq_name, _is_fn_queue) = Self::resolve_dlq_name(topic);
let (dlq_name, _is_fn_queue) = self.resolve_dlq_name_for_topic(topic).await;
let channel = self.operation_channel("DLQ count").await?;

let queue = self
.channel
let queue = channel
.queue_declare(
&dlq_name,
lapin::options::QueueDeclareOptions {
Expand All @@ -697,10 +730,10 @@ impl QueueAdapter for RabbitMQAdapter {
/// the same fields the engine's `DlqMessage` struct carries as a JSON
/// object per entry instead of a dedicated type.
async fn dlq_peek(&self, topic: &str, offset: u64, limit: u64) -> anyhow::Result<Vec<Value>> {
let (dlq_name, is_fn_queue) = Self::resolve_dlq_name(topic);
let (dlq_name, is_fn_queue) = self.resolve_dlq_name_for_topic(topic).await;
let channel = self.operation_channel("DLQ browse").await?;

let queue_depth = match self
.channel
let queue_depth = match channel
.queue_declare(
&dlq_name,
lapin::options::QueueDeclareOptions {
Expand All @@ -720,8 +753,7 @@ impl QueueAdapter for RabbitMQAdapter {
let mut deliveries_to_nack = Vec::new();

for i in 0..fetch_count {
let get_result = self
.channel
let get_result = channel
.basic_get(&dlq_name, BasicGetOptions { no_ack: false })
.await
.map_err(|e| anyhow::anyhow!("Failed to get message from DLQ: {}", e))?;
Expand Down Expand Up @@ -1322,53 +1354,22 @@ impl QueueAdapter for RabbitMQAdapter {
/// dropped (no field for it in this worker's `TopicStats`);
/// `delivered`/`failed` are always `0` (not tracked by this adapter).
///
/// Inherits one engine quirk verbatim: the "main queue depth" is looked
/// up via `RabbitNames::new(topic).queue()` (`iii.{topic}.queue`), a
/// queue name that topic-fanout subscriptions never actually declare
/// (`subscribe` declares per-function queues via
/// `RabbitNames::function_queue`, never the bare `.queue()` name) -- so
/// for topic-based (non-`__fn_queue::`) topics this always resolves to
/// depth `0` in both the engine and this port.
async fn topic_stats(&self, topic: &str) -> anyhow::Result<TopicStats> {
let is_configured_function_queue =
self.function_queue_configs.read().await.contains_key(topic);
let (queue_name, dlq_name) = if let Some(name) = topic.strip_prefix("__fn_queue::") {
let names = FnQueueNames::new(name);
(names.queue(), names.dlq())
} else if is_configured_function_queue {
let names = FnQueueNames::new(topic);
(names.queue(), names.dlq())
} else {
let names = RabbitNames::new(topic);
(names.queue(), names.dlq())
};

let depth = match self
.channel
.queue_declare(
&queue_name,
lapin::options::QueueDeclareOptions {
passive: true,
..Default::default()
},
lapin::types::FieldTable::default(),
)
.await
{
Ok(info) => info.message_count() as u64,
Err(_) => 0,
};

let dlq_depth = match self
.channel
.queue_declare(
&dlq_name,
lapin::options::QueueDeclareOptions {
passive: true,
..Default::default()
},
lapin::types::FieldTable::default(),
)
.await
{
Ok(info) => info.message_count() as u64,
Err(_) => 0,
};
let depth = self.passive_queue_count(&queue_name).await;
let dlq_depth = self.passive_queue_count(&dlq_name).await;

Ok(TopicStats {
depth,
Expand Down Expand Up @@ -1498,3 +1499,16 @@ async fn requeue_rabbitmq_deliveries(
anyhow::bail!(errors.join("; "))
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn bare_function_queue_name_resolves_to_function_dlq() {
assert_eq!(
RabbitMQAdapter::resolve_dlq_name("tax-returns", true),
("iii.__fn_queue::tax-returns::dlq.queue".to_string(), true,)
);
}
}
53 changes: 50 additions & 3 deletions queue/src/functions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -120,7 +120,7 @@ fn default_dlq_limit() -> u64 {
50
}

#[derive(Debug, Clone, Serialize, JsonSchema, PartialEq)]
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq)]
pub struct DlqMessage {
pub id: String,
pub payload: Value,
Expand All @@ -130,6 +130,13 @@ pub struct DlqMessage {
pub size_bytes: u64,
}

#[derive(Debug, Deserialize)]
#[serde(untagged)]
enum AdapterDlqMessage {
Normalized(DlqMessage),
Job(Job),
}

pub fn register_all(
iii: &Arc<IIIClient>,
adapter: Arc<SwappableAdapter>,
Expand Down Expand Up @@ -378,8 +385,12 @@ pub async fn dlq_messages(
.into_iter()
.skip(input.offset as usize)
.take(input.limit as usize)
.filter_map(|value| serde_json::from_value::<Job>(value).ok())
.map(dlq_message_from_job)
.filter_map(
|value| match serde_json::from_value::<AdapterDlqMessage>(value).ok()? {
AdapterDlqMessage::Normalized(message) => Some(message),
AdapterDlqMessage::Job(job) => Some(dlq_message_from_job(job)),
},
)
.collect();
Ok(messages)
}
Expand Down Expand Up @@ -783,6 +794,42 @@ mod tests {
);
}

#[tokio::test]
async fn dlq_browse_preserves_normalized_messages() {
let (adapter, mock) = adapter();
*mock.dlq_messages_result.lock().unwrap() = Some(Ok(vec![json!({
"id": "m1",
"payload": {"dead": true},
"error": "Function tax::process exhausted retries",
"failed_at": 123,
"retries": 1,
"size_bytes": 13
})]));

let messages = dlq_messages(
adapter,
DlqMessagesInput {
topic: "__fn_queue::tax-returns".to_string(),
offset: 0,
limit: 50,
},
)
.await
.unwrap();

assert_eq!(
messages,
vec![DlqMessage {
id: "m1".to_string(),
payload: json!({"dead": true}),
error: "Function tax::process exhausted retries".to_string(),
failed_at: 123,
retries: 1,
size_bytes: 13,
}]
);
}

#[tokio::test]
async fn dlq_messages_rejects_empty_topic() {
let (adapter, _mock) = adapter();
Expand Down
Loading
Loading