use crate::sync::error::{SyncError, SyncResult};
use crate::sync::types::SyncMessage;
use parking_lot::Mutex;
use std::collections::VecDeque;
pub use crate::sync::outbox_journal::{FileSyncOutbox, SYNC_OUTBOX_FORMAT_V2};
pub trait SyncOutbox: Send + Sync {
fn try_enqueue(&self, message: SyncMessage, max_len: usize) -> SyncResult<bool>;
fn front(&self) -> SyncResult<Option<SyncMessage>>;
fn acknowledge_front(&self, batch_id: &str) -> SyncResult<()>;
fn messages(&self) -> SyncResult<Vec<SyncMessage>>;
fn len(&self) -> SyncResult<usize>;
fn is_empty(&self) -> SyncResult<bool> {
Ok(self.len()? == 0)
}
}
#[derive(Debug, Default)]
pub struct InMemorySyncOutbox {
messages: Mutex<VecDeque<SyncMessage>>,
}
impl InMemorySyncOutbox {
pub fn new() -> Self {
Self::default()
}
}
impl SyncOutbox for InMemorySyncOutbox {
fn try_enqueue(&self, message: SyncMessage, max_len: usize) -> SyncResult<bool> {
let mut messages = self.messages.lock();
if messages.len() >= max_len {
return Ok(false);
}
messages.push_back(message);
Ok(true)
}
fn front(&self) -> SyncResult<Option<SyncMessage>> {
Ok(self.messages.lock().front().cloned())
}
fn acknowledge_front(&self, batch_id: &str) -> SyncResult<()> {
let mut messages = self.messages.lock();
if messages.front().map(|message| message.batch_id.as_str()) != Some(batch_id) {
return Err(SyncError::InvalidSyncMessage(
"outbox acknowledgement mismatch",
));
}
messages.pop_front();
Ok(())
}
fn messages(&self) -> SyncResult<Vec<SyncMessage>> {
Ok(self.messages.lock().iter().cloned().collect())
}
fn len(&self) -> SyncResult<usize> {
Ok(self.messages.lock().len())
}
}