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
51 changes: 51 additions & 0 deletions src/types/cuckoo_filter_page.cc
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,47 @@ rocksdb::Status CuckooPageCache::SetBucketSlot(uint16_t filter_index, uint32_t n
return rocksdb::Status::OK();
}

rocksdb::Status CuckooPageCache::SetBucketSlotWithUndo(uint16_t filter_index, uint32_t num_buckets,
uint32_t bucket_index, uint32_t slot_idx, uint8_t fingerprint,
SlotMutation *mutation) {
if (slot_idx >= bucket_size_) return rocksdb::Status::InvalidArgument("invalid cuckoo filter bucket slot");

BucketRef bucket;
auto s = ensureBucketLoaded(filter_index, num_buckets, bucket_index, &bucket);
if (!s.ok()) return s;

mutation->filter_index = filter_index;
mutation->num_buckets = num_buckets;
mutation->bucket_index = bucket_index;
mutation->slot_idx = slot_idx;
mutation->old_fingerprint = getBucketRefSlot(bucket, slot_idx);
mutation->page_was_dirty = bucket.page->is_dirty;
setBucketRefSlot(bucket, slot_idx, fingerprint);
return rocksdb::Status::OK();
}

rocksdb::Status CuckooPageCache::RestoreBucketSlot(const SlotMutation &mutation) {
if (mutation.slot_idx >= bucket_size_) {
return rocksdb::Status::Corruption("invalid cuckoo filter slot mutation");
}

BucketLocation location;
auto s = resolveBucketLocation(mutation.filter_index, mutation.num_buckets, mutation.bucket_index, &location);
if (!s.ok()) return s;

auto it = pages_.find(location.page_key);
if (it == pages_.end()) return rocksdb::Status::Corruption("cuckoo filter slot mutation page is not cached");

auto offset = location.offset + mutation.slot_idx;
if (offset >= it->second.data.size()) {
return rocksdb::Status::Corruption("invalid cuckoo filter slot mutation");
}

it->second.data[offset] = static_cast<char>(mutation.old_fingerprint);
it->second.is_dirty = mutation.page_was_dirty;
return rocksdb::Status::OK();
}

rocksdb::Status CuckooPageCache::WriteBackDirtyPages(rocksdb::WriteBatchBase *batch) {
for (const auto &entry : pages_) {
if (!entry.second.is_dirty) continue;
Expand All @@ -125,6 +166,16 @@ rocksdb::Status CuckooPageCache::WriteBackDirtyPages(rocksdb::WriteBatchBase *ba
return rocksdb::Status::OK();
}

void CuckooPageCache::DiscardCleanPages() {
for (auto it = pages_.begin(); it != pages_.end();) {
if (it->second.is_dirty) {
++it;
} else {
it = pages_.erase(it);
}
}
}

void CuckooPageCache::DiscardCachedPages() { pages_.clear(); }

rocksdb::Status CuckooPageCache::resolveBucketLocation(uint16_t filter_index, uint32_t num_buckets,
Expand Down
17 changes: 17 additions & 0 deletions src/types/cuckoo_filter_page.h
Original file line number Diff line number Diff line change
Expand Up @@ -34,9 +34,22 @@ namespace redis {

class CuckooPageCache {
public:
struct SlotMutation {
uint8_t OldFingerprint() const { return old_fingerprint; }

uint16_t filter_index = 0;
uint32_t num_buckets = 0;
uint32_t bucket_index = 0;
uint32_t slot_idx = 0;
uint8_t old_fingerprint = 0;
bool page_was_dirty = false;
};

CuckooPageCache(engine::Storage *storage, engine::Context &ctx, const Slice &ns_key, bool slot_id_encoded,
uint64_t version, uint8_t bucket_size, uint32_t page_size);

uint8_t BucketSize() const { return bucket_size_; }

rocksdb::Status PrefetchBuckets(uint16_t filter_index, uint32_t num_buckets, uint32_t bucket1_index,
uint32_t bucket2_index);
rocksdb::Status TryInsertInBucket(uint16_t filter_index, uint32_t num_buckets, uint32_t bucket_index,
Expand All @@ -45,8 +58,12 @@ class CuckooPageCache {
uint8_t *fingerprint);
rocksdb::Status SetBucketSlot(uint16_t filter_index, uint32_t num_buckets, uint32_t bucket_index, uint32_t slot_idx,
uint8_t fingerprint);
rocksdb::Status SetBucketSlotWithUndo(uint16_t filter_index, uint32_t num_buckets, uint32_t bucket_index,
uint32_t slot_idx, uint8_t fingerprint, SlotMutation *mutation);
rocksdb::Status RestoreBucketSlot(const SlotMutation &mutation);
rocksdb::Status WriteBackDirtyPages(rocksdb::WriteBatchBase *batch);

void DiscardCleanPages();
void DiscardCachedPages();

private:
Expand Down
57 changes: 28 additions & 29 deletions src/types/cuckoo_filter_sub_filter.cc
Original file line number Diff line number Diff line change
Expand Up @@ -20,17 +20,15 @@

#include "cuckoo_filter_sub_filter.h"

#include <vector>

#include "cuckoo_filter.h"
#include "cuckoo_filter_page.h"

namespace redis {

CuckooSubFilter::CuckooSubFilter(engine::Storage *storage, engine::Context &ctx, const Slice &ns_key,
bool slot_id_encoded, uint64_t version, uint8_t bucket_size, uint32_t page_size,
uint16_t filter_index, uint32_t num_buckets)
: bucket_size_(bucket_size),
filter_index_(filter_index),
num_buckets_(num_buckets),
pages_(storage, ctx, ns_key, slot_id_encoded, version, bucket_size, page_size) {}
CuckooSubFilter::CuckooSubFilter(CuckooPageCache &pages, uint16_t filter_index, uint32_t num_buckets)
: bucket_size_(pages.BucketSize()), filter_index_(filter_index), num_buckets_(num_buckets), pages_(pages) {}

rocksdb::Status CuckooSubFilter::TryInsert(uint64_t hash, uint8_t fingerprint, bool *inserted) {
*inserted = false;
Expand All @@ -52,20 +50,29 @@ rocksdb::Status CuckooSubFilter::TryKickOutInsert(uint64_t hash, uint8_t fingerp
uint32_t current_bucket_idx = getPrimaryBucketIndex(hash);
uint8_t current_fp = fingerprint;
uint32_t victim_slot = 0;
std::vector<CuckooPageCache::SlotMutation> mutations;
mutations.reserve(max_iterations);

for (uint16_t iteration = 0; iteration < max_iterations; ++iteration) {
uint8_t old_fp = 0;
auto s = pages_.GetBucketSlot(filter_index_, num_buckets_, current_bucket_idx, victim_slot, &old_fp);
if (!s.ok()) {
pages_.DiscardCachedPages();
return s;
}
s = pages_.SetBucketSlot(filter_index_, num_buckets_, current_bucket_idx, victim_slot, current_fp);
if (!s.ok()) {
pages_.DiscardCachedPages();
return s;
auto rollback = [&]() {
for (auto it = mutations.rbegin(); it != mutations.rend(); ++it) {
auto s = pages_.RestoreBucketSlot(*it);
if (!s.ok()) return s;
}
current_fp = old_fp;
return rocksdb::Status::OK();
};

auto rollback_and_return = [&](const rocksdb::Status &status) {
auto rollback_status = rollback();
return rollback_status.ok() ? status : rollback_status;
};

for (uint16_t iteration = 0; iteration < max_iterations; ++iteration) {
CuckooPageCache::SlotMutation mutation;
auto s = pages_.SetBucketSlotWithUndo(filter_index_, num_buckets_, current_bucket_idx, victim_slot, current_fp,
&mutation);
if (!s.ok()) return rollback_and_return(s);
current_fp = mutation.OldFingerprint();
mutations.push_back(mutation);

if (current_fp == 0) {
*inserted = true;
Expand All @@ -76,10 +83,7 @@ rocksdb::Status CuckooSubFilter::TryKickOutInsert(uint64_t hash, uint8_t fingerp

bool inserted_in_alt_bucket = false;
s = pages_.TryInsertInBucket(filter_index_, num_buckets_, alt_bucket_idx, current_fp, &inserted_in_alt_bucket);
if (!s.ok()) {
pages_.DiscardCachedPages();
return s;
}
if (!s.ok()) return rollback_and_return(s);
if (inserted_in_alt_bucket) {
*inserted = true;
return rocksdb::Status::OK();
Expand All @@ -89,12 +93,7 @@ rocksdb::Status CuckooSubFilter::TryKickOutInsert(uint64_t hash, uint8_t fingerp
victim_slot = (victim_slot + 1) % bucket_size_;
}

pages_.DiscardCachedPages();
return rocksdb::Status::OK();
}

rocksdb::Status CuckooSubFilter::WriteToBatch(rocksdb::WriteBatchBase *batch) {
return pages_.WriteBackDirtyPages(batch);
return rollback();
}

uint32_t CuckooSubFilter::getPrimaryBucketIndex(uint64_t hash) const { return hash % num_buckets_; }
Expand Down
16 changes: 6 additions & 10 deletions src/types/cuckoo_filter_sub_filter.h
Original file line number Diff line number Diff line change
Expand Up @@ -21,28 +21,24 @@
#pragma once

#include <rocksdb/status.h>
#include <rocksdb/write_batch.h>

#include <cstdint>

#include "cuckoo_filter_page.h"

namespace redis {

class CuckooPageCache;

class CuckooSubFilter {
public:
CuckooSubFilter(engine::Storage *storage, engine::Context &ctx, const Slice &ns_key, bool slot_id_encoded,
uint64_t version, uint8_t bucket_size, uint32_t page_size, uint16_t filter_index,
uint32_t num_buckets);
CuckooSubFilter(CuckooPageCache &pages, 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);
// 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.
// Performs speculative kick-out mutations in the page cache. Failed attempts restore only the slots changed by
// this call, preserving mutations staged by the owning chain operation.
rocksdb::Status TryKickOutInsert(uint64_t hash, uint8_t fingerprint, uint16_t max_iterations, bool *inserted);
rocksdb::Status WriteToBatch(rocksdb::WriteBatchBase *batch);

private:
uint32_t getPrimaryBucketIndex(uint64_t hash) const;
Expand All @@ -51,7 +47,7 @@ class CuckooSubFilter {
uint8_t bucket_size_ = 0;
uint16_t filter_index_ = 0;
uint32_t num_buckets_ = 0;
CuckooPageCache pages_;
CuckooPageCache &pages_;
};

} // namespace redis
Loading
Loading