use std::{collections::VecDeque, path::PathBuf, sync::Arc};
use miette::{Result, miette};
use parking_lot::Mutex;
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use crate::persistence::PersistenceStore;
const PENDING_WORK_FILE_NAME: &str = "pending_work_queue";
#[derive(Clone)]
pub struct PendingWorkQueue {
inner: Arc<Mutex<PendingWorkQueueInner>>,
}
struct PendingWorkQueueInner {
path: PathBuf,
state: PersistedPendingWorkQueue,
}
#[derive(Default, Serialize, Deserialize)]
struct PersistedPendingWorkQueue {
queue: VecDeque<PendingWorkEntry>,
}
#[derive(Clone, Serialize, Deserialize)]
struct PendingWorkEntry {
work: PendingWork,
state: PendingWorkEntryState,
}
#[derive(Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
enum PendingWorkEntryState {
Pending,
Claimed,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum PendingEventMoveDirection {
Up,
Down,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub enum PendingWork {
Event { event_id: Uuid },
}
impl PartialEq for PendingWork {
fn eq(&self, other: &Self) -> bool {
match (self, other) {
(Self::Event { event_id: a }, Self::Event { event_id: b }) => a == b,
}
}
}
impl Eq for PendingWork {}
impl PendingWork {
const fn priority(&self) -> u8 {
match self {
Self::Event { .. } => 0,
}
}
}
impl PendingWorkQueue {
pub async fn new() -> Self {
let persistence = PersistenceStore::runtime().await;
Self::open_with_persistence(persistence).await
}
pub async fn with_session(session_id: &str) -> Self {
let persistence = PersistenceStore::for_session(Some(session_id)).await;
Self::open_with_persistence(persistence).await
}
async fn open_with_persistence(persistence: PersistenceStore) -> Self {
let path = persistence.state_file(PENDING_WORK_FILE_NAME);
let mut state: PersistedPendingWorkQueue = persistence
.read_postcard_state_or_default(PENDING_WORK_FILE_NAME, "pending work queue")
.await;
reset_claimed_entries_on_startup(&mut state);
Self {
inner: Arc::new(Mutex::new(PendingWorkQueueInner { path, state })),
}
}
pub fn pending_count(&self) -> usize {
self.inner
.lock()
.state
.queue
.iter()
.filter(|entry| matches!(entry.state, PendingWorkEntryState::Pending))
.count()
}
pub fn enqueue(&self, work: PendingWork) -> Result<bool> {
let mut inner = self.inner.lock();
if let Some(existing) = inner
.state
.queue
.iter_mut()
.find(|entry| entry.work == work)
{
let changed = !same_work_payload(&existing.work, &work);
existing.work = work;
if changed {
persist_locked(&inner)?;
}
return Ok(false);
}
inner.state.queue.push_back(PendingWorkEntry {
work,
state: PendingWorkEntryState::Pending,
});
persist_locked(&inner)?;
drop(inner);
Ok(true)
}
pub fn claim_batch(&self, max_items: usize) -> Result<Vec<PendingWork>> {
if max_items == 0 {
return Ok(Vec::new());
}
let mut inner = self.inner.lock();
let mut claimed = Vec::new();
for _ in 0..max_items {
let Some(index) = select_next_pending_index(&inner.state.queue) else {
break;
};
let entry = &mut inner.state.queue[index];
entry.state = PendingWorkEntryState::Claimed;
claimed.push(entry.work.clone());
}
if !claimed.is_empty() {
persist_locked(&inner)?;
}
drop(inner);
Ok(claimed)
}
pub fn consume(&self, work: &PendingWork) -> Result<bool> {
let mut inner = self.inner.lock();
let Some(index) = inner
.state
.queue
.iter()
.position(|entry| entry.work == *work)
else {
return Ok(false);
};
inner.state.queue.remove(index);
persist_locked(&inner)?;
drop(inner);
Ok(true)
}
pub fn requeue_front(&self, work: PendingWork) -> Result<bool> {
let mut inner = self.inner.lock();
if let Some(index) = inner
.state
.queue
.iter()
.position(|entry| entry.work == work)
{
let mut entry = inner
.state
.queue
.remove(index)
.expect("pending work index should be valid");
entry.work = work;
entry.state = PendingWorkEntryState::Pending;
inner.state.queue.push_front(entry);
persist_locked(&inner)?;
return Ok(true);
}
inner.state.queue.push_front(PendingWorkEntry {
work,
state: PendingWorkEntryState::Pending,
});
persist_locked(&inner)?;
drop(inner);
Ok(true)
}
pub fn pending_event_ids(&self) -> Vec<Uuid> {
let inner = self.inner.lock();
inner
.state
.queue
.iter()
.filter_map(|entry| match (&entry.work, entry.state) {
(PendingWork::Event { event_id }, PendingWorkEntryState::Pending) => {
Some(*event_id)
}
_ => None,
})
.collect()
}
pub fn move_pending_event(
&self,
event_id: Uuid,
direction: &PendingEventMoveDirection,
) -> Result<bool> {
let mut inner = self.inner.lock();
let event_indices = pending_event_indices(&inner.state.queue);
let Some(position) = event_indices.iter().position(|index| {
matches!(
inner.state.queue[*index].work,
PendingWork::Event { event_id: current } if current == event_id
)
}) else {
return Ok(false);
};
let target_position = match direction {
PendingEventMoveDirection::Up => position.checked_sub(1),
PendingEventMoveDirection::Down => {
(position + 1 < event_indices.len()).then_some(position + 1)
}
};
let Some(target_position) = target_position else {
return Ok(false);
};
inner
.state
.queue
.swap(event_indices[position], event_indices[target_position]);
persist_locked(&inner)?;
drop(inner);
Ok(true)
}
pub fn move_pending_event_to_position(
&self,
event_id: Uuid,
target_position: usize,
) -> Result<bool> {
let mut inner = self.inner.lock();
let event_indices = pending_event_indices(&inner.state.queue);
let Some(position) = event_indices.iter().position(|index| {
matches!(
inner.state.queue[*index].work,
PendingWork::Event { event_id: current } if current == event_id
)
}) else {
return Ok(false);
};
if event_indices.len() <= 1 {
return Ok(false);
}
let target_position = target_position.min(event_indices.len() - 1);
if position == target_position {
return Ok(false);
}
let source_index = event_indices[position];
let Some(entry) = inner.state.queue.remove(source_index) else {
return Ok(false);
};
let adjusted_event_indices = pending_event_indices(&inner.state.queue);
let insert_index = if target_position >= adjusted_event_indices.len() {
adjusted_event_indices
.last()
.map_or(inner.state.queue.len(), |index| index + 1)
} else {
adjusted_event_indices[target_position]
};
inner.state.queue.insert(insert_index, entry);
persist_locked(&inner)?;
drop(inner);
Ok(true)
}
pub fn move_pending_event_to_front(&self, event_id: Uuid) -> Result<bool> {
let mut inner = self.inner.lock();
let event_indices = pending_event_indices(&inner.state.queue);
let Some(position) = event_indices.iter().position(|index| {
matches!(
inner.state.queue[*index].work,
PendingWork::Event { event_id: current } if current == event_id
)
}) else {
return Ok(false);
};
if position == 0 {
return Ok(false);
}
let source_index = event_indices[position];
let Some(entry) = inner.state.queue.remove(source_index) else {
return Ok(false);
};
let insert_index = pending_event_indices(&inner.state.queue)
.first()
.copied()
.unwrap_or(inner.state.queue.len());
inner.state.queue.insert(insert_index, entry);
persist_locked(&inner)?;
drop(inner);
Ok(true)
}
pub fn clear_events(&self) -> Result<usize> {
let mut inner = self.inner.lock();
let before = inner.state.queue.len();
inner
.state
.queue
.retain(|entry| !matches!(entry.work, PendingWork::Event { .. }));
let cleared = before.saturating_sub(inner.state.queue.len());
if cleared > 0 {
persist_locked(&inner)?;
}
drop(inner);
Ok(cleared)
}
pub fn shutdown(self) {
let inner = self.inner.lock();
let _ = persist_locked(&inner);
}
}
fn pending_event_indices(queue: &VecDeque<PendingWorkEntry>) -> Vec<usize> {
queue
.iter()
.enumerate()
.filter_map(|(index, entry)| match (&entry.work, entry.state) {
(PendingWork::Event { .. }, PendingWorkEntryState::Pending) => Some(index),
_ => None,
})
.collect()
}
fn select_next_pending_index(queue: &VecDeque<PendingWorkEntry>) -> Option<usize> {
queue
.iter()
.enumerate()
.filter(|(_, entry)| matches!(entry.state, PendingWorkEntryState::Pending))
.min_by_key(|(index, entry)| (entry.work.priority(), *index))
.map(|(index, _)| index)
}
fn same_work_payload(left: &PendingWork, right: &PendingWork) -> bool {
match (left, right) {
(PendingWork::Event { event_id: a }, PendingWork::Event { event_id: b }) => a == b,
}
}
fn persist_locked(inner: &PendingWorkQueueInner) -> Result<()> {
crate::persistence::write_postcard_atomic_sync(
&inner.path,
&inner.state,
crate::persistence::PersistenceFileMode::Default,
)
.map_err(|err| miette!("persist pending work queue failed: {err}"))
}
fn reset_claimed_entries_on_startup(state: &mut PersistedPendingWorkQueue) {
for entry in &mut state.queue {
if matches!(entry.state, PendingWorkEntryState::Claimed) {
entry.state = PendingWorkEntryState::Pending;
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicU64, Ordering};
static TEST_QUEUE_COUNTER: AtomicU64 = AtomicU64::new(1);
fn test_queue() -> PendingWorkQueue {
let unique = TEST_QUEUE_COUNTER.fetch_add(1, Ordering::Relaxed);
let path = std::env::temp_dir().join(format!(
"daat-locus-pending-work-test-{}-{}.bin",
std::process::id(),
unique
));
let _ = std::fs::remove_file(&path);
PendingWorkQueue {
inner: Arc::new(Mutex::new(PendingWorkQueueInner {
path,
state: PersistedPendingWorkQueue::default(),
})),
}
}
#[test]
fn persisted_pending_work_queue_postcard_round_trips() {
let event_id = Uuid::parse_str("33333333-3333-4333-8333-333333333333").expect("event uuid");
let state = PersistedPendingWorkQueue {
queue: VecDeque::from([PendingWorkEntry {
work: PendingWork::Event { event_id },
state: PendingWorkEntryState::Pending,
}]),
};
let bytes = postcard::to_allocvec(&state).expect("encode pending work");
let restored: PersistedPendingWorkQueue =
postcard::from_bytes(&bytes).expect("decode pending work");
assert_eq!(restored.queue.len(), 1);
let PendingWork::Event {
event_id: restored_event_id,
} = &restored.queue[0].work;
assert_eq!(*restored_event_id, event_id);
assert!(matches!(
restored.queue[0].state,
PendingWorkEntryState::Pending
));
}
#[test]
fn requeue_front_reactivates_claimed_event_driver() {
let queue = test_queue();
let event_id = Uuid::new_v4();
let work = PendingWork::Event { event_id };
queue.enqueue(work.clone()).expect("enqueue event");
let claimed = queue.claim_batch(1).expect("claim event");
assert_eq!(claimed.len(), 1);
assert!(matches!(claimed[0], PendingWork::Event { .. }));
assert_eq!(queue.pending_count(), 0);
queue.requeue_front(work).expect("requeue claimed event");
assert_eq!(queue.pending_count(), 1);
let reclaimed = queue.claim_batch(1).expect("claim requeued event");
assert_eq!(reclaimed.len(), 1);
let PendingWork::Event {
event_id: reclaimed_event_id,
} = &reclaimed[0];
assert_eq!(*reclaimed_event_id, event_id);
}
#[test]
fn startup_releases_claimed_pending_work_entries() {
let event_id = Uuid::new_v4();
let mut state = PersistedPendingWorkQueue::default();
state.queue.push_back(PendingWorkEntry {
work: PendingWork::Event { event_id },
state: PendingWorkEntryState::Claimed,
});
reset_claimed_entries_on_startup(&mut state);
assert!(
state
.queue
.iter()
.all(|entry| matches!(entry.state, PendingWorkEntryState::Pending))
);
}
#[test]
fn pending_event_ids_follow_manual_reordering() {
let queue = test_queue();
let first_event_id = Uuid::new_v4();
let second_event_id = Uuid::new_v4();
let third_event_id = Uuid::new_v4();
queue
.enqueue(PendingWork::Event {
event_id: first_event_id,
})
.expect("enqueue first event");
queue
.enqueue(PendingWork::Event {
event_id: second_event_id,
})
.expect("enqueue second event");
queue
.enqueue(PendingWork::Event {
event_id: third_event_id,
})
.expect("enqueue third event");
assert_eq!(
queue.pending_event_ids(),
vec![first_event_id, second_event_id, third_event_id]
);
assert!(
queue
.move_pending_event(third_event_id, &PendingEventMoveDirection::Up)
.expect("move third event up")
);
assert_eq!(
queue.pending_event_ids(),
vec![first_event_id, third_event_id, second_event_id]
);
assert!(
queue
.move_pending_event(first_event_id, &PendingEventMoveDirection::Down)
.expect("move first event down")
);
assert_eq!(
queue.pending_event_ids(),
vec![third_event_id, first_event_id, second_event_id]
);
}
#[test]
fn pending_events_move_to_absolute_position() {
let queue = test_queue();
let first_event_id = Uuid::new_v4();
let second_event_id = Uuid::new_v4();
let third_event_id = Uuid::new_v4();
queue
.enqueue(PendingWork::Event {
event_id: first_event_id,
})
.expect("enqueue first event");
queue
.enqueue(PendingWork::Event {
event_id: second_event_id,
})
.expect("enqueue second event");
queue
.enqueue(PendingWork::Event {
event_id: third_event_id,
})
.expect("enqueue third event");
assert!(
queue
.move_pending_event_to_position(first_event_id, 2)
.expect("move first event to end")
);
assert_eq!(
queue.pending_event_ids(),
vec![second_event_id, third_event_id, first_event_id]
);
assert!(
queue
.move_pending_event_to_position(first_event_id, 0)
.expect("move first event to start")
);
assert_eq!(
queue.pending_event_ids(),
vec![first_event_id, second_event_id, third_event_id]
);
assert!(
!queue
.move_pending_event_to_position(second_event_id, 1)
.expect("move event to current position")
);
}
}