use crate::common::interface::storage::{BlobStorage, Offloadable};
use crate::common::model::config::Config;
use crate::common::model::message::{TaskErrorEvent, TaskEvent, TaskParserEvent};
use crate::common::model::{Prioritizable, Priority, Request, Response};
use crate::common::policy::PolicyResolver;
use crate::errors::ErrorKind;
use crate::queue::batcher::Batcher;
use crate::queue::channel::Channel;
use crate::queue::compensation::{Compensator, Identifiable};
use crate::queue::compression::{compress_payload_owned, decompress_payload};
#[cfg(feature = "queue-kafka")]
use crate::queue::kafka::KafkaQueue;
#[cfg(feature = "queue-nats")]
use crate::queue::nats::NatsQueue;
use crate::queue::{HEADER_ATTEMPT, HEADER_CREATED_AT, MqBackend, NackPolicy, QueuedItem};
use crate::utils::logger::LogModel;
use crate::utils::storage::FileSystemBlobStorage;
use futures::StreamExt;
use futures::future::join_all;
use log::{error, info};
use metrics::counter;
use once_cell::sync::OnceCell;
use rmp_serde as rmps;
use serde_path_to_error;
use std::collections::HashMap;
use std::sync::Arc;
use tokio::sync::mpsc::{Receiver, Sender};
use tokio::sync::{Mutex, Semaphore};
const DEFAULT_COMPRESSION_THRESHOLD: usize = 1024;
const BLOCKING_PAYLOAD_BYTES: usize = 64 * 1024;
fn default_headers() -> HashMap<String, String> {
let mut headers = HashMap::new();
headers.insert(HEADER_ATTEMPT.to_string(), "0".to_string());
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis()
.to_string();
headers.insert(HEADER_CREATED_AT.to_string(), now_ms);
headers
}
fn msgpack_encode<T: serde::Serialize>(value: &T) -> Result<Vec<u8>, rmps::encode::Error> {
rmps::to_vec(value)
}
fn msgpack_decode<T: serde::de::DeserializeOwned>(
bytes: &[u8],
) -> Result<T, serde_path_to_error::Error<rmps::decode::Error>> {
let mut deserializer = rmps::Deserializer::new(bytes);
serde_path_to_error::deserialize(&mut deserializer)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum QueueCodec {
Json,
Msgpack,
}
static QUEUE_CODEC_OVERRIDE: OnceCell<QueueCodec> = OnceCell::new();
fn queue_codec() -> QueueCodec {
if let Some(codec) = QUEUE_CODEC_OVERRIDE.get() {
return *codec;
}
QueueCodec::Msgpack
}
fn set_queue_codec_from_config(cfg: &Config) {
let codec = cfg
.channel_config
.queue_codec
.as_deref()
.unwrap_or("msgpack")
.to_lowercase();
let mapped = match codec.as_str() {
"json" => QueueCodec::Json,
"msgpack" | "rmp" => QueueCodec::Msgpack,
_ => QueueCodec::Msgpack,
};
let _ = QUEUE_CODEC_OVERRIDE.set(mapped);
}
pub struct QueueManager {
pub channel: Arc<Channel>,
pub backend: Option<Arc<dyn MqBackend>>,
pub compensator: Option<Arc<dyn Compensator>>,
pub blob_storage: Option<Arc<dyn BlobStorage>>,
pub batch_concurrency: usize,
pub compression_threshold: usize,
pub nack_policy: NackPolicy,
pub log_topic: String,
}
impl QueueManager {
pub fn new(backend: Option<Arc<dyn MqBackend>>, capacity: usize) -> Self {
QueueManager {
channel: Arc::new(Channel::new(capacity)),
backend,
compensator: None,
blob_storage: None,
batch_concurrency: 50,
compression_threshold: DEFAULT_COMPRESSION_THRESHOLD,
nack_policy: NackPolicy::default(),
log_topic: "log".to_string(),
}
}
pub fn from_config(cfg: &Config) -> Arc<Self> {
Self::from_config_with_log_topic(cfg, None)
}
pub fn from_config_with_log_topic(cfg: &Config, log_topic: Option<&str>) -> Arc<Self> {
set_queue_codec_from_config(cfg);
let channel_config = &cfg.channel_config;
#[allow(unused_variables)] let namespace = &cfg.name;
let nack_policy = if let Some(policy_cfg) = &cfg.policy
&& !policy_cfg.overrides.is_empty()
{
let resolver = PolicyResolver::new(Some(policy_cfg));
let decision =
resolver.resolve_with_kind("queue", Some("nack"), Some("failed"), ErrorKind::Queue);
NackPolicy {
max_retries: decision.policy.max_retries,
backoff_ms: decision.policy.backoff_ms,
}
} else {
NackPolicy {
max_retries: channel_config.nack_max_retries.unwrap_or(0),
backoff_ms: channel_config.nack_backoff_ms.unwrap_or(0),
}
};
let mut queue_manager = if let Some(kafka_config) = &channel_config.kafka {
#[cfg(feature = "queue-kafka")]
let qm = match KafkaQueue::new(
kafka_config,
channel_config.minid_time,
namespace,
nack_policy,
) {
Ok(kafka_queue) => {
info!("KafkaQueue initialized successfully");
QueueManager::new(Some(Arc::new(kafka_queue)), channel_config.capacity)
}
Err(e) => {
error!("KafkaQueue init failed, fallback to in-memory queue: {}", e);
QueueManager::new(None, channel_config.capacity)
}
};
#[cfg(not(feature = "queue-kafka"))]
let qm = {
let _ = kafka_config;
error!(
"channel_config.kafka is set but the `queue-kafka` feature is disabled; \
falling back to in-memory queue"
);
QueueManager::new(None, channel_config.capacity)
};
qm
} else if let Some(nats_config) = &channel_config.nats {
#[cfg(feature = "queue-nats")]
let qm = match NatsQueue::new(
nats_config,
channel_config.minid_time,
namespace,
nack_policy,
) {
Ok(nats_queue) => {
info!("NatsQueue (JetStream) initialized successfully");
QueueManager::new(Some(Arc::new(nats_queue)), channel_config.capacity)
}
Err(e) => {
error!("NatsQueue init failed, fallback to in-memory queue: {}", e);
QueueManager::new(None, channel_config.capacity)
}
};
#[cfg(not(feature = "queue-nats"))]
let qm = {
let _ = nats_config;
error!(
"channel_config.nats is set but the `queue-nats` feature is disabled; \
falling back to in-memory queue"
);
QueueManager::new(None, channel_config.capacity)
};
qm
} else {
info!("In-Memory Queue initialized (Single Node Mode)");
QueueManager::new(None, 10000)
};
if let Some(topic) = log_topic {
queue_manager.log_topic = topic.to_string();
}
queue_manager.nack_policy = nack_policy;
if let Some(concurrency) = channel_config.batch_concurrency {
queue_manager.with_concurrency(concurrency);
}
if let Some(threshold) = channel_config.compression_threshold {
queue_manager.with_compression_threshold(threshold);
}
if let Some(blob_config) = &channel_config.blob_storage {
if let Some(path) = &blob_config.path {
let storage = Arc::new(FileSystemBlobStorage::new(path));
queue_manager.with_blob_storage(storage);
info!("BlobStorage initialized at: {}", path);
}
}
queue_manager.subscribe();
Arc::new(queue_manager)
}
pub fn with_backend(&mut self, backend: Arc<dyn MqBackend>) {
self.backend = Some(backend);
}
pub fn with_compensator(&mut self, compensator: Arc<dyn Compensator>) {
self.compensator = Some(compensator);
}
pub fn with_blob_storage(&mut self, storage: Arc<dyn BlobStorage>) {
self.blob_storage = Some(storage);
}
pub fn with_concurrency(&mut self, concurrency: usize) {
self.batch_concurrency = concurrency;
}
pub fn with_compression_threshold(&mut self, threshold: usize) {
self.compression_threshold = threshold;
}
pub fn with_log_topic(&mut self, topic: impl Into<String>) {
self.log_topic = topic.into();
}
pub fn subscribe(&self) {
if let Some(backend) = &self.backend {
let backend = backend.clone();
let channel = self.channel.clone();
let compensator = self.compensator.clone();
let blob_storage = self.blob_storage.clone();
let concurrency = self.batch_concurrency;
let compression_threshold = self.compression_threshold;
self.spawn_forwarder(
"task",
channel.remote_task_receiver.clone(),
backend.clone(),
blob_storage.clone(),
concurrency,
compression_threshold,
);
self.spawn_forwarder(
"request",
channel.request_receiver.clone(),
backend.clone(),
blob_storage.clone(),
concurrency,
compression_threshold,
);
self.spawn_forwarder(
"response",
channel.response_receiver.clone(),
backend.clone(),
blob_storage.clone(),
concurrency,
compression_threshold,
);
self.spawn_forwarder(
"parser_task",
channel.parser_task_receiver.clone(),
backend.clone(),
blob_storage.clone(),
concurrency,
compression_threshold,
);
self.spawn_forwarder(
"error_task",
channel.error_receiver.clone(),
backend.clone(),
blob_storage.clone(),
concurrency,
compression_threshold,
);
self.spawn_forwarder(
self.log_topic.as_str(),
channel.log_receiver.clone(),
backend.clone(),
blob_storage.clone(),
concurrency,
compression_threshold,
);
tokio::spawn(async move {
Self::subscribe_all_priorities(
"task",
channel.task_sender.clone(),
backend.clone(),
compensator.clone(),
blob_storage.clone(),
concurrency,
)
.await;
Self::subscribe_all_priorities(
"request",
channel.download_request_sender.clone(),
backend.clone(),
compensator.clone(),
blob_storage.clone(),
concurrency,
)
.await;
Self::subscribe_all_priorities(
"response",
channel.remote_response_sender.clone(),
backend.clone(),
compensator.clone(),
blob_storage.clone(),
concurrency,
)
.await;
Self::subscribe_all_priorities(
"parser_task",
channel.remote_parser_task_sender.clone(),
backend.clone(),
compensator.clone(),
blob_storage.clone(),
concurrency,
)
.await;
Self::subscribe_all_priorities(
"error_task",
channel.remote_error_sender.clone(),
backend.clone(),
compensator.clone(),
blob_storage.clone(),
concurrency,
)
.await;
});
}
}
async fn subscribe_all_priorities<T>(
topic_base: &str,
sender: Sender<QueuedItem<T>>,
backend: Arc<dyn MqBackend>,
compensator: Option<Arc<dyn Compensator>>,
blob_storage: Option<Arc<dyn BlobStorage>>,
concurrency: usize,
) where
T: serde::de::DeserializeOwned
+ Send
+ 'static
+ std::fmt::Debug
+ Identifiable
+ Offloadable,
{
let priorities = [Priority::High, Priority::Normal, Priority::Low];
futures::stream::iter(priorities)
.for_each_concurrent(None, |priority| {
let sender = sender.clone();
let backend = backend.clone();
let compensator = compensator.clone();
let blob_storage = blob_storage.clone();
async move {
let topic = format!("{}-{}", topic_base, priority.suffix());
Self::subscribe_topic(
&topic,
sender,
backend,
compensator,
blob_storage,
concurrency,
)
.await;
}
})
.await;
}
async fn subscribe_topic<T>(
topic: &str,
sender: Sender<QueuedItem<T>>,
backend: Arc<dyn MqBackend>,
compensator: Option<Arc<dyn Compensator>>,
blob_storage: Option<Arc<dyn BlobStorage>>,
concurrency: usize,
) where
T: serde::de::DeserializeOwned
+ Send
+ 'static
+ std::fmt::Debug
+ Identifiable
+ Offloadable,
{
let (tx, mut rx) = tokio::sync::mpsc::channel(1024);
if let Err(e) = backend.subscribe(topic, tx).await {
error!("Failed to subscribe to topic {}: {}", topic, e);
return;
}
let topic = topic.to_string();
let batch_concurrency = (concurrency / 50).max(1);
let semaphore = Arc::new(Semaphore::new(batch_concurrency));
tokio::spawn(async move {
let topic_clone = topic.clone();
let sender_clone = sender.clone();
let compensator_clone = compensator.clone();
let blob_storage_clone = blob_storage.clone();
Batcher::run(&mut rx, 50, 5, semaphore, move |items| {
let topic = topic_clone.clone();
let sender = sender_clone.clone();
let compensator = compensator_clone.clone();
let blob_storage = blob_storage_clone.clone();
async move {
Self::process_batch_messages(
items,
&topic,
&sender,
&compensator,
&blob_storage,
)
.await;
}
})
.await;
log::warn!("Topic {} subscription closed", topic);
});
}
async fn process_batch_messages<T>(
messages: Vec<crate::queue::Message>,
topic: &str,
sender: &Sender<QueuedItem<T>>,
compensator: &Option<Arc<dyn Compensator>>,
blob_storage: &Option<Arc<dyn BlobStorage>>,
) where
T: serde::de::DeserializeOwned
+ Send
+ 'static
+ std::fmt::Debug
+ Identifiable
+ Offloadable,
{
let mut use_blocking = messages.len() >= 32;
if !use_blocking {
let mut total_bytes = 0usize;
for msg in &messages {
total_bytes += msg.payload.len();
if total_bytes >= BLOCKING_PAYLOAD_BYTES {
use_blocking = true;
break;
}
}
}
let results = if use_blocking {
tokio::task::spawn_blocking(move || {
messages
.into_iter()
.map(|msg| {
let payload_slice = msg.payload.as_slice();
let decoded_payload = decompress_payload(payload_slice);
let item_res = match queue_codec() {
QueueCodec::Json => {
serde_json::from_slice::<T>(decoded_payload.as_ref()).map_err(|e| {
crate::errors::Error::new(
crate::errors::ErrorKind::Queue,
Some(e),
)
})
}
QueueCodec::Msgpack => msgpack_decode::<T>(decoded_payload.as_ref())
.map_err(|e| {
crate::errors::Error::new(
crate::errors::ErrorKind::Queue,
Some(e),
)
}),
};
(msg, item_res)
})
.collect::<Vec<_>>()
})
.await
} else {
let processed_items = messages
.into_iter()
.map(|msg| {
let payload_slice = msg.payload.as_slice();
let decoded_payload = decompress_payload(payload_slice);
let item_res = match queue_codec() {
QueueCodec::Json => serde_json::from_slice::<T>(decoded_payload.as_ref())
.map_err(|e| {
crate::errors::Error::new(crate::errors::ErrorKind::Queue, Some(e))
}),
QueueCodec::Msgpack => msgpack_decode::<T>(decoded_payload.as_ref())
.map_err(|e| {
crate::errors::Error::new(crate::errors::ErrorKind::Queue, Some(e))
}),
};
(msg, item_res)
})
.collect::<Vec<_>>();
Ok(processed_items)
};
match results {
Ok(processed_items) => {
let tasks = processed_items.into_iter().map(|(msg, result)| {
let topic = topic.to_string();
let compensator = compensator.clone();
let storage = blob_storage.clone();
async move {
match result {
Ok(mut item) => {
if let Some(storage) = storage {
if let Err(e) = item.reload(&storage).await {
error!("Failed to reload item content from storage for topic {}: {}", topic, e);
if let Err(e) = msg.nack("Blob reload failed").await {
error!("Failed to NACK message: {}", e);
}
return None;
}
}
let id = item.get_id();
if let Some(comp) = compensator {
if let Err(e) = comp.add_task(&topic, &id, msg.payload.clone()).await {
error!(
"Failed to add task to compensation queue for topic {}: {}",
topic, e
);
}
}
let msg_ack = msg.clone();
let msg_nack = msg.clone();
let queued_item = QueuedItem::with_ack(
item,
move || Box::pin(async move { msg_ack.ack().await }),
move |reason| Box::pin(async move { msg_nack.nack(reason).await })
);
Some((msg, queued_item))
}
Err(e) => {
let payload_len = msg.payload.len();
let codec = match queue_codec() {
QueueCodec::Json => "json",
QueueCodec::Msgpack => "msgpack",
};
error!(
"Failed to deserialize message from topic {} (codec={}, bytes={}): {}",
topic,
codec,
payload_len,
e
);
if let Err(e) = msg.nack("Deserialization failed").await {
error!("Failed to NACK poison message: {}", e);
}
None
}
}
}
});
let results = join_all(tasks).await;
for (msg, item) in results.into_iter().flatten() {
if let Err(_e) = sender.send(item).await {
let _ = msg.nack("Channel closed").await;
}
}
}
Err(e) => {
error!("Batch deserialization task failed: {}", e);
}
}
}
fn spawn_forwarder<T>(
&self,
topic: &str,
receiver: Arc<Mutex<Receiver<QueuedItem<T>>>>,
backend: Arc<dyn MqBackend>,
blob_storage: Option<Arc<dyn BlobStorage>>,
concurrency: usize,
compression_threshold: usize,
) where
T: serde::Serialize + Send + Sync + 'static + Identifiable + Prioritizable + Offloadable,
{
Self::forward_channel(
topic,
receiver,
backend,
blob_storage,
concurrency,
compression_threshold,
)
}
fn forward_channel<T>(
topic: &str,
receiver: Arc<Mutex<Receiver<QueuedItem<T>>>>,
backend: Arc<dyn MqBackend>,
blob_storage: Option<Arc<dyn BlobStorage>>,
concurrency: usize,
compression_threshold: usize,
) where
T: serde::Serialize + Send + Sync + 'static + Identifiable + Prioritizable + Offloadable,
{
let topic = topic.to_string();
let semaphore = Arc::new(Semaphore::new(concurrency));
tokio::spawn(async move {
let mut rx = receiver.lock().await;
let topic_clone = topic.clone();
let backend_clone = backend.clone();
let blob_storage = blob_storage.clone();
let compression_threshold = compression_threshold;
Batcher::run(&mut *rx, 500, 5, semaphore, move |items| {
let topic = topic_clone.clone();
let backend = backend_clone.clone();
let blob_storage = blob_storage.clone();
async move {
Self::flush_batch_grouped(
topic,
backend,
blob_storage,
items,
compression_threshold,
)
.await;
}
})
.await;
});
}
pub fn get_task_push_channel(&self) -> Sender<QueuedItem<TaskEvent>> {
if self.backend.is_none() {
return self.channel.task_sender.clone();
}
self.channel.remote_task_sender.clone()
}
pub fn get_task_pop_channel(&self) -> Arc<Mutex<Receiver<QueuedItem<TaskEvent>>>> {
Arc::clone(&self.channel.task_receiver)
}
pub fn get_request_pop_channel(&self) -> Arc<Mutex<Receiver<QueuedItem<Request>>>> {
self.channel.download_request_receiver.clone()
}
pub fn get_request_push_channel(&self) -> Sender<QueuedItem<Request>> {
if self.backend.is_none() {
return self.channel.download_request_sender.clone();
}
self.channel.request_sender.clone()
}
pub fn get_response_push_channel(&self) -> Sender<QueuedItem<Response>> {
if self.backend.is_none() {
return self.channel.remote_response_sender.clone();
}
self.channel.response_sender.clone()
}
pub fn try_send_local_response(
&self,
item: QueuedItem<Response>,
) -> Result<(), tokio::sync::mpsc::error::TrySendError<QueuedItem<Response>>> {
self.channel.remote_response_sender.try_send(item)
}
pub fn get_response_pop_channel(&self) -> Arc<Mutex<Receiver<QueuedItem<Response>>>> {
Arc::clone(&self.channel.remote_response_receiver)
}
pub fn get_parser_task_pop_channel(&self) -> Arc<Mutex<Receiver<QueuedItem<TaskParserEvent>>>> {
self.channel.remote_parser_task_receiver.clone()
}
pub fn get_parser_task_push_channel(&self) -> Sender<QueuedItem<TaskParserEvent>> {
if self.backend.is_none() {
return self.channel.remote_parser_task_sender.clone();
}
self.channel.parser_task_sender.clone()
}
pub fn get_error_pop_channel(&self) -> Arc<Mutex<Receiver<QueuedItem<TaskErrorEvent>>>> {
Arc::clone(&self.channel.remote_error_receiver)
}
pub fn get_error_push_channel(&self) -> Sender<QueuedItem<TaskErrorEvent>> {
if self.backend.is_none() {
return self.channel.remote_error_sender.clone();
}
self.channel.error_sender.clone()
}
pub fn get_log_push_channel(&self) -> Sender<QueuedItem<LogModel>> {
self.channel.log_sender.clone()
}
pub async fn local_pending_count(&self) -> usize {
let (task, download, response, parser, error, remote_task) =
self.local_pending_breakdown().await;
task + download + response + parser + error + remote_task
}
pub async fn local_pending_breakdown(&self) -> (usize, usize, usize, usize, usize, usize) {
let task = self
.channel
.task_sender
.max_capacity()
.saturating_sub(self.channel.task_sender.capacity());
let download = self
.channel
.download_request_sender
.max_capacity()
.saturating_sub(self.channel.download_request_sender.capacity());
let response = self
.channel
.remote_response_sender
.max_capacity()
.saturating_sub(self.channel.remote_response_sender.capacity());
let parser = self
.channel
.remote_parser_task_sender
.max_capacity()
.saturating_sub(self.channel.remote_parser_task_sender.capacity());
let error = self
.channel
.remote_error_sender
.max_capacity()
.saturating_sub(self.channel.remote_error_sender.capacity());
let remote_task = self
.channel
.remote_task_sender
.max_capacity()
.saturating_sub(self.channel.remote_task_sender.capacity());
(task, download, response, parser, error, remote_task)
}
async fn flush_batch_grouped<T>(
base_topic: String,
backend: Arc<dyn MqBackend>,
blob_storage: Option<Arc<dyn BlobStorage>>,
mut items: Vec<QueuedItem<T>>,
compression_threshold: usize,
) where
T: serde::Serialize + Identifiable + Send + Sync + Prioritizable + Offloadable + 'static,
{
if let Some(storage) = &blob_storage {
for item in &mut items {
if item.inner.should_offload(BLOCKING_PAYLOAD_BYTES) {
if let Err(e) = item.inner.offload(storage).await {
error!("Failed to offload item payload to blob storage: {}", e);
}
}
}
}
if items.is_empty() {
return;
}
let first_priority = items[0].get_priority();
let all_same_priority = items.iter().all(|i| i.get_priority() == first_priority);
if all_same_priority {
let topic = format!("{}-{}", base_topic, first_priority.suffix());
Self::flush_batch(topic, backend, items, compression_threshold).await;
return;
}
let mut groups: HashMap<Priority, Vec<QueuedItem<T>>> = HashMap::new();
for item in items {
groups.entry(item.get_priority()).or_default().push(item);
}
for (priority, group_items) in groups {
let topic = format!("{}-{}", base_topic, priority.suffix());
Self::flush_batch(topic, backend.clone(), group_items, compression_threshold).await;
}
}
async fn flush_batch<T>(
topic: String,
backend: Arc<dyn MqBackend>,
items: Vec<QueuedItem<T>>,
compression_threshold: usize,
) where
T: serde::Serialize + Identifiable + Send + Sync + 'static,
{
let use_blocking = items.len() >= 32;
let payloads_result = if use_blocking {
tokio::task::spawn_blocking(move || Self::encode_items(&items, compression_threshold))
.await
} else {
Ok(Self::encode_items(&items, compression_threshold))
};
match payloads_result {
Ok((payloads, ids, count)) => {
if payloads.is_empty() {
return;
}
if let Err(e) = backend.publish_batch_with_headers(&topic, &payloads).await {
error!("Failed to publish batch to topic {}: {}", topic, e);
} else {
log::info!(
"[QueueManager] forward_channel published batch: topic={} count={}",
topic,
count
);
if let Some(id_list) = ids {
for id in id_list {
log::debug!(
"[QueueManager] forward_channel published: topic={} id={}",
topic,
id
);
}
}
}
}
Err(e) => {
error!("Serialization task join error: {}", e);
}
}
}
pub async fn clean_storage(&self) -> crate::errors::Result<()> {
if let Some(backend) = &self.backend {
backend.clean_storage().await?;
}
Ok(())
}
fn encode_items<T>(
items: &[QueuedItem<T>],
compression_threshold: usize,
) -> (
Vec<(Option<String>, Vec<u8>, HashMap<String, String>)>,
Option<Vec<String>>,
usize,
)
where
T: serde::Serialize + Identifiable,
{
let mut payloads: Vec<(Option<String>, Vec<u8>, HashMap<String, String>)> =
Vec::with_capacity(items.len());
let mut ids = if log::log_enabled!(log::Level::Debug) {
Some(Vec::with_capacity(items.len()))
} else {
None
};
for item in items {
let id = item.get_id();
let partition_key = item.partition_key();
if let Some(id_list) = ids.as_mut() {
id_list.push(id.clone());
}
let encoded = match queue_codec() {
QueueCodec::Json => serde_json::to_vec(&item.inner).map_err(|e| {
crate::errors::Error::new(crate::errors::ErrorKind::Queue, Some(e))
}),
QueueCodec::Msgpack => msgpack_encode(&item.inner).map_err(|e| {
crate::errors::Error::new(crate::errors::ErrorKind::Queue, Some(e))
}),
};
match encoded {
Ok(p) => {
let final_payload = compress_payload_owned(p, compression_threshold);
let headers = default_headers();
payloads.push((Some(partition_key), final_payload, headers));
}
Err(e) => {
error!(
"Failed to serialize item id={}: {}. Item will be skipped (unrecoverable).",
id, e
);
counter!("mocra_queue_encode_errors_total").increment(1);
}
}
}
(payloads, ids, items.len())
}
pub async fn send_to_dlq<T>(
&self,
topic: &str,
item: &T,
reason: &str,
) -> crate::errors::Result<()>
where
T: serde::Serialize + Identifiable + Send + Sync,
{
if let Some(backend) = &self.backend {
let payload = match queue_codec() {
QueueCodec::Json => serde_json::to_vec(item)
.map_err(|e| crate::errors::error::QueueError::OperationFailed(Box::new(e)))?,
QueueCodec::Msgpack => msgpack_encode(item)
.map_err(|e| crate::errors::error::QueueError::OperationFailed(Box::new(e)))?,
};
backend
.send_to_dlq(topic, &item.get_id(), &payload, reason)
.await?;
}
Ok(())
}
pub async fn read_dlq(
&self,
topic: &str,
count: usize,
) -> crate::errors::Result<Vec<(String, Vec<u8>, String, String)>> {
if let Some(backend) = &self.backend {
backend.read_dlq(topic, count).await
} else {
Ok(Vec::new())
}
}
}