use siphasher::sip::SipHasher13;
use std::{hash::Hasher, mem::size_of, ops::Range};
use crate::{
internal::{KeyNamespace, RangeMetadata, aligned_data_entry_size},
store::CandyStore,
types::{Error, MAX_USER_VALUE_SIZE, Result},
};
#[derive(Clone, Copy)]
pub(super) struct QueueNamespaces {
pub(super) meta: KeyNamespace,
pub(super) data: KeyNamespace,
}
const QUEUE_NS: QueueNamespaces = QueueNamespaces {
meta: KeyNamespace::QueueMeta,
data: KeyNamespace::QueueData,
};
const BIG_NS: QueueNamespaces = QueueNamespaces {
meta: KeyNamespace::BigMeta,
data: KeyNamespace::BigData,
};
pub struct QueueIterator<'a> {
store: &'a CandyStore,
queue: Vec<u8>,
ns: QueueNamespaces,
next_idx: u64,
end_idx: u64,
initial_next_idx: u64,
initial_end_idx: u64,
}
type QueueMetadata = RangeMetadata;
impl<'a> QueueIterator<'a> {
fn try_heal_head(&self, new_head: u64) -> Result<()> {
self.store.try_heal_range_head(
self.ns.meta,
&self.queue,
self.initial_next_idx,
new_head,
|store, queue| get_queue_meta(store, self.ns, queue),
|store, queue, meta| set_queue_meta(store, self.ns, queue, meta),
)
}
fn try_heal_tail(&self, new_tail: u64) -> Result<()> {
self.store.try_heal_range_tail(
self.ns.meta,
&self.queue,
self.initial_end_idx,
new_tail,
|store, queue| get_queue_meta(store, self.ns, queue),
|store, queue, meta| set_queue_meta(store, self.ns, queue, meta),
)
}
}
impl Iterator for QueueIterator<'_> {
type Item = Result<(usize, Vec<u8>)>;
fn next(&mut self) -> Option<Self::Item> {
while self.next_idx <= self.end_idx {
let idx = self.next_idx;
self.next_idx += 1;
if idx > self.initial_next_idx + 1000 {
let _ = self.try_heal_head(idx);
self.initial_next_idx = idx;
}
let key = make_queue_data_key(&self.queue, idx);
match self.store.get_ns(self.ns.data, &key) {
Ok(Some(v)) => return Some(Ok((idx as usize, v))),
Ok(None) => continue,
Err(e) => return Some(Err(e)),
}
}
None
}
}
impl DoubleEndedIterator for QueueIterator<'_> {
fn next_back(&mut self) -> Option<<Self as Iterator>::Item> {
while self.next_idx <= self.end_idx {
let idx = self.end_idx;
if self.end_idx == 0 {
self.next_idx = 1;
} else {
self.end_idx -= 1;
}
if idx + 1000 < self.initial_end_idx {
let _ = self.try_heal_tail(idx);
self.initial_end_idx = idx;
}
let key = make_queue_data_key(&self.queue, idx);
match self.store.get_ns(self.ns.data, &key) {
Ok(Some(v)) => return Some(Ok((idx as usize, v))),
Ok(None) => continue,
Err(e) => return Some(Err(e)),
}
}
None
}
}
impl CandyStore {
pub fn push_to_queue_head<B1: AsRef<[u8]> + ?Sized, B2: AsRef<[u8]> + ?Sized>(
&self,
queue_key: &B1,
val: &B2,
) -> Result<usize> {
self.queue_push_head_with_ns(QUEUE_NS, queue_key.as_ref(), val.as_ref())
.map(|idx| idx as usize)
}
pub fn push_to_queue_tail<B1: AsRef<[u8]> + ?Sized, B2: AsRef<[u8]> + ?Sized>(
&self,
queue_key: &B1,
val: &B2,
) -> Result<usize> {
self.queue_push_tail_with_ns(QUEUE_NS, queue_key.as_ref(), val.as_ref())
.map(|idx| idx as usize)
}
pub fn pop_queue_head<B: AsRef<[u8]> + ?Sized>(
&self,
queue_key: &B,
) -> Result<Option<Vec<u8>>> {
Ok(self
.queue_pop_head_with_ns(QUEUE_NS, queue_key.as_ref())?
.map(|(_, value)| value))
}
pub fn pop_queue_head_with_idx<B: AsRef<[u8]> + ?Sized>(
&self,
queue_key: &B,
) -> Result<Option<(usize, Vec<u8>)>> {
Ok(self
.queue_pop_head_with_ns(QUEUE_NS, queue_key.as_ref())?
.map(|(idx, value)| (idx as usize, value)))
}
pub fn pop_queue_tail<B: AsRef<[u8]> + ?Sized>(
&self,
queue_key: &B,
) -> Result<Option<Vec<u8>>> {
Ok(self
.queue_pop_tail_with_ns(QUEUE_NS, queue_key.as_ref())?
.map(|(_, value)| value))
}
pub fn pop_queue_tail_with_idx<B: AsRef<[u8]> + ?Sized>(
&self,
queue_key: &B,
) -> Result<Option<(usize, Vec<u8>)>> {
Ok(self
.queue_pop_tail_with_ns(QUEUE_NS, queue_key.as_ref())?
.map(|(idx, value)| (idx as usize, value)))
}
pub fn peek_queue_head<B: AsRef<[u8]> + ?Sized>(
&self,
queue_key: &B,
) -> Result<Option<Vec<u8>>> {
Ok(self
.queue_peek_head_with_ns(QUEUE_NS, queue_key.as_ref())?
.map(|(_, value)| value))
}
pub fn peek_queue_head_with_idx<B: AsRef<[u8]> + ?Sized>(
&self,
queue_key: &B,
) -> Result<Option<(usize, Vec<u8>)>> {
Ok(self
.queue_peek_head_with_ns(QUEUE_NS, queue_key.as_ref())?
.map(|(idx, value)| (idx as usize, value)))
}
pub fn peek_queue_tail<B: AsRef<[u8]> + ?Sized>(
&self,
queue_key: &B,
) -> Result<Option<Vec<u8>>> {
Ok(self
.queue_peek_tail_with_ns(QUEUE_NS, queue_key.as_ref())?
.map(|(_, value)| value))
}
pub fn peek_queue_tail_with_idx<B: AsRef<[u8]> + ?Sized>(
&self,
queue_key: &B,
) -> Result<Option<(usize, Vec<u8>)>> {
Ok(self
.queue_peek_tail_with_ns(QUEUE_NS, queue_key.as_ref())?
.map(|(idx, value)| (idx as usize, value)))
}
pub fn remove_from_queue<B: AsRef<[u8]> + ?Sized>(
&self,
queue_key: &B,
idx: usize,
) -> Result<Option<Vec<u8>>> {
self.queue_remove_with_ns(QUEUE_NS, queue_key.as_ref(), idx as u64)
}
pub fn discard_queue<B: AsRef<[u8]> + ?Sized>(&self, queue_key: &B) -> Result<bool> {
self.queue_discard_with_ns(QUEUE_NS, queue_key.as_ref())
}
pub fn extend_queue<B: AsRef<[u8]> + ?Sized>(
&self,
queue_key: &B,
items: impl IntoIterator<Item = impl AsRef<[u8]>>,
) -> Result<Range<usize>> {
let mut start = None;
let mut end = None;
for item in items {
let idx = self.push_to_queue_tail(queue_key, &item)?;
if start.is_none() {
start = Some(idx);
}
end = Some(idx + 1);
}
Ok(match (start, end) {
(Some(start), Some(end)) => start..end,
_ => {
let range = self.queue_range(queue_key)?;
range.start..range.start
}
})
}
pub fn queue_len<B: AsRef<[u8]> + ?Sized>(&self, queue_key: &B) -> Result<usize> {
Ok(self.queue_len_with_ns(QUEUE_NS, queue_key.as_ref())? as usize)
}
pub fn queue_range<B: AsRef<[u8]> + ?Sized>(&self, queue_key: &B) -> Result<Range<usize>> {
self.queue_range_with_ns(QUEUE_NS, queue_key.as_ref())
}
pub fn iter_queue<'a, B: AsRef<[u8]> + ?Sized>(&'a self, queue_key: &B) -> QueueIterator<'a> {
self.queue_iter_with_ns(QUEUE_NS, queue_key.as_ref())
}
pub fn set_big<B1: AsRef<[u8]> + ?Sized, B2: AsRef<[u8]> + ?Sized>(
&self,
key: &B1,
value: &B2,
) -> Result<bool> {
self.queue_set_big_with_ns(BIG_NS, key.as_ref(), value.as_ref())
}
pub fn get_big<B: AsRef<[u8]> + ?Sized>(&self, key: &B) -> Result<Option<Vec<u8>>> {
self.queue_get_big_with_ns(BIG_NS, key.as_ref())
}
pub fn remove_big<B: AsRef<[u8]> + ?Sized>(&self, key: &B) -> Result<bool> {
self.queue_discard_with_ns(BIG_NS, key.as_ref())
}
pub(super) fn queue_push_tail_with_ns(
&self,
ns: QueueNamespaces,
queue: &[u8],
value: &[u8],
) -> Result<u64> {
let _lock = self.list_write_guard(ns.meta, queue);
self._queue_push_tail_with_ns(ns, queue, value)
}
fn _queue_push_tail_with_ns(
&self,
ns: QueueNamespaces,
queue: &[u8],
value: &[u8],
) -> Result<u64> {
let mut meta = get_queue_meta(self, ns, queue)?;
let new_tail = meta.tail + 1;
let key = make_queue_data_key(queue, new_tail);
self.set_ns(ns.data, &key, value)?;
meta.tail = new_tail;
meta.count += 1;
if meta.head > meta.tail {
meta.head = new_tail;
}
set_queue_meta(self, ns, queue, meta)?;
Ok(new_tail)
}
pub(super) fn queue_push_head_with_ns(
&self,
ns: QueueNamespaces,
queue: &[u8],
value: &[u8],
) -> Result<u64> {
let _lock = self.list_write_guard(ns.meta, queue);
let mut meta = get_queue_meta(self, ns, queue)?;
let new_head = meta.head - 1;
let key = make_queue_data_key(queue, new_head);
self.set_ns(ns.data, &key, value)?;
meta.head = new_head;
meta.count += 1;
if meta.tail < meta.head {
meta.tail = new_head;
}
set_queue_meta(self, ns, queue, meta)?;
Ok(new_head)
}
pub(super) fn queue_pop_head_with_ns(
&self,
ns: QueueNamespaces,
queue: &[u8],
) -> Result<Option<(u64, Vec<u8>)>> {
let _lock = self.list_write_guard(ns.meta, queue);
let mut meta = get_queue_meta(self, ns, queue)?;
loop {
if meta.head > meta.tail {
return Ok(None);
}
let idx = meta.head;
let key = make_queue_data_key(queue, idx);
let value = self.remove_ns(ns.data, &key)?;
meta.head += 1;
if let Some(value) = value {
meta.count = meta.count.saturating_sub(1);
if meta.head > meta.tail {
meta = QueueMetadata::new();
}
set_queue_meta(self, ns, queue, meta)?;
return Ok(Some((idx, value)));
}
set_queue_meta(self, ns, queue, meta)?;
}
}
pub(super) fn queue_pop_tail_with_ns(
&self,
ns: QueueNamespaces,
queue: &[u8],
) -> Result<Option<(u64, Vec<u8>)>> {
let _lock = self.list_write_guard(ns.meta, queue);
let mut meta = get_queue_meta(self, ns, queue)?;
loop {
if meta.head > meta.tail {
return Ok(None);
}
let idx = meta.tail;
let key = make_queue_data_key(queue, idx);
let value = self.remove_ns(ns.data, &key)?;
meta.tail = meta.tail.saturating_sub(1);
if let Some(value) = value {
meta.count = meta.count.saturating_sub(1);
if meta.head > meta.tail {
meta = QueueMetadata::new();
}
set_queue_meta(self, ns, queue, meta)?;
return Ok(Some((idx, value)));
}
set_queue_meta(self, ns, queue, meta)?;
}
}
pub(super) fn queue_peek_head_with_ns(
&self,
ns: QueueNamespaces,
queue: &[u8],
) -> Result<Option<(u64, Vec<u8>)>> {
let _lock = self.list_read_guard(ns.meta, queue);
let meta = get_queue_meta(self, ns, queue)?;
if meta.head > meta.tail {
return Ok(None);
}
for idx in meta.head..=meta.tail {
let key = make_queue_data_key(queue, idx);
if let Some(value) = self.get_ns(ns.data, &key)? {
return Ok(Some((idx, value)));
}
}
Ok(None)
}
pub(super) fn queue_peek_tail_with_ns(
&self,
ns: QueueNamespaces,
queue: &[u8],
) -> Result<Option<(u64, Vec<u8>)>> {
let _lock = self.list_read_guard(ns.meta, queue);
let meta = get_queue_meta(self, ns, queue)?;
if meta.head > meta.tail {
return Ok(None);
}
for idx in (meta.head..=meta.tail).rev() {
let key = make_queue_data_key(queue, idx);
if let Some(value) = self.get_ns(ns.data, &key)? {
return Ok(Some((idx, value)));
}
}
Ok(None)
}
pub(super) fn queue_len_with_ns(&self, ns: QueueNamespaces, queue: &[u8]) -> Result<u64> {
Ok(get_queue_meta(self, ns, queue)?.count)
}
pub(super) fn queue_discard_with_ns(&self, ns: QueueNamespaces, queue: &[u8]) -> Result<bool> {
let _lock = self.list_write_guard(ns.meta, queue);
self._queue_discard_with_ns(ns, queue)
}
fn _queue_discard_with_ns(&self, ns: QueueNamespaces, queue: &[u8]) -> Result<bool> {
let mut meta = get_queue_meta(self, ns, queue)?;
let had_items = meta.head <= meta.tail;
while meta.head <= meta.tail {
let key = make_queue_data_key(queue, meta.head);
_ = self.remove_ns(ns.data, &key)?;
meta.head += 1;
}
self.remove_ns(ns.meta, queue)?;
Ok(had_items)
}
fn queue_remove_with_ns(
&self,
ns: QueueNamespaces,
queue: &[u8],
idx: u64,
) -> Result<Option<Vec<u8>>> {
let _lock = self.list_write_guard(ns.meta, queue);
let mut meta = get_queue_meta(self, ns, queue)?;
let key = make_queue_data_key(queue, idx);
let removed = match self.remove_ns(ns.data, &key)? {
Some(value) => value,
None => return Ok(None),
};
meta.count = meta.count.saturating_sub(1);
if idx == meta.head {
meta.head += 1;
}
if meta.tail == idx {
meta.tail = meta.tail.saturating_sub(1);
}
if meta.head > meta.tail {
meta = QueueMetadata::new();
}
set_queue_meta(self, ns, queue, meta)?;
Ok(Some(removed))
}
pub(super) fn queue_iter_with_ns<'a>(
&'a self,
ns: QueueNamespaces,
queue: &[u8],
) -> QueueIterator<'a> {
let meta = get_queue_meta(self, ns, queue).unwrap_or_else(|_| QueueMetadata::new());
QueueIterator {
store: self,
queue: queue.to_vec(),
ns,
next_idx: meta.head,
end_idx: meta.tail,
initial_next_idx: meta.head,
initial_end_idx: meta.tail,
}
}
pub(super) fn queue_range_with_ns(
&self,
ns: QueueNamespaces,
queue: &[u8],
) -> Result<Range<usize>> {
let meta = get_queue_meta(self, ns, queue)?;
if meta.count == 0 || meta.head > meta.tail {
return Ok(0..0);
}
Ok(meta.head as usize..meta.tail.saturating_add(1) as usize)
}
pub(super) fn queue_set_big_with_ns(
&self,
ns: QueueNamespaces,
key: &[u8],
value: &[u8],
) -> Result<bool> {
let _lock = self.list_write_guard(ns.meta, key);
let existed = self._queue_discard_with_ns(ns, key)?;
let max_chunk_len = self.max_big_chunk_len(key)?;
for chunk in value.chunks(max_chunk_len) {
self._queue_push_tail_with_ns(ns, key, chunk)?;
}
self._queue_push_tail_with_ns(ns, key, &value.len().to_le_bytes())?;
Ok(existed)
}
pub(super) fn queue_get_big_with_ns(
&self,
ns: QueueNamespaces,
key: &[u8],
) -> Result<Option<Vec<u8>>> {
let _lock = self.list_read_guard(ns.meta, key);
let meta = get_queue_meta(self, ns, key)?;
let expected_chunks = meta.count;
if expected_chunks == 0 {
return Ok(None);
}
let mut collected = Vec::new();
let mut seen = 0u64;
for idx in meta.head..=meta.tail {
let item_key = make_queue_data_key(key, idx);
let Some(chunk) = self.get_ns(ns.data, &item_key)? else {
continue;
};
seen += 1;
if seen == expected_chunks && chunk.len() == size_of::<usize>() {
let recorded_len = usize::from_le_bytes(chunk.as_slice().try_into().unwrap());
if recorded_len == collected.len() {
return Ok(Some(collected));
}
return Ok(None);
}
collected.extend_from_slice(&chunk);
if seen == expected_chunks {
return Ok(None);
}
}
Ok(None)
}
fn max_big_chunk_len(&self, key: &[u8]) -> Result<usize> {
let data_key_len = make_queue_data_key(key, 0).len();
if aligned_data_entry_size(data_key_len, size_of::<usize>()) as usize
> self.inner.config.max_data_file_size as usize
{
return Err(Error::PayloadTooLarge(aligned_data_entry_size(
data_key_len,
size_of::<usize>(),
) as usize));
}
let mut max_chunk_len = MAX_USER_VALUE_SIZE;
while max_chunk_len > 0
&& aligned_data_entry_size(data_key_len, max_chunk_len) as usize
> self.inner.config.max_data_file_size as usize
{
max_chunk_len -= 1;
}
if max_chunk_len == 0 {
return Err(Error::PayloadTooLarge(
aligned_data_entry_size(data_key_len, 1) as usize,
));
}
Ok(max_chunk_len)
}
}
fn get_queue_meta(store: &CandyStore, ns: QueueNamespaces, queue: &[u8]) -> Result<QueueMetadata> {
if let Some(value) = store.get_ns(ns.meta, queue)?
&& let Some(meta) = QueueMetadata::from_bytes(&value)
{
return Ok(meta);
}
Ok(QueueMetadata::new())
}
fn set_queue_meta(
store: &CandyStore,
ns: QueueNamespaces,
queue: &[u8],
meta: QueueMetadata,
) -> Result<()> {
store.set_ns(ns.meta, queue, &meta.to_bytes())?;
Ok(())
}
fn hash_queue_key(queue: &[u8]) -> u64 {
let mut hasher = SipHasher13::new_with_keys(0xb1ccc559a9924eaa, 0x1b1a682059c2d599);
hasher.write(queue);
hasher.finish()
}
fn make_queue_data_key(queue: &[u8], seq: u64) -> [u8; 16] {
let hash = hash_queue_key(queue);
let mut key = [0u8; 16];
key[..8].copy_from_slice(&hash.to_le_bytes());
key[8..].copy_from_slice(&seq.to_be_bytes());
key
}