Skip to content
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