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
64 changes: 62 additions & 2 deletions src/commands/cmd_cuckoo_filter.cc
Original file line number Diff line number Diff line change
Expand Up @@ -131,8 +131,68 @@ class CommandCFAdd : public Commander {
}
};

// Register the CF.RESERVE and CF.ADD commands
class CommandCFExists : public Commander {
public:
Status Parse(const std::vector<std::string> &args) override {
// CF.EXISTS 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 exists = false;
auto s = cuckoo_db.Exists(ctx, args_[1], args_[2], &exists);

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

*output = conn->Bool(exists);
return Status::OK();
}
};

class CommandCFMExists : public Commander {
public:
Status Parse(const std::vector<std::string> &args) override {
// CF.MEXISTS key item [item ...]
if (args.size() < 3) {
return {Status::RedisParseErr, errWrongNumOfArguments};
}
items_.reserve(args.size() - 2);
for (size_t i = 2; i < args.size(); ++i) {
items_.emplace_back(args[i]);
}
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());
std::vector<bool> exists(items_.size(), false);
auto s = cuckoo_db.MExists(ctx, args_[1], items_, &exists);

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

*output = redis::MultiLen(items_.size());
for (bool exist : exists) {
*output += conn->Bool(exist);
}
return Status::OK();
}

private:
std::vector<std::string> items_;
};

// Register the CF.RESERVE, CF.ADD, CF.EXISTS, and CF.MEXISTS 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<CommandCFExists>("cf.exists", 3, "read-only", 1, 1, 1),
MakeCmdAttr<CommandCFMExists>("cf.mexists", -3, "read-only", 1, 1, 1))

} // namespace redis
3 changes: 2 additions & 1 deletion src/storage/redis_metadata.h
Original file line number Diff line number Diff line change
Expand Up @@ -346,7 +346,8 @@ class CuckooChainMetadata : public Metadata {
/// When a filter is full, a new one is created with capacity = base_capacity * expansion^n
uint16_t expansion;

/// The capacity of the first filter.
/// The capacity of the first filter. Auto-created filters use kCFDefaultCapacity; CF.RESERVE stores the requested
/// capacity.
uint64_t base_capacity;

/// Number of fingerprints per bucket
Expand Down
31 changes: 31 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,37 @@ 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::Contains(uint64_t hash, uint8_t fingerprint, bool *exists) {
*exists = 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;

uint8_t slot = 0;
for (size_t i = 0; i < bucket_size_; ++i) {
s = pages_.GetBucketSlot(filter_index_, num_buckets_, bucket1_idx, static_cast<uint32_t>(i), &slot);
if (!s.ok()) return s;
if (slot == fingerprint) {
*exists = true;
return rocksdb::Status::OK();
}
}

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

for (size_t i = 0; i < bucket_size_; ++i) {
s = pages_.GetBucketSlot(filter_index_, num_buckets_, bucket2_idx, static_cast<uint32_t>(i), &slot);
if (!s.ok()) return s;
if (slot == fingerprint) {
*exists = true;
return rocksdb::Status::OK();
}
}

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

rocksdb::Status CuckooSubFilter::TryKickOutInsert(uint64_t hash, uint8_t fingerprint, uint16_t max_iterations,
bool *inserted) {
*inserted = false;
Expand Down
1 change: 1 addition & 0 deletions src/types/cuckoo_filter_sub_filter.h
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ class CuckooSubFilter {
uint32_t NumBuckets() const { return num_buckets_; }

rocksdb::Status TryInsert(uint64_t hash, uint8_t fingerprint, bool *inserted);
rocksdb::Status Contains(uint64_t hash, uint8_t fingerprint, bool *exists);
// 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
61 changes: 60 additions & 1 deletion src/types/redis_cuckoo_chain.cc
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,7 @@ rocksdb::Status CuckooChain::Reserve(engine::Context &ctx, const Slice &user_key
return rocksdb::Status::InvalidArgument("capacity must be larger than 0");
}

// RedisBloom requires minimum capacity to ensure at least one bucket can be created
// Require minimum capacity to ensure at least one bucket can be created
// With load factor 0.955, capacity=1 and bucket_size=4 results in 0 buckets
if (capacity < 2) {
return rocksdb::Status::InvalidArgument("capacity must be at least 2");
Expand Down Expand Up @@ -291,4 +291,63 @@ rocksdb::Status CuckooChain::commitSubFilterAndMetadata(engine::Context &ctx, co
return storage_->Write(ctx, storage_->DefaultWriteOptions(), batch->GetWriteBatch());
}

rocksdb::Status CuckooChain::Exists(engine::Context &ctx, const Slice &user_key, const Slice &item, bool *exists) {
std::vector<bool> result;
auto s = MExists(ctx, user_key, std::vector<std::string>{item.ToString()}, &result);
if (!s.ok()) return s;
*exists = result[0];
return rocksdb::Status::OK();
}

rocksdb::Status CuckooChain::MExists(engine::Context &ctx, const Slice &user_key, const std::vector<std::string> &items,
std::vector<bool> *exists) {
exists->assign(items.size(), false);
std::string ns_key = AppendNamespacePrefix(user_key);

CuckooChainMetadata metadata(false);
auto s = getCuckooChainMetadata(ctx, ns_key, &metadata);
if (s.IsNotFound()) return rocksdb::Status::OK();
if (!s.ok()) return s;

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

std::vector<uint64_t> hashes(items.size());
std::vector<uint8_t> fingerprints(items.size());
for (size_t i = 0; i < items.size(); ++i) {
hashes[i] = CuckooFilterHelper::Hash(items[i].data(), items[i].size());
fingerprints[i] = CuckooFilterHelper::GenerateFingerprint(hashes[i]);
CHECK(fingerprints[i] != 0);
}

size_t found_count = 0;
for (int filter_idx = static_cast<int>(metadata.n_filters) - 1; filter_idx >= 0; --filter_idx) {
auto current_filter_idx = static_cast<uint16_t>(filter_idx);
uint32_t num_buckets = 0;
s = CuckooFilterHelper::GetFilterNumBuckets(metadata.base_capacity, metadata.expansion, metadata.bucket_size,
current_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, current_filter_idx, num_buckets);
for (size_t i = 0; i < items.size(); ++i) {
if ((*exists)[i]) continue;

bool item_exists = false;
s = sub_filter.Contains(hashes[i], fingerprints[i], &item_exists);
if (!s.ok()) return s;
if (item_exists) {
(*exists)[i] = true;
++found_count;
}
}

if (found_count == items.size()) {
break;
}
}

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

} // namespace redis
9 changes: 9 additions & 0 deletions src/types/redis_cuckoo_chain.h
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@

#pragma once

#include <vector>

#include "cuckoo_filter.h"
#include "storage/redis_db.h"
#include "storage/redis_metadata.h"
Expand Down Expand Up @@ -47,6 +49,13 @@ 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);

// Returns true if the item might exist, and false if it definitely does not.
rocksdb::Status Exists(engine::Context &ctx, const Slice &user_key, const Slice &item, bool *exists);

// Returns whether each item might exist in the cuckoo filter.
rocksdb::Status MExists(engine::Context &ctx, const Slice &user_key, const std::vector<std::string> &items,
std::vector<bool> *exists);

private:
// Loads metadata for a cuckoo filter key.
rocksdb::Status getCuckooChainMetadata(engine::Context &ctx, const Slice &ns_key, CuckooChainMetadata *metadata);
Expand Down
Loading
Loading