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)] 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 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 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 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 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 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 }
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 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 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 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 if let Some(storage) = &blob_storage {
784 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 }
791 }
792 }
793 }
794
795 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 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 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 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 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}