Coverage Report

Created: 2026-09-28 07:52

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/src/rocksdb/db/external_sst_file_ingestion_job.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 "db/external_sst_file_ingestion_job.h"
7
8
#include <algorithm>
9
#include <cinttypes>
10
#include <string>
11
#include <unordered_set>
12
#include <vector>
13
14
#include "db/builder.h"
15
#include "db/db_impl/db_impl.h"
16
#include "db/version_edit.h"
17
#include "file/file_util.h"
18
#include "file/random_access_file_reader.h"
19
#include "logging/logging.h"
20
#include "monitoring/statistics_impl.h"
21
#include "options/options_helper.h"
22
#include "table/merging_iterator.h"
23
#include "table/prepared_file_info.h"
24
#include "table/sst_file_writer_collectors.h"
25
#include "table/table_builder.h"
26
#include "table/unique_id_impl.h"
27
#include "test_util/sync_point.h"
28
#include "util/udt_util.h"
29
30
namespace ROCKSDB_NAMESPACE {
31
32
0
bool ExternalSstFileIngestionJob::SupportsAtomicReplaceRangeTombstone() const {
33
0
  return cfd_->is_delete_range_supported() &&
34
0
         cfd_->ioptions().compaction_style == kCompactionStyleUniversal &&
35
0
         ingestion_options_.allow_global_seqno &&
36
0
         !ingestion_options_.allow_db_generated_files &&
37
0
         !ingestion_options_.snapshot_consistency &&
38
0
         ucmp_->timestamp_size() == 0;
39
0
}
40
41
bool ExternalSstFileIngestionJob::CanUseAtomicReplaceRangeTombstone(
42
0
    const std::optional<RangeOpt>& atomic_replace_range) const {
43
0
  return atomic_replace_range.has_value() &&
44
0
         atomic_replace_range->start.has_value() &&
45
0
         atomic_replace_range->limit.has_value() &&
46
0
         SupportsAtomicReplaceRangeTombstone() &&
47
0
         !ingestion_options_.fail_if_not_bottommost_level;
48
0
}
49
50
bool ExternalSstFileIngestionJob::HasPartialOverlap(
51
0
    const VersionStorageInfo* vstorage) const {
52
0
  assert(vstorage != nullptr);
53
0
  assert(atomic_replace_range_.has_value());
54
0
  assert(!atomic_replace_range_->unset());
55
0
  for (int level = 0; level < cfd_->NumberLevels(); ++level) {
56
0
    for (const auto* file : vstorage->LevelFiles(level)) {
57
0
      if (file_range_checker_.Overlaps(*atomic_replace_range_, file->smallest,
58
0
                                       file->largest) &&
59
0
          !file_range_checker_.Contains(*atomic_replace_range_, file->smallest,
60
0
                                        file->largest)) {
61
0
        return true;
62
0
      }
63
0
    }
64
0
  }
65
0
  return false;
66
0
}
67
68
size_t ExternalSstFileIngestionJob::NumFilesToPrepare(
69
    size_t num_external_files,
70
0
    const std::optional<RangeOpt>& atomic_replace_range) const {
71
0
  return num_external_files +
72
0
         static_cast<size_t>(
73
0
             CanUseAtomicReplaceRangeTombstone(atomic_replace_range));
74
0
}
75
76
Status ExternalSstFileIngestionJob::PrepareAtomicReplaceRangeTombstone(
77
    const Slice& start, const Slice& limit, uint64_t file_number,
78
0
    SuperVersion* super_version) {
79
0
  assert(cfd_->is_delete_range_supported());
80
0
  assert(cfd_->ioptions().compaction_style == kCompactionStyleUniversal);
81
0
  assert(ingestion_options_.allow_global_seqno);
82
0
  assert(!ingestion_options_.allow_db_generated_files);
83
0
  assert(!ingestion_options_.fail_if_not_bottommost_level);
84
0
  assert(!ingestion_options_.snapshot_consistency);
85
0
  assert(ucmp_->timestamp_size() == 0);
86
87
0
  const uint32_t path_id = 0;
88
0
  const std::string file_path =
89
0
      TableFileName(cfd_->ioptions().cf_paths, file_number, path_id);
90
0
  ExternalSstFileInfo external_file_info;
91
0
  Status status;
92
0
  {
93
0
    Options options(
94
0
        BuildDBOptions(db_options_, mutable_db_options_),
95
0
        BuildColumnFamilyOptions(cfd_->initial_cf_options(),
96
0
                                 super_version->mutable_cf_options));
97
0
    SstFileWriter writer(env_options_, options, ucmp_);
98
0
    status = writer.Open(
99
0
        file_path, super_version->mutable_cf_options.default_write_temperature);
100
0
    if (status.ok()) {
101
0
      status = writer.DeleteRange(start, limit);
102
0
    }
103
0
    if (status.ok()) {
104
0
      status = writer.Finish(&external_file_info);
105
0
    }
106
0
  }
107
0
  if (!status.ok()) {
108
0
    fs_->DeleteFile(file_path, IOOptions(), nullptr).PermitUncheckedError();
109
0
    return status;
110
0
  }
111
112
0
  IngestedFileInfo tombstone;
113
0
  tombstone.file_temperature =
114
0
      super_version->mutable_cf_options.default_write_temperature;
115
0
  tombstone.prefetch_lmax_index_and_filter_blocks =
116
0
      ingestion_options_.prefetch_lmax_index_and_filter_blocks;
117
0
  const PreparedFileInfo* prepared_file_info =
118
0
      ingestion_options_.write_global_seqno
119
0
          ? nullptr
120
0
          : external_file_info.prepared_file_info.get();
121
0
  status = GetIngestedFileInfo(file_path, file_number, prepared_file_info,
122
0
                               &tombstone, super_version);
123
0
  if (!status.ok()) {
124
0
    fs_->DeleteFile(file_path, IOOptions(), nullptr).PermitUncheckedError();
125
0
    return status;
126
0
  }
127
128
0
  tombstone.internal_file_path = file_path;
129
0
  tombstone.copy_file = true;
130
0
  tombstone.generated_for_ingestion = true;
131
0
  tombstone.file_checksum = external_file_info.file_checksum;
132
0
  tombstone.file_checksum_func_name =
133
0
      external_file_info.file_checksum_func_name;
134
0
  atomic_replace_range_tombstone_.emplace(std::move(tombstone));
135
0
  return Status::OK();
136
0
}
137
138
0
void ExternalSstFileIngestionJob::ActivateAtomicReplaceRangeTombstone() {
139
0
  assert(atomic_replace_range_tombstone_.has_value());
140
  // Prepare rejects an empty external file list and verifies every input file
141
  // is contained in the replacement range. The covering tombstone must
142
  // therefore overlap at least one input file.
143
0
  assert(!files_to_ingest_.empty());
144
0
  file_batches_to_ingest_.clear();
145
0
  files_to_ingest_.emplace_back(std::move(*atomic_replace_range_tombstone_));
146
0
  std::rotate(files_to_ingest_.begin(), files_to_ingest_.end() - 1,
147
0
              files_to_ingest_.end());
148
0
  atomic_replace_range_tombstone_.reset();
149
0
  atomic_replace_range_tombstone_active_ = true;
150
0
  files_overlap_ = ComputeFilesOverlap(files_to_ingest_);
151
0
  assert(files_overlap_);
152
0
  DivideInputFilesIntoBatches();
153
0
}
154
155
Status ExternalSstFileIngestionJob::Prepare(
156
    const std::vector<std::string>& external_files_paths,
157
    const std::vector<std::string>& files_checksums,
158
    const std::vector<std::string>& files_checksum_func_names,
159
    const std::vector<const PreparedFileInfo*>& file_infos,
160
    const std::optional<RangeOpt>& atomic_replace_range,
161
    const Temperature& file_temperature, uint64_t next_file_number,
162
0
    SuperVersion* sv) {
163
0
  Status status;
164
165
  // Read the information of files we are ingesting.
166
0
  for (size_t i = 0; i < external_files_paths.size(); i++) {
167
0
    const std::string& file_path = external_files_paths[i];
168
0
    IngestedFileInfo file_to_ingest;
169
    // For temperature, first assume it matches provided hint
170
0
    file_to_ingest.file_temperature = file_temperature;
171
0
    file_to_ingest.prefetch_lmax_index_and_filter_blocks =
172
0
        ingestion_options_.prefetch_lmax_index_and_filter_blocks;
173
0
    const PreparedFileInfo* prepared_file_info =
174
0
        file_infos.empty() ? nullptr : file_infos[i];
175
0
    status = GetIngestedFileInfo(file_path, next_file_number++,
176
0
                                 prepared_file_info, &file_to_ingest, sv);
177
0
    if (!status.ok()) {
178
0
      ROCKS_LOG_WARN(db_options_.info_log,
179
0
                     "Failed to get ingested file info: %s: %s",
180
0
                     file_path.c_str(), status.ToString().c_str());
181
0
      return status;
182
0
    }
183
184
    // Files generated in another DB or CF may have a different column family
185
    // ID, so we let it pass here.
186
0
    if (file_to_ingest.cf_id !=
187
0
            TablePropertiesCollectorFactory::Context::kUnknownColumnFamily &&
188
0
        file_to_ingest.cf_id != cfd_->GetID() &&
189
0
        !ingestion_options_.allow_db_generated_files) {
190
0
      return Status::InvalidArgument(
191
0
          "External file column family id don't match");
192
0
    }
193
194
0
    if (file_to_ingest.num_entries == 0 &&
195
0
        file_to_ingest.num_range_deletions == 0) {
196
0
      return Status::InvalidArgument("File contain no entries");
197
0
    }
198
199
0
    if (!file_to_ingest.smallest_internal_key.Valid() ||
200
0
        !file_to_ingest.largest_internal_key.Valid()) {
201
0
      return Status::Corruption("Generated table have corrupted keys");
202
0
    }
203
204
0
    files_to_ingest_.emplace_back(std::move(file_to_ingest));
205
0
  }
206
207
0
  auto num_files = files_to_ingest_.size();
208
0
  if (num_files == 0) {
209
0
    return Status::InvalidArgument("The list of files is empty");
210
0
  }
211
  // Detect whether the input files overlap one another; this drives how they
212
  // are divided into batches below.
213
0
  files_overlap_ = ComputeFilesOverlap(files_to_ingest_);
214
215
0
  if (atomic_replace_range.has_value()) {
216
0
    atomic_replace_range_.emplace();
217
218
0
    if (atomic_replace_range->start && atomic_replace_range->limit) {
219
0
      if (ucmp_->CompareWithoutTimestamp(*atomic_replace_range->start,
220
0
                                         *atomic_replace_range->limit) >= 0) {
221
0
        return Status::InvalidArgument(
222
0
            "Atomic replace range limit must be greater than start");
223
0
      }
224
      // User keys to internal keys (with timestamps)
225
0
      const size_t ts_sz = ucmp_->timestamp_size();
226
0
      std::string start_with_ts, limit_with_ts;
227
0
      auto [start, limit] = MaybeAddTimestampsToRange(
228
0
          atomic_replace_range->start, atomic_replace_range->limit, ts_sz,
229
0
          &start_with_ts, &limit_with_ts);
230
0
      assert(start.has_value());
231
0
      assert(limit.has_value());
232
0
      atomic_replace_range_->smallest_internal_key.Set(
233
0
          *start, kMaxSequenceNumber, kValueTypeForSeek);
234
0
      atomic_replace_range_->largest_internal_key =
235
0
          RangeTombstone(*start, *limit, kMaxSequenceNumber).SerializeEndKey();
236
      // Check files to ingest against replace range
237
0
      for (size_t i = 0; i < num_files; i++) {
238
0
        if (!file_range_checker_.Contains(*atomic_replace_range_,
239
0
                                          files_to_ingest_[i])) {
240
0
          return Status::InvalidArgument(
241
0
              "Atomic replace range does not contain all files");
242
0
        }
243
0
      }
244
0
    } else {
245
      // Currently if either bound is not present, both must be
246
0
      assert(atomic_replace_range->start.has_value() == false);
247
0
      assert(atomic_replace_range->limit.has_value() == false);
248
0
      assert(atomic_replace_range_->smallest_internal_key.unset());
249
0
      assert(atomic_replace_range_->largest_internal_key.unset());
250
0
    }
251
0
  }
252
253
0
  if (files_overlap_) {
254
0
    if (ingestion_options_.ingest_behind) {
255
0
      return Status::NotSupported(
256
0
          "Files with overlapping ranges cannot be ingested with ingestion "
257
0
          "behind mode.");
258
0
    }
259
260
    // Overlapping files need at least two different sequence numbers. If
261
    // settings disables global seqno, ingestion will fail anyway, so fail
262
    // fast in prepare.
263
0
    if (!ingestion_options_.allow_global_seqno &&
264
0
        !ingestion_options_.allow_db_generated_files) {
265
0
      return Status::InvalidArgument(
266
0
          "Global seqno is required, but disabled (because external files key "
267
0
          "range overlaps).");
268
0
    }
269
270
0
    if (ucmp_->timestamp_size() > 0) {
271
0
      return Status::NotSupported(
272
0
          "Files with overlapping ranges cannot be ingested to column "
273
0
          "family with user-defined timestamp enabled.");
274
0
    }
275
0
  }
276
277
  // Copy/Move external files into DB
278
0
  std::unordered_set<size_t> ingestion_path_ids;
279
0
  for (IngestedFileInfo& f : files_to_ingest_) {
280
0
    f.copy_file = false;
281
0
    const std::string path_outside_db = f.external_file_path;
282
0
    const std::string path_inside_db = TableFileName(
283
0
        cfd_->ioptions().cf_paths, f.fd.GetNumber(), f.fd.GetPathId());
284
0
    if (ingestion_options_.move_files || ingestion_options_.link_files) {
285
0
      status =
286
0
          fs_->LinkFile(path_outside_db, path_inside_db, IOOptions(), nullptr);
287
0
      if (status.ok()) {
288
        // It is unsafe to assume application had sync the file and file
289
        // directory before ingest the file. For integrity of RocksDB we need
290
        // to sync the file.
291
0
        TEST_SYNC_POINT("ExternalSstFileIngestionJob::BeforeSyncIngestedFile");
292
0
        Status s = fs_->SyncFile(path_inside_db, env_options_, IOOptions(),
293
0
                                 db_options_.use_fsync, nullptr);
294
0
        TEST_SYNC_POINT("ExternalSstFileIngestionJob::AfterSyncIngestedFile");
295
0
        TEST_SYNC_POINT_CALLBACK(
296
0
            "ExternalSstFileIngestionJob::CheckSyncReturnCode", &s);
297
        // Some file systems (especially remote/distributed) don't support
298
        // explicitly syncing the file and don't require it. Ignore the
299
        // NotSupported error in that case.
300
0
        if (!s.IsNotSupported()) {
301
0
          status = s;
302
0
          if (!status.ok()) {
303
0
            ROCKS_LOG_WARN(db_options_.info_log,
304
0
                           "Failed to sync ingested file %s: %s",
305
0
                           path_inside_db.c_str(), status.ToString().c_str());
306
0
          }
307
0
        }
308
0
      } else if (status.IsNotSupported() &&
309
0
                 ingestion_options_.failed_move_fall_back_to_copy) {
310
        // Original file is on a different FS, use copy instead of hard linking.
311
0
        f.copy_file = true;
312
0
        ROCKS_LOG_INFO(db_options_.info_log,
313
0
                       "Tried to link file %s but it's not supported : %s",
314
0
                       path_outside_db.c_str(), status.ToString().c_str());
315
0
      } else {
316
0
        ROCKS_LOG_WARN(db_options_.info_log, "Failed to link file %s to %s: %s",
317
0
                       path_outside_db.c_str(), path_inside_db.c_str(),
318
0
                       status.ToString().c_str());
319
0
      }
320
0
    } else {
321
0
      f.copy_file = true;
322
0
    }
323
324
0
    if (f.copy_file) {
325
0
      TEST_SYNC_POINT_CALLBACK("ExternalSstFileIngestionJob::Prepare:CopyFile",
326
0
                               nullptr);
327
      // Always determining the destination temperature from the ingested-to
328
      // level would be difficult because in general we only find out the level
329
      // ingested to later, during Run().
330
      // However, we can guarantee "last level" temperature for when the user
331
      // requires ingestion to the last level.
332
0
      Temperature dst_temp =
333
0
          (ingestion_options_.ingest_behind ||
334
0
           ingestion_options_.fail_if_not_bottommost_level)
335
0
              ? sv->mutable_cf_options.last_level_temperature
336
0
              : sv->mutable_cf_options.default_write_temperature;
337
      // Note: CopyFile also syncs the new file.
338
0
      status = CopyFile(fs_.get(), path_outside_db, f.file_temperature,
339
0
                        path_inside_db, dst_temp, 0, db_options_.use_fsync,
340
0
                        io_tracer_);
341
      // The destination of the copy will be ingested
342
0
      f.file_temperature = dst_temp;
343
344
0
      if (!status.ok()) {
345
0
        ROCKS_LOG_WARN(db_options_.info_log, "Failed to copy file %s to %s: %s",
346
0
                       path_outside_db.c_str(), path_inside_db.c_str(),
347
0
                       status.ToString().c_str());
348
0
      }
349
0
    } else {
350
      // Note: we currently assume that linking files does not cross
351
      // temperatures, so no need to change f.file_temperature
352
0
    }
353
0
    TEST_SYNC_POINT("ExternalSstFileIngestionJob::Prepare:FileAdded");
354
0
    if (!status.ok()) {
355
0
      break;
356
0
    }
357
0
    f.internal_file_path = path_inside_db;
358
    // Initialize the checksum information of ingested files.
359
0
    f.file_checksum = kUnknownFileChecksum;
360
0
    f.file_checksum_func_name = kUnknownFileChecksumFuncName;
361
0
    ingestion_path_ids.insert(f.fd.GetPathId());
362
0
  }
363
364
0
  if (status.ok() && CanUseAtomicReplaceRangeTombstone(atomic_replace_range) &&
365
0
      HasPartialOverlap(sv->current->storage_info())) {
366
0
    assert(atomic_replace_range->start.has_value());
367
0
    assert(atomic_replace_range->limit.has_value());
368
0
    status = PrepareAtomicReplaceRangeTombstone(*atomic_replace_range->start,
369
0
                                                *atomic_replace_range->limit,
370
0
                                                next_file_number, sv);
371
0
    if (status.ok()) {
372
0
      assert(atomic_replace_range_tombstone_.has_value());
373
0
      ingestion_path_ids.insert(
374
0
          atomic_replace_range_tombstone_->fd.GetPathId());
375
0
    }
376
0
  }
377
378
0
  TEST_SYNC_POINT("ExternalSstFileIngestionJob::BeforeSyncDir");
379
0
  if (status.ok()) {
380
0
    for (auto path_id : ingestion_path_ids) {
381
0
      status = directories_->GetDataDir(path_id)->FsyncWithDirOptions(
382
0
          IOOptions(), nullptr,
383
0
          DirFsyncOptions(DirFsyncOptions::FsyncReason::kNewFileSynced));
384
0
      if (!status.ok()) {
385
0
        ROCKS_LOG_WARN(db_options_.info_log,
386
0
                       "Failed to sync directory %" ROCKSDB_PRIszt
387
0
                       " while ingest file: %s",
388
0
                       path_id, status.ToString().c_str());
389
0
        break;
390
0
      }
391
0
    }
392
0
  }
393
0
  TEST_SYNC_POINT("ExternalSstFileIngestionJob::AfterSyncDir");
394
395
  // Generate and check the sst file checksum. Note that, if
396
  // IngestExternalFileOptions::write_global_seqno is true, we will not update
397
  // the checksum information in the files_to_ingests_ here, since the file is
398
  // updated with the new global_seqno. After global_seqno is updated, DB will
399
  // generate the new checksum and store it in the Manifest. In all other cases
400
  // if ingestion_options_.write_global_seqno == true and
401
  // verify_file_checksum is false, we only check the checksum function name.
402
0
  if (status.ok() && db_options_.file_checksum_gen_factory != nullptr) {
403
0
    if (ingestion_options_.verify_file_checksum == false &&
404
0
        files_checksums.size() == files_to_ingest_.size() &&
405
0
        files_checksum_func_names.size() == files_to_ingest_.size()) {
406
      // Only when verify_file_checksum == false and the checksum for ingested
407
      // files are provided, DB will use the provided checksum and does not
408
      // generate the checksum for ingested files.
409
0
      need_generate_file_checksum_ = false;
410
0
    } else {
411
0
      need_generate_file_checksum_ = true;
412
0
    }
413
0
    std::vector<std::string> generated_checksums;
414
0
    std::vector<std::string> generated_checksum_func_names;
415
    // Step 1: generate the checksum for ingested sst file.
416
0
    if (need_generate_file_checksum_) {
417
0
      for (size_t i = 0; i < files_to_ingest_.size(); i++) {
418
0
        std::string generated_checksum;
419
0
        std::string generated_checksum_func_name;
420
0
        std::string requested_checksum_func_name =
421
0
            i < files_checksum_func_names.size() ? files_checksum_func_names[i]
422
0
                                                 : "";
423
        // TODO: rate limit file reads for checksum calculation during file
424
        // ingestion.
425
        // TODO: plumb Env::IOActivity
426
0
        ReadOptions ro;
427
        // Pass user-provided checksums through FileOptions when available.
428
        // The caller may not have provided checksums at all (empty vectors),
429
        // so we guard with a bounds check.
430
0
        FileOptions fopts;
431
0
        if (i < files_checksums.size()) {
432
0
          fopts.file_checksum = files_checksums[i];
433
0
        }
434
0
        if (i < files_checksum_func_names.size()) {
435
0
          fopts.file_checksum_func_name = files_checksum_func_names[i];
436
0
        } else {
437
0
          fopts.file_checksum_func_name = kNoFileChecksumFuncName;
438
0
        }
439
0
        IOStatus io_s = GenerateOneFileChecksum(
440
0
            fs_.get(), files_to_ingest_[i].internal_file_path,
441
0
            db_options_.file_checksum_gen_factory.get(),
442
0
            requested_checksum_func_name, &generated_checksum,
443
0
            &generated_checksum_func_name,
444
0
            ingestion_options_.verify_checksums_readahead_size,
445
0
            db_options_.allow_mmap_reads, io_tracer_,
446
0
            db_options_.rate_limiter.get(), ro, db_options_.stats,
447
0
            db_options_.clock, fopts);
448
0
        if (!io_s.ok()) {
449
0
          status = io_s;
450
0
          ROCKS_LOG_WARN(db_options_.info_log,
451
0
                         "Sst file checksum generation of file: %s failed: %s",
452
0
                         files_to_ingest_[i].internal_file_path.c_str(),
453
0
                         status.ToString().c_str());
454
0
          break;
455
0
        }
456
0
        if (ingestion_options_.write_global_seqno == false) {
457
0
          files_to_ingest_[i].file_checksum = generated_checksum;
458
0
          files_to_ingest_[i].file_checksum_func_name =
459
0
              generated_checksum_func_name;
460
0
        }
461
0
        generated_checksums.push_back(generated_checksum);
462
0
        generated_checksum_func_names.push_back(generated_checksum_func_name);
463
0
      }
464
0
    }
465
466
    // Step 2: based on the verify_file_checksum and ingested checksum
467
    // information, do the verification.
468
0
    if (status.ok()) {
469
0
      if (files_checksums.size() == files_to_ingest_.size() &&
470
0
          files_checksum_func_names.size() == files_to_ingest_.size()) {
471
        // Verify the checksum and checksum function name.
472
0
        if (ingestion_options_.verify_file_checksum) {
473
0
          for (size_t i = 0; i < files_to_ingest_.size(); i++) {
474
0
            if (files_checksum_func_names[i] !=
475
0
                generated_checksum_func_names[i]) {
476
0
              status = Status::InvalidArgument(
477
0
                  "DB file checksum gen factory " +
478
0
                  std::string(db_options_.file_checksum_gen_factory->Name()) +
479
0
                  " generated checksum function name " +
480
0
                  generated_checksum_func_names[i] + " for file " +
481
0
                  external_files_paths[i] +
482
0
                  " which does not match requested/provided " +
483
0
                  files_checksum_func_names[i]);
484
0
              break;
485
0
            }
486
0
            if (files_checksums[i] != generated_checksums[i]) {
487
0
              status = Status::Corruption(
488
0
                  "Checksum verification mismatch for ingestion file " +
489
0
                  external_files_paths[i] + " using function " +
490
0
                  generated_checksum_func_names[i] + ". Expected: " +
491
0
                  Slice(files_checksums[i]).ToString(/*hex=*/true) +
492
0
                  " Computed: " +
493
0
                  Slice(generated_checksums[i]).ToString(/*hex=*/true));
494
0
              break;
495
0
            }
496
0
          }
497
0
        } else {
498
          // If verify_file_checksum is not enabled, we only verify the factory
499
          // recognizes the checksum function name. If it does not match, fail
500
          // the ingestion. If matches, we trust the ingested checksum
501
          // information and store in the Manifest.
502
0
          for (size_t i = 0; i < files_to_ingest_.size(); i++) {
503
0
            FileChecksumGenContext gen_context;
504
0
            gen_context.file_name = files_to_ingest_[i].internal_file_path;
505
0
            gen_context.requested_checksum_func_name =
506
0
                files_checksum_func_names[i];
507
0
            auto file_checksum_gen =
508
0
                db_options_.file_checksum_gen_factory
509
0
                    ->CreateFileChecksumGenerator(gen_context);
510
511
0
            if (file_checksum_gen == nullptr ||
512
0
                files_checksum_func_names[i] != file_checksum_gen->Name()) {
513
0
              status = Status::InvalidArgument(
514
0
                  "Checksum function name " + files_checksum_func_names[i] +
515
0
                  " for file " + external_files_paths[i] +
516
0
                  " not recognized by DB checksum gen factory" +
517
0
                  db_options_.file_checksum_gen_factory->Name() +
518
0
                  (file_checksum_gen ? (" Returned function " +
519
0
                                        std::string(file_checksum_gen->Name()))
520
0
                                     : ""));
521
0
              break;
522
0
            }
523
0
            files_to_ingest_[i].file_checksum = files_checksums[i];
524
0
            files_to_ingest_[i].file_checksum_func_name =
525
0
                files_checksum_func_names[i];
526
0
          }
527
0
        }
528
0
      } else if (files_checksums.size() != files_checksum_func_names.size() ||
529
0
                 files_checksums.size() != 0) {
530
        // The checksum or checksum function name vector are not both empty
531
        // and they are incomplete.
532
0
        status = Status::InvalidArgument(
533
0
            "The checksum information of ingested sst files are nonempty and "
534
0
            "the size of checksums or the size of the checksum function "
535
0
            "names does not match with the number of ingested sst files");
536
0
      }
537
0
      if (!status.ok()) {
538
0
        ROCKS_LOG_WARN(db_options_.info_log, "Ingestion failed: %s",
539
0
                       status.ToString().c_str());
540
0
      }
541
0
    }
542
0
  }
543
544
0
  if (status.ok()) {
545
0
    DivideInputFilesIntoBatches();
546
0
  }
547
548
0
  return status;
549
0
}
550
551
0
void ExternalSstFileIngestionJob::DivideInputFilesIntoBatches() {
552
0
  if (!files_overlap_) {
553
    // No overlap, treat as one batch without the need of tracking overall batch
554
    // range.
555
0
    file_batches_to_ingest_.emplace_back(/* _track_batch_range= */ false);
556
0
    for (auto& file : files_to_ingest_) {
557
0
      file_batches_to_ingest_.back().AddFile(&file, file_range_checker_);
558
0
    }
559
0
    return;
560
0
  }
561
562
0
  file_batches_to_ingest_.emplace_back(/* _track_batch_range= */ true);
563
0
  for (auto& file : files_to_ingest_) {
564
0
    if (!file_batches_to_ingest_.back().unset() &&
565
0
        file_range_checker_.Overlaps(file_batches_to_ingest_.back(), file,
566
0
                                     /* known_sorted= */ false)) {
567
0
      file_batches_to_ingest_.emplace_back(/* _track_batch_range= */ true);
568
0
    }
569
0
    file_batches_to_ingest_.back().AddFile(&file, file_range_checker_);
570
0
  }
571
0
}
572
573
bool ExternalSstFileIngestionJob::ComputeFilesOverlap(
574
0
    const autovector<IngestedFileInfo>& files) const {
575
0
  const size_t num_files = files.size();
576
0
  if (num_files <= 1) {
577
0
    return false;
578
0
  }
579
  // Verify whether the files have overlapping ranges by sorting copies of the
580
  // file ranges and checking adjacent pairs.
581
0
  autovector<const IngestedFileInfo*> sorted_files;
582
0
  for (size_t i = 0; i < num_files; i++) {
583
0
    sorted_files.push_back(&files[i]);
584
0
  }
585
0
  std::sort(sorted_files.begin(), sorted_files.end(), file_range_checker_);
586
0
  for (size_t i = 0; i + 1 < num_files; i++) {
587
0
    if (file_range_checker_.Overlaps(*sorted_files[i], *sorted_files[i + 1],
588
0
                                     /* known_sorted= */ true)) {
589
0
      return true;
590
0
    }
591
0
  }
592
0
  return false;
593
0
}
594
595
Status ExternalSstFileIngestionJob::MergeForSameColumnFamily(
596
0
    ExternalSstFileIngestionJob* other) {
597
0
  assert(other != nullptr);
598
0
  assert(other != this);
599
0
  assert(cfd_ == other->cfd_);
600
0
  if (atomic_replace_range_.has_value() ||
601
0
      other->atomic_replace_range_.has_value()) {
602
0
    return Status::NotSupported(
603
0
        "cannot merge file ingestion handles for the same column family when "
604
0
        "atomic_replace_range is used");
605
0
  }
606
0
  if (!(ingestion_options_ == other->ingestion_options_)) {
607
0
    return Status::InvalidArgument(
608
0
        "file ingestion handles for the same column family must be prepared "
609
0
        "with the same IngestExternalFileOptions");
610
0
  }
611
  // Append the other job's prepared files after this job's so that, for any
612
  // overlapping keys, the other job's data wins via a higher assigned sequence
613
  // number -- the same semantics as passing all the files to a single ingestion
614
  // call in this order. Recompute overlap and rebuild the batches over the
615
  // union.
616
0
  for (IngestedFileInfo& file : other->files_to_ingest_) {
617
0
    files_to_ingest_.push_back(std::move(file));
618
0
  }
619
0
  other->files_to_ingest_.clear();
620
0
  other->file_batches_to_ingest_.clear();
621
0
  files_overlap_ = ComputeFilesOverlap(files_to_ingest_);
622
0
  file_batches_to_ingest_.clear();
623
0
  DivideInputFilesIntoBatches();
624
0
  return Status::OK();
625
0
}
626
627
Status ExternalSstFileIngestionJob::NeedsFlush(bool* flush_needed,
628
0
                                               SuperVersion* super_version) {
629
0
  Status status;
630
0
  if (atomic_replace_range_.has_value() && atomic_replace_range_->unset()) {
631
    // For replacing whole CF, we can simply check whether memtable is empty
632
0
    *flush_needed = !super_version->mem->IsEmpty();
633
0
  } else {
634
0
    autovector<UserKeyRange> ranges;
635
0
    if (atomic_replace_range_.has_value()) {
636
0
      assert(!atomic_replace_range_->smallest_internal_key.unset());
637
0
      assert(!atomic_replace_range_->largest_internal_key.unset());
638
      // NOTE: we already checked in Prepare() that the atomic_replace_range
639
      // covers all the files_to_ingest.
640
0
      ranges.emplace_back(
641
0
          atomic_replace_range_->smallest_internal_key.user_key(),
642
0
          atomic_replace_range_->largest_internal_key.user_key());
643
0
    } else {
644
0
      ranges.reserve(files_to_ingest_.size());
645
0
      for (const IngestedFileInfo& file_to_ingest : files_to_ingest_) {
646
0
        ranges.emplace_back(file_to_ingest.start_ukey,
647
0
                            file_to_ingest.limit_ukey);
648
0
      }
649
0
    }
650
0
    status = cfd_->RangesOverlapWithMemtables(
651
0
        ranges, super_version, db_options_.allow_data_in_errors, flush_needed,
652
0
        atomic_replace_range_.has_value() && !atomic_replace_range_->unset());
653
0
    if (!status.ok()) {
654
0
      ROCKS_LOG_WARN(db_options_.info_log,
655
0
                     "Failed to check ranges overlap with memtables: %s",
656
0
                     status.ToString().c_str());
657
0
    }
658
0
  }
659
0
  if (status.ok() && *flush_needed) {
660
0
    if (!ingestion_options_.allow_blocking_flush) {
661
0
      status = Status::InvalidArgument("External file requires flush");
662
0
    }
663
0
    if (ucmp_->timestamp_size() > 0) {
664
0
      status = Status::InvalidArgument(
665
0
          "Column family enables user-defined timestamps, please make "
666
0
          "sure the key range (without timestamp) of external file does not "
667
0
          "overlap with key range in the memtables.");
668
0
    }
669
0
  }
670
0
  return status;
671
0
}
672
673
// REQUIRES: we have become the only writer by entering both write_thread_ and
674
// nonmem_write_thread_
675
0
Status ExternalSstFileIngestionJob::Run() {
676
0
  SuperVersion* super_version = cfd_->GetSuperVersion();
677
  // If column family is flushed after Prepare and before Run, we should have a
678
  // specific state of Memtables. The mutable Memtable should be empty, and the
679
  // immutable Memtable list should be empty.
680
0
  if (flushed_before_run_ && (super_version->imm->NumNotFlushed() != 0 ||
681
0
                              !super_version->mem->IsEmpty())) {
682
0
    return Status::TryAgain(
683
0
        "Inconsistent memtable state detected when flushed before run.");
684
0
  }
685
0
  Status status;
686
#ifndef NDEBUG
687
  // We should never run the job with a memtable that is overlapping
688
  // with the files we are ingesting
689
  bool need_flush = false;
690
  status = NeedsFlush(&need_flush, super_version);
691
  if (!status.ok()) {
692
    ROCKS_LOG_WARN(db_options_.info_log,
693
                   "Failed to check if flush is needed: %s",
694
                   status.ToString().c_str());
695
    return status;
696
  }
697
  if (need_flush) {
698
    return Status::TryAgain("need_flush");
699
  }
700
  assert(status.ok() && need_flush == false);
701
#endif
702
703
0
  bool force_global_seqno = false;
704
705
0
  if (ingestion_options_.snapshot_consistency && !db_snapshots_->empty()) {
706
    // We need to assign a global sequence number to all the files even
707
    // if the don't overlap with any ranges since we have snapshots
708
0
    force_global_seqno = true;
709
0
  }
710
  // It is safe to use this instead of LastAllocatedSequence since we are
711
  // the only active writer, and hence they are equal
712
0
  SequenceNumber last_seqno = versions_->LastSequence();
713
0
  edit_.SetColumnFamily(cfd_->GetID());
714
715
0
  if (atomic_replace_range_.has_value()) {
716
0
    auto* vstorage = super_version->current->storage_info();
717
0
    if (atomic_replace_range_->unset()) {
718
0
      if (cfd_->compaction_picker()->IsCompactionInProgress()) {
719
0
        return Status::InvalidArgument(
720
0
            "Atomic replace range (full) overlaps with pending compaction");
721
0
      }
722
0
      for (int lvl = 0; lvl < cfd_->NumberLevels(); lvl++) {
723
0
        for (auto file : vstorage->LevelFiles(lvl)) {
724
          // Set up to delete file to be replaced
725
0
          edit_.DeleteFile(lvl, file->fd.GetNumber());
726
0
        }
727
0
      }
728
0
    } else {
729
0
      assert(!atomic_replace_range_->smallest_internal_key.unset());
730
0
      assert(!atomic_replace_range_->largest_internal_key.unset());
731
0
      bool has_partial_overlap = false;
732
0
      for (int lvl = 0; lvl < cfd_->NumberLevels(); lvl++) {
733
0
        if (cfd_->RangeOverlapWithCompaction(
734
0
                atomic_replace_range_->smallest_internal_key.user_key(),
735
0
                atomic_replace_range_->largest_internal_key.user_key(), lvl,
736
0
                /*range_limit_exclusive=*/true)) {
737
0
          return Status::InvalidArgument(
738
0
              "Atomic replace range overlaps with pending compaction");
739
0
        }
740
0
        for (auto file : vstorage->LevelFiles(lvl)) {
741
0
          if (file_range_checker_.Overlaps(*atomic_replace_range_,
742
0
                                           file->smallest, file->largest)) {
743
0
            if (file_range_checker_.Contains(*atomic_replace_range_,
744
0
                                             file->smallest, file->largest)) {
745
              // Set up to delete file to be replaced
746
0
              edit_.DeleteFile(lvl, file->fd.GetNumber());
747
0
            } else {
748
0
              has_partial_overlap = true;
749
0
            }
750
0
          }
751
0
        }
752
0
      }
753
0
      if (has_partial_overlap) {
754
0
        if (ingestion_options_.fail_if_not_bottommost_level &&
755
0
            SupportsAtomicReplaceRangeTombstone()) {
756
0
          return Status::TryAgain(
757
0
              "Atomic replace range partially overlaps with an existing file, "
758
0
              "so replacement files cannot all be ingested to Lmax");
759
0
        }
760
0
        if (!atomic_replace_range_tombstone_.has_value()) {
761
0
          if (SupportsAtomicReplaceRangeTombstone()) {
762
0
            return Status::TryAgain(
763
0
                "Atomic replace range acquired a partial overlap after "
764
0
                "preparation; retry ingestion");
765
0
          }
766
0
          return Status::InvalidArgument(
767
0
              "Atomic replace range partially overlaps with existing file");
768
0
        }
769
0
        ActivateAtomicReplaceRangeTombstone();
770
0
      }
771
0
    }
772
0
  }
773
774
  // Find levels to ingest into
775
0
  std::optional<int> prev_batch_uppermost_level;
776
  // Batches at the front contain older updates and are placed deeper in the
777
  // LSM tree than later overlapping batches.
778
0
  for (auto& batch : file_batches_to_ingest_) {
779
0
    int batch_uppermost_level = 0;
780
0
    status = AssignLevelsForOneBatch(batch, super_version, force_global_seqno,
781
0
                                     &last_seqno, &batch_uppermost_level,
782
0
                                     prev_batch_uppermost_level);
783
0
    if (!status.ok()) {
784
0
      ROCKS_LOG_WARN(db_options_.info_log,
785
0
                     "Failed to assign levels for one batch: %s",
786
0
                     status.ToString().c_str());
787
0
      return status;
788
0
    }
789
790
0
    prev_batch_uppermost_level = batch_uppermost_level;
791
0
  }
792
793
0
  CreateEquivalentFileIngestingCompactions();
794
0
  return status;
795
0
}
796
797
Status ExternalSstFileIngestionJob::AssignLevelsForOneBatch(
798
    FileBatchInfo& batch, SuperVersion* super_version, bool force_global_seqno,
799
    SequenceNumber* last_seqno, int* batch_uppermost_level,
800
0
    std::optional<int> prev_batch_uppermost_level) {
801
0
  Status status;
802
0
  assert(batch_uppermost_level);
803
0
  *batch_uppermost_level = std::numeric_limits<int>::max();
804
0
  for (IngestedFileInfo* file : batch.files) {
805
0
    assert(file);
806
0
    SequenceNumber assigned_seqno = 0;
807
0
    if (ingestion_options_.ingest_behind) {
808
0
      status = CheckLevelForIngestedBehindFile(file);
809
0
    } else {
810
0
      status = AssignLevelAndSeqnoForIngestedFile(
811
0
          super_version, force_global_seqno, cfd_->ioptions().compaction_style,
812
0
          *last_seqno, file, &assigned_seqno, prev_batch_uppermost_level);
813
0
    }
814
815
    // Modify the smallest/largest internal key to include the sequence number
816
    // that we just learned. Only overwrite sequence number zero. There could
817
    // be a nonzero sequence number already to indicate a range tombstone's
818
    // exclusive endpoint.
819
0
    ParsedInternalKey smallest_parsed, largest_parsed;
820
0
    if (status.ok()) {
821
0
      status = ParseInternalKey(*(file->smallest_internal_key.rep()),
822
0
                                &smallest_parsed, false /* log_err_key */);
823
0
    }
824
0
    if (status.ok()) {
825
0
      status = ParseInternalKey(*(file->largest_internal_key.rep()),
826
0
                                &largest_parsed, false /* log_err_key */);
827
0
    }
828
0
    if (!status.ok()) {
829
0
      ROCKS_LOG_WARN(db_options_.info_log, "Failed to parse internal key: %s",
830
0
                     status.ToString().c_str());
831
0
      return status;
832
0
    }
833
834
    // If any ingested file overlaps with the DB, it will fail here.
835
0
    if (ingestion_options_.allow_db_generated_files && assigned_seqno != 0) {
836
0
      return Status::InvalidArgument(
837
0
          "An ingested file overlaps with existing data in the DB and has been "
838
0
          "assigned a non-zero sequence number, which is not allowed when "
839
0
          "'allow_db_generated_files' is enabled.");
840
0
    }
841
842
0
    if (smallest_parsed.sequence == 0 && assigned_seqno != 0) {
843
0
      UpdateInternalKey(file->smallest_internal_key.rep(), assigned_seqno,
844
0
                        smallest_parsed.type);
845
0
    }
846
0
    if (largest_parsed.sequence == 0 && assigned_seqno != 0) {
847
0
      UpdateInternalKey(file->largest_internal_key.rep(), assigned_seqno,
848
0
                        largest_parsed.type);
849
0
    }
850
851
0
    status = AssignGlobalSeqnoForIngestedFile(file, assigned_seqno);
852
0
    if (!status.ok()) {
853
0
      ROCKS_LOG_WARN(
854
0
          db_options_.info_log,
855
0
          "Failed to assign global sequence number for ingested file: %s",
856
0
          status.ToString().c_str());
857
0
      return status;
858
0
    }
859
0
    TEST_SYNC_POINT_CALLBACK("ExternalSstFileIngestionJob::Run",
860
0
                             &assigned_seqno);
861
0
    assert(assigned_seqno == 0 || assigned_seqno == *last_seqno + 1);
862
0
    if (assigned_seqno > *last_seqno) {
863
0
      *last_seqno = assigned_seqno;
864
0
    }
865
0
    max_assigned_seqno_ = std::max(max_assigned_seqno_, assigned_seqno);
866
867
0
    status = GenerateChecksumForIngestedFile(file);
868
0
    if (!status.ok()) {
869
0
      ROCKS_LOG_WARN(db_options_.info_log,
870
0
                     "Failed to generate checksum for ingested file: %s",
871
0
                     status.ToString().c_str());
872
0
      return status;
873
0
    }
874
875
    // We use the import time as the ancester time. This is the time the data
876
    // is written to the database.
877
0
    int64_t temp_current_time = 0;
878
0
    uint64_t current_time = kUnknownFileCreationTime;
879
0
    uint64_t oldest_ancester_time = kUnknownOldestAncesterTime;
880
0
    if (clock_->GetCurrentTime(&temp_current_time).ok()) {
881
0
      current_time = oldest_ancester_time =
882
0
          static_cast<uint64_t>(temp_current_time);
883
0
    }
884
0
    uint64_t tail_size = FileMetaData::CalculateTailSize(
885
0
        file->fd.GetFileSize(), file->table_properties);
886
887
0
    bool marked_for_compaction =
888
0
        file->table_properties.num_range_deletions == 1 &&
889
0
        (file->table_properties.num_entries ==
890
0
         file->table_properties.num_range_deletions);
891
0
    SequenceNumber smallest_seqno = file->assigned_seqno;
892
0
    SequenceNumber largest_seqno = file->assigned_seqno;
893
0
    if (ingestion_options_.allow_db_generated_files) {
894
0
      assert(file->assigned_seqno == 0);
895
0
      assert(file->smallest_seqno != kMaxSequenceNumber);
896
0
      assert(file->largest_seqno != kMaxSequenceNumber);
897
0
      smallest_seqno = file->smallest_seqno;
898
0
      largest_seqno = file->largest_seqno;
899
0
      max_assigned_seqno_ = std::max(max_assigned_seqno_, file->largest_seqno);
900
0
    }
901
0
    FileMetaData f_metadata(
902
0
        file->fd.GetNumber(), file->fd.GetPathId(), file->fd.GetFileSize(),
903
0
        file->smallest_internal_key, file->largest_internal_key, smallest_seqno,
904
0
        largest_seqno, false, file->file_temperature, kInvalidBlobFileNumber,
905
0
        oldest_ancester_time, current_time,
906
0
        ingestion_options_.ingest_behind
907
0
            ? kReservedEpochNumberForFileIngestedBehind
908
0
            : cfd_->NewEpochNumber(),  // orders files ingested to L0
909
0
        file->file_checksum, file->file_checksum_func_name, file->unique_id, 0,
910
0
        tail_size, file->user_defined_timestamps_persisted, "", "");
911
0
    f_metadata.temperature = file->file_temperature;
912
0
    f_metadata.marked_for_compaction = marked_for_compaction;
913
0
    f_metadata.skip_index_and_filter_blocks_prefetch =
914
0
        !file->prefetch_lmax_index_and_filter_blocks &&
915
0
        file->picked_level == cfd_->NumberLevels() - 1;
916
    // Extract min/max timestamps from table properties for UDT support.
917
    // This ensures ingested files have proper timestamp ranges in FileMetaData,
918
    // similar to files created by flush and compaction.
919
0
    ExtractTimestampFromTableProperties(file->table_properties, &f_metadata);
920
    // Retrieve file open metadata for fast SST open
921
0
    if (mutable_db_options_.fast_sst_open) {
922
0
      std::unique_ptr<FSRandomAccessFile> readable_file;
923
0
      FileOptions fopts{env_options_};
924
0
      fopts.file_checksum = f_metadata.file_checksum;
925
0
      fopts.file_checksum_func_name = f_metadata.file_checksum_func_name;
926
0
      IOStatus io_s = fs_->NewRandomAccessFile(file->internal_file_path, fopts,
927
0
                                               &readable_file, nullptr);
928
0
      if (io_s.ok()) {
929
0
        io_s =
930
0
            readable_file->GetFileOpenMetadata(&f_metadata.file_open_metadata);
931
0
        if (io_s.ok() && !f_metadata.file_open_metadata.empty() &&
932
0
            f_metadata.file_open_metadata.size() <=
933
0
                FSRandomAccessFile::kMaxFileOpenMetadataSize) {
934
0
          RecordTick(db_options_.stats, FILE_OPEN_METADATA_RETRIEVED);
935
0
        } else {
936
0
          if (io_s.ok() && f_metadata.file_open_metadata.size() >
937
0
                               FSRandomAccessFile::kMaxFileOpenMetadataSize) {
938
0
            ROCKS_LOG_WARN(db_options_.info_log,
939
0
                           "File open metadata for %s too large (%zu bytes), "
940
0
                           "ignoring",
941
0
                           file->internal_file_path.c_str(),
942
0
                           f_metadata.file_open_metadata.size());
943
0
          }
944
0
          f_metadata.file_open_metadata.clear();
945
0
        }
946
0
      }
947
0
    }
948
0
    edit_.AddFile(file->picked_level, f_metadata);
949
950
0
    *batch_uppermost_level =
951
0
        std::min(*batch_uppermost_level, file->picked_level);
952
0
  }
953
954
0
  return Status::OK();
955
0
}
956
957
0
void ExternalSstFileIngestionJob::CreateEquivalentFileIngestingCompactions() {
958
  // A map from output level to input of compactions equivalent to this
959
  // ingestion job.
960
  // TODO: simplify below logic to creating compaction per ingested file
961
  // instead of per output level, once we figure out how to treat ingested files
962
  // with adjacent range deletion tombstones to same output level in the same
963
  // job as non-overlapping compactions.
964
0
  std::map<int, CompactionInputFiles>
965
0
      output_level_to_file_ingesting_compaction_input;
966
967
0
  for (const auto& pair : edit_.GetNewFiles()) {
968
0
    int output_level = pair.first;
969
0
    const FileMetaData& f_metadata = pair.second;
970
971
0
    CompactionInputFiles& input =
972
0
        output_level_to_file_ingesting_compaction_input[output_level];
973
0
    if (input.files.empty()) {
974
      // Treat the source level of ingested files to be level 0
975
0
      input.level = 0;
976
0
    }
977
978
0
    compaction_input_metdatas_.push_back(new FileMetaData(f_metadata));
979
0
    input.files.push_back(compaction_input_metdatas_.back());
980
0
  }
981
982
0
  for (const auto& pair : output_level_to_file_ingesting_compaction_input) {
983
0
    int output_level = pair.first;
984
0
    const CompactionInputFiles& input = pair.second;
985
986
0
    const auto& mutable_cf_options = cfd_->GetLatestMutableCFOptions();
987
0
    file_ingesting_compactions_.push_back(new Compaction(
988
0
        cfd_->current()->storage_info(), cfd_->ioptions(), mutable_cf_options,
989
0
        mutable_db_options_, {input}, output_level,
990
        /* output file size limit not applicable */
991
0
        MaxFileSizeForLevel(mutable_cf_options, output_level,
992
0
                            cfd_->ioptions().compaction_style),
993
0
        LLONG_MAX /* max compaction bytes, not applicable */,
994
0
        0 /* output path ID, not applicable */, mutable_cf_options.compression,
995
0
        mutable_cf_options.compression_opts, Temperature::kUnknown,
996
0
        0 /* max_subcompaction, not applicable */,
997
0
        {} /* grandparents, not applicable */,
998
0
        std::nullopt /* earliest_snapshot */, nullptr /* snapshot_checker */,
999
0
        CompactionReason::kExternalSstIngestion, "" /* trim_ts */,
1000
0
        -1 /* score, not applicable */,
1001
0
        files_overlap_ /* l0_files_might_overlap, not applicable */));
1002
0
  }
1003
0
}
1004
1005
0
void ExternalSstFileIngestionJob::RegisterRange() {
1006
0
  for (const auto& c : file_ingesting_compactions_) {
1007
0
    cfd_->compaction_picker()->RegisterCompaction(c);
1008
0
  }
1009
0
}
1010
1011
0
void ExternalSstFileIngestionJob::UnregisterRange() {
1012
0
  for (const auto& c : file_ingesting_compactions_) {
1013
0
    cfd_->compaction_picker()->UnregisterCompaction(c);
1014
0
    delete c;
1015
0
  }
1016
0
  file_ingesting_compactions_.clear();
1017
1018
0
  for (const auto& f : compaction_input_metdatas_) {
1019
0
    delete f;
1020
0
  }
1021
0
  compaction_input_metdatas_.clear();
1022
0
}
1023
1024
0
void ExternalSstFileIngestionJob::UpdateStats() {
1025
  // Update internal stats for new ingested files
1026
0
  uint64_t total_keys = 0;
1027
0
  uint64_t total_l0_files = 0;
1028
0
  uint64_t total_time = clock_->NowMicros() - job_start_time_;
1029
1030
0
  EventLoggerStream stream = event_logger_->Log();
1031
0
  stream << "event" << "ingest_finished";
1032
0
  stream << "files_ingested";
1033
0
  stream.StartArray();
1034
1035
0
  for (IngestedFileInfo& f : files_to_ingest_) {
1036
0
    InternalStats::CompactionStats stats(
1037
0
        CompactionReason::kExternalSstIngestion, 1);
1038
0
    stats.micros = total_time;
1039
    // If actual copy occurred for this file, then we need to count the file
1040
    // size as the actual bytes written. If the file was linked, then we ignore
1041
    // the bytes written for file metadata.
1042
    // TODO (yanqin) maybe account for file metadata bytes for exact accuracy?
1043
0
    if (f.copy_file) {
1044
0
      stats.bytes_written = f.fd.GetFileSize();
1045
0
    } else {
1046
0
      stats.bytes_moved = f.fd.GetFileSize();
1047
0
    }
1048
0
    stats.num_output_files = 1;
1049
0
    cfd_->internal_stats()->AddCompactionStats(f.picked_level,
1050
0
                                               Env::Priority::USER, stats);
1051
0
    cfd_->internal_stats()->AddCFStats(InternalStats::BYTES_INGESTED_ADD_FILE,
1052
0
                                       f.fd.GetFileSize());
1053
0
    total_keys += f.num_entries;
1054
0
    if (f.picked_level == 0) {
1055
0
      total_l0_files += 1;
1056
0
    }
1057
0
    ROCKS_LOG_INFO(
1058
0
        db_options_.info_log,
1059
0
        "[AddFile] External SST file %s was ingested in L%d with path %s "
1060
0
        "(global_seqno=%" PRIu64 ")\n",
1061
0
        f.external_file_path.c_str(), f.picked_level,
1062
0
        f.internal_file_path.c_str(), f.assigned_seqno);
1063
0
    stream << "file" << f.internal_file_path << "level" << f.picked_level;
1064
0
  }
1065
0
  stream.EndArray();
1066
1067
0
  stream << "lsm_state";
1068
0
  stream.StartArray();
1069
0
  auto vstorage = cfd_->current()->storage_info();
1070
0
  for (int level = 0; level < vstorage->num_levels(); ++level) {
1071
0
    stream << vstorage->NumLevelFiles(level);
1072
0
  }
1073
0
  stream.EndArray();
1074
1075
0
  cfd_->internal_stats()->AddCFStats(InternalStats::INGESTED_NUM_KEYS_TOTAL,
1076
0
                                     total_keys);
1077
0
  cfd_->internal_stats()->AddCFStats(InternalStats::INGESTED_NUM_FILES_TOTAL,
1078
0
                                     files_to_ingest_.size());
1079
0
  cfd_->internal_stats()->AddCFStats(
1080
0
      InternalStats::INGESTED_LEVEL0_NUM_FILES_TOTAL, total_l0_files);
1081
0
}
1082
1083
0
void ExternalSstFileIngestionJob::Cleanup(const Status& status) {
1084
0
  IOOptions io_opts;
1085
0
  if (!status.ok()) {
1086
    // We failed to add the files to the database
1087
    // remove all the files we copied
1088
0
    DeleteInternalFiles();
1089
0
    files_overlap_ = false;
1090
0
  } else {
1091
0
    if (atomic_replace_range_tombstone_.has_value()) {
1092
0
      Status s =
1093
0
          fs_->DeleteFile(atomic_replace_range_tombstone_->internal_file_path,
1094
0
                          io_opts, nullptr);
1095
0
      if (!s.ok()) {
1096
0
        ROCKS_LOG_WARN(
1097
0
            db_options_.info_log,
1098
0
            "Failed to remove unused atomic replace range tombstone %s: %s",
1099
0
            atomic_replace_range_tombstone_->internal_file_path.c_str(),
1100
0
            s.ToString().c_str());
1101
0
      }
1102
0
      atomic_replace_range_tombstone_.reset();
1103
0
    }
1104
0
  }
1105
0
  if (status.ok() && ingestion_options_.move_files) {
1106
    // The files were moved and added successfully, remove original file links
1107
0
    for (IngestedFileInfo& f : files_to_ingest_) {
1108
0
      if (f.generated_for_ingestion) {
1109
0
        continue;
1110
0
      }
1111
0
      Status s = fs_->DeleteFile(f.external_file_path, io_opts, nullptr);
1112
0
      if (!s.ok()) {
1113
0
        ROCKS_LOG_WARN(
1114
0
            db_options_.info_log,
1115
0
            "%s was added to DB successfully but failed to remove original "
1116
0
            "file link : %s",
1117
0
            f.external_file_path.c_str(), s.ToString().c_str());
1118
0
      }
1119
0
    }
1120
0
  }
1121
0
}
1122
1123
0
void ExternalSstFileIngestionJob::DeleteInternalFiles() {
1124
0
  IOOptions io_opts;
1125
0
  for (IngestedFileInfo& f : files_to_ingest_) {
1126
0
    if (f.internal_file_path.empty()) {
1127
0
      continue;
1128
0
    }
1129
0
    Status s = fs_->DeleteFile(f.internal_file_path, io_opts, nullptr);
1130
0
    if (!s.ok()) {
1131
0
      ROCKS_LOG_WARN(db_options_.info_log,
1132
0
                     "AddFile() clean up for file %s failed : %s",
1133
0
                     f.internal_file_path.c_str(), s.ToString().c_str());
1134
0
    }
1135
0
  }
1136
0
  if (atomic_replace_range_tombstone_.has_value() &&
1137
0
      !atomic_replace_range_tombstone_->internal_file_path.empty()) {
1138
0
    Status s = fs_->DeleteFile(
1139
0
        atomic_replace_range_tombstone_->internal_file_path, io_opts, nullptr);
1140
0
    if (!s.ok()) {
1141
0
      ROCKS_LOG_WARN(
1142
0
          db_options_.info_log,
1143
0
          "AddFile() clean up for generated tombstone file %s failed : %s",
1144
0
          atomic_replace_range_tombstone_->internal_file_path.c_str(),
1145
0
          s.ToString().c_str());
1146
0
    }
1147
0
    atomic_replace_range_tombstone_.reset();
1148
0
  }
1149
0
}
1150
1151
Status ExternalSstFileIngestionJob::ResetTableReader(
1152
    const std::string& external_file, uint64_t new_file_number,
1153
    bool user_defined_timestamps_persisted, SuperVersion* sv,
1154
    IngestedFileInfo* file_to_ingest,
1155
0
    std::unique_ptr<TableReader>* table_reader) {
1156
0
  std::unique_ptr<FSRandomAccessFile> sst_file;
1157
0
  FileOptions fo{env_options_};
1158
0
  fo.temperature = file_to_ingest->file_temperature;
1159
0
  fo.file_checksum_func_name = kNoFileChecksumFuncName;
1160
0
  Status status =
1161
0
      fs_->NewRandomAccessFile(external_file, fo, &sst_file, nullptr);
1162
0
  if (!status.ok()) {
1163
0
    ROCKS_LOG_WARN(
1164
0
        db_options_.info_log,
1165
0
        "Failed to create random access file for external file %s: %s",
1166
0
        external_file.c_str(), status.ToString().c_str());
1167
0
    return status;
1168
0
  }
1169
0
  Temperature updated_temp = sst_file->GetTemperature();
1170
0
  if (updated_temp != Temperature::kUnknown &&
1171
0
      updated_temp != file_to_ingest->file_temperature) {
1172
    // The hint was missing or wrong. Track temperature reported by storage.
1173
0
    file_to_ingest->file_temperature = updated_temp;
1174
0
  }
1175
0
  std::unique_ptr<RandomAccessFileReader> sst_file_reader(
1176
0
      new RandomAccessFileReader(std::move(sst_file), external_file,
1177
0
                                 nullptr /*Env*/, io_tracer_));
1178
0
  table_reader->reset();
1179
0
  ReadOptions ro;
1180
0
  ro.fill_cache = ingestion_options_.fill_cache;
1181
0
  status = sv->mutable_cf_options.table_factory->NewTableReader(
1182
0
      ro,
1183
0
      TableReaderOptions(
1184
0
          cfd_->ioptions(), sv->mutable_cf_options.prefix_extractor,
1185
0
          sv->mutable_cf_options.compression_manager.get(), env_options_,
1186
0
          cfd_->internal_comparator(),
1187
0
          sv->mutable_cf_options.block_protection_bytes_per_key,
1188
0
          /*skip_filters*/ false, /*immortal*/ false,
1189
0
          /*force_direct_prefetch*/ false, /*level*/ -1,
1190
0
          /*block_cache_tracer*/ nullptr,
1191
0
          /*max_file_size_for_l0_meta_pin*/ 0, versions_->DbSessionId(),
1192
0
          /*cur_file_num*/ new_file_number,
1193
0
          /* unique_id */ {}, /* largest_seqno */ 0,
1194
0
          /* tail_size */ 0, user_defined_timestamps_persisted),
1195
0
      std::move(sst_file_reader), file_to_ingest->file_size, table_reader,
1196
      // No need to prefetch index/filter if caching is not needed.
1197
0
      /*prefetch_index_and_filter_in_cache=*/ingestion_options_.fill_cache);
1198
0
  return status;
1199
0
}
1200
1201
Status ExternalSstFileIngestionJob::SanityCheckTableProperties(
1202
    const std::string& external_file, const TableProperties& props,
1203
0
    IngestedFileInfo* file_to_ingest) {
1204
0
  const auto& uprops = props.user_collected_properties;
1205
1206
  // Get table version
1207
0
  auto version_iter = uprops.find(ExternalSstFilePropertyNames::kVersion);
1208
0
  if (version_iter == uprops.end()) {
1209
0
    assert(!SstFileWriter::CreatedBySstFileWriter(props));
1210
0
    if (!ingestion_options_.allow_db_generated_files) {
1211
0
      return Status::Corruption("External file version not found");
1212
0
    } else {
1213
      // 0 is special version for when a file from live DB does not have the
1214
      // version table property
1215
0
      file_to_ingest->version = 0;
1216
0
    }
1217
0
  } else {
1218
0
    assert(SstFileWriter::CreatedBySstFileWriter(props));
1219
0
    file_to_ingest->version = DecodeFixed32(version_iter->second.c_str());
1220
0
  }
1221
1222
0
  auto seqno_iter = uprops.find(ExternalSstFilePropertyNames::kGlobalSeqno);
1223
0
  if (file_to_ingest->version == 2) {
1224
    // version 2 imply that we have global sequence number
1225
0
    if (seqno_iter == uprops.end()) {
1226
0
      return Status::Corruption(
1227
0
          "External file global sequence number not found");
1228
0
    }
1229
1230
    // Set the global sequence number
1231
0
    file_to_ingest->original_seqno = DecodeFixed64(seqno_iter->second.c_str());
1232
0
    file_to_ingest->global_seqno_offset =
1233
0
        static_cast<size_t>(props.external_sst_file_global_seqno_offset);
1234
    // The on-disk offset is only needed if we will write the global seqno back
1235
    // into the file (write_global_seqno). The metadata fast-path does not open
1236
    // the file, and its in-memory table properties do not carry the offset
1237
    // (it is only computed while reading the file back); that is fine as long
1238
    // as write_global_seqno is not requested.
1239
0
    if (ingestion_options_.write_global_seqno &&
1240
0
        file_to_ingest->global_seqno_offset == 0) {
1241
0
      return Status::Corruption("Was not able to find file global seqno field");
1242
0
    }
1243
0
  } else if (file_to_ingest->version == 1) {
1244
    // SST file V1 should not have global seqno field
1245
0
    assert(seqno_iter == uprops.end());
1246
0
    file_to_ingest->original_seqno = 0;
1247
0
    if (ingestion_options_.allow_blocking_flush ||
1248
0
        ingestion_options_.allow_global_seqno) {
1249
0
      return Status::InvalidArgument(
1250
0
          "External SST file V1 does not support global seqno");
1251
0
    }
1252
0
  } else if (file_to_ingest->version == 0) {
1253
    // allow_db_generated_files is true
1254
0
    assert(seqno_iter == uprops.end());
1255
0
    file_to_ingest->original_seqno = 0;
1256
0
    file_to_ingest->global_seqno_offset = 0;
1257
0
  } else {
1258
0
    return Status::InvalidArgument("External file version " +
1259
0
                                   std::to_string(file_to_ingest->version) +
1260
0
                                   " is not supported");
1261
0
  }
1262
1263
0
  file_to_ingest->cf_id = static_cast<uint32_t>(props.column_family_id);
1264
0
  file_to_ingest->table_properties = props;
1265
1266
  // Get number of entries in table
1267
0
  file_to_ingest->num_entries = props.num_entries;
1268
0
  file_to_ingest->num_range_deletions = props.num_range_deletions;
1269
1270
  // Validate table properties related to comparator name and user defined
1271
  // timestamps persisted flag.
1272
0
  file_to_ingest->user_defined_timestamps_persisted =
1273
0
      static_cast<bool>(props.user_defined_timestamps_persisted);
1274
0
  bool mark_sst_file_has_no_udt = false;
1275
0
  Status s = ValidateUserDefinedTimestampsOptions(
1276
0
      cfd_->user_comparator(), props.comparator_name,
1277
0
      cfd_->ioptions().persist_user_defined_timestamps,
1278
0
      file_to_ingest->user_defined_timestamps_persisted,
1279
0
      &mark_sst_file_has_no_udt);
1280
0
  if (s.ok() && mark_sst_file_has_no_udt) {
1281
    // A column family that enables user-defined timestamps in Memtable only
1282
    // feature can also ingest external files created by a setting that disables
1283
    // user-defined timestamps. In that case, we need to re-mark the
1284
    // user_defined_timestamps_persisted flag for the file. The open-and-scan
1285
    // caller is then responsible for reopening its `TableReader` with the
1286
    // updated flag.
1287
0
    file_to_ingest->user_defined_timestamps_persisted = false;
1288
0
  } else if (!s.ok()) {
1289
0
    ROCKS_LOG_WARN(
1290
0
        db_options_.info_log,
1291
0
        "ValidateUserDefinedTimestampsOptions failed for external file %s: %s",
1292
0
        external_file.c_str(), s.ToString().c_str());
1293
0
    return s;
1294
0
  }
1295
1296
0
  return s;
1297
0
}
1298
1299
Status ExternalSstFileIngestionJob::GetIngestedFileInfoFromFileInfo(
1300
    const std::string& external_file,
1301
    const PreparedFileInfo& prepared_file_info,
1302
0
    IngestedFileInfo* file_to_ingest) {
1303
  // Boundaries, size, and table properties were obtained without opening the
1304
  // file (e.g. produced by SstFileWriter::Finish). We reuse them directly.
1305
0
  file_to_ingest->file_size = prepared_file_info.file_size;
1306
0
  Status status = SanityCheckTableProperties(
1307
0
      external_file, prepared_file_info.table_properties, file_to_ingest);
1308
0
  if (!status.ok()) {
1309
0
    ROCKS_LOG_WARN(
1310
0
        db_options_.info_log,
1311
0
        "Failed to sanity check table properties for external file %s: %s",
1312
0
        external_file.c_str(), status.ToString().c_str());
1313
0
    return status;
1314
0
  }
1315
1316
0
  const size_t ts_sz = ucmp_->timestamp_size();
1317
0
  if (ts_sz > 0 && !file_to_ingest->user_defined_timestamps_persisted) {
1318
0
    auto pad_timestamp = [ts_sz](std::string* result, const Slice& key) {
1319
0
      assert(result->empty());
1320
0
      if (ExtractValueType(key) == kTypeRangeDeletion) {
1321
0
        PadInternalKeyWithMaxTimestamp(result, key, ts_sz);
1322
0
      } else {
1323
0
        PadInternalKeyWithMinTimestamp(result, key, ts_sz);
1324
0
      }
1325
0
    };
1326
0
    pad_timestamp(file_to_ingest->smallest_internal_key.rep(),
1327
0
                  prepared_file_info.smallest.Encode());
1328
0
    pad_timestamp(file_to_ingest->largest_internal_key.rep(),
1329
0
                  prepared_file_info.largest.Encode());
1330
0
  } else {
1331
0
    file_to_ingest->smallest_internal_key = prepared_file_info.smallest;
1332
0
    file_to_ingest->largest_internal_key = prepared_file_info.largest;
1333
0
  }
1334
1335
0
  if (ingestion_options_.allow_db_generated_files) {
1336
    // Sequence numbers are preserved (not reassigned), so the bounds must be
1337
    // known before we skip the GetSeqnoBoundaryForFile scan.
1338
0
    if (!file_to_ingest->table_properties.HasKeyLargestSeqno()) {
1339
0
      return Status::Corruption("Unknown largest seqno for db generated file.");
1340
0
    }
1341
0
    file_to_ingest->largest_seqno =
1342
0
        file_to_ingest->table_properties.key_largest_seqno;
1343
0
    if (file_to_ingest->largest_seqno == 0) {
1344
0
      file_to_ingest->smallest_seqno = 0;
1345
0
    } else {
1346
0
      if (!file_to_ingest->table_properties.HasKeySmallestSeqno()) {
1347
0
        return Status::Corruption(
1348
0
            "Unknown smallest seqno for db generated file.");
1349
0
      }
1350
0
      file_to_ingest->smallest_seqno =
1351
0
          file_to_ingest->table_properties.key_smallest_seqno;
1352
0
    }
1353
0
  } else {
1354
    // Normal ingestion reassigns a global sequence number later, so the file's
1355
    // keys must currently be at seqno 0 (mirror of the open-and-scan check).
1356
0
    SequenceNumber largest_seqno =
1357
0
        file_to_ingest->table_properties.key_largest_seqno;
1358
    // UINT64_MAX means unknown and the file is generated before table property
1359
    // `key_largest_seqno` is introduced.
1360
0
    if (largest_seqno != UINT64_MAX && largest_seqno > 0) {
1361
0
      return Status::Corruption(
1362
0
          "External file has non zero largest sequence number " +
1363
0
          std::to_string(largest_seqno));
1364
0
    }
1365
0
  }
1366
0
  return Status::OK();
1367
0
}
1368
1369
Status ExternalSstFileIngestionJob::GetIngestedFileInfoFromFile(
1370
    const std::string& external_file, uint64_t new_file_number,
1371
    IngestedFileInfo* file_to_ingest, SuperVersion* sv,
1372
0
    std::unique_ptr<TableReader>* out_table_reader) {
1373
0
  TEST_SYNC_POINT("ExternalSstFileIngestionJob::GetIngestedFileInfo:ReadPath");
1374
1375
  // Get external file size
1376
0
  Status status = fs_->GetFileSize(external_file, IOOptions(),
1377
0
                                   &file_to_ingest->file_size, nullptr);
1378
0
  if (!status.ok()) {
1379
0
    ROCKS_LOG_WARN(db_options_.info_log,
1380
0
                   "Failed to get file size for external file %s: %s",
1381
0
                   external_file.c_str(), status.ToString().c_str());
1382
0
    return status;
1383
0
  }
1384
1385
  // Create TableReader for external file.
1386
0
  std::unique_ptr<TableReader> table_reader;
1387
  // Initially create the `TableReader` with flag
1388
  // `user_defined_timestamps_persisted` to be true since that's the most
1389
  // common case
1390
0
  status = ResetTableReader(external_file, new_file_number,
1391
0
                            /*user_defined_timestamps_persisted=*/true, sv,
1392
0
                            file_to_ingest, &table_reader);
1393
0
  if (!status.ok()) {
1394
0
    ROCKS_LOG_WARN(db_options_.info_log,
1395
0
                   "Failed to reset table reader for external file %s: %s",
1396
0
                   external_file.c_str(), status.ToString().c_str());
1397
0
    return status;
1398
0
  }
1399
1400
0
  status = SanityCheckTableProperties(
1401
0
      external_file, *table_reader->GetTableProperties(), file_to_ingest);
1402
0
  if (!status.ok()) {
1403
0
    ROCKS_LOG_WARN(
1404
0
        db_options_.info_log,
1405
0
        "Failed to sanity check table properties for external file %s: %s",
1406
0
        external_file.c_str(), status.ToString().c_str());
1407
0
    return status;
1408
0
  }
1409
1410
  // The `TableReader` above was opened with `user_defined_timestamps_persisted`
1411
  // assumed true. If the sanity check determined the file has no persisted
1412
  // timestamps (UDT-in-Memtable-only feature), reopen it with the corrected
1413
  // flag so keys are parsed properly by the scan below.
1414
0
  if (ucmp_->timestamp_size() > 0 &&
1415
0
      !file_to_ingest->user_defined_timestamps_persisted) {
1416
0
    status = ResetTableReader(external_file, new_file_number,
1417
0
                              file_to_ingest->user_defined_timestamps_persisted,
1418
0
                              sv, file_to_ingest, &table_reader);
1419
0
    if (!status.ok()) {
1420
0
      ROCKS_LOG_WARN(db_options_.info_log,
1421
0
                     "Failed to reset table reader for external file %s: %s",
1422
0
                     external_file.c_str(), status.ToString().c_str());
1423
0
      return status;
1424
0
    }
1425
0
  }
1426
1427
0
  const bool allow_data_in_errors = db_options_.allow_data_in_errors;
1428
0
  ParsedInternalKey key;
1429
0
  if (ingestion_options_.allow_db_generated_files) {
1430
    // We are ingesting a DB generated SST file for which we don't reassign
1431
    // sequence numbers. We need its smallest sequence number and largest
1432
    // sequence number for FileMetaData.
1433
0
    Status seqno_status = GetSeqnoBoundaryForFile(
1434
0
        table_reader.get(), sv, file_to_ingest, allow_data_in_errors);
1435
1436
0
    if (!seqno_status.ok()) {
1437
0
      ROCKS_LOG_WARN(
1438
0
          db_options_.info_log,
1439
0
          "Failed to get sequence number boundary for external file %s: %s",
1440
0
          external_file.c_str(), seqno_status.ToString().c_str());
1441
0
      return seqno_status;
1442
0
    }
1443
0
    assert(file_to_ingest->smallest_seqno <= file_to_ingest->largest_seqno);
1444
0
    assert(file_to_ingest->largest_seqno < kMaxSequenceNumber);
1445
0
  } else {
1446
0
    SequenceNumber largest_seqno =
1447
0
        table_reader.get()->GetTableProperties()->key_largest_seqno;
1448
    // UINT64_MAX means unknown and the file is generated before table
1449
    // property `key_largest_seqno` is introduced.
1450
0
    if (largest_seqno != UINT64_MAX && largest_seqno > 0) {
1451
0
      return Status::Corruption(
1452
0
          "External file has non zero largest sequence number " +
1453
0
          std::to_string(largest_seqno));
1454
0
    }
1455
0
  }
1456
1457
  // TODO: plumb Env::IOActivity, Env::IOPriority
1458
0
  ReadOptions ro;
1459
0
  ro.fill_cache = ingestion_options_.fill_cache;
1460
0
  std::unique_ptr<InternalIterator> iter(table_reader->NewIterator(
1461
0
      ro, sv->mutable_cf_options.prefix_extractor.get(), /*arena=*/nullptr,
1462
0
      /*skip_filters=*/false, TableReaderCaller::kExternalSSTIngestion));
1463
1464
  // Get first (smallest) and last (largest) key from file.
1465
0
  iter->SeekToFirst();
1466
0
  if (iter->Valid()) {
1467
0
    Status pik_status =
1468
0
        ParseInternalKey(iter->key(), &key, allow_data_in_errors);
1469
0
    if (!pik_status.ok()) {
1470
0
      return Status::Corruption("Corrupted key in external file. ",
1471
0
                                pik_status.getState());
1472
0
    }
1473
0
    if (key.sequence != 0 && !ingestion_options_.allow_db_generated_files) {
1474
0
      return Status::Corruption("External file has non zero sequence number");
1475
0
    }
1476
0
    file_to_ingest->smallest_internal_key.SetFrom(key);
1477
1478
0
    Slice largest;
1479
0
    if (strcmp(sv->mutable_cf_options.table_factory->Name(), "PlainTable") ==
1480
0
        0) {
1481
      // PlainTable iterator does not support SeekToLast().
1482
0
      largest = iter->key();
1483
0
      for (; iter->Valid(); iter->Next()) {
1484
0
        if (cfd_->internal_comparator().Compare(iter->key(), largest) > 0) {
1485
0
          largest = iter->key();
1486
0
        }
1487
0
      }
1488
0
      if (!iter->status().ok()) {
1489
0
        return iter->status();
1490
0
      }
1491
0
    } else {
1492
0
      iter->SeekToLast();
1493
0
      if (!iter->Valid()) {
1494
0
        if (iter->status().ok()) {
1495
          // The file contains at least 1 key since iter is valid after
1496
          // SeekToFirst().
1497
0
          return Status::Corruption("Can not find largest key in sst file");
1498
0
        } else {
1499
0
          return iter->status();
1500
0
        }
1501
0
      }
1502
0
      largest = iter->key();
1503
0
    }
1504
1505
0
    pik_status = ParseInternalKey(largest, &key, allow_data_in_errors);
1506
0
    if (!pik_status.ok()) {
1507
0
      return Status::Corruption("Corrupted key in external file. ",
1508
0
                                pik_status.getState());
1509
0
    }
1510
0
    if (key.sequence != 0 && !ingestion_options_.allow_db_generated_files) {
1511
0
      return Status::Corruption("External file has non zero sequence number");
1512
0
    }
1513
0
    file_to_ingest->largest_internal_key.SetFrom(key);
1514
0
  } else if (!iter->status().ok()) {
1515
0
    return iter->status();
1516
0
  }
1517
1518
0
  std::unique_ptr<InternalIterator> range_del_iter(
1519
0
      table_reader->NewRangeTombstoneIterator(ro));
1520
  // We may need to adjust these key bounds, depending on whether any range
1521
  // deletion tombstones extend past them.
1522
0
  if (range_del_iter != nullptr) {
1523
0
    for (range_del_iter->SeekToFirst(); range_del_iter->Valid();
1524
0
         range_del_iter->Next()) {
1525
0
      Status pik_status =
1526
0
          ParseInternalKey(range_del_iter->key(), &key, allow_data_in_errors);
1527
0
      if (!pik_status.ok()) {
1528
0
        return Status::Corruption("Corrupted key in external file. ",
1529
0
                                  pik_status.getState());
1530
0
      }
1531
0
      if (key.sequence != 0 && !ingestion_options_.allow_db_generated_files) {
1532
0
        return Status::Corruption(
1533
0
            "External file has a range deletion with non zero sequence "
1534
0
            "number.");
1535
0
      }
1536
#ifndef NDEBUG
1537
      // To keep aligned with the fast path that expects range deletion keys to
1538
      // always use max ts.
1539
      const size_t ts_sz = ucmp_->timestamp_size();
1540
      if (ts_sz > 0) {
1541
        const std::string max_ts(ts_sz, '\xff');
1542
        assert(key.user_key.size() >= ts_sz);
1543
        assert(ucmp_->CompareTimestamp(
1544
                   ExtractTimestampFromUserKey(key.user_key, ts_sz), max_ts) ==
1545
               0);
1546
        assert(range_del_iter->value().size() >= ts_sz);
1547
        assert(ucmp_->CompareTimestamp(
1548
                   ExtractTimestampFromUserKey(range_del_iter->value(), ts_sz),
1549
                   max_ts) == 0);
1550
      }
1551
#endif
1552
0
      RangeTombstone tombstone(key, range_del_iter->value());
1553
0
      file_range_checker_.MaybeUpdateRange(tombstone.SerializeKey(),
1554
0
                                           tombstone.SerializeEndKey(),
1555
0
                                           file_to_ingest);
1556
0
    }
1557
0
  }
1558
1559
0
  *out_table_reader = std::move(table_reader);
1560
0
  return Status::OK();
1561
0
}
1562
1563
Status ExternalSstFileIngestionJob::GetIngestedFileInfo(
1564
    const std::string& external_file, uint64_t new_file_number,
1565
    const PreparedFileInfo* prepared_file_info,
1566
0
    IngestedFileInfo* file_to_ingest, SuperVersion* sv) {
1567
0
  file_to_ingest->external_file_path = external_file;
1568
1569
0
  std::unique_ptr<TableReader> table_reader;
1570
0
  Status status =
1571
0
      prepared_file_info != nullptr
1572
0
          ? GetIngestedFileInfoFromFileInfo(external_file, *prepared_file_info,
1573
0
                                            file_to_ingest)
1574
0
          : GetIngestedFileInfoFromFile(external_file, new_file_number,
1575
0
                                        file_to_ingest, sv, &table_reader);
1576
0
  if (!status.ok()) {
1577
0
    return status;
1578
0
  }
1579
1580
0
  assert(file_to_ingest->file_size > 0);
1581
0
  assert(!file_to_ingest->unset());
1582
0
  assert(file_to_ingest->table_properties.num_entries > 0 ||
1583
0
         file_to_ingest->table_properties.num_range_deletions > 0);
1584
0
  if (ingestion_options_.allow_db_generated_files) {
1585
    // These files keep their original sequence numbers (derived, not
1586
    // reassigned), so the bounds must be valid here.
1587
0
    assert(file_to_ingest->smallest_seqno <= file_to_ingest->largest_seqno);
1588
0
    assert(file_to_ingest->largest_seqno < kMaxSequenceNumber);
1589
0
  }
1590
1591
  // Verify the file checksum if requested. The open-and-scan path already has a
1592
  // `TableReader` open; the fast-path opens one here only when verification is
1593
  // requested -- otherwise it performs no file I/O.
1594
0
  if (ingestion_options_.verify_checksums_before_ingest) {
1595
0
    if (table_reader == nullptr) {
1596
0
      status =
1597
0
          ResetTableReader(external_file, new_file_number,
1598
0
                           file_to_ingest->user_defined_timestamps_persisted,
1599
0
                           sv, file_to_ingest, &table_reader);
1600
0
      if (!status.ok()) {
1601
0
        ROCKS_LOG_WARN(db_options_.info_log,
1602
0
                       "Failed to reset table reader for external file %s: %s",
1603
0
                       external_file.c_str(), status.ToString().c_str());
1604
0
        return status;
1605
0
      }
1606
0
    }
1607
    // If customized readahead size is needed, we can pass a user option all the
1608
    // way to here. Right now we just rely on the default readahead to keep
1609
    // things simple.
1610
    // TODO: plumb Env::IOActivity, Env::IOPriority
1611
0
    ReadOptions ro;
1612
0
    ro.readahead_size = ingestion_options_.verify_checksums_readahead_size;
1613
0
    ro.fill_cache = ingestion_options_.fill_cache;
1614
0
    status = table_reader->VerifyChecksum(
1615
0
        ro, TableReaderCaller::kExternalSSTIngestion);
1616
0
    if (!status.ok()) {
1617
0
      ROCKS_LOG_WARN(db_options_.info_log,
1618
0
                     "Failed to verify checksum for external file %s: %s",
1619
0
                     external_file.c_str(), status.ToString().c_str());
1620
0
      return status;
1621
0
    }
1622
0
  }
1623
1624
  // Assign FD with number.
1625
0
  file_to_ingest->fd =
1626
0
      FileDescriptor(new_file_number, 0, file_to_ingest->file_size);
1627
1628
0
  const size_t ts_sz = ucmp_->timestamp_size();
1629
0
  Slice smallest = file_to_ingest->smallest_internal_key.user_key();
1630
0
  Slice largest = file_to_ingest->largest_internal_key.user_key();
1631
0
  if (ts_sz > 0) {
1632
0
    AppendUserKeyWithMaxTimestamp(&file_to_ingest->start_ukey, smallest, ts_sz);
1633
0
    AppendUserKeyWithMinTimestamp(&file_to_ingest->limit_ukey, largest, ts_sz);
1634
0
  } else {
1635
0
    file_to_ingest->start_ukey.assign(smallest.data(), smallest.size());
1636
0
    file_to_ingest->limit_ukey.assign(largest.data(), largest.size());
1637
0
  }
1638
1639
0
  auto s =
1640
0
      GetSstInternalUniqueId(file_to_ingest->table_properties.db_id,
1641
0
                             file_to_ingest->table_properties.db_session_id,
1642
0
                             file_to_ingest->table_properties.orig_file_number,
1643
0
                             &(file_to_ingest->unique_id));
1644
0
  if (!s.ok()) {
1645
0
    ROCKS_LOG_WARN(db_options_.info_log,
1646
0
                   "Failed to get SST unique id for file %s",
1647
0
                   file_to_ingest->internal_file_path.c_str());
1648
0
    file_to_ingest->unique_id = kNullUniqueId64x2;
1649
0
  }
1650
1651
0
  return status;
1652
0
}
1653
1654
Status ExternalSstFileIngestionJob::AssignLevelAndSeqnoForIngestedFile(
1655
    SuperVersion* sv, bool force_global_seqno, CompactionStyle compaction_style,
1656
    SequenceNumber last_seqno, IngestedFileInfo* file_to_ingest,
1657
    SequenceNumber* assigned_seqno,
1658
0
    std::optional<int> prev_batch_uppermost_level) {
1659
0
  Status status;
1660
0
  *assigned_seqno = 0;
1661
0
  const size_t ts_sz = ucmp_->timestamp_size();
1662
0
  assert(!prev_batch_uppermost_level.has_value() ||
1663
0
         prev_batch_uppermost_level.value() < cfd_->NumberLevels());
1664
0
  bool must_assign_to_l0 = (prev_batch_uppermost_level.has_value() &&
1665
0
                            prev_batch_uppermost_level.value() == 0) ||
1666
0
                           compaction_style == kCompactionStyleFIFO;
1667
1668
0
  if (force_global_seqno || (!ingestion_options_.allow_db_generated_files &&
1669
0
                             (files_overlap_ || must_assign_to_l0))) {
1670
0
    *assigned_seqno = last_seqno + 1;
1671
0
    if (must_assign_to_l0) {
1672
0
      assert(ts_sz == 0);
1673
0
      file_to_ingest->picked_level = 0;
1674
0
      if (ingestion_options_.fail_if_not_bottommost_level &&
1675
0
          cfd_->NumberLevels() > 1) {
1676
0
        status = Status::TryAgain(
1677
0
            "Files cannot be ingested to Lmax. Please make sure key range of "
1678
0
            "Lmax does not overlap with files to ingest.");
1679
0
      }
1680
0
      return status;
1681
0
    }
1682
0
  }
1683
1684
0
  bool overlap_with_db = false;
1685
0
  Arena arena;
1686
  // TODO: plumb Env::IOActivity, Env::IOPriority
1687
0
  ReadOptions ro;
1688
0
  ro.fill_cache = ingestion_options_.fill_cache;
1689
0
  ro.total_order_seek = true;
1690
0
  int target_level = 0;
1691
0
  auto* vstorage = cfd_->current()->storage_info();
1692
0
  assert(!must_assign_to_l0 || ingestion_options_.allow_db_generated_files);
1693
0
  int assigned_level_exclusive_end = cfd_->NumberLevels();
1694
0
  if (must_assign_to_l0) {
1695
0
    assigned_level_exclusive_end = 0;
1696
0
  } else if (prev_batch_uppermost_level.has_value()) {
1697
0
    assigned_level_exclusive_end = prev_batch_uppermost_level.value();
1698
0
  }
1699
1700
  // When ingesting db generated files, we require that ingested files do not
1701
  // overlap with any file in the DB. So we need to check all levels.
1702
0
  int overlap_checking_exclusive_end =
1703
0
      ingestion_options_.allow_db_generated_files
1704
0
          ? cfd_->NumberLevels()
1705
0
          : assigned_level_exclusive_end;
1706
0
  const bool is_generated_range_tombstone =
1707
0
      atomic_replace_range_tombstone_active_ &&
1708
0
      file_to_ingest->generated_for_ingestion;
1709
0
  for (int lvl = 0; lvl < overlap_checking_exclusive_end; lvl++) {
1710
0
    if (lvl > 0 && lvl < vstorage->base_level()) {
1711
0
      continue;
1712
0
    }
1713
0
    if (lvl < assigned_level_exclusive_end &&
1714
0
        atomic_replace_range_.has_value() &&
1715
0
        !atomic_replace_range_tombstone_active_) {
1716
0
      target_level = lvl;
1717
0
      continue;
1718
0
    }
1719
0
    if (cfd_->RangeOverlapWithCompaction(file_to_ingest->start_ukey,
1720
0
                                         file_to_ingest->limit_ukey, lvl,
1721
                                         /*range_limit_exclusive=*/
1722
0
                                         is_generated_range_tombstone)) {
1723
      // We must use L0 or any level higher than `lvl` to be able to overwrite
1724
      // the compaction output keys that we overlap with in this level, We also
1725
      // need to assign this file a seqno to overwrite the compaction output
1726
      // keys in level `lvl`
1727
0
      overlap_with_db = true;
1728
0
      break;
1729
0
    } else if (vstorage->NumLevelFiles(lvl) > 0) {
1730
0
      bool overlap_with_level = false;
1731
0
      status = sv->current->OverlapWithLevelIterator(
1732
0
          ro, env_options_, file_to_ingest->start_ukey,
1733
0
          file_to_ingest->limit_ukey, lvl, &overlap_with_level);
1734
0
      if (!status.ok()) {
1735
0
        ROCKS_LOG_WARN(db_options_.info_log,
1736
0
                       "Failed to check overlap with level iterator: %s",
1737
0
                       status.ToString().c_str());
1738
0
        return status;
1739
0
      }
1740
0
      if (overlap_with_level) {
1741
        // We must use L0 or any level higher than `lvl` to be able to overwrite
1742
        // the keys that we overlap with in this level, We also need to assign
1743
        // this file a seqno to overwrite the existing keys in level `lvl`
1744
0
        overlap_with_db = true;
1745
0
        break;
1746
0
      }
1747
0
    }
1748
1749
    // We don't overlap with any keys in this level, but we still need to check
1750
    // if our file can fit in it
1751
0
    if (lvl < assigned_level_exclusive_end &&
1752
0
        IngestedFileFitInLevel(file_to_ingest, lvl)) {
1753
0
      target_level = lvl;
1754
0
    }
1755
0
  }
1756
1757
0
  if (ingestion_options_.fail_if_not_bottommost_level &&
1758
0
      target_level < cfd_->NumberLevels() - 1) {
1759
0
    status = Status::TryAgain(
1760
0
        "Files cannot be ingested to Lmax. Please make sure key range of Lmax "
1761
0
        "and ongoing compaction's output to Lmax does not overlap with files "
1762
0
        "to ingest. Input files overlapping with each other can cause some "
1763
0
        "file to be assigned to non Lmax level.");
1764
0
    return status;
1765
0
  }
1766
1767
0
  TEST_SYNC_POINT_CALLBACK(
1768
0
      "ExternalSstFileIngestionJob::AssignLevelAndSeqnoForIngestedFile",
1769
0
      &overlap_with_db);
1770
0
  file_to_ingest->picked_level = target_level;
1771
0
  if (overlap_with_db) {
1772
0
    if (ts_sz > 0) {
1773
0
      status = Status::InvalidArgument(
1774
0
          "Column family enables user-defined timestamps, please make sure the "
1775
0
          "key range (without timestamp) of external file does not overlap "
1776
0
          "with key range (without timestamp) in the db");
1777
0
      return status;
1778
0
    }
1779
0
    if (*assigned_seqno == 0) {
1780
0
      *assigned_seqno = last_seqno + 1;
1781
0
    }
1782
0
  }
1783
1784
0
  return status;
1785
0
}
1786
1787
Status ExternalSstFileIngestionJob::CheckLevelForIngestedBehindFile(
1788
0
    IngestedFileInfo* file_to_ingest) {
1789
0
  assert(!atomic_replace_range_.has_value());
1790
1791
0
  auto* vstorage = cfd_->current()->storage_info();
1792
  // First, check if new files fit in the last level
1793
0
  int last_lvl = cfd_->NumberLevels() - 1;
1794
0
  if (!IngestedFileFitInLevel(file_to_ingest, last_lvl)) {
1795
0
    return Status::InvalidArgument(
1796
0
        "Can't ingest_behind file as it doesn't fit "
1797
0
        "at the last level!");
1798
0
  }
1799
1800
  // Second, check if despite cf_allow_ingest_behind=true we still have 0
1801
  // seqnums at some upper level
1802
0
  for (int lvl = 0; lvl < cfd_->NumberLevels() - 1; lvl++) {
1803
0
    for (auto file : vstorage->LevelFiles(lvl)) {
1804
0
      if (file->fd.smallest_seqno == 0) {
1805
0
        return Status::InvalidArgument(
1806
0
            "Can't ingest_behind file as despite cf_allow_ingest_behind=true "
1807
0
            "there are files with 0 seqno in database at upper levels!");
1808
0
      }
1809
0
    }
1810
0
  }
1811
1812
0
  file_to_ingest->picked_level = last_lvl;
1813
0
  return Status::OK();
1814
0
}
1815
1816
Status ExternalSstFileIngestionJob::AssignGlobalSeqnoForIngestedFile(
1817
0
    IngestedFileInfo* file_to_ingest, SequenceNumber seqno) {
1818
0
  if (ingestion_options_.allow_db_generated_files) {
1819
0
    assert(seqno == 0);
1820
0
    assert(file_to_ingest->original_seqno == 0);
1821
0
  }
1822
0
  if (file_to_ingest->original_seqno == seqno) {
1823
    // This file already has the correct global seqno.
1824
0
    return Status::OK();
1825
0
  } else if (!ingestion_options_.allow_global_seqno) {
1826
0
    return Status::InvalidArgument("Global seqno is required, but disabled");
1827
0
  } else if (ingestion_options_.write_global_seqno &&
1828
0
             file_to_ingest->global_seqno_offset == 0) {
1829
0
    return Status::InvalidArgument(
1830
0
        "Trying to set global seqno for a file that don't have a global seqno "
1831
0
        "field");
1832
0
  }
1833
1834
0
  if (ingestion_options_.write_global_seqno) {
1835
    // Determine if we can write global_seqno to a given offset of file.
1836
    // If the file system does not support random write, then we should not.
1837
    // Otherwise we should.
1838
0
    std::unique_ptr<FSRandomRWFile> rwfile;
1839
0
    Status status = fs_->NewRandomRWFile(file_to_ingest->internal_file_path,
1840
0
                                         env_options_, &rwfile, nullptr);
1841
0
    TEST_SYNC_POINT_CALLBACK("ExternalSstFileIngestionJob::NewRandomRWFile",
1842
0
                             &status);
1843
0
    if (status.ok()) {
1844
0
      FSRandomRWFilePtr fsptr(std::move(rwfile), io_tracer_,
1845
0
                              file_to_ingest->internal_file_path);
1846
0
      std::string seqno_val;
1847
0
      PutFixed64(&seqno_val, seqno);
1848
0
      status = fsptr->Write(file_to_ingest->global_seqno_offset, seqno_val,
1849
0
                            IOOptions(), nullptr);
1850
0
      if (!status.ok()) {
1851
0
        ROCKS_LOG_WARN(db_options_.info_log,
1852
0
                       "Failed to write global seqno to %s: %s",
1853
0
                       file_to_ingest->internal_file_path.c_str(),
1854
0
                       status.ToString().c_str());
1855
0
        return status;
1856
0
      }
1857
1858
0
      if (status.ok()) {
1859
0
        TEST_SYNC_POINT("ExternalSstFileIngestionJob::BeforeSyncGlobalSeqno");
1860
0
        status = SyncIngestedFile(fsptr.get());
1861
0
        TEST_SYNC_POINT("ExternalSstFileIngestionJob::AfterSyncGlobalSeqno");
1862
0
        if (!status.ok()) {
1863
0
          ROCKS_LOG_WARN(db_options_.info_log,
1864
0
                         "Failed to sync ingested file %s after writing global "
1865
0
                         "sequence number: %s",
1866
0
                         file_to_ingest->internal_file_path.c_str(),
1867
0
                         status.ToString().c_str());
1868
0
        }
1869
0
      }
1870
0
      if (!status.ok()) {
1871
0
        return status;
1872
0
      }
1873
0
    } else if (!status.IsNotSupported()) {
1874
0
      ROCKS_LOG_WARN(
1875
0
          db_options_.info_log,
1876
0
          "Failed to open ingested file %s for random read/write: %s",
1877
0
          file_to_ingest->internal_file_path.c_str(),
1878
0
          status.ToString().c_str());
1879
0
      return status;
1880
0
    }
1881
0
  }
1882
1883
0
  file_to_ingest->assigned_seqno = seqno;
1884
0
  return Status::OK();
1885
0
}
1886
1887
IOStatus ExternalSstFileIngestionJob::GenerateChecksumForIngestedFile(
1888
0
    IngestedFileInfo* file_to_ingest) {
1889
0
  if (db_options_.file_checksum_gen_factory == nullptr ||
1890
0
      need_generate_file_checksum_ == false ||
1891
0
      ingestion_options_.write_global_seqno == false) {
1892
    // If file_checksum_gen_factory is not set, we are not able to generate
1893
    // the checksum. if write_global_seqno is false, it means we will use
1894
    // file checksum generated during Prepare(). This step will be skipped.
1895
0
    return IOStatus::OK();
1896
0
  }
1897
0
  std::string file_checksum;
1898
0
  std::string file_checksum_func_name;
1899
0
  std::string requested_checksum_func_name;
1900
  // TODO: rate limit file reads for checksum calculation during file ingestion.
1901
  // TODO: plumb Env::IOActivity
1902
0
  ReadOptions ro;
1903
0
  FileOptions gen_fopts;
1904
0
  gen_fopts.file_checksum_func_name = kNoFileChecksumFuncName;
1905
0
  IOStatus io_s = GenerateOneFileChecksum(
1906
0
      fs_.get(), file_to_ingest->internal_file_path,
1907
0
      db_options_.file_checksum_gen_factory.get(), requested_checksum_func_name,
1908
0
      &file_checksum, &file_checksum_func_name,
1909
0
      ingestion_options_.verify_checksums_readahead_size,
1910
0
      db_options_.allow_mmap_reads, io_tracer_, db_options_.rate_limiter.get(),
1911
0
      ro, db_options_.stats, db_options_.clock, gen_fopts);
1912
0
  if (!io_s.ok()) {
1913
0
    ROCKS_LOG_WARN(
1914
0
        db_options_.info_log, "Failed to generate checksum for %s: %s",
1915
0
        file_to_ingest->internal_file_path.c_str(), io_s.ToString().c_str());
1916
0
    return io_s;
1917
0
  }
1918
0
  file_to_ingest->file_checksum = std::move(file_checksum);
1919
0
  file_to_ingest->file_checksum_func_name = std::move(file_checksum_func_name);
1920
0
  return IOStatus::OK();
1921
0
}
1922
1923
bool ExternalSstFileIngestionJob::IngestedFileFitInLevel(
1924
0
    const IngestedFileInfo* file_to_ingest, int level) {
1925
0
  if (level == 0) {
1926
    // Files can always fit in L0
1927
0
    return true;
1928
0
  }
1929
1930
0
  auto* vstorage = cfd_->current()->storage_info();
1931
0
  Slice file_smallest_user_key(file_to_ingest->start_ukey);
1932
0
  Slice file_largest_user_key(file_to_ingest->limit_ukey);
1933
1934
0
  if (vstorage->OverlapInLevel(level, &file_smallest_user_key,
1935
0
                               &file_largest_user_key)) {
1936
    // File overlap with another files in this level, we cannot
1937
    // add it to this level
1938
0
    return false;
1939
0
  }
1940
1941
  // File did not overlap with level files, nor compaction output
1942
0
  return true;
1943
0
}
1944
1945
template <typename TWritableFile>
1946
0
Status ExternalSstFileIngestionJob::SyncIngestedFile(TWritableFile* file) {
1947
0
  assert(file != nullptr);
1948
0
  if (db_options_.use_fsync) {
1949
0
    return file->Fsync(IOOptions(), nullptr);
1950
0
  } else {
1951
0
    return file->Sync(IOOptions(), nullptr);
1952
0
  }
1953
0
}
1954
1955
Status ExternalSstFileIngestionJob::GetSeqnoBoundaryForFile(
1956
    TableReader* table_reader, SuperVersion* sv,
1957
0
    IngestedFileInfo* file_to_ingest, bool allow_data_in_errors) {
1958
0
  const auto tp = table_reader->GetTableProperties();
1959
0
  const bool has_largest_seqno = tp->HasKeyLargestSeqno();
1960
0
  SequenceNumber largest_seqno = tp->key_largest_seqno;
1961
0
  if (has_largest_seqno) {
1962
0
    file_to_ingest->largest_seqno = largest_seqno;
1963
0
    if (largest_seqno == 0) {
1964
0
      file_to_ingest->smallest_seqno = 0;
1965
0
      return Status::OK();
1966
0
    }
1967
0
    if (tp->HasKeySmallestSeqno()) {
1968
0
      file_to_ingest->smallest_seqno = tp->key_smallest_seqno;
1969
0
      return Status::OK();
1970
0
    }
1971
0
  }
1972
1973
  // For older SST files they may not be recorded in table properties, so
1974
  // we scan the file to find out.
1975
0
  TEST_SYNC_POINT(
1976
0
      "ExternalSstFileIngestionJob::GetSeqnoBoundaryForFile:FileScan");
1977
0
  SequenceNumber smallest_seqno = kMaxSequenceNumber;
1978
0
  SequenceNumber largest_seqno_from_iter = 0;
1979
0
  ReadOptions ro;
1980
0
  ro.fill_cache = ingestion_options_.fill_cache;
1981
0
  std::unique_ptr<InternalIterator> iter(table_reader->NewIterator(
1982
0
      ro, sv->mutable_cf_options.prefix_extractor.get(), /*arena=*/nullptr,
1983
0
      /*skip_filters=*/false, TableReaderCaller::kExternalSSTIngestion));
1984
0
  ParsedInternalKey key;
1985
0
  iter->SeekToFirst();
1986
0
  while (iter->Valid()) {
1987
0
    Status pik_status =
1988
0
        ParseInternalKey(iter->key(), &key, allow_data_in_errors);
1989
0
    if (!pik_status.ok()) {
1990
0
      return Status::Corruption("Corrupted key in external file. ",
1991
0
                                pik_status.getState());
1992
0
    }
1993
0
    smallest_seqno = std::min(smallest_seqno, key.sequence);
1994
0
    largest_seqno_from_iter = std::max(largest_seqno_from_iter, key.sequence);
1995
0
    iter->Next();
1996
0
  }
1997
0
  if (!iter->status().ok()) {
1998
0
    return iter->status();
1999
0
  }
2000
2001
0
  if (table_reader->GetTableProperties()->num_range_deletions > 0) {
2002
0
    std::unique_ptr<InternalIterator> range_del_iter(
2003
0
        table_reader->NewRangeTombstoneIterator(ro));
2004
0
    if (range_del_iter != nullptr) {
2005
0
      for (range_del_iter->SeekToFirst(); range_del_iter->Valid();
2006
0
           range_del_iter->Next()) {
2007
0
        Status pik_status =
2008
0
            ParseInternalKey(range_del_iter->key(), &key, allow_data_in_errors);
2009
0
        if (!pik_status.ok()) {
2010
0
          return Status::Corruption("Corrupted key in external file. ",
2011
0
                                    pik_status.getState());
2012
0
        }
2013
0
        smallest_seqno = std::min(smallest_seqno, key.sequence);
2014
0
        largest_seqno_from_iter =
2015
0
            std::max(largest_seqno_from_iter, key.sequence);
2016
0
      }
2017
0
      if (!range_del_iter->status().ok()) {
2018
0
        return range_del_iter->status();
2019
0
      }
2020
0
    }
2021
0
  }
2022
2023
0
  file_to_ingest->smallest_seqno = smallest_seqno;
2024
0
  if (!has_largest_seqno) {
2025
0
    file_to_ingest->largest_seqno = largest_seqno_from_iter;
2026
0
  } else {
2027
0
    assert(largest_seqno == largest_seqno_from_iter);
2028
0
    file_to_ingest->largest_seqno = largest_seqno;
2029
0
  }
2030
2031
0
  if (file_to_ingest->largest_seqno == kMaxSequenceNumber) {
2032
0
    return Status::InvalidArgument(
2033
0
        "Unknown smallest seqno for db generated file.");
2034
0
  }
2035
0
  if (file_to_ingest->smallest_seqno == kMaxSequenceNumber) {
2036
0
    return Status::InvalidArgument(
2037
0
        "Unknown largest seqno for db generated file.");
2038
0
  }
2039
0
  return Status::OK();
2040
0
}
2041
2042
}  // namespace ROCKSDB_NAMESPACE