use wkv::TtlOpt;
use wresp::length::{try_read_length, try_write_length};
use super::{
cmd_strings as cs,
cmd_strings::{
abort_with_error_message, abort_with_wrong_number_of_arguments, write_error_raw, write_raw,
},
parser::{
resp_ext::{RespSliceExt, RespVecExt},
session_parse_state::{strict_i32, strict_i64},
},
rdb_crc64,
resp_server_session::RespServerSession,
ttl_sync::{
del_ttl_sync, now_unix_ms, probe_alive, put_ttl_sync, read_adjudicated_sync, ttl_of_sync,
},
};
const RDB_VERSION: u16 = 11;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ExpireCmd {
Expire,
Pexpire,
Expireat,
Pexpireat,
}
impl ExpireCmd {
pub const fn as_str(self) -> &'static str {
match self {
Self::Expire => "EXPIRE",
Self::Pexpire => "PEXPIRE",
Self::Expireat => "EXPIREAT",
Self::Pexpireat => "PEXPIREAT",
}
}
const fn is_relative(self) -> bool {
matches!(self, Self::Expire | Self::Pexpire)
}
const fn is_millis(self) -> bool {
matches!(self, Self::Pexpire | Self::Pexpireat)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TtlCmd {
Ttl,
Pttl,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ExpireTimeCmd {
Expiretime,
Pexpiretime,
}
fn try_parse_expire_option(raw: &[u8]) -> bool {
raw.eq_ignore_ascii_case(b"NX")
|| raw.eq_ignore_ascii_case(b"XX")
|| raw.eq_ignore_ascii_case(b"GT")
|| raw.eq_ignore_ascii_case(b"LT")
}
impl RespServerSession {
pub fn network_restore<'a, D: wdev::Device>(
&mut self,
parse_state: &[&[u8]],
store: &wkv::BatchStoreSession<'a, D>,
output: &mut Vec<u8>,
) -> wresp::Result<bool> {
if parse_state.len() != 3 {
abort_with_wrong_number_of_arguments(output, "RESTORE");
return Ok(true);
}
let key = parse_state[0];
let Some(expiry) = strict_i32(parse_state[1]) else {
abort_with_error_message(output, cs::RESP_ERR_TIMEOUT_NOT_VALID_FLOAT);
return Ok(true);
};
let expiry = i64::from(expiry);
let value = parse_state[2];
if value.first() != Some(&0x00) {
write_error_raw(output, "ERR RESTORE currently only supports string types");
return Ok(true);
}
if value.len() < 10 {
write_error_raw(output, "ERR DUMP payload version or checksum are wrong");
return Ok(true);
}
let footer = &value[value.len() - 10..];
let rdb_version = u16::from_le_bytes([footer[0], footer[1]]);
if rdb_version > RDB_VERSION {
write_error_raw(output, "ERR DUMP payload version or checksum are wrong");
return Ok(true);
}
let calculated_crc = rdb_crc64::hash(&value[..value.len() - 8]);
if calculated_crc != footer[2..] {
write_error_raw(output, "ERR DUMP payload version or checksum are wrong");
return Ok(true);
}
let Some((length, payload_start)) = try_read_length(&value[1..]) else {
write_error_raw(output, "ERR DUMP payload length format is invalid");
return Ok(true);
};
let Some(val) = value
.get(payload_start + 1..payload_start + 1 + length as usize)
.filter(|_| payload_start as u64 + 1 + length as u64 <= value.len() as u64)
else {
write_error_raw(output, "ERR DUMP payload length format is invalid");
return Ok(true);
};
let exists = match probe_alive(store, key) {
Ok(Some(alive)) => alive,
Ok(None) => return Ok(false),
Err(_) => {
output.write_resp_error("generic error");
return Ok(true);
}
};
if exists {
write_error_raw(output, cs::RESP_ERR_BUSSYKEY);
return Ok(true);
}
match store.try_upsert_sync(key, val) {
Ok(Ok(_)) => {}
Ok(Err(_)) => return Ok(false),
Err(_) => {
output.write_resp_error("generic error");
return Ok(true);
}
}
if expiry > 0 {
let expire_at_ms = now_unix_ms() as i64 + expiry.saturating_mul(1000);
match put_ttl_sync(store, key, expire_at_ms.max(0) as u64) {
Ok(true) => {}
Ok(false) => return Ok(false),
Err(_) => {
output.write_resp_error("generic error");
return Ok(true);
}
}
}
write_raw(output, cs::RESP_OK);
Ok(true)
}
pub fn network_dump<'a, D: wdev::Device>(
&mut self,
parse_state: &[&[u8]],
store: &wkv::BatchStoreSession<'a, D>,
output: &mut Vec<u8>,
) -> wresp::Result<bool> {
if parse_state.len() != 1 {
abort_with_wrong_number_of_arguments(output, "DUMP");
return Ok(true);
}
let key = parse_state[0];
match read_adjudicated_sync(store, key, |v| v.to_vec()) {
Ok(Some(Some(value))) => {
let mut encoded_len = [0u8; 5];
let Some(bytes_written) = try_write_length(value.len() as u32, &mut encoded_len) else {
write_error_raw(output, "ERR DUMP payload length is invalid");
return Ok(true);
};
let encoded_len = &encoded_len[..bytes_written];
let payload_len = 1 + encoded_len.len() + value.len() + 2 + 8;
output.push(b'$');
let mut buf = itoa::Buffer::new();
output.extend_from_slice(buf.format(payload_len).as_bytes());
output.extend_from_slice(b"\r\n");
output.push(0x00);
output.extend_from_slice(encoded_len);
output.extend_from_slice(&value);
output.extend_from_slice(&RDB_VERSION.to_le_bytes());
let framed = output.len() - (payload_len - 8);
let crc = rdb_crc64::hash(&output[framed..]);
output.extend_from_slice(&crc);
output.extend_from_slice(b"\r\n");
}
Ok(Some(None)) => output.write_resp_null(),
Ok(None) => return Ok(false),
Err(_) => output.write_resp_error("generic error"),
}
Ok(true)
}
pub fn network_rename<'a, D: wdev::Device>(
&mut self,
parse_state: &[&[u8]],
store: &wkv::BatchStoreSession<'a, D>,
output: &mut Vec<u8>,
) -> wresp::Result<bool> {
if parse_state.len() != 2 {
abort_with_wrong_number_of_arguments(output, "RENAME");
return Ok(true);
}
rename_sync(store, parse_state[0], parse_state[1], false, output)
}
pub fn network_renamenx<'a, D: wdev::Device>(
&mut self,
parse_state: &[&[u8]],
store: &wkv::BatchStoreSession<'a, D>,
output: &mut Vec<u8>,
) -> wresp::Result<bool> {
if parse_state.len() != 2 {
abort_with_wrong_number_of_arguments(output, "RENAMENX");
return Ok(true);
}
rename_sync(store, parse_state[0], parse_state[1], true, output)
}
pub fn network_getdel<'a, D: wdev::Device>(
&mut self,
parse_state: &[&[u8]],
store: &wkv::BatchStoreSession<'a, D>,
output: &mut Vec<u8>,
) -> wresp::Result<bool> {
if parse_state.len() != 1 {
abort_with_wrong_number_of_arguments(output, "GETDEL");
return Ok(true);
}
let key = parse_state[0];
match store.try_read_sync(key, |v| v.to_vec()) {
Ok(Some(Some(val))) => {
match store.try_delete_sync(key) {
Ok(Ok(_)) => output.write_resp_bulk_string(&val),
Ok(Err(_)) => return Ok(false),
Err(_) => output.write_resp_error("generic error"),
}
}
Ok(Some(None)) => {
output.write_resp_null();
}
Ok(None) => return Ok(false),
Err(_) => output.write_resp_error("generic error"),
}
Ok(true)
}
pub fn network_exists<'a, D: wdev::Device>(
&mut self,
parse_state: &[&[u8]],
store: &wkv::BatchStoreSession<'a, D>,
output: &mut Vec<u8>,
) -> wresp::Result<bool> {
if parse_state.is_empty() {
abort_with_wrong_number_of_arguments(output, "EXISTS");
return Ok(true);
}
let mut exists_count = 0i64;
for key in parse_state {
match store.try_read_sync(key, |_| ()) {
Ok(Some(Some(_))) => exists_count += 1,
Ok(Some(None)) => {}
Ok(None) => return Ok(false),
Err(_) => {
output.write_resp_error("generic error");
return Ok(true);
}
}
}
output.write_resp_int(exists_count);
Ok(true)
}
pub fn network_expire<'a, D: wdev::Device>(
&mut self,
command: ExpireCmd,
parse_state: &[&[u8]],
store: &wkv::BatchStoreSession<'a, D>,
output: &mut Vec<u8>,
) -> wresp::Result<bool> {
let count = parse_state.len();
if !(2..=4).contains(&count) {
abort_with_wrong_number_of_arguments(output, command.as_str());
return Ok(true);
}
let key = parse_state[0];
let Some(expiration) = strict_i64(parse_state[1]) else {
abort_with_error_message(output, cs::RESP_ERR_GENERIC_VALUE_IS_NOT_INTEGER);
return Ok(true);
};
if expiration < 0 {
abort_with_error_message(output, cs::RESP_ERR_INVALID_EXPIRE_TIME);
return Ok(true);
}
let mut opt = TtlOpt::NONE;
if count > 2 {
if !try_parse_expire_option(parse_state[2]) {
abort_with_error_message(
output,
&cs::GENERIC_ERR_UNSUPPORTED_OPTION.replace("{0}", parse_state[2].as_str_safe()),
);
return Ok(true);
}
opt = parse_expire_option(parse_state[2], TtlOpt::NONE);
if count > 3 {
if !try_parse_expire_option(parse_state[3]) {
abort_with_error_message(
output,
&cs::GENERIC_ERR_UNSUPPORTED_OPTION.replace("{0}", parse_state[3].as_str_safe()),
);
return Ok(true);
}
let first = parse_state[2];
let compatible = (first.eq_ignore_ascii_case(b"XX")
&& (parse_state[3].eq_ignore_ascii_case(b"GT")
|| parse_state[3].eq_ignore_ascii_case(b"LT")))
|| ((first.eq_ignore_ascii_case(b"GT") || first.eq_ignore_ascii_case(b"LT"))
&& parse_state[3].eq_ignore_ascii_case(b"XX"));
if !compatible {
abort_with_error_message(
output,
"ERR NX and XX, GT or LT options at the same time are not compatible",
);
return Ok(true);
}
opt = merge_expire_options(first, parse_state[3]);
}
}
let now = now_unix_ms() as i64;
let expire_at_ms = if command.is_relative() {
let span = if command.is_millis() {
expiration
} else {
expiration.saturating_mul(1000)
};
now.saturating_add(span)
} else if command.is_millis() {
expiration
} else {
expiration.saturating_mul(1000)
};
match expire_apply_sync(store, key, expire_at_ms.max(0) as u64, opt) {
Ok(Some(applied)) => {
if applied != 0 {
write_raw(output, cs::RESP_RETURN_VAL_1);
} else {
write_raw(output, cs::RESP_RETURN_VAL_0);
}
Ok(true)
}
Ok(None) => Ok(false),
Err(_) => {
output.write_resp_error("generic error");
Ok(true)
}
}
}
pub fn network_persist<'a, D: wdev::Device>(
&mut self,
parse_state: &[&[u8]],
store: &wkv::BatchStoreSession<'a, D>,
output: &mut Vec<u8>,
) -> wresp::Result<bool> {
if parse_state.len() != 1 {
abort_with_wrong_number_of_arguments(output, "PERSIST");
return Ok(true);
}
let key = parse_state[0];
match persist_apply_sync(store, key) {
Ok(Some(removed)) => {
output.write_resp_int(i64::from(removed));
Ok(true)
}
Ok(None) => Ok(false),
Err(_) => {
output.write_resp_error("generic error");
Ok(true)
}
}
}
pub fn network_ttl<'a, D: wdev::Device>(
&mut self,
command: TtlCmd,
parse_state: &[&[u8]],
store: &wkv::BatchStoreSession<'a, D>,
output: &mut Vec<u8>,
) -> wresp::Result<bool> {
if parse_state.len() != 1 {
abort_with_wrong_number_of_arguments(output, command_name_of_ttl(command));
return Ok(true);
}
let key = parse_state[0];
match ttl_read_sync(store, key) {
Ok(Some(read)) => {
let value = match read {
ExpiryRead::Missing => -2,
ExpiryRead::NoExpiry => -1,
ExpiryRead::At(exp) => {
let remaining = exp.saturating_sub(now_unix_ms());
if command == TtlCmd::Pttl {
remaining as i64
} else {
(remaining as i64 + 500) / 1000
}
}
};
output.write_resp_int(value);
Ok(true)
}
Ok(None) => Ok(false),
Err(_) => {
output.write_resp_error("generic error");
Ok(true)
}
}
}
pub fn network_expiretime<'a, D: wdev::Device>(
&mut self,
command: ExpireTimeCmd,
parse_state: &[&[u8]],
store: &wkv::BatchStoreSession<'a, D>,
output: &mut Vec<u8>,
) -> wresp::Result<bool> {
if parse_state.len() != 1 {
abort_with_wrong_number_of_arguments(output, "EXPIRETIME");
return Ok(true);
}
let key = parse_state[0];
match expiretime_read_sync(store, key) {
Ok(Some(read)) => {
let value = match read {
ExpiryRead::Missing => -2,
ExpiryRead::NoExpiry => -1,
ExpiryRead::At(exp) => {
if command == ExpireTimeCmd::Pexpiretime {
exp as i64
} else {
exp as i64 / 1000
}
}
};
output.write_resp_int(value);
Ok(true)
}
Ok(None) => Ok(false),
Err(_) => {
output.write_resp_error("generic error");
Ok(true)
}
}
}
}
const fn command_name_of_ttl(command: TtlCmd) -> &'static str {
match command {
TtlCmd::Ttl => "TTL",
TtlCmd::Pttl => "PTTL",
}
}
fn rename_sync<'a, D: wdev::Device>(
store: &wkv::BatchStoreSession<'a, D>,
old_key: &[u8],
new_key: &[u8],
nx: bool,
output: &mut Vec<u8>,
) -> wresp::Result<bool> {
if old_key == new_key {
if nx {
output.write_resp_int(1);
} else {
write_raw(output, cs::RESP_OK);
}
return Ok(true);
}
let old_val = match read_adjudicated_sync(store, old_key, |v| v.to_vec()) {
Ok(Some(Some(val))) => val,
Ok(Some(None)) => {
abort_with_error_message(output, cs::RESP_ERR_GENERIC_NOSUCHKEY);
return Ok(true);
}
Ok(None) => return Ok(false),
Err(_) => {
output.write_resp_error("generic error");
return Ok(true);
}
};
let old_ttl = match ttl_of_sync(store, old_key) {
Ok(Some(ttl)) => ttl,
Ok(None) => return Ok(false),
Err(_) => {
output.write_resp_error("generic error");
return Ok(true);
}
};
if nx {
match probe_alive(store, new_key) {
Ok(Some(true)) => {
output.write_resp_int(0);
return Ok(true);
}
Ok(Some(false)) => {}
Ok(None) => return Ok(false),
Err(_) => {
output.write_resp_error("generic error");
return Ok(true);
}
}
}
match store.try_upsert_sync(new_key, old_val.as_slice()) {
Ok(Ok(_)) => {}
Ok(Err(_)) => return Ok(false),
Err(_) => {
output.write_resp_error("generic error");
return Ok(true);
}
}
if let Some(exp) = old_ttl {
match put_ttl_sync(store, new_key, exp) {
Ok(true) => {}
Ok(false) => return Ok(false),
Err(_) => {
output.write_resp_error("generic error");
return Ok(true);
}
}
}
if old_ttl.is_some() {
match del_ttl_sync(store, old_key) {
Ok(true) => {}
Ok(false) => return Ok(false),
Err(_) => {
output.write_resp_error("generic error");
return Ok(true);
}
}
}
match store.try_delete_sync(old_key) {
Ok(Ok(_)) => {}
Ok(Err(_)) => return Ok(false),
Err(_) => {
output.write_resp_error("generic error");
return Ok(true);
}
}
if nx {
output.write_resp_int(1);
} else {
write_raw(output, cs::RESP_OK);
}
Ok(true)
}
fn parse_expire_option(raw: &[u8], base: TtlOpt) -> TtlOpt {
let mut opt = base;
if raw.eq_ignore_ascii_case(b"NX") {
opt.nx = true;
} else if raw.eq_ignore_ascii_case(b"XX") {
opt.xx = true;
} else if raw.eq_ignore_ascii_case(b"GT") {
opt.gt = true;
} else if raw.eq_ignore_ascii_case(b"LT") {
opt.lt = true;
}
opt
}
fn merge_expire_options(first: &[u8], second: &[u8]) -> TtlOpt {
let opt = parse_expire_option(first, TtlOpt::NONE);
parse_expire_option(second, opt)
}
fn expire_apply_sync<'a, D: wdev::Device>(
store: &wkv::BatchStoreSession<'a, D>,
key: &[u8],
expire_at_ms: u64,
opt: TtlOpt,
) -> Result<Option<i32>, wkv::Error> {
match probe_alive(store, key)? {
None => Ok(None),
Some(false) => Ok(Some(0)),
Some(true) => {
let current = match ttl_of_sync(store, key)? {
None => return Ok(None),
Some(cur) => cur,
};
let denied = match current {
None => opt.xx || opt.gt,
Some(c) => opt.nx || (opt.gt && expire_at_ms <= c) || (opt.lt && expire_at_ms >= c),
};
if denied {
return Ok(Some(0));
}
if expire_at_ms <= now_unix_ms() {
if !del_ttl_sync(store, key)? {
return Ok(None);
}
return store
.try_delete_sync(key)
.map(|done| if done.is_ok() { Some(1) } else { None });
}
put_ttl_sync(store, key, expire_at_ms).map(|done| if done { Some(1) } else { None })
}
}
}
fn persist_apply_sync<'a, D: wdev::Device>(
store: &wkv::BatchStoreSession<'a, D>,
key: &[u8],
) -> Result<Option<i32>, wkv::Error> {
match probe_alive(store, key)? {
None => Ok(None),
Some(false) => Ok(Some(0)),
Some(true) => match ttl_of_sync(store, key)? {
None => Ok(None),
Some(None) => Ok(Some(0)),
Some(Some(_)) => del_ttl_sync(store, key).map(|done| if done { Some(1) } else { None }),
},
}
}
enum ExpiryRead {
Missing,
NoExpiry,
At(u64),
}
fn ttl_read_sync<'a, D: wdev::Device>(
store: &wkv::BatchStoreSession<'a, D>,
key: &[u8],
) -> Result<Option<ExpiryRead>, wkv::Error> {
match probe_alive(store, key)? {
None => Ok(None),
Some(false) => Ok(Some(ExpiryRead::Missing)),
Some(true) => match ttl_of_sync(store, key)? {
None => Ok(None),
Some(None) => Ok(Some(ExpiryRead::NoExpiry)),
Some(Some(exp)) => Ok(Some(ExpiryRead::At(exp))),
},
}
}
fn expiretime_read_sync<'a, D: wdev::Device>(
store: &wkv::BatchStoreSession<'a, D>,
key: &[u8],
) -> Result<Option<ExpiryRead>, wkv::Error> {
ttl_read_sync(store, key)
}
#[cfg(test)]
mod tests {
use super::{
super::{
batch_harness::with_batch,
ttl_sync::{data_alive_sync, now_unix_ms, put_ttl_sync, ttl_of_sync},
},
*,
};
fn dump_payload(val: &[u8]) -> Vec<u8> {
let mut payload = vec![0x00];
let mut len_buf = [0u8; 5];
let n = try_write_length(val.len() as u32, &mut len_buf).unwrap();
payload.extend_from_slice(&len_buf[..n]);
payload.extend_from_slice(val);
payload.extend_from_slice(&RDB_VERSION.to_le_bytes());
let crc = rdb_crc64::hash(&payload);
payload.extend_from_slice(&crc);
payload
}
#[test]
fn restore_writes_string_value() {
with_batch(|s, batch| {
let payload = dump_payload(b"hello");
let mut out = Vec::new();
let _ = s
.network_restore(&[b"k", b"0", &payload], batch, &mut out)
.unwrap();
assert_eq!(out, b"+OK\r\n");
let mut out = Vec::new();
let _ = s.network_getdel(&[b"k"], batch, &mut out).unwrap();
assert_eq!(out, b"$5\r\nhello\r\n");
});
}
#[test]
fn restore_busykey_and_validation_errors() {
with_batch(|s, batch| {
let payload = dump_payload(b"v");
let _ = s
.network_restore(&[b"k", b"0", &payload], batch, &mut Vec::new())
.unwrap();
let mut out = Vec::new();
let _ = s
.network_restore(&[b"k", b"0", &payload], batch, &mut out)
.unwrap();
assert_eq!(out, b"-BUSYKEY Target key name already exists.\r\n");
let mut out = Vec::new();
let _ = s.network_restore(&[b"k", b"0"], batch, &mut out).unwrap();
assert_eq!(
out,
b"-ERR wrong number of arguments for 'RESTORE' command\r\n"
);
let mut out = Vec::new();
let _ = s
.network_restore(&[b"k", b"abc", &payload], batch, &mut out)
.unwrap();
assert_eq!(out, b"-ERR timeout is not a float or out of range\r\n");
let mut bad_type = payload.clone();
bad_type[0] = 0x0c;
let mut out = Vec::new();
let _ = s
.network_restore(&[b"k2", b"0", &bad_type], batch, &mut out)
.unwrap();
assert_eq!(
out,
b"-ERR RESTORE currently only supports string types\r\n"
);
let mut out = Vec::new();
let _ = s
.network_restore(&[b"k2", b"0", &[0x00, 1, 2, 3]], batch, &mut out)
.unwrap();
assert_eq!(out, b"-ERR DUMP payload version or checksum are wrong\r\n");
let mut corrupted = dump_payload(b"v");
let last = corrupted.len() - 1;
corrupted[last] ^= 0xff;
let mut out = Vec::new();
let _ = s
.network_restore(&[b"k3", b"0", &corrupted], batch, &mut out)
.unwrap();
assert_eq!(out, b"-ERR DUMP payload version or checksum are wrong\r\n");
});
}
#[test]
fn restore_with_expiry() {
with_batch(|s, batch| {
let payload = dump_payload(b"v");
let mut out = Vec::new();
let _ = s
.network_restore(&[b"k", b"100", &payload], batch, &mut out)
.unwrap();
assert_eq!(out, b"+OK\r\n");
let ttl = ttl_of_sync(batch, b"k").unwrap().unwrap().unwrap();
assert!(ttl > now_unix_ms() + 99_000);
});
}
#[test]
fn dump_frame_roundtrips_through_restore() {
with_batch(|s, batch| {
let _ = s
.network_set(&[b"k", b"hello"], batch, &mut Vec::new())
.unwrap();
let mut out = Vec::new();
let _ = s.network_dump(&[b"k"], batch, &mut out).unwrap();
assert!(out.starts_with(b"$"));
let type_pos = out.iter().position(|b| *b == 0).unwrap();
let payload = &out[type_pos..out.len() - 2];
let mut out2 = Vec::new();
let _ = s
.network_restore(&[b"k2", b"0", payload], batch, &mut out2)
.unwrap();
assert_eq!(out2, b"+OK\r\n");
let mut out3 = Vec::new();
let _ = s.network_getdel(&[b"k2"], batch, &mut out3).unwrap();
assert_eq!(out3, b"$5\r\nhello\r\n");
let mut out4 = Vec::new();
let _ = s.network_dump(&[b"missing"], batch, &mut out4).unwrap();
assert_eq!(out4, b"$-1\r\n");
let mut out5 = Vec::new();
let _ = s.network_dump(&[], batch, &mut out5).unwrap();
assert_eq!(
out5,
b"-ERR wrong number of arguments for 'DUMP' command\r\n"
);
});
}
#[test]
fn rename_and_renamenx() {
with_batch(|s, batch| {
let _ = s
.network_set(&[b"a", b"va"], batch, &mut Vec::new())
.unwrap();
let mut out = Vec::new();
let _ = s.network_rename(&[b"a", b"b"], batch, &mut out).unwrap();
assert_eq!(out, b"+OK\r\n");
let mut out = Vec::new();
let _ = s.network_exists(&[b"a"], batch, &mut out).unwrap();
assert_eq!(out, b":0\r\n");
let mut out = Vec::new();
let _ = s.network_rename(&[b"a", b"c"], batch, &mut out).unwrap();
assert_eq!(out, b"-ERR no such key\r\n");
let mut out = Vec::new();
let _ = s.network_renamenx(&[b"b", b"c"], batch, &mut out).unwrap();
assert_eq!(out, b":1\r\n");
let _ = s
.network_set(&[b"x", b"vx"], batch, &mut Vec::new())
.unwrap();
let mut out = Vec::new();
let _ = s.network_renamenx(&[b"c", b"x"], batch, &mut out).unwrap();
assert_eq!(out, b":0\r\n");
let mut out = Vec::new();
let _ = s.network_rename(&[b"c", b"c"], batch, &mut out).unwrap();
assert_eq!(out, b"+OK\r\n");
let mut out = Vec::new();
let _ = s.network_renamenx(&[b"c", b"c"], batch, &mut out).unwrap();
assert_eq!(out, b":1\r\n");
});
}
#[test]
fn rename_transfers_ttl() {
with_batch(|s, batch| {
let _ = s
.network_set(&[b"k", b"v"], batch, &mut Vec::new())
.unwrap();
let mut out = Vec::new();
let _ = s
.network_expire(ExpireCmd::Expire, &[b"k", b"100"], batch, &mut out)
.unwrap();
assert_eq!(out, b":1\r\n");
let mut out = Vec::new();
let _ = s.network_rename(&[b"k", b"k2"], batch, &mut out).unwrap();
assert_eq!(out, b"+OK\r\n");
let ttl = ttl_of_sync(batch, b"k2").unwrap().unwrap().unwrap();
assert!(ttl > now_unix_ms() + 99_000);
assert_eq!(ttl_of_sync(batch, b"k").unwrap(), Some(None));
assert_eq!(data_alive_sync(batch, b"k").unwrap(), Some(false));
let _ = s
.network_set(&[b"k3", b"v3"], batch, &mut Vec::new())
.unwrap();
let mut out = Vec::new();
let _ = s
.network_renamenx(&[b"k3", b"k4"], batch, &mut out)
.unwrap();
assert_eq!(out, b":1\r\n");
assert_eq!(ttl_of_sync(batch, b"k4").unwrap(), Some(None));
let _ = s
.network_set(&[b"a", b"va"], batch, &mut Vec::new())
.unwrap();
let _ = put_ttl_sync(batch, b"a", now_unix_ms() + 60_000).unwrap();
let mut out = Vec::new();
let _ = s.network_renamenx(&[b"k4", b"a"], batch, &mut out).unwrap();
assert_eq!(out, b":0\r\n");
let ttl = ttl_of_sync(batch, b"a").unwrap().unwrap().unwrap();
assert!(ttl > now_unix_ms() + 59_000 && ttl < now_unix_ms() + 61_000);
});
}
#[test]
fn expire_validation_and_effects() {
with_batch(|s, batch| {
let _ = s
.network_set(&[b"k", b"v"], batch, &mut Vec::new())
.unwrap();
let mut out = Vec::new();
let _ = s
.network_expire(ExpireCmd::Expire, &[b"k", b"100"], batch, &mut out)
.unwrap();
assert_eq!(out, b":1\r\n");
let ttl = ttl_of_sync(batch, b"k").unwrap().unwrap().unwrap();
assert!(ttl > now_unix_ms() + 99_000);
let mut out = Vec::new();
let _ = s
.network_expire(ExpireCmd::Expire, &[b"k", b"100", b"NX"], batch, &mut out)
.unwrap();
assert_eq!(out, b":0\r\n");
let mut out = Vec::new();
let _ = s
.network_expire(ExpireCmd::Expire, &[b"k", b"1", b"GT"], batch, &mut out)
.unwrap();
assert_eq!(out, b":0\r\n");
let mut out = Vec::new();
let _ = s
.network_expire(
ExpireCmd::Pexpire,
&[b"k", b"50000", b"XX"],
batch,
&mut out,
)
.unwrap();
assert_eq!(out, b":1\r\n");
let mut out = Vec::new();
let _ = s
.network_expire(
ExpireCmd::Expire,
&[b"k", b"1", b"XX", b"GT"],
batch,
&mut out,
)
.unwrap();
assert_eq!(out, b":0\r\n");
let mut out = Vec::new();
let _ = s
.network_expire(
ExpireCmd::Expire,
&[b"k", b"1", b"NX", b"GT"],
batch,
&mut out,
)
.unwrap();
assert_eq!(
out,
b"-ERR NX and XX, GT or LT options at the same time are not compatible\r\n"
);
let mut out = Vec::new();
let _ = s
.network_expire(ExpireCmd::Expire, &[b"k", b"1", b"WHAT"], batch, &mut out)
.unwrap();
assert_eq!(out, b"-ERR Unsupported option WHAT\r\n");
let mut out = Vec::new();
let _ = s
.network_expire(ExpireCmd::Expire, &[b"k", b"-1"], batch, &mut out)
.unwrap();
assert_eq!(out, b"-ERR invalid expire time, must be >= 0\r\n");
let mut out = Vec::new();
let _ = s
.network_expire(ExpireCmd::Expire, &[b"k", b"abc"], batch, &mut out)
.unwrap();
assert_eq!(out, b"-ERR value is not an integer or out of range.\r\n");
let mut out = Vec::new();
let _ = s
.network_expire(ExpireCmd::Expire, &[b"nope", b"100"], batch, &mut out)
.unwrap();
assert_eq!(out, b":0\r\n");
let mut out = Vec::new();
let abs = (now_unix_ms() / 1000 + 100).to_string();
let _ = s
.network_expire(
ExpireCmd::Expireat,
&[b"k", abs.as_bytes()],
batch,
&mut out,
)
.unwrap();
assert_eq!(out, b":1\r\n");
});
}
#[test]
fn persist_ttl_expiretime_semantics() {
with_batch(|s, batch| {
let _ = s
.network_set(&[b"k", b"v"], batch, &mut Vec::new())
.unwrap();
let mut out = Vec::new();
let _ = s.network_persist(&[b"k"], batch, &mut out).unwrap();
assert_eq!(out, b":0\r\n");
let mut out = Vec::new();
let _ = s
.network_ttl(TtlCmd::Ttl, &[b"k"], batch, &mut out)
.unwrap();
assert_eq!(out, b":-1\r\n");
let mut out = Vec::new();
let _ = s
.network_expiretime(ExpireTimeCmd::Expiretime, &[b"k"], batch, &mut out)
.unwrap();
assert_eq!(out, b":-1\r\n");
let _ = put_ttl_sync(batch, b"k", now_unix_ms() + 60_000).unwrap();
let mut out = Vec::new();
let _ = s
.network_ttl(TtlCmd::Ttl, &[b"k"], batch, &mut out)
.unwrap();
assert!(out != b":-1\r\n" && out != b":-2\r\n");
let mut out = Vec::new();
let _ = s
.network_ttl(TtlCmd::Pttl, &[b"k"], batch, &mut out)
.unwrap();
assert!(out != b":-1\r\n" && out != b":-2\r\n");
let mut out = Vec::new();
let _ = s
.network_expiretime(ExpireTimeCmd::Pexpiretime, &[b"k"], batch, &mut out)
.unwrap();
assert!(out != b":-1\r\n");
let mut out = Vec::new();
let _ = s.network_persist(&[b"k"], batch, &mut out).unwrap();
assert_eq!(out, b":1\r\n");
assert_eq!(ttl_of_sync(batch, b"k").unwrap(), Some(None));
let mut out = Vec::new();
let _ = s
.network_ttl(TtlCmd::Ttl, &[b"nope"], batch, &mut out)
.unwrap();
assert_eq!(out, b":-2\r\n");
let mut out = Vec::new();
let _ = s.network_persist(&[b"nope"], batch, &mut out).unwrap();
assert_eq!(out, b":0\r\n");
let mut out = Vec::new();
let _ = s
.network_ttl(TtlCmd::Ttl, &[b"a", b"b"], batch, &mut out)
.unwrap();
assert_eq!(out, b"-ERR wrong number of arguments for 'TTL' command\r\n");
});
}
#[test]
fn getdel_and_exists() {
with_batch(|s, batch| {
let _ = s
.network_set(&[b"k", b"v"], batch, &mut Vec::new())
.unwrap();
let mut out = Vec::new();
let _ = s.network_getdel(&[b"k"], batch, &mut out).unwrap();
assert_eq!(out, b"$1\r\nv\r\n");
let mut out = Vec::new();
let _ = s.network_getdel(&[b"k"], batch, &mut out).unwrap();
assert_eq!(out, b"$-1\r\n");
let _ = s
.network_set(&[b"x", b"vx"], batch, &mut Vec::new())
.unwrap();
let mut out = Vec::new();
let _ = s
.network_exists(&[b"x", b"x", b"nope"], batch, &mut out)
.unwrap();
assert_eq!(out, b":2\r\n");
});
}
}