use std::marker::PhantomData;
use std::sync::Arc;
use bytes::Bytes;
use serde::Serialize;
use crate::queue::backend::SenderBackend;
use crate::queue::error::WorkQueueSendError;
pub struct WorkQueueSender<T: Serialize> {
backend: Arc<dyn SenderBackend>,
_marker: PhantomData<fn(T)>,
}
impl<T: Serialize> Clone for WorkQueueSender<T> {
fn clone(&self) -> Self {
Self {
backend: Arc::clone(&self.backend),
_marker: PhantomData,
}
}
}
impl<T: Serialize> WorkQueueSender<T> {
pub(crate) fn new(backend: Arc<dyn SenderBackend>) -> Self {
Self {
backend,
_marker: PhantomData,
}
}
pub async fn enqueue(&self, data: &T) -> Result<(), WorkQueueSendError> {
let bytes = serialize(data)?;
self.backend.send(bytes).await
}
pub fn try_enqueue(&self, data: &T) -> Result<(), WorkQueueSendError> {
let bytes = serialize(data)?;
self.backend.try_send(bytes)
}
pub async fn close(&self) {
self.backend.close().await;
}
}
fn serialize<T: Serialize>(data: &T) -> Result<Bytes, WorkQueueSendError> {
rmp_serde::to_vec(data)
.map(Bytes::from)
.map_err(|e| WorkQueueSendError::Serialization(e.to_string()))
}