Coverage Report

Created: 2026-09-28 07:52

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/src/rocksdb/db/version_edit_handler.h
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
// Copyright (c) 2011 The LevelDB Authors. All rights reserved.
7
// Use of this source code is governed by a BSD-style license that can be
8
// found in the LICENSE file. See the AUTHORS file for names of contributors.
9
10
#pragma once
11
12
#include "db/version_builder.h"
13
#include "db/version_edit.h"
14
#include "db/version_set.h"
15
16
namespace ROCKSDB_NAMESPACE {
17
18
struct FileMetaData;
19
20
class VersionEditHandlerBase {
21
 public:
22
  explicit VersionEditHandlerBase(const ReadOptions& read_options)
23
53.3k
      : read_options_(read_options),
24
53.3k
        max_manifest_read_size_(std::numeric_limits<uint64_t>::max()) {}
25
26
53.3k
  virtual ~VersionEditHandlerBase() {}
27
28
  void Iterate(log::Reader& reader, Status* log_read_status);
29
30
53.3k
  const Status& status() const { return status_; }
31
32
0
  AtomicGroupReadBuffer& GetReadBuffer() { return read_buffer_; }
33
34
0
  uint64_t GetLastValidRecordEnd() const { return last_valid_record_end_; }
35
36
 protected:
37
  explicit VersionEditHandlerBase(const ReadOptions& read_options,
38
                                  uint64_t max_read_size)
39
0
      : read_options_(read_options), max_manifest_read_size_(max_read_size) {}
40
0
  virtual Status Initialize() { return Status::OK(); }
41
42
  virtual Status ApplyVersionEdit(VersionEdit& edit,
43
                                  ColumnFamilyData** cfd) = 0;
44
45
0
  virtual Status OnAtomicGroupReplayBegin() { return Status::OK(); }
46
0
  virtual Status OnAtomicGroupReplayEnd() { return Status::OK(); }
47
48
  virtual void CheckIterationResult(const log::Reader& /*reader*/,
49
0
                                    Status* /*s*/) {}
50
51
0
  void ResetReadState() {
52
0
    read_buffer_.Clear();
53
0
    last_valid_record_end_ = 0;
54
0
  }
55
56
  Status status_;
57
58
  const ReadOptions& read_options_;
59
60
  // File offset at the end of the last fully applied logical unit.
61
  // For non-atomic edits, this advances after each edit is decoded.
62
  // For atomic groups, this only advances after the entire group is
63
  // decoded, so an incomplete atomic group at the tail is excluded.
64
  uint64_t last_valid_record_end_ = 0;
65
66
 private:
67
  AtomicGroupReadBuffer read_buffer_;
68
  const uint64_t max_manifest_read_size_;
69
};
70
71
class ListColumnFamiliesHandler : public VersionEditHandlerBase {
72
 public:
73
  explicit ListColumnFamiliesHandler(const ReadOptions& read_options)
74
0
      : VersionEditHandlerBase(read_options) {}
75
76
0
  ~ListColumnFamiliesHandler() override {}
77
78
0
  const std::map<uint32_t, std::string> GetColumnFamilyNames() const {
79
0
    return column_family_names_;
80
0
  }
81
82
 protected:
83
  Status ApplyVersionEdit(VersionEdit& edit,
84
                          ColumnFamilyData** /*unused*/) override;
85
86
 private:
87
  // default column family is always implicitly there
88
  std::map<uint32_t, std::string> column_family_names_{
89
      {0, kDefaultColumnFamilyName}};
90
};
91
92
class FileChecksumRetriever : public VersionEditHandlerBase {
93
 public:
94
  FileChecksumRetriever(const ReadOptions& read_options, uint64_t max_read_size)
95
0
      : VersionEditHandlerBase(read_options, max_read_size) {}
96
97
0
  ~FileChecksumRetriever() override {}
98
99
  Status FetchFileChecksumList(FileChecksumList& file_checksum_list);
100
101
 protected:
102
  Status ApplyVersionEdit(VersionEdit& edit,
103
                          ColumnFamilyData** /*unused*/) override;
104
105
 private:
106
  // Map from CF to file # to string pair, where first portion of the value
107
  // is checksum, and second portion of the value is checksum function name.
108
  //
109
  // [column family id A]
110
  //      |
111
  //      |-- [file #1] -> [checksum #1, checksum function name #1]
112
  //      |-- [file #2] -> [checksum #2, checksum function name #2]
113
  //      |
114
  //     ...
115
  //      |
116
  //      |-- [file #N] -> [checksum #N, checksum function name #N]
117
  // [column family id B]
118
  //      |
119
  //      |-- [file #1] -> [checksum #1, checksum function name #1]
120
  //      |
121
  //     ...
122
  //      |
123
  //      |-- [file #M] -> [checksum #M, checksum function name #M]
124
  //      |
125
  //     ...
126
  std::unordered_map<
127
      uint32_t,
128
      std::unordered_map<uint64_t, std::pair<std::string, std::string>>>
129
      cf_file_checksums_;
130
};
131
132
using VersionBuilderUPtr = std::unique_ptr<BaseReferencedVersionBuilder>;
133
134
// A class used for scanning MANIFEST file.
135
// VersionEditHandler reads a MANIFEST file, parses the version edits, and
136
// builds the version set's in-memory state, e.g. the version storage info for
137
// the versions of column families. It replays all the version edits in one
138
// MANIFEST file to build the end version.
139
//
140
// To use this class and its subclasses,
141
// 1. Create an object of VersionEditHandler or its subclasses.
142
//    VersionEditHandler handler(read_only, column_families, version_set,
143
//                               track_found_and_missing_files,
144
//                               no_error_if_files_missing);
145
// 2. Status s = handler.Iterate(reader, &db_id);
146
// 3. Check s and handle possible errors.
147
//
148
// Not thread-safe, external synchronization is necessary if an object of
149
// VersionEditHandler is shared by multiple threads.
150
class VersionEditHandler : public VersionEditHandlerBase {
151
 public:
152
  explicit VersionEditHandler(
153
      bool read_only,
154
      const std::vector<ColumnFamilyDescriptor>& column_families,
155
      VersionSet* version_set, bool track_found_and_missing_files,
156
      bool no_error_if_files_missing,
157
      const std::shared_ptr<IOTracer>& io_tracer,
158
      const ReadOptions& read_options, bool allow_incomplete_valid_version,
159
      EpochNumberRequirement epoch_number_requirement =
160
          EpochNumberRequirement::kMustPresent,
161
      bool skip_load_table_files = false)
162
53.3k
      : VersionEditHandler(
163
53.3k
            read_only, column_families, version_set,
164
53.3k
            track_found_and_missing_files, no_error_if_files_missing, io_tracer,
165
53.3k
            read_options, skip_load_table_files, allow_incomplete_valid_version,
166
53.3k
            epoch_number_requirement) {}
167
168
53.3k
  ~VersionEditHandler() override {}
169
170
53.3k
  const VersionEditParams& GetVersionEditParams() const {
171
53.3k
    return version_edit_params_;
172
53.3k
  }
173
174
53.3k
  void GetDbId(std::string* db_id) const {
175
53.3k
    if (db_id && version_edit_params_.HasDbId()) {
176
53.3k
      *db_id = version_edit_params_.GetDbId();
177
53.3k
    }
178
53.3k
  }
179
180
  virtual Status VerifyFile(ColumnFamilyData* /*cfd*/,
181
                            const std::string& /*fpath*/, int /*level*/,
182
0
                            const FileMetaData& /*fmeta*/) {
183
0
    return Status::OK();
184
0
  }
185
186
  virtual Status VerifyBlobFile(ColumnFamilyData* /*cfd*/,
187
                                uint64_t /*blob_file_num*/,
188
0
                                const BlobFileAddition& /*blob_addition*/) {
189
0
    return Status::OK();
190
0
  }
191
192
 protected:
193
  explicit VersionEditHandler(
194
      bool read_only, std::vector<ColumnFamilyDescriptor> column_families,
195
      VersionSet* version_set, bool track_found_and_missing_files,
196
      bool no_error_if_files_missing,
197
      const std::shared_ptr<IOTracer>& io_tracer,
198
      const ReadOptions& read_options, bool skip_load_table_files,
199
      bool allow_incomplete_valid_version,
200
      EpochNumberRequirement epoch_number_requirement =
201
          EpochNumberRequirement::kMustPresent);
202
203
  Status ApplyVersionEdit(VersionEdit& edit, ColumnFamilyData** cfd) override;
204
205
  virtual Status OnColumnFamilyAdd(VersionEdit& edit, ColumnFamilyData** cfd);
206
207
  Status OnColumnFamilyDrop(VersionEdit& edit, ColumnFamilyData** cfd);
208
209
  Status OnNonCfOperation(VersionEdit& edit, ColumnFamilyData** cfd);
210
211
  Status OnWalAddition(VersionEdit& edit);
212
213
  Status OnWalDeletion(VersionEdit& edit);
214
215
  Status Initialize() override;
216
217
  void CheckColumnFamilyId(const VersionEdit& edit, bool* do_not_open_cf,
218
                           bool* cf_in_builders) const;
219
220
  void CheckIterationResult(const log::Reader& reader, Status* s) override;
221
222
  ColumnFamilyData* CreateCfAndInit(const ColumnFamilyOptions& cf_options,
223
                                    const VersionEdit& edit);
224
225
  virtual ColumnFamilyData* DestroyCfAndCleanup(const VersionEdit& edit);
226
227
  virtual Status MaybeCreateVersionBeforeApplyEdit(const VersionEdit& edit,
228
                                                   ColumnFamilyData* cfd,
229
                                                   bool force_create_version);
230
231
  virtual Status LoadTables(ColumnFamilyData* cfd,
232
                            bool prefetch_index_and_filter_in_cache,
233
                            bool is_initial_load);
234
235
53.3k
  virtual bool MustOpenAllColumnFamilies() const {
236
53.3k
    return !version_set_->unchanging();
237
53.3k
  }
238
239
  const bool read_only_;
240
  std::vector<ColumnFamilyDescriptor> column_families_;
241
  VersionSet* version_set_;
242
  std::unordered_map<uint32_t, VersionBuilderUPtr> builders_;
243
  std::unordered_map<std::string, ColumnFamilyOptions> name_to_options_;
244
  const bool track_found_and_missing_files_;
245
  // Keeps track of column families in manifest that were not found in
246
  // column families parameters. Namely, the user asks to not open these column
247
  // families. In non read only mode, if those column families are not dropped
248
  // by subsequent manifest records, Recover() will return failure status.
249
  std::unordered_map<uint32_t, std::string> do_not_open_column_families_;
250
  VersionEditParams version_edit_params_;
251
  bool no_error_if_files_missing_;
252
  std::shared_ptr<IOTracer> io_tracer_;
253
  bool skip_load_table_files_;
254
  bool initialized_;
255
  std::unique_ptr<std::unordered_map<uint32_t, std::string>> cf_to_cmp_names_;
256
  // If false, only a complete Version for which all files consisting it can be
257
  // found is considered a valid Version. If true, besides complete Version, an
258
  // incomplete Version with only a suffix of L0 files missing is also
259
  // considered valid if the Version is never edited in an atomic group.
260
  const bool allow_incomplete_valid_version_;
261
  EpochNumberRequirement epoch_number_requirement_;
262
  std::unordered_set<uint32_t> cfds_to_mark_no_udt_;
263
264
 private:
265
  Status ExtractInfoFromVersionEdit(ColumnFamilyData* cfd,
266
                                    const VersionEdit& edit);
267
268
  // When `FileMetaData.user_defined_timestamps_persisted` is false and
269
  // user-defined timestamp size is non-zero. User-defined timestamps are
270
  // stripped from file boundaries: `smallest`, `largest` in
271
  // `VersionEdit.DecodeFrom` before they were written to Manifest.
272
  // This is the mirroring change to handle file boundaries on the Manifest read
273
  // path for this scenario: to pad a minimum timestamp to the user key in
274
  // `smallest` and `largest` so their format are consistent with the running
275
  // user comparator.
276
  Status MaybeHandleFileBoundariesForNewFiles(VersionEdit& edit,
277
                                              const ColumnFamilyData* cfd);
278
};
279
280
// A class similar to its base class, i.e. VersionEditHandler.
281
// Unlike VersionEditHandler that only aims to build the end version, this class
282
// supports building the most recent point in time version. A point in time
283
// version is a version for which no files are missing, or if
284
// `allow_incomplete_valid_version` is true, only a suffix of L0 files (and
285
// their associated blob files) are missing.
286
//
287
// Building a point in time version when end version is not available can
288
// be useful for best efforts recovery (options.best_efforts_recovery), which
289
// uses this class and sets `allow_incomplete_valid_version` to true.
290
// It's also useful for secondary instances/follower instances for which end
291
// version could be transiently unavailable. These two cases use subclass
292
// `ManifestTailer` and sets `allow_incomplete_valid_version` to false.
293
//
294
// Not thread-safe, external synchronization is necessary if an object of
295
// VersionEditHandlerPointInTime is shared by multiple threads.
296
class VersionEditHandlerPointInTime : public VersionEditHandler {
297
 public:
298
  VersionEditHandlerPointInTime(
299
      bool read_only, std::vector<ColumnFamilyDescriptor> column_families,
300
      VersionSet* version_set, const std::shared_ptr<IOTracer>& io_tracer,
301
      const ReadOptions& read_options, bool allow_incomplete_valid_version,
302
      bool trust_manifest_recovery,
303
      EpochNumberRequirement epoch_number_requirement =
304
          EpochNumberRequirement::kMustPresent);
305
  ~VersionEditHandlerPointInTime() override;
306
307
  bool HasMissingFiles() const;
308
309
  // Returns the column family's log number as of the Version most recently
310
  // installed for it by this handler, i.e. the log number that the MANIFEST
311
  // records reflected by that Version had put in effect. Data in WALs older
312
  // than the returned number is readable from the files that Version
313
  // references, without those WALs.
314
  //
315
  // Returns 0 when no Version has been installed for `cf_id`.
316
  //
317
  // REQUIRES: db mutex
318
  uint64_t GetInstalledVersionLogNumber(uint32_t cf_id) const;
319
320
  virtual Status VerifyFile(ColumnFamilyData* cfd, const std::string& fpath,
321
                            int level, const FileMetaData& fmeta) override;
322
  virtual Status VerifyBlobFile(ColumnFamilyData* cfd, uint64_t blob_file_num,
323
                                const BlobFileAddition& blob_addition) override;
324
325
 protected:
326
  Status OnAtomicGroupReplayBegin() override;
327
  Status OnAtomicGroupReplayEnd() override;
328
  void CheckIterationResult(const log::Reader& reader, Status* s) override;
329
330
  ColumnFamilyData* DestroyCfAndCleanup(const VersionEdit& edit) override;
331
  // `MaybeCreateVersionBeforeApplyEdit(..., false)` creates a version upon a
332
  // negative edge trigger (transition from valid to invalid).
333
  //
334
  // `MaybeCreateVersionBeforeApplyEdit(..., true)` creates a version on a
335
  // positive level trigger (state is valid).
336
  Status MaybeCreateVersionBeforeApplyEdit(const VersionEdit& edit,
337
                                           ColumnFamilyData* cfd,
338
                                           bool force_create_version) override;
339
340
  Status LoadTables(ColumnFamilyData* cfd,
341
                    bool prefetch_index_and_filter_in_cache,
342
                    bool is_initial_load) override;
343
344
  // A Version built from the MANIFEST records read up to some point in time,
345
  // together with the column family's log number that those records had put in
346
  // effect. Keeping the two together is what lets an installed Version report
347
  // the log number it covers.
348
  struct PointInTimeVersion {
349
    Version* version = nullptr;
350
    uint64_t log_number = 0;
351
  };
352
353
  std::unordered_map<uint32_t, PointInTimeVersion> versions_;
354
355
  // `atomic_update_versions_` is for ensuring all-or-nothing AtomicGroup
356
  // recoveries.  When `atomic_update_versions_` is nonempty, it serves as a
357
  // barrier to updating `versions_` until all its values are populated.
358
  std::unordered_map<uint32_t, PointInTimeVersion> atomic_update_versions_;
359
  // `atomic_update_versions_missing_` counts the null `version`s in
360
  // `atomic_update_versions_`.
361
  size_t atomic_update_versions_missing_;
362
363
  bool in_atomic_group_ = false;
364
365
  // When true (set only via the OpenAndCompact remote-compaction path),
366
  // recovery trusts the MANIFEST when reconstructing the LSM version: it does
367
  // not stat/open SST or blob files to classify them found/missing (VerifyFile
368
  // / VerifyBlobFile return OK) and does not open candidate-version table
369
  // handlers (LoadTableHandlers is skipped). This prevents both a fatal open of
370
  // a transient/obsolete file (added then later deleted in the MANIFEST,
371
  // already physically removed by a live primary) and a rollback to an earlier,
372
  // wrong LSM shape. Safe because the compaction's input files are
373
  // ref-protected from deletion by the primary, so any missing file is
374
  // necessarily a non-input file the compaction never reads; inputs are opened
375
  // (and unique-id verified) on demand by the compaction itself.
376
  //
377
  // NOTE: This does not yet suppress the bounded file-property sampling in
378
  // Version::PrepareAppend -> UpdateAccumulatedStats, which still reads some
379
  // non-input files' properties (tolerating any that are missing). That path is
380
  // needed today to populate input FileMetaData::num_entries for compaction
381
  // input-record-count verification. A strict "open only compaction inputs"
382
  // recovery is a planned follow-up that also skips that sampling and sources
383
  // input num_entries from the worker's own input-table-properties read.
384
  const bool trust_manifest_recovery_ = false;
385
386
 private:
387
  bool AtomicUpdateVersionsCompleted();
388
  bool AtomicUpdateVersionsContains(uint32_t cfid);
389
  void AtomicUpdateVersionsDropCf(uint32_t cfid);
390
391
  // This function is called for Version updates for column families in an
392
  // incomplete atomic update. It buffers them in `atomic_update_versions_`.
393
  void AtomicUpdateVersionsPut(PointInTimeVersion pit_version);
394
395
  // This function is called upon completion of an atomic update. It applies the
396
  // updates buffered in `atomic_update_versions_` to `versions_`.
397
  void AtomicUpdateVersionsApply();
398
399
  // The log number of the Version last installed for each column family. See
400
  // GetInstalledVersionLogNumber().
401
  std::unordered_map<uint32_t, uint64_t> installed_version_log_numbers_;
402
};
403
404
// A class similar to `VersionEditHandlerPointInTime` that parse MANIFEST and
405
// builds point in time version.
406
// `ManifestTailer` supports reading one MANIFEST file in multiple tailing
407
// attempts and supports switching to a different MANIFEST after
408
// `PrepareToReadNewManifest` is called. This class is used by secondary and
409
// follower instance.
410
class ManifestTailer : public VersionEditHandlerPointInTime {
411
 public:
412
  explicit ManifestTailer(std::vector<ColumnFamilyDescriptor> column_families,
413
                          VersionSet* version_set,
414
                          const std::shared_ptr<IOTracer>& io_tracer,
415
                          const ReadOptions& read_options,
416
                          bool trust_manifest_recovery,
417
                          EpochNumberRequirement epoch_number_requirement =
418
                              EpochNumberRequirement::kMustPresent)
419
0
      : VersionEditHandlerPointInTime(
420
0
            /*read_only=*/true, column_families, version_set, io_tracer,
421
0
            read_options,
422
0
            /*allow_incomplete_valid_version=*/false, trust_manifest_recovery,
423
0
            epoch_number_requirement),
424
0
        mode_(Mode::kRecovery) {}
425
426
  Status VerifyFile(ColumnFamilyData* cfd, const std::string& fpath, int level,
427
                    const FileMetaData& fmeta) override;
428
429
0
  void PrepareToReadNewManifest() {
430
0
    initialized_ = false;
431
0
    ResetReadState();
432
0
  }
433
434
0
  std::unordered_set<ColumnFamilyData*>& GetUpdatedColumnFamilies() {
435
0
    return cfds_changed_;
436
0
  }
437
438
  std::vector<std::string> GetAndClearIntermediateFiles();
439
440
 protected:
441
  Status Initialize() override;
442
443
0
  bool MustOpenAllColumnFamilies() const override { return false; }
444
445
  Status ApplyVersionEdit(VersionEdit& edit, ColumnFamilyData** cfd) override;
446
447
  Status OnColumnFamilyAdd(VersionEdit& edit, ColumnFamilyData** cfd) override;
448
449
  void CheckIterationResult(const log::Reader& reader, Status* s) override;
450
451
  enum Mode : uint8_t {
452
    kRecovery = 0,
453
    kCatchUp = 1,
454
  };
455
456
  Mode mode_;
457
  std::unordered_set<ColumnFamilyData*> cfds_changed_;
458
};
459
460
class DumpManifestHandler : public VersionEditHandler {
461
 public:
462
  DumpManifestHandler(std::vector<ColumnFamilyDescriptor> column_families,
463
                      VersionSet* version_set,
464
                      const std::shared_ptr<IOTracer>& io_tracer,
465
                      const ReadOptions& read_options, bool verbose, bool hex,
466
                      bool json)
467
0
      : VersionEditHandler(
468
0
            /*read_only=*/true, column_families, version_set,
469
0
            /*track_found_and_missing_files=*/false,
470
0
            /*no_error_if_files_missing=*/false, io_tracer, read_options,
471
0
            /*skip_load_table_files=*/true,
472
0
            /*allow_incomplete_valid_version=*/false,
473
0
            /*epoch_number_requirement=*/EpochNumberRequirement::kMustPresent),
474
0
        verbose_(verbose),
475
0
        hex_(hex),
476
0
        json_(json),
477
0
        count_(0) {
478
0
    cf_to_cmp_names_.reset(new std::unordered_map<uint32_t, std::string>());
479
0
  }
480
481
0
  ~DumpManifestHandler() override {}
482
483
0
  Status ApplyVersionEdit(VersionEdit& edit, ColumnFamilyData** cfd) override {
484
    // Write out each individual edit
485
0
    if (json_) {
486
      // Print out DebugStrings. Can include non-terminating null characters.
487
0
      std::string edit_dump_str = edit.DebugJSON(count_, hex_);
488
0
      fwrite(edit_dump_str.data(), sizeof(char), edit_dump_str.size(), stdout);
489
0
      fwrite("\n", sizeof(char), 1, stdout);
490
0
    } else if (verbose_) {
491
      // Print out DebugStrings. Can include non-terminating null characters.
492
0
      std::string edit_dump_str = edit.DebugString(hex_);
493
      fwrite(edit_dump_str.data(), sizeof(char), edit_dump_str.size(), stdout);
494
0
    }
495
0
    ++count_;
496
0
    return VersionEditHandler::ApplyVersionEdit(edit, cfd);
497
0
  }
498
499
  void CheckIterationResult(const log::Reader& reader, Status* s) override;
500
501
 private:
502
  const bool verbose_;
503
  const bool hex_;
504
  const bool json_;
505
  int count_;
506
};
507
508
}  // namespace ROCKSDB_NAMESPACE