use std::sync::atomic::Ordering;
use std::sync::Arc;
use tracing::error;
use crate::ilink::types::{HubExt, WeixinMessage};
use super::super::*;
pub async fn push_to_queue_pub(
queue: &Arc<dyn MessageQueue>,
metrics: &Metrics,
vtoken: &str,
msg: WeixinMessage,
) {
push_to_queue(queue, metrics, vtoken, msg).await;
}
pub(super) async fn push_to_queue(
queue: &Arc<dyn MessageQueue>,
metrics: &Metrics,
vtoken: &str,
msg: WeixinMessage,
) {
match queue.push(vtoken, msg).await {
Ok(false) => {
metrics.messages_dispatched.fetch_add(1, Ordering::Relaxed);
}
Ok(true) => {
metrics.messages_dropped.fetch_add(1, Ordering::Relaxed);
}
Err(e) => {
error!(error = %e, vtoken = %crate::redact_token(vtoken), "failed to push message to queue");
metrics.messages_dropped.fetch_add(1, Ordering::Relaxed);
}
}
}
pub(super) async fn push_shared_to_queue(
queue: &Arc<dyn MessageQueue>,
metrics: &Metrics,
vtoken: &str,
base: Arc<WeixinMessage>,
context_token: Option<String>,
hub_ext: Option<HubExt>,
) {
match queue
.push_shared(vtoken, base, context_token, hub_ext)
.await
{
Ok(false) => {
metrics.messages_dispatched.fetch_add(1, Ordering::Relaxed);
}
Ok(true) => {
metrics.messages_dropped.fetch_add(1, Ordering::Relaxed);
}
Err(e) => {
error!(error = %e, vtoken = %crate::redact_token(vtoken), "failed to push shared message to queue");
metrics.messages_dropped.fetch_add(1, Ordering::Relaxed);
}
}
}