Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 27 additions & 2 deletions src/commands/cmd_cuckoo_filter.cc
Original file line number Diff line number Diff line change
Expand Up @@ -131,8 +131,33 @@ class CommandCFAdd : public Commander {
}
};

// Register the CF.RESERVE and CF.ADD commands
class CommandCFDel : public Commander {
public:
Status Parse(const std::vector<std::string> &args) override {
// CF.DEL key item
if (args.size() != 3) {
return {Status::RedisParseErr, errWrongNumOfArguments};
}
return Commander::Parse(args);
}

Status Execute(engine::Context &ctx, Server *srv, Connection *conn, std::string *output) override {
redis::CuckooChain cuckoo_db(srv->storage, conn->GetNamespace());
bool deleted = false;
auto s = cuckoo_db.Delete(ctx, args_[1], args_[2], &deleted);

if (!s.ok()) {
return {Status::RedisExecErr, s.ToString()};
}

*output = redis::Integer(deleted ? 1 : 0);
return Status::OK();
}
};

// Register the CF.RESERVE, CF.ADD and CF.DEL commands
REDIS_REGISTER_COMMANDS(CuckooFilter, MakeCmdAttr<CommandCFReserve>("cf.reserve", -3, "write", 1, 1, 1),
MakeCmdAttr<CommandCFAdd>("cf.add", 3, "write", 1, 1, 1))
MakeCmdAttr<CommandCFAdd>("cf.add", 3, "write", 1, 1, 1),
MakeCmdAttr<CommandCFDel>("cf.del", 3, "write", 1, 1, 1))

} // namespace redis
1 change: 1 addition & 0 deletions src/types/cuckoo_filter.h
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ constexpr long double kCuckooFilterLoadFactor = 0.955L;
constexpr uint64_t kCuckooFilterMaxSupportedBuckets = std::numeric_limits<uint32_t>::max() / 2 + 1ULL;
constexpr uint64_t kCuckooFilterFingerprintModulus = 255;
constexpr uint64_t kCuckooFilterAltHashMultiplier = 0x5bd1e995ULL;
constexpr uint8_t kEmptyCuckooFingerprint = 0;

// Cuckoo filter implementation from the paper:
// "Cuckoo Filter: Practically Better Than Bloom" by Fan et al.
Expand Down
13 changes: 12 additions & 1 deletion src/types/cuckoo_filter_page.cc
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,13 @@ rocksdb::Status CuckooPageCache::SetBucketSlot(uint16_t filter_index, uint32_t n
rocksdb::Status CuckooPageCache::WriteBackDirtyPages(rocksdb::WriteBatchBase *batch) {
for (const auto &entry : pages_) {
if (!entry.second.is_dirty) continue;
if (std::all_of(entry.second.data.begin(), entry.second.data.end(), [](char value) { return value == 0; })) {
if (entry.second.exists) {
auto s = batch->Delete(entry.first);
if (!s.ok()) return s;
}
continue;
}
auto s = batch->Put(entry.first, entry.second.data);
if (!s.ok()) return s;
}
Expand Down Expand Up @@ -167,6 +174,7 @@ rocksdb::Status CuckooPageCache::loadPage(const BucketLocation &location, PageEn
PageEntry page_entry;
auto s = storage_->Get(ctx_, ctx_.GetReadOptions(), location.page_key, &page_entry.data);
if (!s.ok() && !s.IsNotFound()) return s;
page_entry.exists = s.ok();
s = normalizePage(s, location.expected_page_size, &page_entry);
if (!s.ok()) return s;

Expand All @@ -193,7 +201,10 @@ rocksdb::Status CuckooPageCache::loadPages(const std::vector<BucketLocation> &lo

for (size_t i = 0; i < locations.size(); ++i) {
PageEntry page_entry;
if (statuses[i].ok()) page_entry.data.assign(values[i].data(), values[i].size());
if (statuses[i].ok()) {
page_entry.data.assign(values[i].data(), values[i].size());
page_entry.exists = true;
}
auto s = normalizePage(statuses[i], locations[i].expected_page_size, &page_entry);
if (!s.ok()) return s;
pages_.emplace(locations[i].page_key, std::move(page_entry));
Expand Down
1 change: 1 addition & 0 deletions src/types/cuckoo_filter_page.h
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@ class CuckooPageCache {
private:
struct PageEntry {
std::string data;
bool exists = false;
bool is_dirty = false;
};

Expand Down
32 changes: 32 additions & 0 deletions src/types/cuckoo_filter_sub_filter.cc
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,38 @@ rocksdb::Status CuckooSubFilter::TryInsert(uint64_t hash, uint8_t fingerprint, b
return pages_.TryInsertInBucket(filter_index_, num_buckets_, bucket2_idx, fingerprint, inserted);
}

rocksdb::Status CuckooSubFilter::Delete(uint64_t hash, uint8_t fingerprint, bool *deleted) {
*deleted = false;
uint32_t bucket1_idx = getPrimaryBucketIndex(hash);
uint32_t bucket2_idx = getSecondaryBucketIndex(hash, fingerprint);
auto s = pages_.PrefetchBuckets(filter_index_, num_buckets_, bucket1_idx, bucket2_idx);
if (!s.ok()) return s;

for (uint32_t slot_idx = 0; slot_idx < bucket_size_; ++slot_idx) {
uint8_t current_fingerprint = kEmptyCuckooFingerprint;
s = pages_.GetBucketSlot(filter_index_, num_buckets_, bucket1_idx, slot_idx, &current_fingerprint);
if (!s.ok()) return s;
if (current_fingerprint != fingerprint) continue;

*deleted = true;
return pages_.SetBucketSlot(filter_index_, num_buckets_, bucket1_idx, slot_idx, kEmptyCuckooFingerprint);
}

if (bucket1_idx == bucket2_idx) return rocksdb::Status::OK();

for (uint32_t slot_idx = 0; slot_idx < bucket_size_; ++slot_idx) {
uint8_t current_fingerprint = kEmptyCuckooFingerprint;
s = pages_.GetBucketSlot(filter_index_, num_buckets_, bucket2_idx, slot_idx, &current_fingerprint);
if (!s.ok()) return s;
if (current_fingerprint != fingerprint) continue;

*deleted = true;
return pages_.SetBucketSlot(filter_index_, num_buckets_, bucket2_idx, slot_idx, kEmptyCuckooFingerprint);
}

return rocksdb::Status::OK();
}

rocksdb::Status CuckooSubFilter::TryKickOutInsert(uint64_t hash, uint8_t fingerprint, uint16_t max_iterations,
bool *inserted) {
*inserted = false;
Expand Down
4 changes: 1 addition & 3 deletions src/types/cuckoo_filter_sub_filter.h
Original file line number Diff line number Diff line change
Expand Up @@ -35,10 +35,8 @@ class CuckooSubFilter {
uint64_t version, uint8_t bucket_size, uint32_t page_size, uint16_t filter_index,
uint32_t num_buckets);

uint16_t Index() const { return filter_index_; }
uint32_t NumBuckets() const { return num_buckets_; }

rocksdb::Status TryInsert(uint64_t hash, uint8_t fingerprint, bool *inserted);
rocksdb::Status Delete(uint64_t hash, uint8_t fingerprint, bool *deleted);
// Performs speculative kick-out mutations in the page cache. On success, dirty pages remain staged for
// WriteToBatch(); on inserted=false or non-OK status, cached pages are discarded before returning.
rocksdb::Status TryKickOutInsert(uint64_t hash, uint8_t fingerprint, uint16_t max_iterations, bool *inserted);
Expand Down
57 changes: 57 additions & 0 deletions src/types/redis_cuckoo_chain.cc
Original file line number Diff line number Diff line change
Expand Up @@ -179,6 +179,45 @@ rocksdb::Status CuckooChain::Add(engine::Context &ctx, const Slice &user_key, co
return rocksdb::Status::Aborted("filter is full");
}

rocksdb::Status CuckooChain::Delete(engine::Context &ctx, const Slice &user_key, const Slice &item, bool *deleted) {
*deleted = false;
std::string ns_key = AppendNamespacePrefix(user_key);

CuckooChainMetadata metadata(false);
auto s = getCuckooChainMetadata(ctx, ns_key, &metadata);
if (s.IsNotFound()) return rocksdb::Status::NotFound("Not found");
if (!s.ok()) return s;
Comment on lines +186 to +189

s = validateMetadata(metadata);
if (!s.ok()) return s;

uint64_t hash = CuckooFilterHelper::Hash(item.data(), item.size());
uint8_t fingerprint = CuckooFilterHelper::GenerateFingerprint(hash);

for (int filter_idx = static_cast<int>(metadata.n_filters) - 1; filter_idx >= 0; --filter_idx) {
uint32_t num_buckets = 0;
s = CuckooFilterHelper::GetFilterNumBuckets(metadata.base_capacity, metadata.expansion, metadata.bucket_size,
static_cast<uint16_t>(filter_idx), &num_buckets);
if (!s.ok()) return s;

CuckooSubFilter sub_filter(storage_, ctx, ns_key, storage_->IsSlotIdEncoded(), metadata.version,
metadata.bucket_size, metadata.page_size, static_cast<uint16_t>(filter_idx),
num_buckets);
bool found = false;
s = sub_filter.Delete(hash, fingerprint, &found);
if (!s.ok()) return s;
if (!found) continue;

if (metadata.size == 0) return rocksdb::Status::Corruption("invalid metadata: size is 0");
metadata.size--;
metadata.num_deleted_items++;
*deleted = true;
return commitDelete(ctx, user_key, ns_key, &metadata, &sub_filter);
}

return rocksdb::Status::OK();
}

rocksdb::Status CuckooChain::tryCuckooInsert(engine::Context &ctx, const Slice &user_key, const std::string &ns_key,
CuckooChainMetadata *metadata, uint64_t hash, uint8_t fingerprint,
bool *inserted) {
Expand Down Expand Up @@ -291,4 +330,22 @@ rocksdb::Status CuckooChain::commitSubFilterAndMetadata(engine::Context &ctx, co
return storage_->Write(ctx, storage_->DefaultWriteOptions(), batch->GetWriteBatch());
}

rocksdb::Status CuckooChain::commitDelete(engine::Context &ctx, const Slice &user_key, const std::string &ns_key,
CuckooChainMetadata *metadata, CuckooSubFilter *sub_filter) {
auto batch = storage_->GetWriteBatchBase();
WriteBatchLogData log_data(kRedisCuckooFilter, std::vector<std::string>{"del", user_key.ToString()});
auto s = batch->PutLogData(log_data.Encode());
if (!s.ok()) return s;

s = sub_filter->WriteToBatch(batch.Get());
if (!s.ok()) return s;

std::string metadata_bytes;
metadata->Encode(&metadata_bytes);
s = batch->Put(metadata_cf_handle_, ns_key, metadata_bytes);
if (!s.ok()) return s;

return storage_->Write(ctx, storage_->DefaultWriteOptions(), batch->GetWriteBatch());
}

} // namespace redis
5 changes: 5 additions & 0 deletions src/types/redis_cuckoo_chain.h
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,9 @@ class CuckooChain : public Database {
// Duplicate items are allowed, so added is true whenever insertion succeeds.
rocksdb::Status Add(engine::Context &ctx, const Slice &user_key, const Slice &item, bool *added);

// Deletes one matching fingerprint from the cuckoo filter.
rocksdb::Status Delete(engine::Context &ctx, const Slice &user_key, const Slice &item, bool *deleted);

private:
// Loads metadata for a cuckoo filter key.
rocksdb::Status getCuckooChainMetadata(engine::Context &ctx, const Slice &ns_key, CuckooChainMetadata *metadata);
Expand All @@ -62,6 +65,8 @@ class CuckooChain : public Database {
bool *inserted);
rocksdb::Status commitSubFilterAndMetadata(engine::Context &ctx, const Slice &user_key, const std::string &ns_key,
CuckooChainMetadata *metadata, CuckooSubFilter *sub_filter);
rocksdb::Status commitDelete(engine::Context &ctx, const Slice &user_key, const std::string &ns_key,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I feel that commitDelete is a bit odd from an architectural perspective, as it duplicates some of the logic in commitSubFilterAndMetadata. Could we reuse commitSubFilterAndMetadata for the commit-related logic instead?

CuckooChainMetadata *metadata, CuckooSubFilter *sub_filter);
};

} // namespace redis
37 changes: 37 additions & 0 deletions tests/cppunit/types/cuckoo_filter_page_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -238,6 +238,43 @@ TEST_F(RedisCuckooPageCacheTest, SetBucketSlotWritesOnlyTargetSlot) {
EXPECT_EQ(page, expected);
}

TEST_F(RedisCuckooPageCacheTest, SkipsWritingNewPageWhenItBecomesEmpty) {
auto metadata = makeMetadata(1);
redis::CuckooPageCache pages(storage_.get(), *ctx_, ns_key_, storage_->IsSlotIdEncoded(), metadata.version,
metadata.bucket_size, metadata.page_size);

auto s = pages.SetBucketSlot(0, 1, 0, 0, 88);
ASSERT_TRUE(s.ok()) << s.ToString();
s = pages.SetBucketSlot(0, 1, 0, 0, 0);
ASSERT_TRUE(s.ok()) << s.ToString();

auto batch = storage_->GetWriteBatchBase();
s = pages.WriteBackDirtyPages(batch.Get());
ASSERT_TRUE(s.ok()) << s.ToString();
EXPECT_EQ(batch->GetWriteBatch()->Count(), 0);
}

TEST_F(RedisCuckooPageCacheTest, DeletesExistingPageWhenItBecomesEmpty) {
auto metadata = makeMetadata(1);
auto page_key = makePageKey(metadata, 0, 0);
writePage(page_key, std::string{static_cast<char>(88)});
redis::CuckooPageCache pages(storage_.get(), *ctx_, ns_key_, storage_->IsSlotIdEncoded(), metadata.version,
metadata.bucket_size, metadata.page_size);

auto s = pages.SetBucketSlot(0, 1, 0, 0, 0);
ASSERT_TRUE(s.ok()) << s.ToString();

auto batch = storage_->GetWriteBatchBase();
s = pages.WriteBackDirtyPages(batch.Get());
ASSERT_TRUE(s.ok()) << s.ToString();
EXPECT_EQ(batch->GetWriteBatch()->Count(), 1);
commitBatch(batch.Get());

std::string page;
s = readPage(page_key, &page);
EXPECT_TRUE(s.IsNotFound()) << s.ToString();
}

TEST_F(RedisCuckooPageCacheTest, InvalidBucketAndSlotArguments) {
auto metadata = makeMetadata(4);
redis::CuckooPageCache pages(storage_.get(), *ctx_, ns_key_, storage_->IsSlotIdEncoded(), metadata.version,
Expand Down
Loading
Loading