Skip to content

Commit c100eef

Browse files
committed
fix(storage): align DEL with generation lifecycle
1 parent cbc2895 commit c100eef

2 files changed

Lines changed: 145 additions & 69 deletions

File tree

src/storage/src/redis_strings.rs

Lines changed: 69 additions & 68 deletions
Original file line numberDiff line numberDiff line change
@@ -1993,7 +1993,20 @@ impl Redis {
19931993
});
19941994
};
19951995

1996-
if value.is_empty() || self.is_stale(&value)? {
1996+
let is_live = if value.is_empty() {
1997+
false
1998+
} else {
1999+
match DataType::try_from(value[0])? {
2000+
DataType::String => !ParsedStringsValue::new(&value[..])?.is_stale(),
2001+
DataType::Hash | DataType::Set | DataType::ZSet => {
2002+
ParsedBaseMetaValue::new(&value[..])?.is_valid()
2003+
}
2004+
DataType::List => ParsedListsMetaValue::new(&value[..])?.is_valid(),
2005+
_ => false,
2006+
}
2007+
};
2008+
2009+
if !is_live {
19972010
return Err(crate::error::Error::KeyNotFound {
19982011
key: String::from_utf8_lossy(key).to_string(),
19992012
location: snafu::location!(),
@@ -2005,9 +2018,6 @@ impl Redis {
20052018

20062019
/// Delete a key (works for all data types)
20072020
pub fn del_key(&self, key: &[u8]) -> Result<bool> {
2008-
use crate::storage_define::{PREFIX_RESERVE_LENGTH, encode_user_key};
2009-
use bytes::BufMut;
2010-
20112021
let db = self.db.as_ref().context(OptionNoneSnafu {
20122022
message: "db is not initialized".to_string(),
20132023
})?;
@@ -2019,77 +2029,68 @@ impl Redis {
20192029
message: "MetaCF is not initialized".to_string(),
20202030
})?;
20212031

2022-
// Check if key exists in MetaCF (where all metadata is stored)
2023-
let string_key = BaseKey::new(key);
2032+
let key_str = String::from_utf8_lossy(key).to_string();
2033+
let _lock = ScopeRecordLock::new(self.lock_mgr.as_ref(), &key_str);
20242034
let meta_key = BaseMetaKey::new(key);
2025-
2026-
let string_existed = db
2027-
.get_cf_opt(&meta_cf, &string_key.encode()?, &self.read_options)
2028-
.context(RocksSnafu)?
2029-
.is_some();
2030-
let meta_existed = db
2031-
.get_cf_opt(&meta_cf, &meta_key.encode()?, &self.read_options)
2035+
let encoded_meta_key = meta_key.encode()?;
2036+
let Some(value) = db
2037+
.get_cf_opt(&meta_cf, &encoded_meta_key, &self.read_options)
20322038
.context(RocksSnafu)?
2033-
.is_some();
2034-
2035-
if string_existed || meta_existed {
2036-
let encoded = string_key.encode()?;
2037-
2038-
// Build correct prefix for data CF scanning:
2039-
// Data keys format: | reserve1 (8B) | encoded_user_key | version (8B) | data | reserve2 |
2040-
// We need prefix: | reserve1 (8B) | encoded_user_key (with \x00\x00 delimiter) |
2041-
// Note: BaseKey.encode() includes reserve2, which would not match data keys.
2042-
let mut data_key_prefix =
2043-
bytes::BytesMut::with_capacity(PREFIX_RESERVE_LENGTH + key.len() * 2 + 2);
2044-
data_key_prefix.put_slice(&[0u8; PREFIX_RESERVE_LENGTH]);
2045-
encode_user_key(key, &mut data_key_prefix)?;
2046-
2047-
// Collect all keys to delete first (to avoid borrow conflicts)
2048-
let mut keys_to_delete: Vec<(ColumnFamilyIndex, Vec<u8>)> = Vec::new();
2049-
2050-
// Delete from MetaCF
2051-
if string_existed {
2052-
keys_to_delete.push((ColumnFamilyIndex::MetaCF, encoded.to_vec()));
2053-
}
2054-
if meta_existed {
2055-
keys_to_delete.push((ColumnFamilyIndex::MetaCF, meta_key.encode()?.to_vec()));
2056-
}
2039+
else {
2040+
return Ok(false);
2041+
};
2042+
if value.is_empty() {
2043+
return Ok(false);
2044+
}
20572045

2058-
// For composite data types, perform prefix scan to delete all related entries
2059-
for cf_index in [
2060-
ColumnFamilyIndex::HashesDataCF,
2061-
ColumnFamilyIndex::SetsDataCF,
2062-
ColumnFamilyIndex::ListsDataCF,
2063-
ColumnFamilyIndex::ZsetsDataCF,
2064-
ColumnFamilyIndex::ZsetsScoreCF,
2065-
] {
2066-
let cf_handle = self.get_cf_handle(cf_index);
2067-
if let Some(cf) = cf_handle {
2068-
// Prefix-scan data CF and delete all derived keys
2069-
let iter = db.iterator_cf(
2070-
&cf,
2071-
rocksdb::IteratorMode::From(&data_key_prefix, rocksdb::Direction::Forward),
2072-
);
2073-
for item in iter {
2074-
let (k, _) = item.context(RocksSnafu)?;
2075-
if !k.starts_with(&data_key_prefix) {
2076-
break;
2077-
}
2078-
keys_to_delete.push((cf_index, k.to_vec()));
2079-
}
2046+
match DataType::try_from(value[0])? {
2047+
DataType::String => {
2048+
let parsed = ParsedStringsValue::new(&value[..])?;
2049+
if parsed.is_stale() {
2050+
return Ok(false);
20802051
}
2052+
let mut batch = self.create_batch()?;
2053+
batch.delete(ColumnFamilyIndex::MetaCF, &encoded_meta_key)?;
2054+
batch.commit()?;
2055+
self.update_specific_key_statistics(DataType::String, &key_str, 1)?;
20812056
}
2082-
2083-
// Now create batch and delete all collected keys
2084-
let mut batch = self.create_batch()?;
2085-
for (cf_idx, key) in keys_to_delete {
2086-
batch.delete(cf_idx, &key)?;
2057+
DataType::Hash | DataType::Set | DataType::ZSet => {
2058+
let mut parsed = ParsedBaseMetaValue::new(&value[..])?;
2059+
if !parsed.is_valid() {
2060+
return Ok(false);
2061+
}
2062+
let count = parsed.count();
2063+
parsed.set_count(0);
2064+
parsed.update_version();
2065+
parsed.set_etime(0);
2066+
let data_type = DataType::try_from(value[0])?;
2067+
let mut batch = self.create_batch()?;
2068+
batch.put(
2069+
ColumnFamilyIndex::MetaCF,
2070+
&encoded_meta_key,
2071+
parsed.encoded(),
2072+
)?;
2073+
batch.commit()?;
2074+
self.update_specific_key_statistics(data_type, &key_str, count)?;
20872075
}
2088-
batch.commit()?;
2089-
Ok(true)
2090-
} else {
2091-
Ok(false)
2076+
DataType::List => {
2077+
let mut parsed = ParsedListsMetaValue::new(&value[..])?;
2078+
if !parsed.is_valid() {
2079+
return Ok(false);
2080+
}
2081+
let count = parsed.count();
2082+
parsed.set_count(0);
2083+
parsed.update_version();
2084+
parsed.set_etime(0);
2085+
let mut batch = self.create_batch()?;
2086+
batch.put(ColumnFamilyIndex::MetaCF, &encoded_meta_key, parsed.value())?;
2087+
batch.commit()?;
2088+
self.update_specific_key_statistics(DataType::List, &key_str, count)?;
2089+
}
2090+
_ => return Ok(false),
20922091
}
2092+
2093+
Ok(true)
20932094
}
20942095

20952096
/// Scan for keys matching a pattern

src/storage/tests/ttl_test.rs

Lines changed: 76 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,10 @@
2020
use std::sync::Arc;
2121

2222
use storage::ZsetScoreMember;
23-
use storage::{StorageOptions, storage::Storage, unique_test_db_path};
23+
use storage::{
24+
BaseMetaKey, ColumnFamilyIndex, StorageOptions, format_base_meta_value::ParsedBaseMetaValue,
25+
storage::Storage, unique_test_db_path,
26+
};
2427

2528
#[tokio::test]
2629
async fn test_ttl_basic_operations() {
@@ -337,6 +340,78 @@ async fn test_del_command() {
337340
storage.shutdown().await;
338341
}
339342

343+
#[tokio::test]
344+
async fn test_del_returns_zero_for_stale_string() {
345+
let db_path = unique_test_db_path();
346+
let mut storage = Storage::new(1, 0);
347+
let options = Arc::new(StorageOptions::default());
348+
let _receiver = storage.open(options, &db_path).unwrap();
349+
350+
let key = b"stale_del_key";
351+
storage.set(key, b"value").unwrap();
352+
assert!(storage.pexpire(key, 1).unwrap());
353+
tokio::time::sleep(std::time::Duration::from_millis(5)).await;
354+
assert_eq!(storage.del(&[key.to_vec()]).unwrap(), 0);
355+
356+
storage.shutdown().await;
357+
}
358+
359+
#[tokio::test]
360+
async fn test_del_composite_recreate_hides_old_members() {
361+
let db_path = unique_test_db_path();
362+
let mut storage = Storage::new(1, 0);
363+
let options = Arc::new(StorageOptions::default());
364+
let _receiver = storage.open(options, &db_path).unwrap();
365+
366+
let key = b"del_recreate_hash";
367+
storage.hset(key, b"old-field", b"old-value").unwrap();
368+
assert_eq!(storage.del(&[key.to_vec()]).unwrap(), 1);
369+
assert_eq!(storage.exists(&[key.to_vec()]).unwrap(), 0);
370+
assert_eq!(storage.key_type(key).unwrap(), "none");
371+
assert_eq!(storage.ttl(key).unwrap(), -2);
372+
assert!(storage.keys(b"del_recreate_hash").unwrap().is_empty());
373+
storage.hset(key, b"new-field", b"new-value").unwrap();
374+
assert_eq!(storage.hget(key, b"old-field").unwrap(), None);
375+
assert_eq!(
376+
storage.hget(key, b"new-field").unwrap(),
377+
Some("new-value".to_string())
378+
);
379+
380+
storage.shutdown().await;
381+
}
382+
383+
#[tokio::test]
384+
async fn test_del_writes_composite_meta_tombstone_with_new_version() {
385+
let db_path = unique_test_db_path();
386+
let mut storage = Storage::new(1, 0);
387+
let options = Arc::new(StorageOptions::default());
388+
let _receiver = storage.open(options, &db_path).unwrap();
389+
390+
let key = b"del_meta_tombstone";
391+
storage.hset(key, b"field", b"value").unwrap();
392+
let meta_key = BaseMetaKey::new(key).encode().unwrap();
393+
let original_version = {
394+
let instance = &storage.insts[0];
395+
let meta_cf = instance.get_cf_handle(ColumnFamilyIndex::MetaCF).unwrap();
396+
let db = instance.db().unwrap();
397+
let original = db.get_cf(&meta_cf, &meta_key).unwrap().unwrap();
398+
ParsedBaseMetaValue::new(&original[..]).unwrap().version()
399+
};
400+
401+
assert_eq!(storage.del(&[key.to_vec()]).unwrap(), 1);
402+
{
403+
let instance = &storage.insts[0];
404+
let meta_cf = instance.get_cf_handle(ColumnFamilyIndex::MetaCF).unwrap();
405+
let db = instance.db().unwrap();
406+
let tombstone = db.get_cf(&meta_cf, &meta_key).unwrap().unwrap();
407+
let parsed = ParsedBaseMetaValue::new(&tombstone[..]).unwrap();
408+
assert_eq!(parsed.count(), 0);
409+
assert!(parsed.version() > original_version);
410+
}
411+
412+
storage.shutdown().await;
413+
}
414+
340415
#[tokio::test]
341416
async fn test_keys_command() {
342417
let db_path = unique_test_db_path();

0 commit comments

Comments
 (0)