use crate::{
api::list::{
ListItemKeyComposer, compose_list_meta_key, compose_list_prefix_stack,
conf::{ERR_INDEX_OUT_OF_RANGE, ERR_RANK_ZERO, PosSpec},
meta::ListMeta,
traits::List,
},
error::{Error, Result},
key::{clear_prefix_in_batch, get_meta_checked},
meta::current_now_ms,
traits::DbLike,
};
impl<T: DbLike> List for T {
#[inline]
fn lpush<K: AsRef<[u8]>, V: AsRef<[u8]>>(&self, key: K, values: &[V]) -> Result<u64> {
list_push_internal(self, key.as_ref(), values, true, true)
}
#[inline]
fn rpush<K: AsRef<[u8]>, V: AsRef<[u8]>>(&self, key: K, values: &[V]) -> Result<u64> {
list_push_internal(self, key.as_ref(), values, true, false)
}
#[inline]
fn lpushx<K: AsRef<[u8]>, V: AsRef<[u8]>>(&self, key: K, values: &[V]) -> Result<u64> {
list_push_internal(self, key.as_ref(), values, false, true)
}
#[inline]
fn rpushx<K: AsRef<[u8]>, V: AsRef<[u8]>>(&self, key: K, values: &[V]) -> Result<u64> {
list_push_internal(self, key.as_ref(), values, false, false)
}
#[inline]
fn lpop<K: AsRef<[u8]>>(&self, key: K, count: usize) -> Result<Vec<Vec<u8>>> {
list_pop_internal(self, key.as_ref(), count, true)
}
#[inline]
fn rpop<K: AsRef<[u8]>>(&self, key: K, count: usize) -> Result<Vec<Vec<u8>>> {
list_pop_internal(self, key.as_ref(), count, false)
}
#[inline]
fn llen<K: AsRef<[u8]>>(&self, key: K) -> Result<u64> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_list_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
match get_meta_checked::<ListMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) => Ok(m.base.size),
None => Ok(0),
}
}
fn lrange<K: AsRef<[u8]>>(&self, key: K, start: i64, stop: i64) -> Result<Vec<Vec<u8>>> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_list_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let meta = match get_meta_checked::<ListMeta>(self, key_bytes, &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, key_bytes);
let data_ks = self.data();
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) = data_ks.get(item_k)? {
results.push(val.to_vec());
}
}
Ok(results)
}
fn lindex<K: AsRef<[u8]>>(&self, key: K, index: i64) -> Result<Option<Vec<u8>>> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_list_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let meta = match get_meta_checked::<ListMeta>(self, key_bytes, &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, key_bytes);
let actual_idx = meta.head.wrapping_add(actual_offset as u64);
let item_k = composer.key_for_idx(actual_idx);
let val = self.data().get(item_k)?;
Ok(val.map(|v| v.to_vec()))
}
fn lset<K: AsRef<[u8]>, V: AsRef<[u8]>>(&self, key: K, index: i64, value: V) -> Result<()> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_list_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let meta = match get_meta_checked::<ListMeta>(self, key_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Err(Error::invalid_data(ERR_INDEX_OUT_OF_RANGE)),
};
if meta.base.size == 0 {
return Err(Error::invalid_data(ERR_INDEX_OUT_OF_RANGE));
}
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, key_bytes);
let actual_idx = meta.head.wrapping_add(actual_offset as u64);
let item_k = composer.key_for_idx(actual_idx);
let mut batch = self.batch();
batch.insert(self.data(), item_k, value.as_ref());
batch.commit()?;
Ok(())
}
fn ltrim<K: AsRef<[u8]>>(&self, key: K, start: i64, stop: i64) -> Result<()> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_list_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let mut meta = match get_meta_checked::<ListMeta>(self, key_bytes, &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.batch();
let mut composer = ListItemKeyComposer::new(&kc, key_bytes);
let data_ks = self.data();
let meta_ks = self.meta();
if s >= len || e < 0 || s > e {
for offset in 0..meta.base.size {
let idx = meta.head.wrapping_add(offset);
let item_k = composer.key_for_idx(idx);
batch.remove(data_ks, item_k);
}
batch.remove(meta_ks, &meta_k);
batch.commit()?;
return Ok(());
}
if s < 0 {
s = 0;
}
if e >= len {
e = len - 1;
}
for offset in 0..(s as u64) {
let idx = meta.head.wrapping_add(offset);
let item_k = composer.key_for_idx(idx);
batch.remove(data_ks, item_k);
}
for offset in ((e + 1) as u64)..meta.base.size {
let idx = meta.head.wrapping_add(offset);
let item_k = composer.key_for_idx(idx);
batch.remove(data_ks, item_k);
}
let new_size = (e - s + 1) as u64;
let new_head = meta.head.wrapping_add(s as u64);
let new_tail = new_head.wrapping_add(new_size);
meta.base.size = new_size;
meta.head = new_head;
meta.tail = new_tail;
batch.insert(meta_ks, &meta_k, meta.encode());
batch.commit()?;
Ok(())
}
fn linsert<K: AsRef<[u8]>, P: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
key: K,
before: bool,
pivot: P,
elem: V,
) -> Result<i64> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_list_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let mut meta = match get_meta_checked::<ListMeta>(self, key_bytes, &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, key_bytes);
let data_ks = self.data();
let meta_ks = self.meta();
let mut pivot_offset = None;
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) = data_ks.get(item_k)?
&& v.as_ref() == pivot_bytes
{
pivot_offset = Some(offset);
break;
}
}
let pivot_offset = match pivot_offset {
Some(o) => o,
None => return Ok(-1),
};
let insert_offset = if before {
pivot_offset
} else {
pivot_offset + 1
};
let mut batch = self.batch();
if insert_offset < len / 2 {
let new_head = meta.head.wrapping_sub(1);
for offset in 0..insert_offset {
let from_idx = meta.head.wrapping_add(offset as u64);
let to_idx = new_head.wrapping_add(offset as u64);
if let Some(val) = data_ks.get(composer.key_for_idx(from_idx))? {
batch.insert(data_ks, composer.key_for_idx(to_idx), val.as_ref());
}
}
let target_idx = new_head.wrapping_add(insert_offset as u64);
batch.insert(data_ks, composer.key_for_idx(target_idx), elem.as_ref());
meta.head = new_head;
} else {
let old_tail = meta.tail;
let new_tail = old_tail.wrapping_add(1);
for offset in (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(val) = data_ks.get(composer.key_for_idx(from_idx))? {
batch.insert(data_ks, composer.key_for_idx(to_idx), val.as_ref());
}
}
let target_idx = meta.head.wrapping_add(insert_offset as u64);
batch.insert(data_ks, composer.key_for_idx(target_idx), elem.as_ref());
meta.tail = new_tail;
}
meta.base.size += 1;
batch.insert(meta_ks, &meta_k, meta.encode());
batch.commit()?;
Ok(meta.base.size as i64)
}
fn lrem<K: AsRef<[u8]>, V: AsRef<[u8]>>(&self, key: K, count: i64, elem: V) -> Result<u64> {
let key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_list_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let mut meta = match get_meta_checked::<ListMeta>(self, key_bytes, &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, key_bytes);
let data_ks = self.data();
let meta_ks = self.meta();
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) = 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 offset in (0..len).rev() {
let idx = meta.head.wrapping_add(offset as u64);
let item_k = composer.key_for_idx(idx);
if let Some(v) = 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 deleted_count = to_delete_offsets.len() as u64;
let mut batch = self.batch();
if deleted_count == meta.base.size {
for offset in 0..len {
let idx = meta.head.wrapping_add(offset as u64);
batch.remove(data_ks, composer.key_for_idx(idx));
}
batch.remove(meta_ks, &meta_k);
batch.commit()?;
return Ok(deleted_count);
}
let mut write_idx = 0usize;
let mut del_idx = 0usize;
for read_idx in 0..len {
if del_idx < to_delete_offsets.len() && to_delete_offsets[del_idx] == read_idx {
del_idx += 1;
continue;
}
if write_idx != read_idx {
let from_key_idx = meta.head.wrapping_add(read_idx as u64);
let to_key_idx = meta.head.wrapping_add(write_idx as u64);
if let Some(val) = data_ks.get(composer.key_for_idx(from_key_idx))? {
batch.insert(data_ks, composer.key_for_idx(to_key_idx), val.as_ref());
}
}
write_idx += 1;
}
for extra_idx in write_idx..len {
let key_idx = meta.head.wrapping_add(extra_idx as u64);
batch.remove(data_ks, composer.key_for_idx(key_idx));
}
meta.base.size -= deleted_count;
meta.tail = meta.head.wrapping_add(meta.base.size);
batch.insert(meta_ks, &meta_k, meta.encode());
batch.commit()?;
Ok(deleted_count)
}
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 = current_now_ms();
let kc = self.kc();
let src_meta_k = compose_list_meta_key(&kc, src_bytes);
let mut src_meta = match get_meta_checked::<ListMeta>(self, src_bytes, &src_meta_k, now_ms)? {
Some(m) => m,
None => return Ok(None),
};
if src_meta.base.size == 0 {
return Ok(None);
}
let data_ks = self.data();
let meta_ks = self.meta();
if src_bytes == dst_bytes {
let mut composer = ListItemKeyComposer::new(&kc, src_bytes);
let curr_idx = if src_left {
src_meta.head
} else {
src_meta.tail.wrapping_sub(1)
};
let elem = match data_ks.get(composer.key_for_idx(curr_idx))? {
Some(v) => v.to_vec(),
None => return Ok(None),
};
if src_meta.base.size == 1 && src_left == dst_left {
return Ok(Some(elem));
}
let mut batch = self.batch();
batch.remove(data_ks, composer.key_for_idx(curr_idx));
if src_left {
src_meta.head = src_meta.head.wrapping_add(1);
} else {
src_meta.tail = src_meta.tail.wrapping_sub(1);
}
let target_idx = if dst_left {
src_meta.head = src_meta.head.wrapping_sub(1);
src_meta.head
} else {
let t = src_meta.tail;
src_meta.tail = src_meta.tail.wrapping_add(1);
t
};
batch.insert(data_ks, composer.key_for_idx(target_idx), &elem);
batch.insert(meta_ks, &src_meta_k, src_meta.encode());
batch.commit()?;
return Ok(Some(elem));
}
let mut src_composer = ListItemKeyComposer::new(&kc, src_bytes);
let curr_src_idx = if src_left {
src_meta.head
} else {
src_meta.tail.wrapping_sub(1)
};
let elem = match data_ks.get(src_composer.key_for_idx(curr_src_idx))? {
Some(v) => v.to_vec(),
None => return Ok(None),
};
let dst_meta_k = compose_list_meta_key(&kc, dst_bytes);
let mut batch = self.batch();
let (mut dst_meta, _) =
prepare_list_meta_for_write(self, dst_bytes, &dst_meta_k, now_ms, &mut batch)?;
batch.remove(data_ks, src_composer.key_for_idx(curr_src_idx));
src_meta.base.size -= 1;
if src_left {
src_meta.head = src_meta.head.wrapping_add(1);
} else {
src_meta.tail = src_meta.tail.wrapping_sub(1);
}
if src_meta.base.size == 0 {
batch.remove(meta_ks, &src_meta_k);
} else {
batch.insert(meta_ks, &src_meta_k, src_meta.encode());
}
let mut dst_composer = ListItemKeyComposer::new(&kc, dst_bytes);
let target_dst_idx = if dst_left {
dst_meta.head = dst_meta.head.wrapping_sub(1);
dst_meta.head
} else {
let t = dst_meta.tail;
dst_meta.tail = dst_meta.tail.wrapping_add(1);
t
};
dst_meta.base.size += 1;
batch.insert(data_ks, dst_composer.key_for_idx(target_dst_idx), &elem);
batch.insert(meta_ks, &dst_meta_k, dst_meta.encode());
batch.commit()?;
Ok(Some(elem))
}
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 key_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_list_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let meta = match get_meta_checked::<ListMeta>(self, key_bytes, &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 = match spec.count {
Some(c) if c > 0 => Vec::with_capacity(c.min(limit)),
_ => Vec::with_capacity(1),
};
let count_limit = spec.count.unwrap_or(1);
let is_multi_count = spec.count.is_some();
let mut composer = ListItemKeyComposer::new(&kc, key_bytes);
let data_ks = self.data();
let mut rank_count = 0usize;
if !reversed {
for offset in 0..limit {
let idx = meta.head.wrapping_add(offset as u64);
let item_k = composer.key_for_idx(idx);
if let Some(val) = data_ks.get(item_k)?
&& val.as_ref() == elem_bytes
{
rank_count += 1;
if rank_count >= target_rank {
matches.push(offset as i64);
if is_multi_count && count_limit > 0 && matches.len() >= count_limit {
break;
}
if !is_multi_count {
break;
}
}
}
}
} else {
let start_offset = len - 1;
let end_offset = len.saturating_sub(limit);
for offset in (end_offset..=start_offset).rev() {
let idx = meta.head.wrapping_add(offset as u64);
let item_k = composer.key_for_idx(idx);
if let Some(val) = data_ks.get(item_k)?
&& val.as_ref() == elem_bytes
{
rank_count += 1;
if rank_count >= target_rank {
matches.push(offset as i64);
if is_multi_count && count_limit > 0 && matches.len() >= count_limit {
break;
}
if !is_multi_count {
break;
}
}
}
}
}
Ok(matches)
}
}
pub fn prepare_list_meta_for_write(
db: &impl DbLike,
k_bytes: &[u8],
meta_k: &[u8],
now_ms: u64,
batch: &mut fjall::OwnedWriteBatch,
) -> Result<(ListMeta, bool)> {
let kc = db.kc();
match get_meta_checked::<ListMeta>(db, k_bytes, meta_k, now_ms)? {
Some(meta) => Ok((meta, true)),
None => {
let prefix = compose_list_prefix_stack(&kc, k_bytes);
clear_prefix_in_batch(db.data(), &prefix, batch)?;
Ok((ListMeta::new_with_version(0), false))
}
}
}
fn list_push_internal<T: DbLike, V: AsRef<[u8]>>(
db: &T,
key_bytes: &[u8],
values: &[V],
create_if_missing: bool,
push_left: bool,
) -> Result<u64> {
if values.is_empty() {
return Ok(0);
}
let kc = db.kc();
let meta_k = compose_list_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let mut batch = db.batch();
let (mut meta, metadata_existed) =
prepare_list_meta_for_write(db, key_bytes, &meta_k, now_ms, &mut batch)?;
if !create_if_missing && (!metadata_existed || meta.base.size == 0) {
return Ok(0);
}
let mut composer = ListItemKeyComposer::new(&kc, key_bytes);
let data_ks = db.data();
let meta_ks = db.meta();
for v in values {
let v_bytes = v.as_ref();
let target_idx = if push_left {
meta.head = meta.head.wrapping_sub(1);
meta.head
} else {
let t = meta.tail;
meta.tail = meta.tail.wrapping_add(1);
t
};
let item_k = composer.key_for_idx(target_idx);
batch.insert(data_ks, item_k, v_bytes);
meta.base.size += 1;
}
batch.insert(meta_ks, &meta_k, meta.encode());
batch.commit()?;
Ok(meta.base.size)
}
fn list_pop_internal<T: DbLike>(
db: &T,
key_bytes: &[u8],
count: usize,
pop_left: bool,
) -> Result<Vec<Vec<u8>>> {
if count == 0 {
return Ok(Vec::new());
}
let kc = db.kc();
let meta_k = compose_list_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let mut meta = match get_meta_checked::<ListMeta>(db, key_bytes, &meta_k, now_ms)? {
Some(m) if m.base.size > 0 => m,
_ => return Ok(Vec::new()),
};
let actual_count = (count as u64).min(meta.base.size);
let mut results = Vec::with_capacity(actual_count as usize);
let mut batch = db.batch();
let mut composer = ListItemKeyComposer::new(&kc, key_bytes);
let data_ks = db.data();
let meta_ks = db.meta();
for _ in 0..actual_count {
let target_idx = if pop_left {
let h = meta.head;
meta.head = meta.head.wrapping_add(1);
h
} else {
meta.tail = meta.tail.wrapping_sub(1);
meta.tail
};
let item_k = composer.key_for_idx(target_idx);
if let Some(val) = data_ks.get(item_k)? {
results.push(val.to_vec());
batch.remove(data_ks, item_k);
}
}
meta.base.size -= actual_count;
if meta.base.size == 0 {
batch.remove(meta_ks, &meta_k);
} else {
batch.insert(meta_ks, &meta_k, meta.encode());
}
batch.commit()?;
Ok(results)
}