pub mod conf;
pub mod meta;
pub use conf::PosSpec;
pub use meta::ListMeta;
use std::str;
use crate::db::WeDb;
use crate::error::{Error, Result};
use crate::key_composer::KeyComposer;
const ERR_NO_SUCH_KEY: &str = "ERR no such key";
const ERR_INDEX_OUT_OF_RANGE: &str = "ERR index out of range";
const ERR_RANK_ZERO: &str = "ERR RANK can't be zero: must be 1, 2, 3, ... or -1, -2, -3...";
const HEX_DIGITS: &[u8; 16] = b"0123456789abcdef";
#[inline(always)]
fn write_hex_u64(buf: &mut [u8], val: u64) {
for i in 0..16 {
let shift = (15 - i) * 4;
buf[i] = HEX_DIGITS[((val >> shift) & 0x0f) as usize];
}
}
#[derive(Debug, Clone)]
struct ListItemKeyComposer {
buf: Vec<u8>,
prefix_len: usize,
}
impl ListItemKeyComposer {
#[inline]
fn new(kc: &KeyComposer<'_>, key: &str) -> Self {
let mut buf = kc.list_prefix(key);
let prefix_len = buf.len();
buf.resize(prefix_len + 16, 0);
Self { buf, prefix_len }
}
#[inline(always)]
fn key_for_idx(&mut self, idx: u64) -> &[u8] {
write_hex_u64(&mut self.buf[self.prefix_len..self.prefix_len + 16], idx);
&self.buf
}
}
impl WeDb {
#[inline]
fn load_valid_list_meta(&self, meta_k: &str, now_ms: u64) -> Result<Option<ListMeta>> {
match self.meta_ks.get(meta_k.as_bytes())? {
Some(m_bytes) => Ok(ListMeta::decode(&m_bytes).filter(|m| !m.is_expired(now_ms))),
None => Ok(None),
}
}
pub fn lpush<K: AsRef<[u8]>, V: AsRef<[u8]>>(&self, key: K, values: &[V]) -> Result<u64> {
self.list_push_internal(key.as_ref(), values, true, true)
}
pub fn rpush<K: AsRef<[u8]>, V: AsRef<[u8]>>(&self, key: K, values: &[V]) -> Result<u64> {
self.list_push_internal(key.as_ref(), values, true, false)
}
pub fn lpushx<K: AsRef<[u8]>, V: AsRef<[u8]>>(&self, key: K, values: &[V]) -> Result<u64> {
self.list_push_internal(key.as_ref(), values, false, true)
}
pub fn rpushx<K: AsRef<[u8]>, V: AsRef<[u8]>>(&self, key: K, values: &[V]) -> Result<u64> {
self.list_push_internal(key.as_ref(), values, false, false)
}
fn list_push_internal<V: AsRef<[u8]>>(
&self,
key: &[u8],
values: &[V],
create_if_missing: bool,
left: bool,
) -> Result<u64> {
if values.is_empty() {
return self.llen(key);
}
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key).unwrap_or("");
let meta_k = kc.list_meta(k_str);
let now_ms = ts_::sec() * 1000;
let raw_meta_opt = match self.meta_ks.get(meta_k.as_bytes())? {
Some(m_bytes) => ListMeta::decode(&m_bytes),
None => None,
};
let mut batch = self.db.batch();
let mut composer = ListItemKeyComposer::new(&kc, k_str);
let meta_opt = match raw_meta_opt {
Some(m) => {
if m.is_expired(now_ms) {
for i in 0..m.base.size {
let old_idx = m.head.wrapping_add(i);
batch.remove(&self.data_ks, composer.key_for_idx(old_idx));
}
None
} else {
Some(m)
}
}
None => None,
};
if meta_opt.is_none() && !create_if_missing {
return Ok(0);
}
let mut meta = meta_opt.unwrap_or_else(|| ListMeta::new_with_version(0));
if left {
for val in values {
meta.head = meta.head.wrapping_sub(1);
meta.base.size = meta.base.size.saturating_add(1);
let item_k = composer.key_for_idx(meta.head);
batch.insert(&self.data_ks, item_k, val.as_ref());
}
} else {
for val in values {
let item_k = composer.key_for_idx(meta.tail);
batch.insert(&self.data_ks, item_k, val.as_ref());
meta.tail = meta.tail.wrapping_add(1);
meta.base.size = meta.base.size.saturating_add(1);
}
}
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
batch.commit()?;
Ok(meta.base.size)
}
pub fn lpop<K: AsRef<[u8]>>(&self, key: K, count: usize) -> Result<Vec<Vec<u8>>> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.list_meta(k_str);
let now_ms = ts_::sec() * 1000;
let mut meta = match self.load_valid_list_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(Vec::new()),
};
if meta.base.size == 0 || count == 0 {
return Ok(Vec::new());
}
let num_pop = (count as u64).min(meta.base.size);
let mut results = Vec::with_capacity(num_pop as usize);
let mut batch = self.db.batch();
let mut composer = ListItemKeyComposer::new(&kc, k_str);
for _ in 0..num_pop {
let item_k = composer.key_for_idx(meta.head);
if let Some(val) = self.data_ks.get(item_k)? {
results.push(val.to_vec());
}
batch.remove(&self.data_ks, item_k);
meta.head = meta.head.wrapping_add(1);
meta.base.size = meta.base.size.saturating_sub(1);
}
if meta.base.size == 0 {
batch.remove(&self.meta_ks, meta_k.as_bytes());
} else {
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
}
batch.commit()?;
Ok(results)
}
pub fn rpop<K: AsRef<[u8]>>(&self, key: K, count: usize) -> Result<Vec<Vec<u8>>> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.list_meta(k_str);
let now_ms = ts_::sec() * 1000;
let mut meta = match self.load_valid_list_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(Vec::new()),
};
if meta.base.size == 0 || count == 0 {
return Ok(Vec::new());
}
let num_pop = (count as u64).min(meta.base.size);
let mut results = Vec::with_capacity(num_pop as usize);
let mut batch = self.db.batch();
let mut composer = ListItemKeyComposer::new(&kc, k_str);
for _ in 0..num_pop {
meta.tail = meta.tail.wrapping_sub(1);
meta.base.size = meta.base.size.saturating_sub(1);
let item_k = composer.key_for_idx(meta.tail);
if let Some(val) = self.data_ks.get(item_k)? {
results.push(val.to_vec());
}
batch.remove(&self.data_ks, item_k);
}
if meta.base.size == 0 {
batch.remove(&self.meta_ks, meta_k.as_bytes());
} else {
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
}
batch.commit()?;
Ok(results)
}
pub fn llen<K: AsRef<[u8]>>(&self, key: K) -> Result<u64> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.list_meta(k_str);
let now_ms = ts_::sec() * 1000;
Ok(self
.load_valid_list_meta(&meta_k, now_ms)?
.map_or(0, |m| m.base.size))
}
pub fn lrange<K: AsRef<[u8]>>(&self, key: K, start: i64, stop: i64) -> Result<Vec<Vec<u8>>> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.list_meta(k_str);
let now_ms = ts_::sec() * 1000;
let meta = match self.load_valid_list_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(Vec::new()),
};
if meta.base.size == 0 {
return Ok(Vec::new());
}
let len = meta.base.size as i64;
let mut s = if start < 0 {
len.checked_add(start).unwrap_or(i64::MIN)
} else {
start
};
let mut e = if stop < 0 {
len.checked_add(stop).unwrap_or(i64::MIN)
} else {
stop
};
if s >= len || e < 0 || s > e {
return Ok(Vec::new());
}
if s < 0 {
s = 0;
}
if e >= len {
e = len - 1;
}
let num_elems = (e - s + 1) as usize;
let mut results = Vec::with_capacity(num_elems);
let mut composer = ListItemKeyComposer::new(&kc, k_str);
for idx in s..=e {
let actual_idx = meta.head.wrapping_add(idx as u64);
let item_k = composer.key_for_idx(actual_idx);
if let Some(val) = self.data_ks.get(item_k)? {
results.push(val.to_vec());
}
}
Ok(results)
}
pub fn lindex<K: AsRef<[u8]>>(&self, key: K, index: i64) -> Result<Option<Vec<u8>>> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.list_meta(k_str);
let now_ms = ts_::sec() * 1000;
let meta = match self.load_valid_list_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(None),
};
if meta.base.size == 0 {
return Ok(None);
}
let len = meta.base.size as i64;
let actual_offset = if index < 0 {
len.checked_add(index).unwrap_or(i64::MIN)
} else {
index
};
if actual_offset < 0 || actual_offset >= len {
return Ok(None);
}
let mut composer = ListItemKeyComposer::new(&kc, k_str);
let actual_idx = meta.head.wrapping_add(actual_offset as u64);
let item_k = composer.key_for_idx(actual_idx);
let val = self.data_ks.get(item_k)?;
Ok(val.map(|v| v.to_vec()))
}
pub fn lset<K: AsRef<[u8]>, V: AsRef<[u8]>>(&self, key: K, index: i64, value: V) -> Result<()> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.list_meta(k_str);
let now_ms = ts_::sec() * 1000;
let meta = match self.load_valid_list_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Err(Error::invalid_data(ERR_NO_SUCH_KEY)),
};
let len = meta.base.size as i64;
let actual_offset = if index < 0 {
len.checked_add(index).unwrap_or(i64::MIN)
} else {
index
};
if actual_offset < 0 || actual_offset >= len {
return Err(Error::invalid_data(ERR_INDEX_OUT_OF_RANGE));
}
let mut composer = ListItemKeyComposer::new(&kc, k_str);
let actual_idx = meta.head.wrapping_add(actual_offset as u64);
let item_k = composer.key_for_idx(actual_idx);
let val_bytes = value.as_ref();
if let Some(existing) = self.data_ks.get(item_k)?
&& existing.as_ref() == val_bytes
{
return Ok(());
}
self.data_ks.insert(item_k, val_bytes)?;
Ok(())
}
pub fn linsert<K: AsRef<[u8]>, P: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
key: K,
before: bool,
pivot: P,
elem: V,
) -> Result<i64> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.list_meta(k_str);
let now_ms = ts_::sec() * 1000;
let mut meta = match self.load_valid_list_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(0),
};
if meta.base.size == 0 {
return Ok(0);
}
let len = meta.base.size as usize;
let pivot_bytes = pivot.as_ref();
let mut composer = ListItemKeyComposer::new(&kc, k_str);
let pivot_offset = (0..len)
.find_map(|offset| {
let idx = meta.head.wrapping_add(offset as u64);
let item_k = composer.key_for_idx(idx);
match self.data_ks.get(item_k) {
Ok(Some(v)) if v.as_ref() == pivot_bytes => Some(Ok(offset)),
Ok(_) => None,
Err(e) => Some(Err(e)),
}
})
.transpose()?;
let p_off = match pivot_offset {
Some(off) => off,
None => return Ok(-1),
};
let target_insert_offset = if before { p_off } else { p_off + 1 };
let left_cost = target_insert_offset;
let right_cost = len - target_insert_offset;
let mut batch = self.db.batch();
if left_cost <= right_cost {
for offset in 0..target_insert_offset {
let from_idx = meta.head.wrapping_add(offset as u64);
let to_idx = from_idx.wrapping_sub(1);
if let Some(v) = self.data_ks.get(composer.key_for_idx(from_idx))? {
let val_bytes = v.to_vec();
batch.insert(&self.data_ks, composer.key_for_idx(to_idx), val_bytes);
}
}
let insert_idx = meta
.head
.wrapping_add(target_insert_offset as u64)
.wrapping_sub(1);
batch.insert(
&self.data_ks,
composer.key_for_idx(insert_idx),
elem.as_ref(),
);
meta.head = meta.head.wrapping_sub(1);
} else {
for offset in (target_insert_offset..len).rev() {
let from_idx = meta.head.wrapping_add(offset as u64);
let to_idx = from_idx.wrapping_add(1);
if let Some(v) = self.data_ks.get(composer.key_for_idx(from_idx))? {
let val_bytes = v.to_vec();
batch.insert(&self.data_ks, composer.key_for_idx(to_idx), val_bytes);
}
}
let insert_idx = meta.head.wrapping_add(target_insert_offset as u64);
batch.insert(
&self.data_ks,
composer.key_for_idx(insert_idx),
elem.as_ref(),
);
meta.tail = meta.tail.wrapping_add(1);
}
meta.base.size += 1;
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
batch.commit()?;
Ok(meta.base.size as i64)
}
pub fn lrem<K: AsRef<[u8]>, V: AsRef<[u8]>>(&self, key: K, count: i64, elem: V) -> Result<u64> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.list_meta(k_str);
let now_ms = ts_::sec() * 1000;
let mut meta = match self.load_valid_list_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(0),
};
if meta.base.size == 0 {
return Ok(0);
}
let len = meta.base.size as usize;
let target_del_limit = if count == 0 {
usize::MAX
} else {
count.unsigned_abs() as usize
};
let elem_bytes = elem.as_ref();
let mut to_delete_offsets = Vec::new();
let mut composer = ListItemKeyComposer::new(&kc, k_str);
if count >= 0 {
for offset in 0..len {
let idx = meta.head.wrapping_add(offset as u64);
let item_k = composer.key_for_idx(idx);
if let Some(v) = self.data_ks.get(item_k)?
&& v.as_ref() == elem_bytes
{
to_delete_offsets.push(offset);
if to_delete_offsets.len() >= target_del_limit {
break;
}
}
}
} else {
for step in 0..len {
let offset = len - 1 - step;
let idx = meta.head.wrapping_add(offset as u64);
let item_k = composer.key_for_idx(idx);
if let Some(v) = self.data_ks.get(item_k)?
&& v.as_ref() == elem_bytes
{
to_delete_offsets.push(offset);
if to_delete_offsets.len() >= target_del_limit {
break;
}
}
}
to_delete_offsets.reverse();
}
if to_delete_offsets.is_empty() {
return Ok(0);
}
let del_cnt = to_delete_offsets.len();
let mut batch = self.db.batch();
if del_cnt == len {
for offset in 0..len {
let idx = meta.head.wrapping_add(offset as u64);
batch.remove(&self.data_ks, composer.key_for_idx(idx));
}
batch.remove(&self.meta_ks, meta_k.as_bytes());
batch.commit()?;
return Ok(del_cnt as u64);
}
let min_del_offset = to_delete_offsets[0];
let max_del_offset = to_delete_offsets[del_cnt - 1];
let left_cost = max_del_offset;
let right_cost = len - 1 - min_del_offset;
if left_cost <= right_cost {
let mut target_offset = max_del_offset;
let mut del_idx_cursor = del_cnt;
for offset in (0..=max_del_offset).rev() {
if del_idx_cursor > 0 && to_delete_offsets[del_idx_cursor - 1] == offset {
del_idx_cursor -= 1;
} else {
if target_offset != offset {
let from_idx = meta.head.wrapping_add(offset as u64);
let to_idx = meta.head.wrapping_add(target_offset as u64);
if let Some(v) = self.data_ks.get(composer.key_for_idx(from_idx))? {
let val_bytes = v.to_vec();
batch.insert(&self.data_ks, composer.key_for_idx(to_idx), val_bytes);
}
}
target_offset = target_offset.saturating_sub(1);
}
}
for offset in 0..del_cnt {
let idx = meta.head.wrapping_add(offset as u64);
batch.remove(&self.data_ks, composer.key_for_idx(idx));
}
meta.head = meta.head.wrapping_add(del_cnt as u64);
} else {
let mut target_offset = min_del_offset;
let mut del_idx_cursor = 0;
for offset in min_del_offset..len {
if del_idx_cursor < del_cnt && to_delete_offsets[del_idx_cursor] == offset {
del_idx_cursor += 1;
} else {
if target_offset != offset {
let from_idx = meta.head.wrapping_add(offset as u64);
let to_idx = meta.head.wrapping_add(target_offset as u64);
if let Some(v) = self.data_ks.get(composer.key_for_idx(from_idx))? {
let val_bytes = v.to_vec();
batch.insert(&self.data_ks, composer.key_for_idx(to_idx), val_bytes);
}
}
target_offset += 1;
}
}
for offset in (len - del_cnt)..len {
let idx = meta.head.wrapping_add(offset as u64);
batch.remove(&self.data_ks, composer.key_for_idx(idx));
}
meta.tail = meta.tail.wrapping_sub(del_cnt as u64);
}
meta.base.size -= del_cnt as u64;
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
batch.commit()?;
Ok(del_cnt as u64)
}
pub fn ltrim<K: AsRef<[u8]>>(&self, key: K, start: i64, stop: i64) -> Result<()> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.list_meta(k_str);
let now_ms = ts_::sec() * 1000;
let mut meta = match self.load_valid_list_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(()),
};
if meta.base.size == 0 {
return Ok(());
}
let len = meta.base.size as i64;
let mut s = if start < 0 {
len.checked_add(start).unwrap_or(i64::MIN)
} else {
start
};
let mut e = if stop < 0 {
len.checked_add(stop).unwrap_or(i64::MIN)
} else {
stop
};
let mut batch = self.db.batch();
let mut composer = ListItemKeyComposer::new(&kc, k_str);
if s > e || s >= len || e < 0 {
for i in 0..meta.base.size {
let idx = meta.head.wrapping_add(i);
batch.remove(&self.data_ks, composer.key_for_idx(idx));
}
batch.remove(&self.meta_ks, meta_k.as_bytes());
batch.commit()?;
return Ok(());
}
if s < 0 {
s = 0;
}
if e >= len {
e = len - 1;
}
for i in 0..s {
let idx = meta.head.wrapping_add(i as u64);
batch.remove(&self.data_ks, composer.key_for_idx(idx));
}
for i in (e + 1)..len {
let idx = meta.head.wrapping_add(i as u64);
batch.remove(&self.data_ks, composer.key_for_idx(idx));
}
let new_size = (e - s + 1) as u64;
meta.head = meta.head.wrapping_add(s as u64);
meta.tail = meta.head.wrapping_add(new_size);
meta.base.size = new_size;
if new_size == 0 {
batch.remove(&self.meta_ks, meta_k.as_bytes());
} else {
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
}
batch.commit()?;
Ok(())
}
pub fn lmove<S: AsRef<[u8]>, D: AsRef<[u8]>>(
&self,
src: S,
dst: D,
src_left: bool,
dst_left: bool,
) -> Result<Option<Vec<u8>>> {
let src_bytes = src.as_ref();
let dst_bytes = dst.as_ref();
let now_ms = ts_::sec() * 1000;
let kc = KeyComposer::new("default");
let src_str = str::from_utf8(src_bytes).unwrap_or("");
let dst_str = str::from_utf8(dst_bytes).unwrap_or("");
let src_meta_k = kc.list_meta(src_str);
let mut src_meta = match self.load_valid_list_meta(&src_meta_k, now_ms)? {
Some(m) => m,
None => return Ok(None),
};
if src_meta.base.size == 0 {
return Ok(None);
}
if src_bytes == dst_bytes {
let mut composer = ListItemKeyComposer::new(&kc, src_str);
let curr_idx = if src_left {
src_meta.head
} else {
src_meta.tail.wrapping_sub(1)
};
let elem = match self.data_ks.get(composer.key_for_idx(curr_idx))? {
Some(v) => v.to_vec(),
None => return Ok(None),
};
if src_left == dst_left || src_meta.base.size == 1 {
return Ok(Some(elem));
}
let mut batch = self.db.batch();
batch.remove(&self.data_ks, composer.key_for_idx(curr_idx));
if src_left {
let new_tail_idx = src_meta.tail;
batch.insert(&self.data_ks, composer.key_for_idx(new_tail_idx), &elem);
src_meta.head = src_meta.head.wrapping_add(1);
src_meta.tail = src_meta.tail.wrapping_add(1);
} else {
let new_head_idx = src_meta.head.wrapping_sub(1);
batch.insert(&self.data_ks, composer.key_for_idx(new_head_idx), &elem);
src_meta.head = src_meta.head.wrapping_sub(1);
src_meta.tail = src_meta.tail.wrapping_sub(1);
}
batch.insert(&self.meta_ks, src_meta_k.as_bytes(), src_meta.encode());
batch.commit()?;
return Ok(Some(elem));
}
let dst_meta_k = kc.list_meta(dst_str);
let raw_dst_meta_opt = match self.meta_ks.get(dst_meta_k.as_bytes())? {
Some(m_bytes) => ListMeta::decode(&m_bytes),
None => None,
};
let mut batch = self.db.batch();
let mut src_composer = ListItemKeyComposer::new(&kc, src_str);
let mut dst_composer = ListItemKeyComposer::new(&kc, dst_str);
let dst_meta_opt = match raw_dst_meta_opt {
Some(m) => {
if m.is_expired(now_ms) {
for i in 0..m.base.size {
let old_idx = m.head.wrapping_add(i);
batch.remove(&self.data_ks, dst_composer.key_for_idx(old_idx));
}
None
} else {
Some(m)
}
}
None => None,
};
let mut dst_meta = dst_meta_opt.unwrap_or_else(|| ListMeta::new_with_version(0));
let src_idx = if src_left {
src_meta.head
} else {
src_meta.tail.wrapping_sub(1)
};
let elem = match self.data_ks.get(src_composer.key_for_idx(src_idx))? {
Some(v) => v.to_vec(),
None => return Ok(None),
};
batch.remove(&self.data_ks, src_composer.key_for_idx(src_idx));
if src_left {
src_meta.head = src_meta.head.wrapping_add(1);
} else {
src_meta.tail = src_meta.tail.wrapping_sub(1);
}
src_meta.base.size -= 1;
if src_meta.base.size == 0 {
batch.remove(&self.meta_ks, src_meta_k.as_bytes());
} else {
batch.insert(&self.meta_ks, src_meta_k.as_bytes(), src_meta.encode());
}
let dst_idx = if dst_left {
let idx = dst_meta.head.wrapping_sub(1);
dst_meta.head = idx;
idx
} else {
let idx = dst_meta.tail;
dst_meta.tail = idx.wrapping_add(1);
idx
};
batch.insert(&self.data_ks, dst_composer.key_for_idx(dst_idx), &elem);
dst_meta.base.size += 1;
batch.insert(&self.meta_ks, dst_meta_k.as_bytes(), dst_meta.encode());
batch.commit()?;
Ok(Some(elem))
}
#[inline]
pub fn rpoplpush<S: AsRef<[u8]>, D: AsRef<[u8]>>(
&self,
source: S,
destination: D,
) -> Result<Option<Vec<u8>>> {
self.lmove(source, destination, false, true)
}
pub fn lpos<K: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
key: K,
elem: V,
spec: PosSpec,
) -> Result<Vec<i64>> {
if spec.rank == 0 {
return Err(Error::invalid_data(ERR_RANK_ZERO));
}
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.list_meta(k_str);
let now_ms = ts_::sec() * 1000;
let meta = match self.load_valid_list_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(Vec::new()),
};
if meta.base.size == 0 {
return Ok(Vec::new());
}
let len = meta.base.size as usize;
let reversed = spec.rank < 0;
let target_rank = spec.rank.unsigned_abs() as usize;
let limit = spec
.max_len
.map(|m| if m == 0 { len } else { m.min(len) })
.unwrap_or(len);
let elem_bytes = elem.as_ref();
let mut matches = Vec::new();
let mut match_count = 0;
let mut composer = ListItemKeyComposer::new(&kc, k_str);
for i in 0..limit {
let offset = if !reversed { i } else { len - 1 - i };
let idx = meta.head.wrapping_add(offset as u64);
let item_k = composer.key_for_idx(idx);
if let Some(v) = self.data_ks.get(item_k)?
&& v.as_ref() == elem_bytes
{
match_count += 1;
if match_count >= target_rank {
matches.push(offset as i64);
if let Some(c) = spec.count {
if c > 0 && matches.len() >= c {
break;
}
} else {
break;
}
}
}
}
Ok(matches)
}
pub fn lexpireat<K: AsRef<[u8]>>(&self, key: K, expire_at_ms: u64) -> Result<bool> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.list_meta(k_str);
let now_ms = ts_::sec() * 1000;
let mut meta = match self.load_valid_list_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(false),
};
meta.base.expire_at = expire_at_ms;
self.meta_ks.insert(meta_k.as_bytes(), meta.encode())?;
Ok(true)
}
pub fn lexpire<K: AsRef<[u8]>>(&self, key: K, seconds: u64) -> Result<bool> {
let expire_at_ms = (ts_::sec() + seconds) * 1000;
self.lexpireat(key, expire_at_ms)
}
pub fn lttl<K: AsRef<[u8]>>(&self, key: K) -> Result<i64> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.list_meta(k_str);
let now_ms = ts_::sec() * 1000;
match self.load_valid_list_meta(&meta_k, now_ms)? {
Some(meta) => Ok(meta.ttl(now_ms)),
None => Ok(-2),
}
}
pub fn lpersist<K: AsRef<[u8]>>(&self, key: K) -> Result<bool> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.list_meta(k_str);
let now_ms = ts_::sec() * 1000;
let mut meta = match self.load_valid_list_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(false),
};
if meta.base.expire_at == 0 {
return Ok(false);
}
meta.base.expire_at = 0;
self.meta_ks.insert(meta_k.as_bytes(), meta.encode())?;
Ok(true)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_list_basic_operations() -> Result<()> {
let temp_dir = tempfile::tempdir().map_err(Error::io)?;
let db = WeDb::open(temp_dir.path())?;
assert_eq!(db.lpush("mylist", &[b"b", b"a"])?, 2); assert_eq!(db.rpush("mylist", &[b"c", b"d"])?, 4); assert_eq!(db.llen("mylist")?, 4);
let range = db.lrange("mylist", 0, -1)?;
assert_eq!(
range,
vec![b"a".to_vec(), b"b".to_vec(), b"c".to_vec(), b"d".to_vec()]
);
assert_eq!(db.lindex("mylist", 0)?, Some(b"a".to_vec()));
assert_eq!(db.lindex("mylist", -1)?, Some(b"d".to_vec()));
assert_eq!(db.lindex("mylist", i64::MIN)?, None);
assert_eq!(db.lindex("mylist", i64::MAX)?, None);
db.lset("mylist", 1, b"B")?;
assert_eq!(db.lindex("mylist", 1)?, Some(b"B".to_vec()));
db.lset("mylist", 1, b"B")?;
assert_eq!(db.lindex("mylist", 1)?, Some(b"B".to_vec()));
assert_eq!(db.linsert("mylist", true, b"c", b"X")?, 5);
assert_eq!(
db.lrange("mylist", 0, -1)?,
vec![
b"a".to_vec(),
b"B".to_vec(),
b"X".to_vec(),
b"c".to_vec(),
b"d".to_vec()
]
);
let pos = db.lpos("mylist", b"X", PosSpec::default())?;
assert_eq!(pos, vec![2]);
assert_eq!(db.lrem("mylist", 1, b"X")?, 1);
assert_eq!(db.llen("mylist")?, 4);
let moved = db.lmove("mylist", "otherlist", true, false)?;
assert_eq!(moved, Some(b"a".to_vec()));
assert_eq!(db.llen("mylist")?, 3);
assert_eq!(db.llen("otherlist")?, 1);
let rpop_res = db.rpoplpush("otherlist", "mylist")?;
assert_eq!(rpop_res, Some(b"a".to_vec()));
assert_eq!(db.llen("otherlist")?, 0);
assert_eq!(db.llen("mylist")?, 4);
db.ltrim("mylist", 0, 1)?;
assert_eq!(db.llen("mylist")?, 2);
let popped = db.lpop("mylist", 1)?;
assert_eq!(popped, vec![b"a".to_vec()]);
let popped_r = db.rpop("mylist", 1)?;
assert_eq!(popped_r, vec![b"B".to_vec()]);
assert_eq!(db.llen("mylist")?, 0);
Ok(())
}
#[test]
fn test_list_extended_branches() -> Result<()> {
let temp_dir = tempfile::tempdir().map_err(Error::io)?;
let db = WeDb::open(temp_dir.path())?;
assert_eq!(db.lpop("empty", 1)?, Vec::<Vec<u8>>::new());
assert_eq!(db.rpop("empty", 1)?, Vec::<Vec<u8>>::new());
assert_eq!(db.llen("empty")?, 0);
assert_eq!(db.lrange("empty", 0, -1)?, Vec::<Vec<u8>>::new());
assert_eq!(db.lindex("empty", 0)?, None);
assert_eq!(db.linsert("empty", true, "p", "v")?, 0);
assert_eq!(db.lrem("empty", 0, "v")?, 0);
assert_eq!(
db.lpos("empty", "v", PosSpec::default())?,
Vec::<i64>::new()
);
assert_eq!(db.lmove("empty", "dst", true, false)?, None);
db.rpush("single", &["only"])?;
assert_eq!(
db.lmove("single", "single", true, false)?,
Some(b"only".to_vec())
);
assert_eq!(
db.lmove("single", "single", true, true)?,
Some(b"only".to_vec())
);
assert_eq!(db.llen("single")?, 1);
db.rpush("rot", &["1", "2", "3"])?;
assert_eq!(db.lmove("rot", "rot", true, false)?, Some(b"1".to_vec()));
assert_eq!(
db.lrange("rot", 0, -1)?,
vec![b"2".to_vec(), b"3".to_vec(), b"1".to_vec()]
);
assert_eq!(db.lmove("rot", "rot", false, true)?, Some(b"1".to_vec()));
assert_eq!(
db.lrange("rot", 0, -1)?,
vec![b"1".to_vec(), b"2".to_vec(), b"3".to_vec()]
);
assert!(db.lpos("rot", "1", PosSpec::new().with_rank(0)).is_err());
db.rpush("rot", &["1"])?; let all_pos = db.lpos("rot", "1", PosSpec::new().with_count(0))?;
assert_eq!(all_pos, vec![0, 3]);
db.rpush("shift_list", &["e1", "x", "e2", "x", "e3", "e4", "e5"])?;
assert_eq!(db.lrem("shift_list", 1, "x")?, 1);
assert_eq!(
db.lrange("shift_list", 0, -1)?,
vec![
b"e1".to_vec(),
b"e2".to_vec(),
b"x".to_vec(),
b"e3".to_vec(),
b"e4".to_vec(),
b"e5".to_vec()
]
);
assert_eq!(db.lrem("shift_list", -1, "x")?, 1);
assert_eq!(
db.lrange("shift_list", 0, -1)?,
vec![
b"e1".to_vec(),
b"e2".to_vec(),
b"e3".to_vec(),
b"e4".to_vec(),
b"e5".to_vec()
]
);
assert_eq!(db.lttl("shift_list")?, -1);
assert!(db.lexpire("shift_list", 300)?);
assert!(db.lttl("shift_list")? > 0);
assert!(db.lpersist("shift_list")?);
assert_eq!(db.lttl("shift_list")?, -1);
assert_eq!(db.lttl("nonexistent")?, -2);
Ok(())
}
}