Skip to main content

mocra_core/queue/
manager.rs

1use crate::common::interface::storage::{BlobStorage, Offloadable};
2use crate::common::model::config::Config;
3use crate::common::model::message::{TaskErrorEvent, TaskEvent, TaskParserEvent};
4use crate::common::model::{Prioritizable, Priority, Request, Response};
5use crate::common::policy::PolicyResolver;
6use crate::errors::ErrorKind;
7use crate::queue::batcher::Batcher;
8use crate::queue::channel::Channel;
9use crate::queue::compensation::{Compensator, Identifiable};
10use crate::queue::compression::{compress_payload_owned, decompress_payload};
11#[cfg(feature = "queue-kafka")]
12use crate::queue::kafka::KafkaQueue;
13#[cfg(feature = "queue-nats")]
14use crate::queue::nats::NatsQueue;
15use crate::queue::{HEADER_ATTEMPT, HEADER_CREATED_AT, MqBackend, NackPolicy, QueuedItem};
16use crate::utils::logger::LogModel;
17use crate::utils::storage::FileSystemBlobStorage;
18use futures::StreamExt;
19use futures::future::join_all;
20use log::{error, info};
21use metrics::counter;
22use once_cell::sync::OnceCell;
23use rmp_serde as rmps;
24use serde_path_to_error;
25use std::collections::HashMap;
26use std::sync::Arc;
27use tokio::sync::mpsc::{Receiver, Sender};
28use tokio::sync::{Mutex, Semaphore};
29
30const DEFAULT_COMPRESSION_THRESHOLD: usize = 1024;
31const BLOCKING_PAYLOAD_BYTES: usize = 64 * 1024;
32
33fn default_headers() -> HashMap<String, String> {
34    let mut headers = HashMap::new();
35    headers.insert(HEADER_ATTEMPT.to_string(), "0".to_string());
36    let now_ms = std::time::SystemTime::now()
37        .duration_since(std::time::UNIX_EPOCH)
38        .unwrap_or_default()
39        .as_millis()
40        .to_string();
41    headers.insert(HEADER_CREATED_AT.to_string(), now_ms);
42    headers
43}
44
45fn msgpack_encode<T: serde::Serialize>(value: &T) -> Result<Vec<u8>, rmps::encode::Error> {
46    rmps::to_vec(value)
47}
48
49fn msgpack_decode<T: serde::de::DeserializeOwned>(
50    bytes: &[u8],
51) -> Result<T, serde_path_to_error::Error<rmps::decode::Error>> {
52    let mut deserializer = rmps::Deserializer::new(bytes);
53    serde_path_to_error::deserialize(&mut deserializer)
54}
55
56#[derive(Debug, Clone, Copy, PartialEq, Eq)]
57enum QueueCodec {
58    Json,
59    Msgpack,
60}
61
62static QUEUE_CODEC_OVERRIDE: OnceCell<QueueCodec> = OnceCell::new();
63
64fn queue_codec() -> QueueCodec {
65    if let Some(codec) = QUEUE_CODEC_OVERRIDE.get() {
66        return *codec;
67    }
68    QueueCodec::Msgpack
69}
70
71fn set_queue_codec_from_config(cfg: &Config) {
72    let codec = cfg
73        .channel_config
74        .queue_codec
75        .as_deref()
76        .unwrap_or("msgpack")
77        .to_lowercase();
78
79    let mapped = match codec.as_str() {
80        "json" => QueueCodec::Json,
81        "msgpack" | "rmp" => QueueCodec::Msgpack,
82        _ => QueueCodec::Msgpack,
83    };
84
85    let _ = QUEUE_CODEC_OVERRIDE.set(mapped);
86}
87
88pub struct QueueManager {
89    pub channel: Arc<Channel>,
90    pub backend: Option<Arc<dyn MqBackend>>,
91    pub compensator: Option<Arc<dyn Compensator>>,
92    pub blob_storage: Option<Arc<dyn BlobStorage>>,
93    pub batch_concurrency: usize,
94    pub compression_threshold: usize,
95    pub nack_policy: NackPolicy,
96    pub log_topic: String,
97}
98
99impl QueueManager {
100    pub fn new(backend: Option<Arc<dyn MqBackend>>, capacity: usize) -> Self {
101        QueueManager {
102            channel: Arc::new(Channel::new(capacity)),
103            backend,
104            compensator: None,
105            blob_storage: None,
106            batch_concurrency: 50,
107            compression_threshold: DEFAULT_COMPRESSION_THRESHOLD,
108            nack_policy: NackPolicy::default(),
109            log_topic: "log".to_string(),
110        }
111    }
112
113    pub fn from_config(cfg: &Config) -> Arc<Self> {
114        Self::from_config_with_log_topic(cfg, None)
115    }
116
117    pub fn from_config_with_log_topic(cfg: &Config, log_topic: Option<&str>) -> Arc<Self> {
118        set_queue_codec_from_config(cfg);
119        let channel_config = &cfg.channel_config;
120        #[allow(unused_variables)] // used by queue-kafka / queue-nats backends
121        let namespace = &cfg.name;
122
123        let nack_policy = if let Some(policy_cfg) = &cfg.policy
124            && !policy_cfg.overrides.is_empty()
125        {
126            let resolver = PolicyResolver::new(Some(policy_cfg));
127            let decision =
128                resolver.resolve_with_kind("queue", Some("nack"), Some("failed"), ErrorKind::Queue);
129            NackPolicy {
130                max_retries: decision.policy.max_retries,
131                backoff_ms: decision.policy.backoff_ms,
132            }
133        } else {
134            NackPolicy {
135                max_retries: channel_config.nack_max_retries.unwrap_or(0),
136                backoff_ms: channel_config.nack_backoff_ms.unwrap_or(0),
137            }
138        };
139
140        let mut queue_manager = if let Some(kafka_config) = &channel_config.kafka {
141            #[cfg(feature = "queue-kafka")]
142            let qm = match KafkaQueue::new(
143                kafka_config,
144                channel_config.minid_time,
145                namespace,
146                nack_policy,
147            ) {
148                Ok(kafka_queue) => {
149                    info!("KafkaQueue initialized successfully");
150                    QueueManager::new(Some(Arc::new(kafka_queue)), channel_config.capacity)
151                }
152                Err(e) => {
153                    error!("KafkaQueue init failed, fallback to in-memory queue: {}", e);
154                    QueueManager::new(None, channel_config.capacity)
155                }
156            };
157            #[cfg(not(feature = "queue-kafka"))]
158            let qm = {
159                let _ = kafka_config;
160                error!(
161                    "channel_config.kafka is set but the `queue-kafka` feature is disabled; \
162                     falling back to in-memory queue"
163                );
164                QueueManager::new(None, channel_config.capacity)
165            };
166            qm
167        } else if let Some(nats_config) = &channel_config.nats {
168            #[cfg(feature = "queue-nats")]
169            let qm = match NatsQueue::new(
170                nats_config,
171                channel_config.minid_time,
172                namespace,
173                nack_policy,
174            ) {
175                Ok(nats_queue) => {
176                    info!("NatsQueue (JetStream) initialized successfully");
177                    QueueManager::new(Some(Arc::new(nats_queue)), channel_config.capacity)
178                }
179                Err(e) => {
180                    error!("NatsQueue init failed, fallback to in-memory queue: {}", e);
181                    QueueManager::new(None, channel_config.capacity)
182                }
183            };
184            #[cfg(not(feature = "queue-nats"))]
185            let qm = {
186                let _ = nats_config;
187                error!(
188                    "channel_config.nats is set but the `queue-nats` feature is disabled; \
189                     falling back to in-memory queue"
190                );
191                QueueManager::new(None, channel_config.capacity)
192            };
193            qm
194        } else {
195            info!("In-Memory Queue initialized (Single Node Mode)");
196            QueueManager::new(None, 10000)
197        };
198
199        if let Some(topic) = log_topic {
200            queue_manager.log_topic = topic.to_string();
201        }
202
203        queue_manager.nack_policy = nack_policy;
204
205        if let Some(concurrency) = channel_config.batch_concurrency {
206            queue_manager.with_concurrency(concurrency);
207        }
208        if let Some(threshold) = channel_config.compression_threshold {
209            queue_manager.with_compression_threshold(threshold);
210        }
211
212        if let Some(blob_config) = &channel_config.blob_storage {
213            if let Some(path) = &blob_config.path {
214                let storage = Arc::new(FileSystemBlobStorage::new(path));
215                queue_manager.with_blob_storage(storage);
216                info!("BlobStorage initialized at: {}", path);
217            }
218        }
219
220        queue_manager.subscribe();
221        Arc::new(queue_manager)
222    }
223
224    pub fn with_backend(&mut self, backend: Arc<dyn MqBackend>) {
225        self.backend = Some(backend);
226    }
227
228    pub fn with_compensator(&mut self, compensator: Arc<dyn Compensator>) {
229        self.compensator = Some(compensator);
230    }
231
232    pub fn with_blob_storage(&mut self, storage: Arc<dyn BlobStorage>) {
233        self.blob_storage = Some(storage);
234    }
235
236    pub fn with_concurrency(&mut self, concurrency: usize) {
237        self.batch_concurrency = concurrency;
238    }
239
240    pub fn with_compression_threshold(&mut self, threshold: usize) {
241        self.compression_threshold = threshold;
242    }
243
244    pub fn with_log_topic(&mut self, topic: impl Into<String>) {
245        self.log_topic = topic.into();
246    }
247
248    pub fn subscribe(&self) {
249        if let Some(backend) = &self.backend {
250            let backend = backend.clone();
251            let channel = self.channel.clone();
252            let compensator = self.compensator.clone();
253            let blob_storage = self.blob_storage.clone();
254            let concurrency = self.batch_concurrency;
255            let compression_threshold = self.compression_threshold;
256
257            // Define outbound channels (Local -> Remote)
258            // format: (topic, receiver_channel)
259            self.spawn_forwarder(
260                "task",
261                channel.remote_task_receiver.clone(),
262                backend.clone(),
263                blob_storage.clone(),
264                concurrency,
265                compression_threshold,
266            );
267            self.spawn_forwarder(
268                "request",
269                channel.request_receiver.clone(),
270                backend.clone(),
271                blob_storage.clone(),
272                concurrency,
273                compression_threshold,
274            );
275            self.spawn_forwarder(
276                "response",
277                channel.response_receiver.clone(),
278                backend.clone(),
279                blob_storage.clone(),
280                concurrency,
281                compression_threshold,
282            );
283            self.spawn_forwarder(
284                "parser_task",
285                channel.parser_task_receiver.clone(),
286                backend.clone(),
287                blob_storage.clone(),
288                concurrency,
289                compression_threshold,
290            );
291            self.spawn_forwarder(
292                "error_task",
293                channel.error_receiver.clone(),
294                backend.clone(),
295                blob_storage.clone(),
296                concurrency,
297                compression_threshold,
298            );
299            self.spawn_forwarder(
300                self.log_topic.as_str(),
301                channel.log_receiver.clone(),
302                backend.clone(),
303                blob_storage.clone(),
304                concurrency,
305                compression_threshold,
306            );
307
308            // Define inbound subscriptions (Remote -> Local)
309            tokio::spawn(async move {
310                Self::subscribe_all_priorities(
311                    "task",
312                    channel.task_sender.clone(),
313                    backend.clone(),
314                    compensator.clone(),
315                    blob_storage.clone(),
316                    concurrency,
317                )
318                .await;
319                Self::subscribe_all_priorities(
320                    "request",
321                    channel.download_request_sender.clone(),
322                    backend.clone(),
323                    compensator.clone(),
324                    blob_storage.clone(),
325                    concurrency,
326                )
327                .await;
328                Self::subscribe_all_priorities(
329                    "response",
330                    channel.remote_response_sender.clone(),
331                    backend.clone(),
332                    compensator.clone(),
333                    blob_storage.clone(),
334                    concurrency,
335                )
336                .await;
337                Self::subscribe_all_priorities(
338                    "parser_task",
339                    channel.remote_parser_task_sender.clone(),
340                    backend.clone(),
341                    compensator.clone(),
342                    blob_storage.clone(),
343                    concurrency,
344                )
345                .await;
346                Self::subscribe_all_priorities(
347                    "error_task",
348                    channel.remote_error_sender.clone(),
349                    backend.clone(),
350                    compensator.clone(),
351                    blob_storage.clone(),
352                    concurrency,
353                )
354                .await;
355            });
356        }
357    }
358
359    async fn subscribe_all_priorities<T>(
360        topic_base: &str,
361        sender: Sender<QueuedItem<T>>,
362        backend: Arc<dyn MqBackend>,
363        compensator: Option<Arc<dyn Compensator>>,
364        blob_storage: Option<Arc<dyn BlobStorage>>,
365        concurrency: usize,
366    ) where
367        T: serde::de::DeserializeOwned
368            + Send
369            + 'static
370            + std::fmt::Debug
371            + Identifiable
372            + Offloadable,
373    {
374        let priorities = [Priority::High, Priority::Normal, Priority::Low];
375        futures::stream::iter(priorities)
376            .for_each_concurrent(None, |priority| {
377                let sender = sender.clone();
378                let backend = backend.clone();
379                let compensator = compensator.clone();
380                let blob_storage = blob_storage.clone();
381                async move {
382                    let topic = format!("{}-{}", topic_base, priority.suffix());
383                    Self::subscribe_topic(
384                        &topic,
385                        sender,
386                        backend,
387                        compensator,
388                        blob_storage,
389                        concurrency,
390                    )
391                    .await;
392                }
393            })
394            .await;
395    }
396
397    async fn subscribe_topic<T>(
398        topic: &str,
399        sender: Sender<QueuedItem<T>>,
400        backend: Arc<dyn MqBackend>,
401        compensator: Option<Arc<dyn Compensator>>,
402        blob_storage: Option<Arc<dyn BlobStorage>>,
403        concurrency: usize,
404    ) where
405        T: serde::de::DeserializeOwned
406            + Send
407            + 'static
408            + std::fmt::Debug
409            + Identifiable
410            + Offloadable,
411    {
412        let (tx, mut rx) = tokio::sync::mpsc::channel(1024);
413        if let Err(e) = backend.subscribe(topic, tx).await {
414            error!("Failed to subscribe to topic {}: {}", topic, e);
415            return;
416        }
417
418        let topic = topic.to_string();
419        // Use semaphore to limit number of concurrent BATCHES roughly
420        let batch_concurrency = (concurrency / 50).max(1);
421        let semaphore = Arc::new(Semaphore::new(batch_concurrency));
422
423        tokio::spawn(async move {
424            let topic_clone = topic.clone();
425            let sender_clone = sender.clone();
426            let compensator_clone = compensator.clone();
427            let blob_storage_clone = blob_storage.clone();
428
429            Batcher::run(&mut rx, 50, 5, semaphore, move |items| {
430                let topic = topic_clone.clone();
431                let sender = sender_clone.clone();
432                let compensator = compensator_clone.clone();
433                let blob_storage = blob_storage_clone.clone();
434                async move {
435                    Self::process_batch_messages(
436                        items,
437                        &topic,
438                        &sender,
439                        &compensator,
440                        &blob_storage,
441                    )
442                    .await;
443                }
444            })
445            .await;
446            log::warn!("Topic {} subscription closed", topic);
447        });
448    }
449
450    async fn process_batch_messages<T>(
451        messages: Vec<crate::queue::Message>,
452        topic: &str,
453        sender: &Sender<QueuedItem<T>>,
454        compensator: &Option<Arc<dyn Compensator>>,
455        blob_storage: &Option<Arc<dyn BlobStorage>>,
456    ) where
457        T: serde::de::DeserializeOwned
458            + Send
459            + 'static
460            + std::fmt::Debug
461            + Identifiable
462            + Offloadable,
463    {
464        // Offload deserialization to blocking thread for larger batches to avoid blocking the async runtime.
465        // For small batches, decode inline to reduce thread switch overhead.
466        let mut use_blocking = messages.len() >= 32;
467        if !use_blocking {
468            let mut total_bytes = 0usize;
469            for msg in &messages {
470                total_bytes += msg.payload.len();
471                if total_bytes >= BLOCKING_PAYLOAD_BYTES {
472                    use_blocking = true;
473                    break;
474                }
475            }
476        }
477        let results = if use_blocking {
478            tokio::task::spawn_blocking(move || {
479                messages
480                    .into_iter()
481                    .map(|msg| {
482                        let payload_slice = msg.payload.as_slice();
483                        let decoded_payload = decompress_payload(payload_slice);
484                        let item_res = match queue_codec() {
485                            QueueCodec::Json => {
486                                serde_json::from_slice::<T>(decoded_payload.as_ref()).map_err(|e| {
487                                    crate::errors::Error::new(
488                                        crate::errors::ErrorKind::Queue,
489                                        Some(e),
490                                    )
491                                })
492                            }
493                            QueueCodec::Msgpack => msgpack_decode::<T>(decoded_payload.as_ref())
494                                .map_err(|e| {
495                                    crate::errors::Error::new(
496                                        crate::errors::ErrorKind::Queue,
497                                        Some(e),
498                                    )
499                                }),
500                        };
501                        (msg, item_res)
502                    })
503                    .collect::<Vec<_>>()
504            })
505            .await
506        } else {
507            let processed_items = messages
508                .into_iter()
509                .map(|msg| {
510                    let payload_slice = msg.payload.as_slice();
511                    let decoded_payload = decompress_payload(payload_slice);
512                    let item_res = match queue_codec() {
513                        QueueCodec::Json => serde_json::from_slice::<T>(decoded_payload.as_ref())
514                            .map_err(|e| {
515                                crate::errors::Error::new(crate::errors::ErrorKind::Queue, Some(e))
516                            }),
517                        QueueCodec::Msgpack => msgpack_decode::<T>(decoded_payload.as_ref())
518                            .map_err(|e| {
519                                crate::errors::Error::new(crate::errors::ErrorKind::Queue, Some(e))
520                            }),
521                    };
522                    (msg, item_res)
523                })
524                .collect::<Vec<_>>();
525            Ok(processed_items)
526        };
527
528        match results {
529            Ok(processed_items) => {
530                let tasks = processed_items.into_iter().map(|(msg, result)| {
531                    let topic = topic.to_string();
532                    let compensator = compensator.clone();
533                    let storage = blob_storage.clone();
534
535                    async move {
536                        match result {
537                            Ok(mut item) => {
538                                // Reload content from blob storage if necessary
539                                if let Some(storage) = storage {
540                                    if let Err(e) = item.reload(&storage).await {
541                                         error!("Failed to reload item content from storage for topic {}: {}", topic, e);
542                                         if let Err(e) = msg.nack("Blob reload failed").await {
543                                             error!("Failed to NACK message: {}", e);
544                                         }
545                                         return None;
546                                    }
547                                }
548
549                                let id = item.get_id();
550
551                                if let Some(comp) = compensator {
552                                     if let Err(e) = comp.add_task(&topic, &id, msg.payload.clone()).await {
553                                        error!(
554                                            "Failed to add task to compensation queue for topic {}: {}",
555                                            topic, e
556                                        );
557                                     }
558                                }
559
560                                let msg_ack = msg.clone();
561                                let msg_nack = msg.clone();
562
563                                let queued_item = QueuedItem::with_ack(
564                                    item,
565                                    move || Box::pin(async move { msg_ack.ack().await }),
566                                    move |reason| Box::pin(async move { msg_nack.nack(reason).await })
567                                );
568
569                                Some((msg, queued_item))
570                            }
571                            Err(e) => {
572                                let payload_len = msg.payload.len();
573                                let codec = match queue_codec() {
574                                    QueueCodec::Json => "json",
575                                    QueueCodec::Msgpack => "msgpack",
576                                };
577
578                                error!(
579                                    "Failed to deserialize message from topic {} (codec={}, bytes={}): {}",
580                                    topic,
581                                    codec,
582                                    payload_len,
583                                    e
584                                );
585                                if let Err(e) = msg.nack("Deserialization failed").await {
586                                    error!("Failed to NACK poison message: {}", e);
587                                }
588                                None
589                            }
590                        }
591                    }
592                });
593
594                let results = join_all(tasks).await;
595
596                for (msg, item) in results.into_iter().flatten() {
597                    if let Err(_e) = sender.send(item).await {
598                        let _ = msg.nack("Channel closed").await;
599                    }
600                }
601            }
602            Err(e) => {
603                error!("Batch deserialization task failed: {}", e);
604                // We can't easily NACK here because we lost ownership of msgs inside the closure if it panicked.
605                // But spawn_blocking usually returns JoinError if panicked.
606            }
607        }
608    }
609
610    fn spawn_forwarder<T>(
611        &self,
612        topic: &str,
613        receiver: Arc<Mutex<Receiver<QueuedItem<T>>>>,
614        backend: Arc<dyn MqBackend>,
615        blob_storage: Option<Arc<dyn BlobStorage>>,
616        concurrency: usize,
617        compression_threshold: usize,
618    ) where
619        T: serde::Serialize + Send + Sync + 'static + Identifiable + Prioritizable + Offloadable,
620    {
621        Self::forward_channel(
622            topic,
623            receiver,
624            backend,
625            blob_storage,
626            concurrency,
627            compression_threshold,
628        )
629    }
630
631    fn forward_channel<T>(
632        topic: &str,
633        receiver: Arc<Mutex<Receiver<QueuedItem<T>>>>,
634        backend: Arc<dyn MqBackend>,
635        blob_storage: Option<Arc<dyn BlobStorage>>,
636        concurrency: usize,
637        compression_threshold: usize,
638    ) where
639        T: serde::Serialize + Send + Sync + 'static + Identifiable + Prioritizable + Offloadable,
640    {
641        let topic = topic.to_string();
642        // Concurrency limit for publishing batches
643        let semaphore = Arc::new(Semaphore::new(concurrency));
644
645        tokio::spawn(async move {
646            let mut rx = receiver.lock().await;
647            let topic_clone = topic.clone();
648            let backend_clone = backend.clone();
649            let blob_storage = blob_storage.clone();
650            let compression_threshold = compression_threshold;
651
652            Batcher::run(&mut *rx, 500, 5, semaphore, move |items| {
653                let topic = topic_clone.clone();
654                let backend = backend_clone.clone();
655                let blob_storage = blob_storage.clone();
656                async move {
657                    Self::flush_batch_grouped(
658                        topic,
659                        backend,
660                        blob_storage,
661                        items,
662                        compression_threshold,
663                    )
664                    .await;
665                }
666            })
667            .await;
668        });
669    }
670
671    pub fn get_task_push_channel(&self) -> Sender<QueuedItem<TaskEvent>> {
672        if self.backend.is_none() {
673            return self.channel.task_sender.clone();
674        }
675        self.channel.remote_task_sender.clone()
676    }
677    pub fn get_task_pop_channel(&self) -> Arc<Mutex<Receiver<QueuedItem<TaskEvent>>>> {
678        Arc::clone(&self.channel.task_receiver)
679    }
680    pub fn get_request_pop_channel(&self) -> Arc<Mutex<Receiver<QueuedItem<Request>>>> {
681        self.channel.download_request_receiver.clone()
682    }
683    pub fn get_request_push_channel(&self) -> Sender<QueuedItem<Request>> {
684        if self.backend.is_none() {
685            return self.channel.download_request_sender.clone();
686        }
687        self.channel.request_sender.clone()
688    }
689    pub fn get_response_push_channel(&self) -> Sender<QueuedItem<Response>> {
690        if self.backend.is_none() {
691            return self.channel.remote_response_sender.clone();
692        }
693        self.channel.response_sender.clone()
694    }
695
696    /// Attempts to send directly to local consumers (bypassing the MQ backend).
697    /// Optimization: if local consumers exist and the channel is not full, send
698    /// directly to avoid serialization and network overhead.
699    pub fn try_send_local_response(
700        &self,
701        item: QueuedItem<Response>,
702    ) -> Result<(), tokio::sync::mpsc::error::TrySendError<QueuedItem<Response>>> {
703        self.channel.remote_response_sender.try_send(item)
704    }
705
706    pub fn get_response_pop_channel(&self) -> Arc<Mutex<Receiver<QueuedItem<Response>>>> {
707        Arc::clone(&self.channel.remote_response_receiver)
708    }
709    pub fn get_parser_task_pop_channel(&self) -> Arc<Mutex<Receiver<QueuedItem<TaskParserEvent>>>> {
710        self.channel.remote_parser_task_receiver.clone()
711    }
712    pub fn get_parser_task_push_channel(&self) -> Sender<QueuedItem<TaskParserEvent>> {
713        if self.backend.is_none() {
714            return self.channel.remote_parser_task_sender.clone();
715        }
716        self.channel.parser_task_sender.clone()
717    }
718    pub fn get_error_pop_channel(&self) -> Arc<Mutex<Receiver<QueuedItem<TaskErrorEvent>>>> {
719        Arc::clone(&self.channel.remote_error_receiver)
720    }
721    pub fn get_error_push_channel(&self) -> Sender<QueuedItem<TaskErrorEvent>> {
722        if self.backend.is_none() {
723            return self.channel.remote_error_sender.clone();
724        }
725        self.channel.error_sender.clone()
726    }
727    pub fn get_log_push_channel(&self) -> Sender<QueuedItem<LogModel>> {
728        self.channel.log_sender.clone()
729    }
730
731    /// Local pending message count across processor queues
732    pub async fn local_pending_count(&self) -> usize {
733        let (task, download, response, parser, error, remote_task) =
734            self.local_pending_breakdown().await;
735        task + download + response + parser + error + remote_task
736    }
737
738    pub async fn local_pending_breakdown(&self) -> (usize, usize, usize, usize, usize, usize) {
739        let task = self
740            .channel
741            .task_sender
742            .max_capacity()
743            .saturating_sub(self.channel.task_sender.capacity());
744        let download = self
745            .channel
746            .download_request_sender
747            .max_capacity()
748            .saturating_sub(self.channel.download_request_sender.capacity());
749        let response = self
750            .channel
751            .remote_response_sender
752            .max_capacity()
753            .saturating_sub(self.channel.remote_response_sender.capacity());
754        let parser = self
755            .channel
756            .remote_parser_task_sender
757            .max_capacity()
758            .saturating_sub(self.channel.remote_parser_task_sender.capacity());
759        let error = self
760            .channel
761            .remote_error_sender
762            .max_capacity()
763            .saturating_sub(self.channel.remote_error_sender.capacity());
764        let remote_task = self
765            .channel
766            .remote_task_sender
767            .max_capacity()
768            .saturating_sub(self.channel.remote_task_sender.capacity());
769
770        (task, download, response, parser, error, remote_task)
771    }
772
773    async fn flush_batch_grouped<T>(
774        base_topic: String,
775        backend: Arc<dyn MqBackend>,
776        blob_storage: Option<Arc<dyn BlobStorage>>,
777        mut items: Vec<QueuedItem<T>>,
778        compression_threshold: usize,
779    ) where
780        T: serde::Serialize + Identifiable + Send + Sync + Prioritizable + Offloadable + 'static,
781    {
782        // Try offloading large payloads if storage is configured
783        if let Some(storage) = &blob_storage {
784            // We iterate mutably to potentially modify items (offload content)
785            for item in &mut items {
786                if item.inner.should_offload(BLOCKING_PAYLOAD_BYTES) {
787                    if let Err(e) = item.inner.offload(storage).await {
788                        error!("Failed to offload item payload to blob storage: {}", e);
789                        // Continue, hoping it might still pass or fail later
790                    }
791                }
792            }
793        }
794
795        // Optimization: Fast path if all items have same priority (likely normal)
796        // or if items is empty
797        if items.is_empty() {
798            return;
799        }
800
801        let first_priority = items[0].get_priority();
802        let all_same_priority = items.iter().all(|i| i.get_priority() == first_priority);
803
804        if all_same_priority {
805            let topic = format!("{}-{}", base_topic, first_priority.suffix());
806            Self::flush_batch(topic, backend, items, compression_threshold).await;
807            return;
808        }
809
810        let mut groups: HashMap<Priority, Vec<QueuedItem<T>>> = HashMap::new();
811        for item in items {
812            groups.entry(item.get_priority()).or_default().push(item);
813        }
814
815        for (priority, group_items) in groups {
816            let topic = format!("{}-{}", base_topic, priority.suffix());
817            Self::flush_batch(topic, backend.clone(), group_items, compression_threshold).await;
818        }
819    }
820
821    async fn flush_batch<T>(
822        topic: String,
823        backend: Arc<dyn MqBackend>,
824        items: Vec<QueuedItem<T>>,
825        compression_threshold: usize,
826    ) where
827        T: serde::Serialize + Identifiable + Send + Sync + 'static,
828    {
829        // Serialize all items once. Decide blocking vs inline based on item count alone
830        // to avoid the previous double-serialization (once to estimate size, once to encode).
831        let use_blocking = items.len() >= 32;
832
833        let payloads_result = if use_blocking {
834            tokio::task::spawn_blocking(move || Self::encode_items(&items, compression_threshold))
835                .await
836        } else {
837            Ok(Self::encode_items(&items, compression_threshold))
838        };
839
840        match payloads_result {
841            Ok((payloads, ids, count)) => {
842                if payloads.is_empty() {
843                    return;
844                }
845
846                if let Err(e) = backend.publish_batch_with_headers(&topic, &payloads).await {
847                    error!("Failed to publish batch to topic {}: {}", topic, e);
848                } else {
849                    log::info!(
850                        "[QueueManager] forward_channel published batch: topic={} count={}",
851                        topic,
852                        count
853                    );
854                    if let Some(id_list) = ids {
855                        for id in id_list {
856                            log::debug!(
857                                "[QueueManager] forward_channel published: topic={} id={}",
858                                topic,
859                                id
860                            );
861                        }
862                    }
863                }
864            }
865            Err(e) => {
866                error!("Serialization task join error: {}", e);
867            }
868        }
869    }
870
871    pub async fn clean_storage(&self) -> crate::errors::Result<()> {
872        if let Some(backend) = &self.backend {
873            backend.clean_storage().await?;
874        }
875        Ok(())
876    }
877
878    /// Encode items into publish-ready payloads (serialize once + compress).
879    fn encode_items<T>(
880        items: &[QueuedItem<T>],
881        compression_threshold: usize,
882    ) -> (
883        Vec<(Option<String>, Vec<u8>, HashMap<String, String>)>,
884        Option<Vec<String>>,
885        usize,
886    )
887    where
888        T: serde::Serialize + Identifiable,
889    {
890        let mut payloads: Vec<(Option<String>, Vec<u8>, HashMap<String, String>)> =
891            Vec::with_capacity(items.len());
892        let mut ids = if log::log_enabled!(log::Level::Debug) {
893            Some(Vec::with_capacity(items.len()))
894        } else {
895            None
896        };
897
898        for item in items {
899            let id = item.get_id();
900            // The partition key drives MQ routing / sharding (account affinity); it is
901            // independent of the id used for deduplication / logging.
902            let partition_key = item.partition_key();
903            if let Some(id_list) = ids.as_mut() {
904                id_list.push(id.clone());
905            }
906            let encoded = match queue_codec() {
907                QueueCodec::Json => serde_json::to_vec(&item.inner).map_err(|e| {
908                    crate::errors::Error::new(crate::errors::ErrorKind::Queue, Some(e))
909                }),
910                QueueCodec::Msgpack => msgpack_encode(&item.inner).map_err(|e| {
911                    crate::errors::Error::new(crate::errors::ErrorKind::Queue, Some(e))
912                }),
913            };
914            match encoded {
915                Ok(p) => {
916                    let final_payload = compress_payload_owned(p, compression_threshold);
917                    let headers = default_headers();
918                    payloads.push((Some(partition_key), final_payload, headers));
919                }
920                Err(e) => {
921                    error!(
922                        "Failed to serialize item id={}: {}. Item will be skipped (unrecoverable).",
923                        id, e
924                    );
925                    counter!("mocra_queue_encode_errors_total").increment(1);
926                }
927            }
928        }
929        (payloads, ids, items.len())
930    }
931
932    /// Send a message to the Dead Letter Queue (DLQ) manually.
933    ///
934    /// This is useful when the application logic decides a message cannot be processed
935    /// even if the message delivery itself was successful (e.g., max logic retries exceeded).
936    pub async fn send_to_dlq<T>(
937        &self,
938        topic: &str,
939        item: &T,
940        reason: &str,
941    ) -> crate::errors::Result<()>
942    where
943        T: serde::Serialize + Identifiable + Send + Sync,
944    {
945        if let Some(backend) = &self.backend {
946            let payload = match queue_codec() {
947                QueueCodec::Json => serde_json::to_vec(item)
948                    .map_err(|e| crate::errors::error::QueueError::OperationFailed(Box::new(e)))?,
949                QueueCodec::Msgpack => msgpack_encode(item)
950                    .map_err(|e| crate::errors::error::QueueError::OperationFailed(Box::new(e)))?,
951            };
952            backend
953                .send_to_dlq(topic, &item.get_id(), &payload, reason)
954                .await?;
955        }
956        Ok(())
957    }
958
959    pub async fn read_dlq(
960        &self,
961        topic: &str,
962        count: usize,
963    ) -> crate::errors::Result<Vec<(String, Vec<u8>, String, String)>> {
964        if let Some(backend) = &self.backend {
965            backend.read_dlq(topic, count).await
966        } else {
967            Ok(Vec::new())
968        }
969    }
970}