use serde::{Serialize, de::DeserializeOwned};
use std::collections::VecDeque;
use std::sync::{Arc, Condvar, Mutex};
use std::time::{Duration, Instant};
use crate::codec::{Bytes, Codec};
use crate::entry::{Entry, ListEntry, ListInner, Meta, ValueType};
use crate::error::Error;
use crate::store::Store;
pub struct ListRef<'a, C: Codec> {
store: &'a Store<C>,
key: &'a str,
}
impl<'a, C: Codec> ListRef<'a, C> {
pub(crate) fn new(store: &'a Store<C>, key: &'a str) -> Self {
Self { store, key }
}
#[inline]
fn enc<T: Serialize>(&self, v: &T) -> Result<Bytes, Error> {
self.store.codec().encode(v)
}
#[inline]
fn dec<T: DeserializeOwned>(&self, b: &[u8]) -> Result<T, Error> {
self.store.codec().decode(b)
}
pub fn lpush<T: Serialize>(&self, value: &T) -> Result<usize, Error> {
let b = self.enc(value)?;
let len = self.store.with_list_mut(self.key, |dq| {
dq.push_front(b);
Ok(dq.len())
})?;
let _ = self.store.list_notify(self.key);
Ok(len)
}
pub fn rpush<T: Serialize>(&self, value: &T) -> Result<usize, Error> {
let b = self.enc(value)?;
let len = self.store.with_list_mut(self.key, |dq| {
dq.push_back(b);
Ok(dq.len())
})?;
let _ = self.store.list_notify(self.key);
Ok(len)
}
pub fn lpop<T: DeserializeOwned>(&self) -> Result<Option<T>, Error> {
self.store.with_list_mut(self.key, |dq| {
let Some(b) = dq.pop_front() else {
return Ok(None);
};
Ok(Some(self.dec::<T>(&b)?))
})
}
pub fn rpop<T: DeserializeOwned>(&self) -> Result<Option<T>, Error> {
self.store.with_list_mut(self.key, |dq| {
let Some(b) = dq.pop_back() else {
return Ok(None);
};
Ok(Some(self.dec::<T>(&b)?))
})
}
pub fn llen(&self) -> Result<usize, Error> {
self.store
.with_list_read(self.key, |opt| Ok(opt.map(|d| d.len()).unwrap_or(0)))
}
pub fn is_empty(&self) -> Result<bool, Error> {
Ok(self.llen()? == 0)
}
pub fn lrange<T: DeserializeOwned>(&self, start: isize, stop: isize) -> Result<Vec<T>, Error> {
self.store.with_list_read(self.key, |opt| {
let Some(dq) = opt else {
return Ok(vec![]);
};
let n = dq.len() as isize;
if n == 0 {
return Ok(vec![]);
}
let mut s = if start < 0 { n + start } else { start };
let mut e = if stop < 0 { n + stop } else { stop };
if s < 0 {
s = 0;
}
if e < 0 {
return Ok(vec![]);
}
if s >= n {
return Ok(vec![]);
}
if e >= n {
e = n - 1;
}
if e < s {
return Ok(vec![]);
}
let want = (e - s + 1) as usize;
let mut out = Vec::with_capacity(want);
for i in s..=e {
out.push(self.dec::<T>(&dq[i as usize])?);
}
Ok(out)
})
}
pub fn lindex<T: DeserializeOwned>(&self, index: isize) -> Result<Option<T>, Error> {
self.store.with_list_read(self.key, |opt| {
let Some(dq) = opt else {
return Ok(None);
};
let n = dq.len() as isize;
if n == 0 {
return Ok(None);
}
let i = if index < 0 { n + index } else { index };
if i < 0 || i >= n {
return Ok(None);
}
Ok(Some(self.dec::<T>(&dq[i as usize])?))
})
}
pub fn lset<T: Serialize>(&self, index: isize, value: &T) -> Result<bool, Error> {
let b = self.enc(value)?;
self.store.with_list_mut(self.key, |dq| {
let n = dq.len() as isize;
if n == 0 {
return Ok(false);
}
let i = if index < 0 { n + index } else { index };
if i < 0 || i >= n {
return Ok(false);
}
dq[i as usize] = b;
Ok(true)
})
}
pub fn ltrim(&self, start: isize, stop: isize) -> Result<(), Error> {
self.store.with_list_mut(self.key, |dq| {
let n = dq.len() as isize;
if n == 0 {
return Ok(());
}
let mut s = if start < 0 { n + start } else { start };
let mut e = if stop < 0 { n + stop } else { stop };
if s < 0 {
s = 0;
}
if e < 0 || s >= n {
dq.clear();
return Ok(());
}
if e >= n {
e = n - 1;
}
if e < s {
dq.clear();
return Ok(());
}
let s = s as usize;
let e = e as usize;
let keep_len = e - s + 1;
let kept: VecDeque<Bytes> = dq.iter().skip(s).take(keep_len).cloned().collect();
*dq = kept;
Ok(())
})
}
pub fn lrem<T: Serialize>(&self, count: isize, value: &T) -> Result<usize, Error> {
let target = self.enc(value)?;
self.store.with_list_mut(self.key, |dq| {
if dq.is_empty() {
return Ok(0);
}
let mut removed = 0usize;
if count == 0 {
let mut i = 0usize;
while i < dq.len() {
if dq[i] == target {
dq.remove(i);
removed += 1;
} else {
i += 1;
}
}
return Ok(removed);
}
if count > 0 {
let limit = count as usize;
let mut i = 0usize;
while i < dq.len() && removed < limit {
if dq[i] == target {
dq.remove(i);
removed += 1;
} else {
i += 1;
}
}
return Ok(removed);
}
let limit = (-count) as usize;
let mut i = dq.len();
while i > 0 && removed < limit {
i -= 1;
if dq[i] == target {
dq.remove(i);
removed += 1;
}
}
Ok(removed)
})
}
pub fn lpos<T: Serialize>(&self, value: &T) -> Result<Option<usize>, Error> {
let target = self.enc(value)?;
self.store.with_list_read(self.key, |opt| {
let Some(dq) = opt else {
return Ok(None);
};
Ok(dq.iter().position(|b| *b == target))
})
}
pub fn blpop<T: DeserializeOwned>(
&self,
timeout: Option<Duration>,
) -> Result<Option<T>, Error> {
self.blocking_pop::<T>(true, timeout)
}
pub fn brpop<T: DeserializeOwned>(
&self,
timeout: Option<Duration>,
) -> Result<Option<T>, Error> {
self.blocking_pop::<T>(false, timeout)
}
fn blocking_pop<T: DeserializeOwned>(
&self,
left: bool,
timeout: Option<Duration>,
) -> Result<Option<T>, Error> {
self.store.purge_if_expired(self.key);
let inner = self.store.with_map_write(|m| {
match m.get(self.key) {
None => {
m.insert(
self.key.to_string(),
Entry::List(ListEntry {
meta: Meta::new(ValueType::List),
inner: Arc::new((
Mutex::new(ListInner {
deque: VecDeque::<Bytes>::new(),
}),
Condvar::new(),
)),
}),
);
}
Some(e) if e.value_type() != ValueType::List => {
return Err(Error::WrongType {
expected: ValueType::List.as_str(),
got: e.value_type().as_str(),
});
}
Some(_) => {}
}
match m.get(self.key).expect("just inserted or validated") {
Entry::List(le) => Ok(le.inner.clone()),
other => Err(Error::WrongType {
expected: ValueType::List.as_str(),
got: other.value_type().as_str(),
}),
}
})?;
let (lock, cv) = &*inner;
let mut guard = lock
.lock()
.map_err(|_| Error::InvalidArgument("list mutex poisoned".to_string()))?;
if let Some(b) = if left {
guard.deque.pop_front()
} else {
guard.deque.pop_back()
} {
return Ok(Some(self.dec::<T>(&b)?));
}
let deadline = timeout.map(|d| Instant::now() + d);
loop {
match deadline {
None => {
guard = cv
.wait(guard)
.map_err(|_| Error::InvalidArgument("condvar wait poisoned".to_string()))?;
}
Some(dl) => {
let now = Instant::now();
if now >= dl {
return Ok(None);
}
let remaining = dl - now;
let (g, _timeout_res) = cv
.wait_timeout(guard, remaining)
.map_err(|_| Error::InvalidArgument("condvar wait poisoned".to_string()))?;
guard = g;
}
}
if let Some(b) = if left {
guard.deque.pop_front()
} else {
guard.deque.pop_back()
} {
return Ok(Some(self.dec::<T>(&b)?));
}
}
}
}
#[cfg(test)]
mod tests {
use crate::Store;
use std::time::Duration;
#[test]
fn list_push_pop_range_trim() {
let store = Store::new();
let l = store.list("l");
l.rpush(&1_i64).unwrap();
l.rpush(&2_i64).unwrap();
l.lpush(&0_i64).unwrap();
assert_eq!(l.llen().unwrap(), 3);
let v: Vec<i64> = l.lrange(0, -1).unwrap();
assert_eq!(v, vec![0, 1, 2]);
let x: Option<i64> = l.lpop().unwrap();
assert_eq!(x, Some(0));
l.ltrim(0, 0).unwrap();
let v2: Vec<i64> = l.lrange(0, -1).unwrap();
assert_eq!(v2, vec![1]);
}
#[test]
fn list_lrem_lpos() {
let store = Store::new();
let l = store.list("r");
l.rpush(&"a").unwrap();
l.rpush(&"b").unwrap();
l.rpush(&"a").unwrap();
assert_eq!(l.lpos(&"a").unwrap(), Some(0));
assert_eq!(l.lrem(0, &"a").unwrap(), 2);
assert_eq!(l.lpos(&"a").unwrap(), None);
}
#[test]
fn list_blocking_timeout() {
let store = Store::new();
let l = store.list("blk");
let v: Option<i64> = l.blpop(Some(Duration::from_millis(10))).unwrap();
assert_eq!(v, None);
}
#[test]
fn list_lset_in_bounds_and_oob() {
let store = Store::new();
let l = store.list("l2");
l.rpush(&1_i64).unwrap();
l.rpush(&2_i64).unwrap();
assert!(l.lset(1, &10_i64).unwrap());
let v: Vec<i64> = l.lrange(0, -1).unwrap();
assert_eq!(v, vec![1, 10]);
assert!(!l.lset(10, &0_i64).unwrap());
}
#[test]
fn list_brpop_wakes_no_timeout() {
use std::sync::Arc;
let store = Arc::new(Store::new());
let store2 = Arc::clone(&store);
let t = std::thread::spawn(move || {
let l = store2.list("q");
let v: Option<String> = l.brpop(None).unwrap();
v
});
std::thread::sleep(Duration::from_millis(20));
store.list("q").rpush(&"only").unwrap();
let got = t.join().unwrap();
assert_eq!(got.as_deref(), Some("only"));
}
#[test]
fn list_rpop_semantics() {
let store = Store::new();
let l = store.list("q2");
l.rpush(&"a").unwrap();
l.rpush(&"b").unwrap();
let got: Option<String> = l.rpop().unwrap();
assert_eq!(got.as_deref(), Some("b"));
}
}