Coverage Report

Created: 2026-09-28 07:52

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/src/rocksdb/cache/compressed_secondary_cache.cc
Line
Count
Source
1
//  Copyright (c) 2011-present, Facebook, Inc.  All rights reserved.
2
//  This source code is licensed under both the GPLv2 (found in the
3
//  COPYING file in the root directory) and Apache 2.0 License
4
//  (found in the LICENSE.Apache file in the root directory).
5
6
#include "cache/compressed_secondary_cache.h"
7
8
#include <algorithm>
9
#include <cstdint>
10
#include <memory>
11
12
#include "memory/memory_allocator_impl.h"
13
#include "monitoring/perf_context_imp.h"
14
#include "util/coding.h"
15
#include "util/compression.h"
16
#include "util/string_util.h"
17
18
namespace ROCKSDB_NAMESPACE {
19
namespace {
20
// Format of values in CompressedSecondaryCache:
21
// If enable_custom_split_merge:
22
//  * A chain of CacheValueChunk representing the sequence of bytes for a tagged
23
//    value. The overall length of the tagged value is determined by the chain
24
//    of CacheValueChunks.
25
// If !enable_custom_split_merge:
26
//  * A LengthPrefixedSlice (starts with varint64 size) of a tagged value.
27
//
28
// A tagged value has a 2-byte header before the "saved" or compressed block
29
// data:
30
//  * 1 byte for "source" CacheTier indicating which tier is responsible for
31
//    compression/decompression.
32
//  * 1 byte for compression type which is generated/used by
33
//    CompressedSecondaryCache iff source == CacheTier::kVolatileCompressedTier
34
//    (original entry passed in was uncompressed). Otherwise, the compression
35
//    type is preserved from the entry passed in.
36
constexpr uint32_t kTagSize = 2;
37
38
// Size of tag + varint size prefix when applicable
39
0
uint32_t GetHeaderSize(size_t data_size, bool enable_split_merge) {
40
0
  return (enable_split_merge ? 0 : VarintLength(kTagSize + data_size)) +
41
0
         kTagSize;
42
0
}
43
}  // namespace
44
45
CompressedSecondaryCache::CompressedSecondaryCache(
46
    const CompressedSecondaryCacheOptions& opts)
47
0
    : cache_(opts.LRUCacheOptions::MakeSharedCache()),
48
0
      cache_options_(opts),
49
0
      cache_res_mgr_(std::make_shared<ConcurrentCacheReservationManager>(
50
0
          std::make_shared<CacheReservationManagerImpl<CacheEntryRole::kMisc>>(
51
0
              cache_))),
52
0
      disable_cache_(opts.capacity == 0) {
53
0
  auto mgr = GetBuiltinV2CompressionManager();
54
0
  compressor_ = mgr->GetCompressor(cache_options_.compression_opts,
55
0
                                   cache_options_.compression_type);
56
0
  decompressor_ =
57
0
      mgr->GetDecompressorOptimizeFor(cache_options_.compression_type);
58
0
}
59
60
0
CompressedSecondaryCache::~CompressedSecondaryCache() = default;
61
62
std::unique_ptr<SecondaryCacheResultHandle> CompressedSecondaryCache::Lookup(
63
    const Slice& key, const Cache::CacheItemHelper* helper,
64
    Cache::CreateContext* create_context, bool /*wait*/, bool advise_erase,
65
0
    Statistics* stats, bool& kept_in_sec_cache) {
66
0
  assert(helper);
67
0
  if (disable_cache_.LoadRelaxed()) {
68
0
    return nullptr;
69
0
  }
70
71
0
  std::unique_ptr<SecondaryCacheResultHandle> handle;
72
0
  kept_in_sec_cache = false;
73
0
  Cache::Handle* lru_handle = cache_->Lookup(key);
74
0
  if (lru_handle == nullptr) {
75
0
    return nullptr;
76
0
  }
77
78
0
  void* handle_value = cache_->Value(lru_handle);
79
0
  if (handle_value == nullptr) {
80
0
    cache_->Release(lru_handle, /*erase_if_last_ref=*/false);
81
0
    RecordTick(stats, COMPRESSED_SECONDARY_CACHE_DUMMY_HITS);
82
0
    return nullptr;
83
0
  }
84
85
0
  std::string merged_value;
86
0
  Slice tagged_data;
87
0
  if (cache_options_.enable_custom_split_merge) {
88
0
    CacheValueChunk* value_chunk_ptr =
89
0
        static_cast<CacheValueChunk*>(handle_value);
90
0
    merged_value = MergeChunksIntoValue(value_chunk_ptr);
91
0
    tagged_data = Slice(merged_value);
92
0
  } else {
93
0
    tagged_data = GetLengthPrefixedSlice(static_cast<char*>(handle_value));
94
0
  }
95
96
0
  auto source = lossless_cast<CacheTier>(tagged_data[0]);
97
0
  auto type = lossless_cast<CompressionType>(tagged_data[1]);
98
99
0
  std::unique_ptr<char[]> uncompressed;
100
0
  Slice saved(tagged_data.data() + kTagSize, tagged_data.size() - kTagSize);
101
0
  if (source == CacheTier::kVolatileCompressedTier) {
102
0
    if (type != kNoCompression) {
103
      // TODO: can we do something to avoid yet another allocation?
104
0
      Decompressor::Args args;
105
0
      args.compressed_data = saved;
106
0
      args.compression_type = type;
107
0
      Status s = decompressor_->ExtractUncompressedSize(args);
108
0
      assert(s.ok());  // in-memory data
109
0
      if (s.ok()) {
110
0
        uncompressed = std::make_unique<char[]>(args.uncompressed_size);
111
0
        s = decompressor_->DecompressBlock(args, uncompressed.get());
112
0
        assert(s.ok());  // in-memory data
113
0
      }
114
0
      if (!s.ok()) {
115
0
        cache_->Release(lru_handle, /*erase_if_last_ref=*/true);
116
0
        return nullptr;
117
0
      }
118
0
      saved = Slice(uncompressed.get(), args.uncompressed_size);
119
0
      type = kNoCompression;
120
      // Free temporary compressed data as early as we can. This could matter
121
      // for unusually large blocks because we also have
122
      // * Another compressed copy above (from lru_cache).
123
      // * The uncompressed copy in `uncompressed`.
124
      // * Another uncompressed copy in `result_value` below.
125
      // Let's try to max out at 3 copies instead of 4.
126
0
      merged_value = std::string();
127
0
    }
128
    // Reduced as if it came from primary cache
129
0
    source = CacheTier::kVolatileTier;
130
0
  }
131
132
0
  Cache::ObjectPtr result_value = nullptr;
133
0
  size_t result_charge = 0;
134
0
  Status s = helper->create_cb(saved, type, source, create_context,
135
0
                               cache_options_.memory_allocator.get(),
136
0
                               &result_value, &result_charge);
137
0
  if (!s.ok()) {
138
0
    cache_->Release(lru_handle, /*erase_if_last_ref=*/true);
139
0
    return nullptr;
140
0
  }
141
142
0
  if (advise_erase) {
143
0
    cache_->Release(lru_handle, /*erase_if_last_ref=*/true);
144
    // Insert a dummy handle.
145
0
    cache_
146
0
        ->Insert(key, /*obj=*/nullptr,
147
0
                 GetHelper(cache_options_.enable_custom_split_merge),
148
0
                 /*charge=*/0)
149
0
        .PermitUncheckedError();
150
0
  } else {
151
0
    kept_in_sec_cache = true;
152
0
    cache_->Release(lru_handle, /*erase_if_last_ref=*/false);
153
0
  }
154
0
  handle.reset(
155
0
      new CompressedSecondaryCacheResultHandle(result_value, result_charge));
156
0
  RecordTick(stats, COMPRESSED_SECONDARY_CACHE_HITS);
157
0
  return handle;
158
0
}
159
160
0
bool CompressedSecondaryCache::MaybeInsertDummy(const Slice& key) {
161
0
  auto internal_helper = GetHelper(cache_options_.enable_custom_split_merge);
162
0
  Cache::Handle* lru_handle = cache_->Lookup(key);
163
0
  if (lru_handle == nullptr) {
164
0
    PERF_COUNTER_ADD(compressed_sec_cache_insert_dummy_count, 1);
165
    // Insert a dummy handle if the handle is evicted for the first time.
166
0
    cache_->Insert(key, /*obj=*/nullptr, internal_helper, /*charge=*/0)
167
0
        .PermitUncheckedError();
168
0
    return true;
169
0
  } else {
170
0
    cache_->Release(lru_handle, /*erase_if_last_ref=*/false);
171
0
  }
172
173
0
  return false;
174
0
}
175
176
Status CompressedSecondaryCache::InsertInternal(
177
    const Slice& key, Cache::ObjectPtr value,
178
    const Cache::CacheItemHelper* helper, CompressionType from_type,
179
0
    CacheTier source) {
180
0
  bool enable_split_merge = cache_options_.enable_custom_split_merge;
181
0
  const Cache::CacheItemHelper* internal_helper = GetHelper(enable_split_merge);
182
183
  // TODO: variant of size_cb that also returns a pointer to the data if
184
  // already available. Saves an allocation if we keep the compressed version.
185
0
  const size_t data_size_original = (*helper->size_cb)(value);
186
187
  // Allocate enough memory for header/tag + original data because (a) we might
188
  // not be attempting compression at all, and (b) we might keep the original if
189
  // compression is insufficient. But we don't need the length prefix with
190
  // enable_split_merge. TODO: be smarter with CacheValueChunk to save an
191
  // allocation in the enable_split_merge case.
192
0
  size_t header_size = GetHeaderSize(data_size_original, enable_split_merge);
193
0
  CacheAllocationPtr allocation = AllocateBlock(
194
0
      header_size + data_size_original, cache_options_.memory_allocator.get());
195
0
  char* data_ptr = allocation.get() + header_size;
196
0
  Slice tagged_data(data_ptr - kTagSize, data_size_original + kTagSize);
197
0
  assert(tagged_data.data() >= allocation.get());
198
199
0
  Status s = (*helper->saveto_cb)(value, 0, data_size_original, data_ptr);
200
0
  if (!s.ok()) {
201
0
    return s;
202
0
  }
203
204
0
  std::unique_ptr<char[]> tagged_compressed_data;
205
0
  CompressionType to_type = kNoCompression;
206
0
  if (compressor_ && from_type == kNoCompression &&
207
0
      !cache_options_.do_not_compress_roles.Contains(helper->role)) {
208
0
    assert(source == CacheTier::kVolatileCompressedTier);
209
210
    // TODO: consider malloc sizes for max acceptable compressed size
211
    // Or maybe max_compressed_bytes_per_kb
212
0
    size_t data_size_compressed = data_size_original - 1;
213
0
    tagged_compressed_data =
214
0
        std::make_unique<char[]>(data_size_compressed + kTagSize);
215
0
    s = compressor_->CompressBlock(Slice(data_ptr, data_size_original),
216
0
                                   tagged_compressed_data.get() + kTagSize,
217
0
                                   &data_size_compressed, &to_type,
218
0
                                   nullptr /*working_area*/);
219
0
    if (!s.ok()) {
220
0
      return s;
221
0
    }
222
0
    PERF_COUNTER_ADD(compressed_sec_cache_uncompressed_bytes,
223
0
                     data_size_original);
224
0
    if (to_type == kNoCompression) {
225
      // Compression rejected or otherwise aborted/failed
226
0
      to_type = kNoCompression;
227
0
      tagged_compressed_data.reset();
228
      // TODO: consider separate counters for rejected compressions
229
0
      PERF_COUNTER_ADD(compressed_sec_cache_compressed_bytes,
230
0
                       data_size_original);
231
0
    } else {
232
0
      PERF_COUNTER_ADD(compressed_sec_cache_compressed_bytes,
233
0
                       data_size_compressed);
234
0
      if (enable_split_merge) {
235
        // Only need tagged_data for copying into CacheValueChunks.
236
0
        tagged_data = Slice(tagged_compressed_data.get(),
237
0
                            data_size_compressed + kTagSize);
238
0
        allocation.reset();
239
0
      } else {
240
        // Replace allocation with compressed version, copied from string
241
0
        header_size = GetHeaderSize(data_size_compressed, enable_split_merge);
242
0
        allocation = AllocateBlock(header_size + data_size_compressed,
243
0
                                   cache_options_.memory_allocator.get());
244
0
        data_ptr = allocation.get() + header_size;
245
        // Ignore unpopulated tag on tagged_compressed_data; will only be
246
        // populated on the new allocation.
247
0
        std::memcpy(data_ptr, tagged_compressed_data.get() + kTagSize,
248
0
                    data_size_compressed);
249
0
        tagged_data =
250
0
            Slice(data_ptr - kTagSize, data_size_compressed + kTagSize);
251
0
        assert(tagged_data.data() >= allocation.get());
252
0
      }
253
0
    }
254
0
  }
255
256
0
  PERF_COUNTER_ADD(compressed_sec_cache_insert_real_count, 1);
257
258
  // Save the tag fields
259
0
  const_cast<char*>(tagged_data.data())[0] = lossless_cast<char>(source);
260
0
  const_cast<char*>(tagged_data.data())[1] = lossless_cast<char>(
261
0
      source == CacheTier::kVolatileCompressedTier ? to_type : from_type);
262
263
0
  if (enable_split_merge) {
264
0
    size_t split_charge{0};
265
0
    CacheValueChunk* value_chunks_head =
266
0
        SplitValueIntoChunks(tagged_data, split_charge);
267
0
    s = cache_->Insert(key, value_chunks_head, internal_helper, split_charge);
268
0
    assert(s.ok());  // LRUCache::Insert() with handle==nullptr always OK
269
0
  } else {
270
    // Save the size prefix
271
0
    char* ptr = allocation.get();
272
0
    ptr = EncodeVarint64(ptr, tagged_data.size());
273
0
    assert(ptr == tagged_data.data());
274
0
#ifdef ROCKSDB_MALLOC_USABLE_SIZE
275
0
    size_t charge = malloc_usable_size(allocation.get());
276
#else
277
    size_t charge = tagged_data.size();
278
#endif
279
0
    s = cache_->Insert(key, allocation.release(), internal_helper, charge);
280
0
    assert(s.ok());  // LRUCache::Insert() with handle==nullptr always OK
281
0
  }
282
0
  return Status::OK();
283
0
}
284
285
Status CompressedSecondaryCache::Insert(const Slice& key,
286
                                        Cache::ObjectPtr value,
287
                                        const Cache::CacheItemHelper* helper,
288
0
                                        bool force_insert) {
289
0
  if (value == nullptr) {
290
0
    return Status::InvalidArgument();
291
0
  }
292
293
0
  if (!force_insert && MaybeInsertDummy(key)) {
294
0
    return Status::OK();
295
0
  }
296
297
0
  return InsertInternal(key, value, helper, kNoCompression,
298
0
                        CacheTier::kVolatileCompressedTier);
299
0
}
300
301
Status CompressedSecondaryCache::InsertSaved(
302
    const Slice& key, const Slice& saved, CompressionType type = kNoCompression,
303
0
    CacheTier source = CacheTier::kVolatileTier) {
304
0
  if (source == CacheTier::kVolatileCompressedTier) {
305
    // Unexpected, would violate InsertInternal preconditions
306
0
    assert(source != CacheTier::kVolatileCompressedTier);
307
0
    return Status::OK();
308
0
  }
309
0
  if (type == kNoCompression) {
310
    // Not currently supported (why?)
311
0
    return Status::OK();
312
0
  }
313
0
  if (cache_options_.enable_custom_split_merge) {
314
    // We don't support custom split/merge for the tiered case (why?)
315
0
    return Status::OK();
316
0
  }
317
318
0
  auto slice_helper = &kSliceCacheItemHelper;
319
0
  if (MaybeInsertDummy(key)) {
320
0
    return Status::OK();
321
0
  }
322
323
0
  return InsertInternal(
324
0
      key, static_cast<Cache::ObjectPtr>(const_cast<Slice*>(&saved)),
325
0
      slice_helper, type, source);
326
0
}
327
328
0
void CompressedSecondaryCache::Erase(const Slice& key) { cache_->Erase(key); }
329
330
0
Status CompressedSecondaryCache::SetCapacity(size_t capacity) {
331
0
  MutexLock l(&capacity_mutex_);
332
0
  cache_options_.capacity = capacity;
333
0
  cache_->SetCapacity(capacity);
334
0
  disable_cache_.StoreRelaxed(capacity == 0);
335
0
  return Status::OK();
336
0
}
337
338
0
Status CompressedSecondaryCache::GetCapacity(size_t& capacity) {
339
0
  MutexLock l(&capacity_mutex_);
340
0
  capacity = cache_options_.capacity;
341
0
  return Status::OK();
342
0
}
343
344
0
std::string CompressedSecondaryCache::GetPrintableOptions() const {
345
0
  std::string ret;
346
0
  ret.reserve(20000);
347
0
  const int kBufferSize{200};
348
0
  char buffer[kBufferSize];
349
0
  ret.append(cache_->GetPrintableOptions());
350
0
  snprintf(buffer, kBufferSize, "    compression_type : %s\n",
351
0
           CompressionTypeToString(cache_options_.compression_type).c_str());
352
0
  ret.append(buffer);
353
0
  snprintf(buffer, kBufferSize, "    compression_opts : %s\n",
354
0
           CompressionOptionsToString(
355
0
               const_cast<CompressionOptions&>(cache_options_.compression_opts))
356
0
               .c_str());
357
0
  ret.append(buffer);
358
0
  return ret;
359
0
}
360
361
// FIXME: this could use a lot of attention, including:
362
// * Use allocator
363
// * We shouldn't be worse than non-split; be more pro-actively aware of
364
// internal fragmentation
365
// * Consider a unified object/chunk structure that may or may not split
366
// * Optimize size overhead of chunks
367
CompressedSecondaryCache::CacheValueChunk*
368
CompressedSecondaryCache::SplitValueIntoChunks(const Slice& value,
369
0
                                               size_t& charge) {
370
0
  assert(!value.empty());
371
0
  const char* src_ptr = value.data();
372
0
  size_t src_size{value.size()};
373
374
0
  CacheValueChunk dummy_head = CacheValueChunk();
375
0
  CacheValueChunk* current_chunk = &dummy_head;
376
  // Do not split when value size is large or there is no compression.
377
0
  size_t predicted_chunk_size{0};
378
0
  size_t actual_chunk_size{0};
379
0
  size_t tmp_size{0};
380
0
  while (src_size > 0) {
381
0
    predicted_chunk_size = sizeof(CacheValueChunk) - 1 + src_size;
382
0
    auto upper =
383
0
        std::upper_bound(malloc_bin_sizes_.begin(), malloc_bin_sizes_.end(),
384
0
                         predicted_chunk_size);
385
    // Do not split when value size is too small, too large, close to a bin
386
    // size, or there is no compression.
387
0
    if (upper == malloc_bin_sizes_.begin() ||
388
0
        upper == malloc_bin_sizes_.end() ||
389
0
        *upper - predicted_chunk_size < malloc_bin_sizes_.front()) {
390
0
      tmp_size = predicted_chunk_size;
391
0
    } else {
392
0
      tmp_size = *(--upper);
393
0
    }
394
395
0
    CacheValueChunk* new_chunk =
396
0
        static_cast<CacheValueChunk*>(static_cast<void*>(new char[tmp_size]));
397
0
    current_chunk->next = new_chunk;
398
0
    current_chunk = current_chunk->next;
399
0
    actual_chunk_size = tmp_size - sizeof(CacheValueChunk) + 1;
400
0
    memcpy(current_chunk->data, src_ptr, actual_chunk_size);
401
0
    current_chunk->size = actual_chunk_size;
402
0
    src_ptr += actual_chunk_size;
403
0
    src_size -= actual_chunk_size;
404
0
    charge += tmp_size;
405
0
  }
406
0
  current_chunk->next = nullptr;
407
408
0
  return dummy_head.next;
409
0
}
410
411
std::string CompressedSecondaryCache::MergeChunksIntoValue(
412
0
    const CacheValueChunk* head) {
413
0
  const CacheValueChunk* current_chunk = head;
414
0
  size_t total_size = 0;
415
0
  while (current_chunk != nullptr) {
416
0
    total_size += current_chunk->size;
417
0
    current_chunk = current_chunk->next;
418
0
  }
419
420
0
  std::string result;
421
0
  result.reserve(total_size);
422
0
  current_chunk = head;
423
0
  while (current_chunk != nullptr) {
424
0
    result.append(current_chunk->data, current_chunk->size);
425
0
    current_chunk = current_chunk->next;
426
0
  }
427
0
  assert(result.size() == total_size);
428
0
  return result;
429
0
}
430
431
const Cache::CacheItemHelper* CompressedSecondaryCache::GetHelper(
432
0
    bool enable_custom_split_merge) const {
433
0
  if (enable_custom_split_merge) {
434
0
    static const Cache::CacheItemHelper kHelper{
435
0
        CacheEntryRole::kMisc,
436
0
        [](Cache::ObjectPtr obj, MemoryAllocator* /*alloc*/) {
437
0
          CacheValueChunk* chunks_head = static_cast<CacheValueChunk*>(obj);
438
0
          while (chunks_head != nullptr) {
439
0
            CacheValueChunk* tmp_chunk = chunks_head;
440
0
            chunks_head = chunks_head->next;
441
0
            tmp_chunk->Free();
442
0
          }
443
0
        }};
444
0
    return &kHelper;
445
0
  } else {
446
0
    static const Cache::CacheItemHelper kHelper{
447
0
        CacheEntryRole::kMisc,
448
0
        [](Cache::ObjectPtr obj, MemoryAllocator* alloc) {
449
0
          if (obj != nullptr) {
450
0
            CacheAllocationDeleter{alloc}(static_cast<char*>(obj));
451
0
          }
452
0
        }};
453
0
    return &kHelper;
454
0
  }
455
0
}
456
457
0
size_t CompressedSecondaryCache::TEST_GetCharge(const Slice& key) {
458
0
  Cache::Handle* lru_handle = cache_->Lookup(key);
459
0
  if (lru_handle == nullptr) {
460
0
    return 0;
461
0
  }
462
0
  size_t charge = cache_->GetCharge(lru_handle);
463
0
  cache_->Release(lru_handle, /*erase_if_last_ref=*/false);
464
0
  return charge;
465
0
}
466
467
std::shared_ptr<SecondaryCache>
468
0
CompressedSecondaryCacheOptions::MakeSharedSecondaryCache() const {
469
0
  return std::make_shared<CompressedSecondaryCache>(*this);
470
0
}
471
472
0
Status CompressedSecondaryCache::Deflate(size_t decrease) {
473
0
  return cache_res_mgr_->UpdateCacheReservation(decrease, /*increase=*/true);
474
0
}
475
476
0
Status CompressedSecondaryCache::Inflate(size_t increase) {
477
0
  return cache_res_mgr_->UpdateCacheReservation(increase, /*increase=*/false);
478
0
}
479
480
}  // namespace ROCKSDB_NAMESPACE