Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
88 changes: 78 additions & 10 deletions fdbclient/FileBackupAgent.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -3803,12 +3803,10 @@ struct BulkDumpTaskFunc : BackupTaskFuncBase {
// Enable BulkDump mode at the DD level before submitting the job
co_await setBulkDumpMode(cx, 1);

// Configure BulkDump job for the full keyspace
// BulkDump/BulkLoad requires the load range to be a subset of the dump range.
// Using normalKeys for both ensures compatibility regardless of user-specified ranges.
// submitBackup rejects a multi-range bulkdump, so the caller's coverage is this one range.
// Store data under data/<container>/bulkdump_data/ to be consistent with backup container layout
std::string bulkDumpRoot = getBackupDataPath(bc->getURL(), "bulkdump_data");
bulkDumpJob = createBulkDumpJob(normalKeys, bulkDumpRoot, BulkLoadType::SST, transportMethod);
bulkDumpJob = createBulkDumpJob(backupRanges[0], bulkDumpRoot, BulkLoadType::SST, transportMethod);

// Submit the BulkDump job
co_await submitBulkDumpJob(cx, bulkDumpJob);
Expand Down Expand Up @@ -4395,9 +4393,16 @@ struct BulkLoadRestoreTaskFunc : RestoreTaskFuncBase {
BulkLoadTransportMethod loadTransportMethod =
isBlobstoreUrl(backupUrl) ? BulkLoadTransportMethod::BLOBSTORE : BulkLoadTransportMethod::CP;

// BulkLoad range must match the BulkDump range (normalKeys)
// The actual restore ranges will be applied via mutation log replay
BulkLoadJobState bulkLoadJob = createBulkLoadJob(dumpJobUid, normalKeys, jobRoot, loadTransportMethod);
// submitRestore rejects a multi-range bulkload, so range 0 is the caller's whole request.
//
// TODO(BulkLoad): reject a restore range the backup never covered. Until the ranges were
// plumbed through, dump and load were both normalKeys, so an out-of-range request could not
// arise; now it submits a job over a range with no manifests and loads nothing rather than
// failing. The dump's actual coverage is the "ranges" array of the keyspace snapshot
// (written at BackupContainerFileSystem.cpp:236-243), i.e. the same JSONDoc parsed above for
// bulkDumpJobId, so the containment check needs no extra read.
BulkLoadJobState bulkLoadJob =
createBulkLoadJob(dumpJobUid, restoreRanges[0], jobRoot, loadTransportMethod);

TraceEvent("BulkLoadRestoreJobCreated")
.detail("RestoreUID", restore.getUid())
Expand Down Expand Up @@ -6848,8 +6853,9 @@ struct StartFullRestoreTaskFunc : RestoreTaskFuncBase {
bulkDumpJobId,
TaskCompletionKey::signal(bulkLoadDone));

// After BulkLoad completes, run RestoreDispatch to apply mutation logs
// Set onlyApplyMutationLogs so it only processes logs, not range files
// Sequencing, not user intent: the SSTs carry the range data, so the dispatch that follows must
// replay logs only. submitRestore rejects --incremental with bulkload, so this can only ever
// overwrite false.
restore.onlyApplyMutationLogs().set(tr, true);

// Add RestoreDispatch task that waits for BulkLoad to complete
Expand Down Expand Up @@ -6882,7 +6888,13 @@ struct StartFullRestoreTaskFunc : RestoreTaskFuncBase {
// If this is an incremental restore, we need to set the applyMutationsMapPrefix
// to the earliest log version so no mutations are missed
Value versionEncoded = BinaryWriter::toValue(Params.firstVersion().get(task), Unversioned());
co_await krmSetRange(tr, restore.applyMutationsMapPrefix(), normalKeys, versionEncoded);
std::vector<KeyRange> ranges = co_await restore.getRestoreRangesOrDefault(tr);
std::vector<Future<Void>> updateMap;
updateMap.reserve(ranges.size());
for (const auto& range : ranges) {
updateMap.push_back(krmSetRange(tr, restore.applyMutationsMapPrefix(), range, versionEncoded));
}
co_await waitForAll(updateMap);
}

co_await taskBucket->finish(tr, task);
Expand Down Expand Up @@ -7115,6 +7127,37 @@ class FileBackupAgentImpl {
}
}

// A bulkdump job carries exactly one key range and the cluster admits one bulk job at a time, so
// several disjoint ranges cannot be expressed as a single snapshot. Rejecting here is what keeps
// the restore honest: the dump would otherwise widen to normalKeys and a later bulkload restore
// would overwrite every key between the caller's ranges. BulkDumpTaskFunc indexes range 0 on the
// strength of this check.
if (snapshotMode != static_cast<int>(SnapshotMode::RANGEFILE) && normalizedRanges.size() != 1) {
TraceEvent(SevWarnAlways, "FBA_SubmitBackupBulkDumpRangeCount")
.detail("TagName", tagName)
.detail("SnapshotMode", snapshotMode)
.detail("RangeCount", normalizedRanges.size());
fprintf(stderr,
"ERROR: --mode bulkdump and --mode both require exactly one key range, but %d remain after "
"coalescing adjacent ranges. Use --mode rangefile, or run one backup per range.\n",
static_cast<int>(normalizedRanges.size()));
throw backup_error();
}

// --incremental suppresses the snapshot entirely (StartFullBackupTaskFunc skips the whole snapshot
// block), so pairing it with a bulkdump mode asks for an SST snapshot and for no snapshot at once.
// Taken silently, it yields a log-only backup that no bulkload restore can use, and the operator
// only discovers that when the restore aborts much later.
if (snapshotMode != static_cast<int>(SnapshotMode::RANGEFILE) && incrementalBackupOnly) {
TraceEvent(SevWarnAlways, "FBA_SubmitBackupBulkDumpIncremental")
.detail("TagName", tagName)
.detail("SnapshotMode", snapshotMode);
fprintf(stderr,
"ERROR: --incremental cannot be combined with --mode bulkdump or --mode both; it suppresses "
"the snapshot those modes exist to produce, leaving a backup no bulkload restore can use.\n");
throw backup_error();
}

config.clear(tr);

Key destUidValue(BinaryWriter::toValue(uid, Unversioned()));
Expand Down Expand Up @@ -7202,6 +7245,31 @@ class FileBackupAgentImpl {
ASSERT(restoreRange.begin.startsWith(removePrefix) && restoreRange.end.startsWith(removePrefix));
}

// Same one-range limit as bulkdump: the bulkload job that ingests the SSTs carries a single range.
// BulkLoadRestoreTaskFunc indexes range 0 on the strength of this check.
if (!useRangeFileRestore && restoreRanges.size() != 1) {
TraceEvent(SevWarnAlways, "FBA_SubmitRestoreBulkLoadRangeCount")
.detail("TagName", tagName)
.detail("RangeCount", restoreRanges.size());
fprintf(stderr,
"ERROR: --mode bulkload requires exactly one key range, but %d remain after coalescing "
"adjacent ranges. Use --mode rangefile, or run one restore per range.\n",
static_cast<int>(restoreRanges.size()));
throw restore_error();
}

// Mirror of the backup-side rejection. A bulkload restore ingests the snapshot via SST, which is
// precisely the range data --incremental declines, and StartFullRestoreTaskFunc picks the bulkload
// branch without consulting the flag. Accepting the pair would restore the whole snapshot the
// caller asked to skip, then overwrite their setting so even status misreports it.
if (!useRangeFileRestore && onlyApplyMutationLogs) {
TraceEvent(SevWarnAlways, "FBA_SubmitRestoreBulkLoadIncremental").detail("TagName", tagName);
fprintf(stderr,
"ERROR: --incremental cannot be combined with --mode bulkload; bulkload restores the "
"snapshot that --incremental skips. Use --mode rangefile for a logs-only restore.\n");
throw restore_error();
}

tr->setOption(FDBTransactionOptions::ACCESS_SYSTEM_KEYS);
tr->setOption(FDBTransactionOptions::LOCK_AWARE);

Expand Down
151 changes: 148 additions & 3 deletions fdbserver/workloads/BackupS3BlobCorrectness.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,13 @@ extern std::atomic<int> g_bulkLoadRestoreTaskCompleteCount;
// S3-specific backup correctness workload - see file header for differences from BackupAndRestoreCorrectness
struct BackupS3BlobCorrectnessWorkload : TestWorkload {
static constexpr auto NAME = "BackupS3BlobCorrectness";

// Inside normalKeys but below every alphanumeric key, so these sort outside any range this workload
// generates while still being ordinary user keys. The '_' after the escape keeps \x01 from swallowing
// a following hex digit.
static inline const StringRef SENTINEL_PREFIX = "\x01_bs3bc_outside/"_sr;
static inline const ValueRef SENTINEL_BEFORE_BACKUP = "before_backup"_sr;
static inline const ValueRef SENTINEL_AFTER_BACKUP = "after_backup"_sr;
double backupAfter, restoreAfter, abortAndRestartAfter;
double minBackupAfter;
double backupStartAt, restoreStartAfterBackupFinished, stopDifferentialAfter;
Expand All @@ -114,6 +121,7 @@ struct BackupS3BlobCorrectnessWorkload : TestWorkload {
bool differentialBackup, performRestore, agentRequest;
Standalone<VectorRef<KeyRangeRef>> backupRanges;
std::vector<KeyRange> skippedRestoreRanges;
std::vector<Key> outOfRangeSentinels;
Standalone<VectorRef<KeyRangeRef>> restoreRanges;
static int backupAgentRequests;
LockDB locked{ false };
Expand All @@ -137,6 +145,9 @@ struct BackupS3BlobCorrectnessWorkload : TestWorkload {
// performValidation: if true, validates backup by restoring with prefix and running audit_storage validate_restore
// This must happen BEFORE clearing the database so we can compare original vs restored data
bool performValidation;
// expectBulkIncrementalRejection: assert that --incremental is refused on both bulk paths. Kept
// behind an option and given its own toml so the probes cannot disturb the key-range regression test.
bool expectBulkIncrementalRejection;
// expectRestoreFailure: the bulkload restore is expected NOT to complete, so the workload asserts the
// failure is reported and recovered from rather than that the data arrived.
bool expectRestoreFailure;
Expand Down Expand Up @@ -198,6 +209,7 @@ struct BackupS3BlobCorrectnessWorkload : TestWorkload {
// performValidation: Validates backup by comparing original data vs restored data
// Uses audit_storage validate_restore - must happen BEFORE clearing database
performValidation = getOption(options, "performValidation"_sr, false);
expectBulkIncrementalRejection = getOption(options, "expectBulkIncrementalRejection"_sr, false);
// Mutually exclusive with performValidation, whose audit compares restored contents that by
// definition do not exist when the restore is expected to fail.
expectRestoreFailure = getOption(options, "expectRestoreFailure"_sr, false);
Expand Down Expand Up @@ -317,9 +329,37 @@ struct BackupS3BlobCorrectnessWorkload : TestWorkload {
}
}

// Backup everything
self->backupRanges.push_back_deep(self->backupRanges.arena(), normalKeys);
self->restoreRanges.push_back_deep(self->restoreRanges.arena(), normalKeys);
// Keys planted outside the backup ranges, each holding SENTINEL_BEFORE_BACKUP. Once the backup has
// stopped they are rewritten to SENTINEL_AFTER_BACKUP, so a restore that widens beyond the backup
// ranges reverts them and is caught. The rewrite has to follow the backup rather than precede it:
// a value the backup captured would be restored to the value the check expects, and the widening
// would pass unnoticed. When the ranges cover all of normalKeys there is no outside and the set is
// empty.
for (int i = 0; i < 3; i++) {
Key candidate(SENTINEL_PREFIX.toString() + std::to_string(i));
bool covered = false;
for (const auto& range : self->backupRanges) {
if (range.contains(candidate)) {
covered = true;
break;
}
}
if (!covered) {
self->outOfRangeSentinels.push_back(candidate);
}
}

if (!self->outOfRangeSentinels.empty()) {
co_await runRYWTransaction(cx, [=](Reference<ReadYourWritesTransaction> tr) -> Future<Void> {
tr->setOption(FDBTransactionOptions::ACCESS_SYSTEM_KEYS);
tr->setOption(FDBTransactionOptions::LOCK_AWARE);
for (const auto& key : self->outOfRangeSentinels) {
tr->set(key, SENTINEL_BEFORE_BACKUP);
}
return Void();
});
TraceEvent("BS3BCW_PlantedOutOfRangeSentinels").detail("Count", self->outOfRangeSentinels.size());
}
}

Future<Void> start(Database const& cx) override {
Expand All @@ -335,6 +375,21 @@ struct BackupS3BlobCorrectnessWorkload : TestWorkload {

void getMetrics(std::vector<PerfMetric>& m) override {}

static Future<Void> checkOutOfRangeSentinels(Reference<ReadYourWritesTransaction> tr, std::vector<Key> sentinels) {
for (const auto& key : sentinels) {
Optional<Value> value = co_await tr->get(key);
bool intact = value.present() && value.get() == SENTINEL_AFTER_BACKUP;
if (!intact) {
TraceEvent(SevError, "BS3BCW_RestoreWroteOutsideBackupRanges")
.detail("Key", printable(key))
.detail("Expected", printable(SENTINEL_AFTER_BACKUP))
.detail("Actual", value.present() ? printable(value.get()) : std::string("<missing>"));
}
ASSERT(intact);
}
co_return;
}

static Future<Void> changePaused(Database cx, FileBackupAgent* backupAgent) {
while (true) {
co_await backupAgent->taskBucket->changePause(cx, deterministicRandom()->coinflip());
Expand Down Expand Up @@ -518,6 +573,38 @@ struct BackupS3BlobCorrectnessWorkload : TestWorkload {

// Testing v1 (non-partitioned) backup approach
// This does not require backup workers
if (expectBulkIncrementalRejection) {
// --incremental suppresses the snapshot that a bulk mode exists to produce, so submitBackup
// must refuse the pair outright. A separate tag keeps a wrongly-accepted submit from
// colliding with the real backup below; the assert fails the test either way.
Error rejected;
try {
co_await backupAgent->submitBackup(cx,
StringRef(backupContainer),
{},
initSnapshotInterval,
snapshotInterval,
tag.toString() + "_incrreject",
backupRanges,
StopWhenDone{ !stopDifferentialDelay },
MutationLogType::DEFAULT,
IncrementalBackupOnly::True,
encryptionKeyFileName,
encryptionKeyFileName.present() ? DEFAULT_ENCRYPTION_BLOCK_SIZE : 0,
snapshotMode);
} catch (Error& e) {
if (e.code() == error_code_actor_cancelled) {
throw;
}
rejected = e;
}
TraceEvent("BS3BCW_IncrementalBackupRejection", randomID)
.error(rejected)
.detail("SnapshotMode", snapshotMode)
.detail("Accepted", rejected.code() == error_code_success);
ASSERT(rejected.code() == error_code_backup_error);
}

try {
co_await backupAgent->submitBackup(cx,
StringRef(backupContainer),
Expand Down Expand Up @@ -753,6 +840,19 @@ struct BackupS3BlobCorrectnessWorkload : TestWorkload {
});
}

// The backup has stopped, so this value cannot reach the snapshot or the mutation logs.
if (!outOfRangeSentinels.empty()) {
co_await runRYWTransaction(cx, [=](Reference<ReadYourWritesTransaction> tr) -> Future<Void> {
tr->setOption(FDBTransactionOptions::ACCESS_SYSTEM_KEYS);
tr->setOption(FDBTransactionOptions::LOCK_AWARE);
for (const auto& key : outOfRangeSentinels) {
tr->set(key, SENTINEL_AFTER_BACKUP);
}
return Void();
});
TraceEvent("BS3BCW_BumpedOutOfRangeSentinels").detail("Count", outOfRangeSentinels.size());
}

// Step 3: Perform the restore (with BulkLoad if configured)
TraceEvent("BS3BCW_Restore")
.detail("LastBackupContainer", lastBackupContainer->getURL())
Expand All @@ -765,6 +865,42 @@ struct BackupS3BlobCorrectnessWorkload : TestWorkload {
// and our lockUID so restore uses the same lock for checkDatabaseLock calls
Version v = ::invalidVersion;
Error restoreError;
if (expectBulkIncrementalRejection) {
// Mirror on the restore side: bulkload ingests the snapshot --incremental skips,
// so submitRestore must refuse the pair. Runs against the backup just taken, so
// the container is describable and a throw can only come from the guard.
Error rejected;
try {
co_await backupAgent.restore(cx,
cx,
Standalone<StringRef>(restoreTag.toString() + "_incrreject"),
KeyRef(lastBackupContainer->getURL()),
lastBackupContainer->getProxy(),
restoreRanges,
WaitForComplete::False,
::invalidVersion,
Verbose::True,
Key(),
Key(),
LockDB::False,
UnlockDB::False,
OnlyApplyMutationLogs::True,
InconsistentSnapshotOnly::False,
::invalidVersion,
lastBackupContainer->getEncryptionKeyFileName(),
lockUID,
/*useRangeFileRestore=*/false);
} catch (Error& e) {
if (e.code() == error_code_actor_cancelled) {
throw;
}
rejected = e;
}
TraceEvent("BS3BCW_IncrementalRestoreRejection")
.error(rejected)
.detail("Accepted", rejected.code() == error_code_success);
ASSERT(rejected.code() == error_code_restore_error);
}
try {
v = co_await backupAgent.restore(cx,
cx,
Expand Down Expand Up @@ -842,6 +978,15 @@ struct BackupS3BlobCorrectnessWorkload : TestWorkload {
ASSERT(bulkLoadCount > 0);
}

// A restore confined to the backup ranges cannot have touched these keys.
if (!outOfRangeSentinels.empty()) {
co_await runRYWTransaction(cx, [=](Reference<ReadYourWritesTransaction> tr) -> Future<Void> {
tr->setOption(FDBTransactionOptions::ACCESS_SYSTEM_KEYS);
tr->setOption(FDBTransactionOptions::LOCK_AWARE);
return checkOutOfRangeSentinels(tr, outOfRangeSentinels);
});
}

// Step 4: Run audit to compare BulkLoad-restored vs traditional-restored
if (performValidation) {
TraceEvent("BS3BCW_ValidationStep4_AuditStarting")
Expand Down
3 changes: 2 additions & 1 deletion tests/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -176,7 +176,8 @@ if(WITH_PYTHON)
add_fdb_test(TEST_FILES slow/BulkDumpingS3WithChaos.toml)
add_fdb_test(TEST_FILES slow/BackupS3BlobBulkLoadRestore.toml)
add_fdb_test(TEST_FILES slow/BackupS3BlobBulkLoadRestoreWithChaos.toml)
add_fdb_test(TEST_FILES slow/BackupS3BlobBulkLoadRestoreMultiRange.toml)
add_fdb_test(TEST_FILES slow/BackupS3BlobBulkLoadRestoreKeyRange.toml)
add_fdb_test(TEST_FILES slow/BackupS3BlobBulkLoadIncrementalRejected.toml)
add_fdb_test(TEST_FILES slow/BackupS3BlobBulkLoadRestoreJobIncomplete.toml)
add_fdb_test(TEST_FILES fast/BulkLoading.toml)
add_fdb_test(TEST_FILES fast/BulkLoadingDestTeamFailure.toml)
Expand Down
Loading
Loading