Skip to content

Commit 874e45d

Browse files
authored
Merge pull request #1 from happy-v587/codex/compact-del-validation
Codex/compact del validation
2 parents 27e9bef + d2a206a commit 874e45d

4 files changed

Lines changed: 201 additions & 4 deletions

File tree

src/storage/src/data_compaction_filter.rs

Lines changed: 42 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,7 @@ const DATA_FILTER_FACTORY_NAME: &std::ffi::CStr = c"DataCompactionFilterFactory"
4848
struct CompactionFilterTestState {
4949
entered: bool,
5050
released: bool,
51+
removed_keys: Vec<Vec<u8>>,
5152
}
5253

5354
#[cfg(test)]
@@ -88,6 +89,27 @@ impl CompactionFilterTestGate {
8889
self.changed.notify_all();
8990
}
9091

92+
pub(crate) fn wait_until_removed(&self, timeout: std::time::Duration) -> Vec<Vec<u8>> {
93+
let state = self
94+
.state
95+
.lock()
96+
.expect("compaction filter test gate mutex should not be poisoned");
97+
let (state, _) = self
98+
.changed
99+
.wait_timeout_while(state, timeout, |state| state.removed_keys.is_empty())
100+
.expect("compaction filter test gate wait should succeed");
101+
state.removed_keys.clone()
102+
}
103+
104+
fn record_removed_key(&self, key: &[u8]) {
105+
let mut state = self
106+
.state
107+
.lock()
108+
.expect("compaction filter test gate mutex should not be poisoned");
109+
state.removed_keys.push(key.to_vec());
110+
self.changed.notify_all();
111+
}
112+
91113
fn enter_and_wait(&self, key: &[u8]) {
92114
if !key.starts_with(&self.key_prefix) {
93115
return;
@@ -161,6 +183,18 @@ fn block_once_for_compaction_filter_test(key: &[u8]) {
161183
}
162184
}
163185

186+
#[cfg(test)]
187+
fn record_compaction_filter_remove(key: &[u8]) {
188+
let gate = compaction_filter_test_gate()
189+
.lock()
190+
.expect("compaction filter test gate registry should not be poisoned")
191+
.as_ref()
192+
.and_then(Weak::upgrade);
193+
if let Some(gate) = gate {
194+
gate.record_removed_key(key);
195+
}
196+
}
197+
164198
#[derive(Debug)]
165199
enum MetaLookup {
166200
Valid,
@@ -326,7 +360,7 @@ impl CompactionFilter for DataCompactionFilter {
326360
return CompactionDecision::Keep;
327361
};
328362

329-
match self.ensure_meta_state(&meta_key) {
363+
let decision = match self.ensure_meta_state(&meta_key) {
330364
MetaLookup::Unavailable => CompactionDecision::Keep,
331365
MetaLookup::NotFound => CompactionDecision::Remove,
332366
MetaLookup::Valid => {
@@ -340,7 +374,14 @@ impl CompactionFilter for DataCompactionFilter {
340374
_ => CompactionDecision::Keep,
341375
}
342376
}
377+
};
378+
379+
#[cfg(test)]
380+
if let CompactionDecision::Remove = &decision {
381+
record_compaction_filter_remove(key);
343382
}
383+
384+
decision
344385
}
345386
}
346387

src/storage/src/redis.rs

Lines changed: 92 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1010,6 +1010,7 @@ macro_rules! get_db_and_cfs {
10101010
#[cfg(test)]
10111011
mod lifecycle_tests {
10121012
use std::sync::{Arc, mpsc};
1013+
use std::thread;
10131014
use std::time::Duration;
10141015

10151016
use super::{
@@ -1018,8 +1019,11 @@ mod lifecycle_tests {
10181019
use crate::data_compaction_filter::{
10191020
CompactionFilterTestGate, install_compaction_filter_test_gate,
10201021
};
1022+
use crate::format_base_key::BaseMetaKey;
1023+
use crate::format_base_meta_value::ParsedSetsMetaValue;
10211024
use crate::{BgTaskHandler, StorageOptions, safe_cleanup_test_db, unique_test_db_path};
10221025
use kstd::lock_mgr::LockMgr;
1026+
use rocksdb::IteratorMode;
10231027

10241028
fn open_compaction_test_redis(path: &std::path::Path) -> Redis {
10251029
let mut storage_options = StorageOptions::default();
@@ -1043,6 +1047,94 @@ mod lifecycle_tests {
10431047
redis
10441048
}
10451049

1050+
#[test]
1051+
fn del_compaction_removes_old_generation_and_preserves_recreated_set() {
1052+
let path = unique_test_db_path();
1053+
safe_cleanup_test_db(&path);
1054+
1055+
let mut storage_options = StorageOptions::default();
1056+
storage_options.options.set_disable_auto_compactions(true);
1057+
let (bg_task_handler, _) = BgTaskHandler::new();
1058+
let mut redis = Redis::new(
1059+
Arc::new(storage_options),
1060+
0,
1061+
Arc::new(bg_task_handler),
1062+
Arc::new(kstd::lock_mgr::LockMgr::new(64)),
1063+
);
1064+
redis
1065+
.open(path.to_str().expect("test DB path should be valid UTF-8"))
1066+
.expect("compaction test Redis should open");
1067+
let redis = Arc::new(redis);
1068+
1069+
let key = b"del_compaction_set";
1070+
redis
1071+
.sadd(key, &[b"old-member"])
1072+
.expect("initial set write should succeed");
1073+
let db = redis.db().expect("Redis should own RocksDB");
1074+
db.flush().expect("initial data should be flushed");
1075+
1076+
let meta_key = BaseMetaKey::new(key).encode().unwrap();
1077+
let meta_cf = redis.get_cf_handle(ColumnFamilyIndex::MetaCF).unwrap();
1078+
let original_meta = db.get_cf(&meta_cf, &meta_key).unwrap().unwrap();
1079+
let original_version = ParsedSetsMetaValue::new(&original_meta[..])
1080+
.unwrap()
1081+
.version();
1082+
let data_cf = redis.get_cf_handle(ColumnFamilyIndex::SetsDataCF).unwrap();
1083+
let old_data_keys: Vec<Vec<u8>> = db
1084+
.iterator_cf(&data_cf, IteratorMode::Start)
1085+
.map(|entry| entry.unwrap().0.to_vec())
1086+
.collect();
1087+
assert_eq!(old_data_keys.len(), 1);
1088+
1089+
assert!(redis.del_key(key).unwrap());
1090+
db.flush().expect("tombstone should be flushed");
1091+
let tombstone = db.get_cf(&meta_cf, &meta_key).unwrap().unwrap();
1092+
let parsed_tombstone = ParsedSetsMetaValue::new(&tombstone[..]).unwrap();
1093+
assert_eq!(parsed_tombstone.count(), 0);
1094+
assert_eq!(parsed_tombstone.etime(), 0);
1095+
assert!(parsed_tombstone.version() > original_version);
1096+
1097+
let gate = CompactionFilterTestGate::new(&[]);
1098+
let _gate_guard = install_compaction_filter_test_gate(Arc::clone(&gate));
1099+
let compactor = Arc::clone(&redis);
1100+
let compaction_thread = thread::spawn(move || {
1101+
compactor
1102+
.compact_range(None, None)
1103+
.expect("manual compaction should succeed");
1104+
});
1105+
assert!(gate.wait_until_entered(Duration::from_secs(10)));
1106+
gate.release();
1107+
compaction_thread.join().unwrap();
1108+
1109+
let removed_keys = gate.wait_until_removed(Duration::from_secs(10));
1110+
assert!(
1111+
old_data_keys
1112+
.iter()
1113+
.any(|old_key| removed_keys.contains(old_key)),
1114+
"DataCompactionFilter did not remove the old generation key"
1115+
);
1116+
let remaining_data_keys: Vec<Vec<u8>> = db
1117+
.iterator_cf(&data_cf, IteratorMode::Start)
1118+
.map(|entry| entry.unwrap().0.to_vec())
1119+
.collect();
1120+
assert!(
1121+
old_data_keys
1122+
.iter()
1123+
.all(|old_key| !remaining_data_keys.contains(old_key)),
1124+
"old generation key survived compaction"
1125+
);
1126+
1127+
redis
1128+
.sadd(key, &[b"new-member"])
1129+
.expect("recreated set write should succeed");
1130+
assert_eq!(redis.smembers(key).unwrap(), vec!["new-member".to_string()]);
1131+
1132+
drop(data_cf);
1133+
drop(meta_cf);
1134+
drop(redis);
1135+
safe_cleanup_test_db(&path);
1136+
}
1137+
10461138
#[test]
10471139
fn dropping_last_owner_waits_for_active_compaction_filter_before_reopen() {
10481140
let path = unique_test_db_path();

src/storage/src/storage.rs

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -402,8 +402,15 @@ impl Storage {
402402
Ok(())
403403
}
404404

405-
fn do_compact_range(&self, _dtype: DataType, _start: &str, _end: &str) -> Result<()> {
406-
log::info!("do_compact_range {_dtype:?} {_start} {_end}");
405+
fn do_compact_range(&self, dtype: DataType, start: &str, end: &str) -> Result<()> {
406+
log::info!("do_compact_range {dtype:?} {start} {end}");
407+
408+
let begin = (!start.is_empty()).then_some(start.as_bytes());
409+
let end = (!end.is_empty()).then_some(end.as_bytes());
410+
for instance in &self.insts {
411+
instance.compact_range(begin, end)?;
412+
}
413+
407414
Ok(())
408415
}
409416

src/storage/tests/storage_basic_test.rs

Lines changed: 58 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,8 +19,9 @@
1919

2020
use std::sync::Arc;
2121

22+
use rocksdb::IteratorMode;
2223
use storage::storage::Storage;
23-
use storage::{BgTask, BgTaskHandler, DataType, StorageOptions};
24+
use storage::{BgTask, BgTaskHandler, ColumnFamilyIndex, DataType, StorageOptions};
2425

2526
// This test ensures:
2627
// - All tasks are sent successfully (no panic)
@@ -69,6 +70,62 @@ async fn test_bg_task_worker_concurrent() {
6970
worker_handle.await.unwrap();
7071
}
7172

73+
#[tokio::test]
74+
async fn test_storage_compact_range_executes_manual_compaction() {
75+
let temp_dir = tempfile::tempdir().unwrap();
76+
let mut storage_options = StorageOptions::default();
77+
storage_options.options.set_disable_auto_compactions(true);
78+
79+
let mut storage = Storage::new(1, 0);
80+
let _receiver = storage
81+
.open(Arc::new(storage_options), temp_dir.path())
82+
.unwrap();
83+
84+
storage.sadd(b"compact_range_key", &[b"member"]).unwrap();
85+
storage.insts[0].db().unwrap().flush().unwrap();
86+
87+
let old_data_keys: Vec<Vec<u8>> = {
88+
let data_cf = storage.insts[0]
89+
.get_cf_handle(ColumnFamilyIndex::SetsDataCF)
90+
.unwrap();
91+
storage.insts[0]
92+
.db()
93+
.unwrap()
94+
.iterator_cf(&data_cf, IteratorMode::Start)
95+
.map(|entry| entry.unwrap().0.to_vec())
96+
.collect()
97+
};
98+
assert_eq!(old_data_keys.len(), 1);
99+
100+
storage.del(&[b"compact_range_key".to_vec()]).unwrap();
101+
storage.insts[0].db().unwrap().flush().unwrap();
102+
103+
storage
104+
.compact_range(DataType::All, "", "", true)
105+
.await
106+
.unwrap();
107+
108+
let remaining_data_keys: Vec<Vec<u8>> = {
109+
let data_cf = storage.insts[0]
110+
.get_cf_handle(ColumnFamilyIndex::SetsDataCF)
111+
.unwrap();
112+
storage.insts[0]
113+
.db()
114+
.unwrap()
115+
.iterator_cf(&data_cf, IteratorMode::Start)
116+
.map(|entry| entry.unwrap().0.to_vec())
117+
.collect()
118+
};
119+
assert!(
120+
old_data_keys
121+
.iter()
122+
.all(|old_key| !remaining_data_keys.contains(old_key)),
123+
"stale set data key survived compaction"
124+
);
125+
126+
storage.shutdown().await;
127+
}
128+
72129
/// Test get_global_smallest_flushed_log_index returns minimum across all instances
73130
#[tokio::test]
74131
async fn test_get_global_smallest_flushed_log_index() {

0 commit comments

Comments
 (0)